diff --git a/Cargo.lock b/Cargo.lock
index fe54ccac6f1..4e48f2cb66f 100644
--- a/Cargo.lock
+++ b/Cargo.lock
@@ -7646,6 +7646,7 @@ dependencies = [
"omicron-uuid-kinds",
"omicron-workspace-hack",
"parallel-task-set",
+ "perfetto-trace",
"regex",
"reqwest 0.13.2",
"serde_json",
@@ -8268,6 +8269,7 @@ dependencies = [
"oxnet",
"oxql-types",
"parse-display",
+ "perfetto-trace",
"proptest",
"regex",
"schemars 0.8.22",
@@ -9243,6 +9245,7 @@ dependencies = [
"parse-display",
"paste",
"pem",
+ "perfetto-trace",
"petgraph 0.8.3",
"pq-sys",
"pretty_assertions",
@@ -11049,6 +11052,16 @@ version = "2.3.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9b4f627cb1b25917193a259e49bdad08f671f8d9708acfd5fe0a8c1455d87220"
+[[package]]
+name = "perfetto-trace"
+version = "0.1.0"
+dependencies = [
+ "chrono",
+ "omicron-workspace-hack",
+ "serde",
+ "serde_json",
+]
+
[[package]]
name = "pest"
version = "2.8.6"
@@ -15147,6 +15160,7 @@ dependencies = [
"omicron-uuid-kinds",
"omicron-workspace-hack",
"parallel-task-set",
+ "perfetto-trace",
"pq-sys",
"serde",
"serde_json",
diff --git a/Cargo.toml b/Cargo.toml
index fc16e0bc2d9..82126e0da9f 100644
--- a/Cargo.toml
+++ b/Cargo.toml
@@ -144,6 +144,7 @@ members = [
"package",
"parallel-task-set",
"passwords",
+ "perfetto-trace",
"range-requests",
"rpaths",
"sled-agent",
@@ -339,6 +340,7 @@ default-members = [
"package",
"parallel-task-set",
"passwords",
+ "perfetto-trace",
"range-requests",
"rpaths",
"sled-agent",
@@ -728,6 +730,7 @@ partial-io = { version = "0.5.4", features = ["proptest1", "tokio1"] }
parse-size = "1.1.0"
paste = "1.0.15"
percent-encoding = "2.3.1"
+perfetto-trace = { path = "perfetto-trace" }
peg = "0.8.5"
pem = "3.0"
# petname's default features pull in clap for CLI parsing, which we don't need.
diff --git a/dev-tools/omdb/src/bin/omdb/nexus.rs b/dev-tools/omdb/src/bin/omdb/nexus.rs
index aedd47acc6c..53a329068a2 100644
--- a/dev-tools/omdb/src/bin/omdb/nexus.rs
+++ b/dev-tools/omdb/src/bin/omdb/nexus.rs
@@ -68,6 +68,7 @@ use nexus_types::internal_api::background::FmRendezvousStatus;
use nexus_types::internal_api::background::IncompleteBootstoreConfigReport;
use nexus_types::internal_api::background::InstanceReincarnationStatus;
use nexus_types::internal_api::background::InstanceUpdaterStatus;
+use nexus_types::internal_api::background::InventoryCollectionStatus;
use nexus_types::internal_api::background::InventoryLoadStatus;
use nexus_types::internal_api::background::LookupRegionPortStatus;
use nexus_types::internal_api::background::PhysicalDiskAdoptionStatus;
@@ -208,10 +209,23 @@ enum BackgroundTasksCommands {
Show(BackgroundTasksShowArgs),
/// Print an event report for a background task if available.
PrintReport(BackgroundTasksPrintReportArgs),
+ /// Save the last inventory collection's timing trace to a file
+ ///
+ /// The file is in Chrome Trace Event format and can be loaded into
+ /// to visualize where collection time was
+ /// spent.
+ InventoryTrace(BackgroundTasksInventoryTraceArgs),
/// Activate one or more background tasks
Activate(BackgroundTasksActivateArgs),
}
+#[derive(Debug, Args)]
+struct BackgroundTasksInventoryTraceArgs {
+ /// where to write the trace JSON
+ #[clap(long)]
+ output: Utf8PathBuf,
+}
+
#[derive(Debug, Args)]
struct BackgroundTasksShowArgs {
/// Names of background tasks to show (default: all)
@@ -746,6 +760,11 @@ impl NexusArgs {
)
.await
}
+ NexusCommands::BackgroundTasks(BackgroundTasksArgs {
+ command: BackgroundTasksCommands::InventoryTrace(args),
+ }) => {
+ cmd_nexus_background_tasks_inventory_trace(&client, args).await
+ }
NexusCommands::BackgroundTasks(BackgroundTasksArgs {
command: BackgroundTasksCommands::Activate(args),
}) => {
@@ -1132,6 +1151,46 @@ async fn cmd_nexus_background_tasks_print_report(
Ok(())
}
+/// Runs `omdb nexus background-tasks inventory-trace`
+async fn cmd_nexus_background_tasks_inventory_trace(
+ client: &nexus_lockstep_client::Client,
+ args: &BackgroundTasksInventoryTraceArgs,
+) -> Result<(), anyhow::Error> {
+ const TASK_NAME: &str = "inventory_collection";
+ let response = client
+ .bgtask_view(TASK_NAME)
+ .await
+ .context("fetching background task")?;
+ let task = response.into_inner();
+ let LastResult::Completed(last) = task.last else {
+ bail!("task {:?} has never completed", TASK_NAME);
+ };
+ let status: InventoryCollectionStatus =
+ serde_json::from_value(last.details.clone()).with_context(|| {
+ format!(
+ "interpreting task details (did the last activation fail?) \
+ -- found {:?}",
+ last.details
+ )
+ })?;
+ let Some(trace) = status.trace else {
+ bail!(
+ "task status has no trace (is this Nexus running a version \
+ that records one?)"
+ );
+ };
+ let json =
+ serde_json::to_string_pretty(&trace).context("serializing trace")?;
+ std::fs::write(&args.output, json)
+ .with_context(|| format!("writing {:?}", args.output))?;
+ println!(
+ "wrote trace for collection {} to {}",
+ status.collection_id, args.output
+ );
+ println!("load it into https://ui.perfetto.dev/ to visualize");
+ Ok(())
+}
+
/// Runs `omdb nexus background-tasks activate`
async fn cmd_nexus_background_tasks_activate(
client: &nexus_lockstep_client::Client,
@@ -2195,14 +2254,7 @@ fn print_task_instance_watcher(details: &serde_json::Value) {
}
fn print_task_inventory_collection(details: &serde_json::Value) {
- #[derive(Deserialize)]
- struct InventorySuccess {
- collection_id: Uuid,
- time_started: DateTime,
- time_done: DateTime,
- }
-
- match serde_json::from_value::(details.clone()) {
+ match serde_json::from_value::(details.clone()) {
Err(error) => eprintln!(
"warning: failed to interpret task details: {:?}: {:?}",
error, details
@@ -2224,6 +2276,20 @@ fn print_task_inventory_collection(details: &serde_json::Value) {
.time_done
.to_rfc3339_opts(SecondsFormat::Secs, true),
);
+ if let Some(trace) = &found_inventory.trace {
+ println!(" phase timings:");
+ for event in
+ trace.trace_events.iter().filter(|e| e.cat == "phase")
+ {
+ // Bare integer milliseconds: the omdb test output
+ // redactor recognizes exactly this form.
+ println!(" {}: {}ms", event.name, event.dur / 1000);
+ }
+ println!(
+ " (fetch the full trace with `omdb nexus \
+ background-tasks inventory-trace`)"
+ );
+ }
}
};
}
diff --git a/dev-tools/omdb/tests/successes.out b/dev-tools/omdb/tests/successes.out
index 129371ef432..967efb2e1b4 100644
--- a/dev-tools/omdb/tests/successes.out
+++ b/dev-tools/omdb/tests/successes.out
@@ -859,6 +859,14 @@ task: "inventory_collection"
last collection id: .....................
last collection started:
last collection done:
+ phase timings:
+ mgs: ms
+ sled_agents: ms
+ keepers: ms
+ cockroach: ms
+ timesync: ms
+ dns_generations: ms
+ (fetch the full trace with `omdb nexus background-tasks inventory-trace`)
task: "inventory_loader"
configured period: every s
@@ -1584,6 +1592,14 @@ task: "inventory_collection"
last collection id: .....................
last collection started:
last collection done:
+ phase timings:
+ mgs: ms
+ sled_agents: ms
+ keepers: ms
+ cockroach: ms
+ timesync: ms
+ dns_generations: ms
+ (fetch the full trace with `omdb nexus background-tasks inventory-trace`)
task: "inventory_loader"
configured period: every s
diff --git a/dev-tools/omdb/tests/usage_errors.out b/dev-tools/omdb/tests/usage_errors.out
index 0cbb3bb455a..f72035266ab 100644
--- a/dev-tools/omdb/tests/usage_errors.out
+++ b/dev-tools/omdb/tests/usage_errors.out
@@ -1419,12 +1419,13 @@ print information about background tasks
Usage: omdb nexus background-tasks [OPTIONS]
Commands:
- doc Show documentation about background tasks
- list Print a summary of the status of all background tasks
- show Print human-readable summary of the status of each background task
- print-report Print an event report for a background task if available
- activate Activate one or more background tasks
- help Print this message or the help of the given subcommand(s)
+ doc Show documentation about background tasks
+ list Print a summary of the status of all background tasks
+ show Print human-readable summary of the status of each background task
+ print-report Print an event report for a background task if available
+ inventory-trace Save the last inventory collection's timing trace to a file
+ activate Activate one or more background tasks
+ help Print this message or the help of the given subcommand(s)
Options:
--log-level log level filter [env: LOG_LEVEL=] [default: warn]
diff --git a/nexus/Cargo.toml b/nexus/Cargo.toml
index 710224338af..01ca6a67787 100644
--- a/nexus/Cargo.toml
+++ b/nexus/Cargo.toml
@@ -90,6 +90,7 @@ oxql-types.workspace = true
parallel-task-set.workspace = true
parse-display.workspace = true
paste.workspace = true
+perfetto-trace.workspace = true
# See omicron-rpaths for more about the "pq-sys" dependency.
pq-sys = "*"
progenitor-client.workspace = true
diff --git a/nexus/inventory/Cargo.toml b/nexus/inventory/Cargo.toml
index 414a184ff89..b79bb985c9b 100644
--- a/nexus/inventory/Cargo.toml
+++ b/nexus/inventory/Cargo.toml
@@ -34,6 +34,7 @@ ntp-admin-client.workspace = true
omicron-common.workspace = true
omicron-uuid-kinds.workspace = true
parallel-task-set.workspace = true
+perfetto-trace.workspace = true
reqwest.workspace = true
serde_json.workspace = true
sled-agent-client.workspace = true
diff --git a/nexus/inventory/src/collector.rs b/nexus/inventory/src/collector.rs
index 12c2b97bcee..ac6981b916f 100644
--- a/nexus/inventory/src/collector.rs
+++ b/nexus/inventory/src/collector.rs
@@ -9,11 +9,13 @@ use crate::builder::CollectionBuilder;
use crate::builder::InventoryError;
use anyhow::Context;
use anyhow::anyhow;
+use chrono::Utc;
use clickhouse_admin_keeper_client::ClientInfo as _;
use clickhouse_admin_types::keeper::ClickhouseKeeperClusterMembership;
use gateway_client::types::GetCfpaParams;
use gateway_client::types::RotCfpaSlot;
use gateway_messages::SpComponent;
+use gateway_types::component::SpIdentifier;
use itertools::Itertools;
use nexus_types::inventory::CabooseWhich;
use nexus_types::inventory::Collection;
@@ -27,6 +29,8 @@ use omicron_common::address::NTP_ADMIN_PORT;
use omicron_common::disk::M2Slot;
use omicron_uuid_kinds::OmicronZoneUuid;
use parallel_task_set::ParallelTaskSet;
+use perfetto_trace::TraceSpan;
+use perfetto_trace::timed;
use sled_agent_types::inventory::Inventory;
use sled_agent_types::inventory::OmicronZoneType;
use sled_agent_types::inventory::ZoneKind;
@@ -52,6 +56,7 @@ pub struct Collector<'a> {
cockroach_admin_client: &'a CockroachClusterAdminClient,
sled_agent_lister: &'a (dyn SledAgentEnumerator + Send + Sync),
in_progress: CollectionBuilder,
+ spans: Vec,
}
impl<'a> Collector<'a> {
@@ -70,6 +75,7 @@ impl<'a> Collector<'a> {
cockroach_admin_client,
sled_agent_lister,
in_progress: CollectionBuilder::new(creator),
+ spans: Vec::new(),
}
}
@@ -80,7 +86,13 @@ impl<'a> Collector<'a> {
/// components. This can take a while and produce any number of errors.
/// Such errors generally don't cause this function to fail. Rather, the
/// returned `Collection` keeps track of these errors.
- pub async fn collect_all(mut self) -> Result {
+ ///
+ /// Also returns a set of timing spans covering each phase of the
+ /// collection and each request made within the concurrent phases,
+ /// suitable for assembling into a trace with `perfetto_trace::assemble`.
+ pub async fn collect_all(
+ mut self,
+ ) -> Result<(Collection, Vec), anyhow::Error> {
// Most of the phases below fan requests out to their targets with
// bounded concurrency (MAX_CONCURRENT_INVENTORY_REQUESTS) and then
// merge the results into the in-progress collection on this task.
@@ -95,26 +107,53 @@ impl<'a> Collector<'a> {
debug!(&self.log, "begin collection");
+ let start = Utc::now();
self.collect_all_mgs().await;
+ self.spans.push(TraceSpan::since("mgs", "phase", start));
+
+ let start = Utc::now();
self.collect_all_sled_agents().await;
+ self.spans.push(TraceSpan::since("sled_agents", "phase", start));
+
+ let start = Utc::now();
self.collect_all_keepers().await;
+ self.spans.push(TraceSpan::since("keepers", "phase", start));
+
+ let start = Utc::now();
self.collect_all_cockroach().await;
+ self.spans.push(TraceSpan::since("cockroach", "phase", start));
// The following must be called after "collect_all_sled_agents",
// or they'll see an empty set of services.
+ let start = Utc::now();
self.collect_all_timesync().await;
+ self.spans.push(TraceSpan::since("timesync", "phase", start));
+
+ let start = Utc::now();
self.collect_all_dns_generations().await;
+ self.spans.push(TraceSpan::since("dns_generations", "phase", start));
debug!(&self.log, "finished collection");
- Ok(self.in_progress.build())
+ Ok((self.in_progress.build(), self.spans))
}
/// Collect inventory from all MGS instances
async fn collect_all_mgs(&mut self) {
for client in &self.mgs_clients {
- Self::collect_one_mgs(client, &self.log, &mut self.in_progress)
- .await;
+ let start = Utc::now();
+ Self::collect_one_mgs(
+ client,
+ &self.log,
+ &mut self.in_progress,
+ &mut self.spans,
+ )
+ .await;
+ self.spans.push(TraceSpan::since(
+ client.baseurl(),
+ "mgs_client",
+ start,
+ ));
}
}
@@ -122,6 +161,7 @@ impl<'a> Collector<'a> {
client: &gateway_client::Client,
log: &Logger,
in_progress: &mut CollectionBuilder,
+ spans: &mut Vec,
) {
debug!(log, "begin collection from MGS";
"mgs_url" => client.baseurl()
@@ -163,277 +203,280 @@ impl<'a> Collector<'a> {
// For each SP that ignition reports up, fetch the state and caboose
// information.
for sp in sps {
- // First, fetch the state of the SP. If that fails, report the
- // error but continue.
- let result =
- client.sp_get(&sp.typ, sp.slot).await.with_context(|| {
- format!(
- "MGS {:?}: fetching state of SP {:?}",
- client.baseurl(),
- sp
- )
- });
- let sp_state = match result {
- Err(error) => {
- in_progress.found_error(InventoryError::from(error));
- continue;
- }
- Ok(response) => response.into_inner(),
- };
+ let start = Utc::now();
+ let name = format!("sp {:?} {}", sp.typ, sp.slot);
+ Self::collect_one_sp(client, log, in_progress, sp).await;
+ spans.push(TraceSpan::since(name, "sp", start));
+ }
+ }
- // Record the state that we found.
- let Some(baseboard_id) = in_progress.found_sp_state(
- client.baseurl(),
- sp.typ,
- sp.slot,
- sp_state,
- ) else {
- // We failed to parse this SP for some reason. The error was
- // reported already. Move on.
- continue;
- };
+ /// Collect inventory reported by one MGS instance for one SP
+ async fn collect_one_sp(
+ client: &gateway_client::Client,
+ log: &Logger,
+ in_progress: &mut CollectionBuilder,
+ sp: SpIdentifier,
+ ) {
+ // First, fetch the state of the SP. If that fails, report the
+ // error but continue.
+ let result = client.sp_get(&sp.typ, sp.slot).await.with_context(|| {
+ format!("MGS {:?}: fetching state of SP {:?}", client.baseurl(), sp)
+ });
+ let sp_state = match result {
+ Err(error) => {
+ in_progress.found_error(InventoryError::from(error));
+ return;
+ }
+ Ok(response) => response.into_inner(),
+ };
- // For sled SPs, collect the currently-active phase 1 slot and the
- // hash of the contents of both slots, if they haven't been
- // collected already. Generally, we'd only get here for the first
- // MGS client. Assuming that one succeeds, the other(s) will skip
- // this loop.
- if matches!(sp.typ, SpType::Sled) {
- if !in_progress
- .found_host_phase_1_active_slot_already(&baseboard_id)
- {
- let result = client
- .sp_component_active_slot_get(
- &sp.typ,
- sp.slot,
- SpComponent::HOST_CPU_BOOT_FLASH.const_as_str(),
+ // Record the state that we found.
+ let Some(baseboard_id) = in_progress.found_sp_state(
+ client.baseurl(),
+ sp.typ,
+ sp.slot,
+ sp_state,
+ ) else {
+ // We failed to parse this SP for some reason. The error was
+ // reported already. Move on.
+ return;
+ };
+
+ // For sled SPs, collect the currently-active phase 1 slot and the
+ // hash of the contents of both slots, if they haven't been
+ // collected already. Generally, we'd only get here for the first
+ // MGS client. Assuming that one succeeds, the other(s) will skip
+ // this loop.
+ if matches!(sp.typ, SpType::Sled) {
+ if !in_progress
+ .found_host_phase_1_active_slot_already(&baseboard_id)
+ {
+ let result = client
+ .sp_component_active_slot_get(
+ &sp.typ,
+ sp.slot,
+ SpComponent::HOST_CPU_BOOT_FLASH.const_as_str(),
+ )
+ .await
+ .with_context(|| {
+ format!(
+ "MGS {:?}: SP {sp:?}: phase 1 active slot",
+ client.baseurl(),
)
- .await
- .with_context(|| {
- format!(
- "MGS {:?}: SP {sp:?}: phase 1 active slot",
- client.baseurl(),
- )
- })
- .and_then(|response| {
- M2Slot::from_mgs_firmware_slot(response.slot)
- .ok_or_else(|| {
- anyhow!(
- "MGS {:?}: SP {sp:?}: \
+ })
+ .and_then(|response| {
+ M2Slot::from_mgs_firmware_slot(response.slot)
+ .ok_or_else(|| {
+ anyhow!(
+ "MGS {:?}: SP {sp:?}: \
invalid host phase 1 slot {}",
- client.baseurl(),
- response.slot
- )
- })
- });
- match result {
- Ok(phase_1_slot) => {
- if let Err(error) = in_progress
- .found_host_phase_1_active_slot(
- &baseboard_id,
client.baseurl(),
- phase_1_slot,
+ response.slot
)
- {
- error!(
- log,
- "error reporting host phase 1 active slot: \
+ })
+ });
+ match result {
+ Ok(phase_1_slot) => {
+ if let Err(error) = in_progress
+ .found_host_phase_1_active_slot(
+ &baseboard_id,
+ client.baseurl(),
+ phase_1_slot,
+ )
+ {
+ error!(
+ log,
+ "error reporting host phase 1 active slot: \
{baseboard_id:?} {phase_1_slot:?} \
{:?}: {error:#}",
- client.baseurl(),
- );
- }
- }
- Err(err) => {
- in_progress.found_error(InventoryError::from(err));
- }
- }
- }
-
- for slot in M2Slot::iter() {
- const PHASE1_HASH_TIMEOUT: Duration =
- Duration::from_secs(30);
-
- if in_progress.found_host_phase_1_flash_hash_already(
- &baseboard_id,
- slot,
- ) {
- continue;
- }
-
- let phase1_slot = match slot {
- M2Slot::A => 0,
- M2Slot::B => 1,
- };
-
- let result = client
- .host_phase_1_flash_hash_calculate_with_timeout(
- sp.typ,
- sp.slot,
- phase1_slot,
- PHASE1_HASH_TIMEOUT,
- )
- .await
- .with_context(|| {
- format!(
- "MGS {:?}: SP {sp:?}: phase 1 slot {slot:?}",
client.baseurl(),
- )
- });
- let hash = match result {
- Err(error) => {
- in_progress
- .found_error(InventoryError::from(error));
- continue;
+ );
}
- Ok(hash) => hash,
- };
- if let Err(error) = in_progress
- .found_host_phase_1_flash_hash(
- &baseboard_id,
- slot,
- client.baseurl(),
- ArtifactHash(hash),
- )
- {
- error!(
- log,
- "error reporting host phase 1 flash hash: \
- {baseboard_id:?} {slot:?} {:?}: {error:#}",
- client.baseurl(),
- );
+ }
+ Err(err) => {
+ in_progress.found_error(InventoryError::from(err));
}
}
}
- // For each kind of caboose that we care about, if it hasn't been
- // fetched already, fetch it and record it. Generally, we'd only
- // get here for the first MGS client. Assuming that one succeeds,
- // the other(s) will skip this loop.
- for which in CabooseWhich::iter() {
- if in_progress.found_caboose_already(&baseboard_id, which) {
+ for slot in M2Slot::iter() {
+ const PHASE1_HASH_TIMEOUT: Duration = Duration::from_secs(30);
+
+ if in_progress
+ .found_host_phase_1_flash_hash_already(&baseboard_id, slot)
+ {
continue;
}
- let (component, slot) = match which {
- CabooseWhich::SpSlot0 => ("sp", 0),
- CabooseWhich::SpSlot1 => ("sp", 1),
- CabooseWhich::RotSlotA => ("rot", 0),
- CabooseWhich::RotSlotB => ("rot", 1),
- CabooseWhich::Stage0 => ("stage0", 0),
- CabooseWhich::Stage0Next => ("stage0", 1),
+ let phase1_slot = match slot {
+ M2Slot::A => 0,
+ M2Slot::B => 1,
};
let result = client
- .sp_component_caboose_get(&sp.typ, sp.slot, component, slot)
+ .host_phase_1_flash_hash_calculate_with_timeout(
+ sp.typ,
+ sp.slot,
+ phase1_slot,
+ PHASE1_HASH_TIMEOUT,
+ )
.await
.with_context(|| {
format!(
- "MGS {:?}: SP {:?}: caboose {:?}",
+ "MGS {:?}: SP {sp:?}: phase 1 slot {slot:?}",
client.baseurl(),
- sp,
- which
)
});
- let caboose = match result {
+ let hash = match result {
Err(error) => {
in_progress.found_error(InventoryError::from(error));
continue;
}
- Ok(response) => response.into_inner(),
+ Ok(hash) => hash,
};
- if let Err(error) = in_progress.found_caboose(
+ if let Err(error) = in_progress.found_host_phase_1_flash_hash(
&baseboard_id,
- which,
+ slot,
client.baseurl(),
- caboose,
+ ArtifactHash(hash),
) {
error!(
log,
- "error reporting caboose: {:?} {:?} {:?}: {:#}",
- baseboard_id,
- which,
+ "error reporting host phase 1 flash hash: \
+ {baseboard_id:?} {slot:?} {:?}: {error:#}",
client.baseurl(),
- error
);
}
}
+ }
+
+ // For each kind of caboose that we care about, if it hasn't been
+ // fetched already, fetch it and record it. Generally, we'd only
+ // get here for the first MGS client. Assuming that one succeeds,
+ // the other(s) will skip this loop.
+ for which in CabooseWhich::iter() {
+ if in_progress.found_caboose_already(&baseboard_id, which) {
+ continue;
+ }
- // For each kind of RoT page that we care about, if it hasn't been
- // fetched already, fetch it and record it. Generally, we'd only
- // get here for the first MGS client. Assuming that one succeeds,
- // the other(s) will skip this loop.
- for which in RotPageWhich::iter() {
- if in_progress.found_rot_page_already(&baseboard_id, which) {
- continue;
- }
+ let (component, slot) = match which {
+ CabooseWhich::SpSlot0 => ("sp", 0),
+ CabooseWhich::SpSlot1 => ("sp", 1),
+ CabooseWhich::RotSlotA => ("rot", 0),
+ CabooseWhich::RotSlotB => ("rot", 1),
+ CabooseWhich::Stage0 => ("stage0", 0),
+ CabooseWhich::Stage0Next => ("stage0", 1),
+ };
- let component = SpComponent::ROT.const_as_str();
-
- let result = match which {
- RotPageWhich::Cmpa => client
- .sp_rot_cmpa_get(&sp.typ, sp.slot, component)
- .await
- .map(|response| response.into_inner().base64_data),
- RotPageWhich::CfpaActive => client
- .sp_rot_cfpa_get(
- &sp.typ,
- sp.slot,
- component,
- &GetCfpaParams { slot: RotCfpaSlot::Active },
- )
- .await
- .map(|response| response.into_inner().base64_data),
- RotPageWhich::CfpaInactive => client
- .sp_rot_cfpa_get(
- &sp.typ,
- sp.slot,
- component,
- &GetCfpaParams { slot: RotCfpaSlot::Inactive },
- )
- .await
- .map(|response| response.into_inner().base64_data),
- RotPageWhich::CfpaScratch => client
- .sp_rot_cfpa_get(
- &sp.typ,
- sp.slot,
- component,
- &GetCfpaParams { slot: RotCfpaSlot::Scratch },
- )
- .await
- .map(|response| response.into_inner().base64_data),
- }
+ let result = client
+ .sp_component_caboose_get(&sp.typ, sp.slot, component, slot)
+ .await
.with_context(|| {
format!(
- "MGS {:?}: SP {:?}: rot page {:?}",
+ "MGS {:?}: SP {:?}: caboose {:?}",
client.baseurl(),
sp,
which
)
});
-
- let page = match result {
- Err(error) => {
- in_progress.found_error(InventoryError::from(error));
- continue;
- }
- Ok(data_base64) => RotPage { data_base64 },
- };
- if let Err(error) = in_progress.found_rot_page(
- &baseboard_id,
+ let caboose = match result {
+ Err(error) => {
+ in_progress.found_error(InventoryError::from(error));
+ continue;
+ }
+ Ok(response) => response.into_inner(),
+ };
+ if let Err(error) = in_progress.found_caboose(
+ &baseboard_id,
+ which,
+ client.baseurl(),
+ caboose,
+ ) {
+ error!(
+ log,
+ "error reporting caboose: {:?} {:?} {:?}: {:#}",
+ baseboard_id,
which,
client.baseurl(),
- page,
- ) {
- error!(
- log,
- "error reporting rot page: {:?} {:?} {:?}: {:#}",
- baseboard_id,
- which,
- client.baseurl(),
- error
- );
+ error
+ );
+ }
+ }
+
+ // For each kind of RoT page that we care about, if it hasn't been
+ // fetched already, fetch it and record it. Generally, we'd only
+ // get here for the first MGS client. Assuming that one succeeds,
+ // the other(s) will skip this loop.
+ for which in RotPageWhich::iter() {
+ if in_progress.found_rot_page_already(&baseboard_id, which) {
+ continue;
+ }
+
+ let component = SpComponent::ROT.const_as_str();
+
+ let result = match which {
+ RotPageWhich::Cmpa => client
+ .sp_rot_cmpa_get(&sp.typ, sp.slot, component)
+ .await
+ .map(|response| response.into_inner().base64_data),
+ RotPageWhich::CfpaActive => client
+ .sp_rot_cfpa_get(
+ &sp.typ,
+ sp.slot,
+ component,
+ &GetCfpaParams { slot: RotCfpaSlot::Active },
+ )
+ .await
+ .map(|response| response.into_inner().base64_data),
+ RotPageWhich::CfpaInactive => client
+ .sp_rot_cfpa_get(
+ &sp.typ,
+ sp.slot,
+ component,
+ &GetCfpaParams { slot: RotCfpaSlot::Inactive },
+ )
+ .await
+ .map(|response| response.into_inner().base64_data),
+ RotPageWhich::CfpaScratch => client
+ .sp_rot_cfpa_get(
+ &sp.typ,
+ sp.slot,
+ component,
+ &GetCfpaParams { slot: RotCfpaSlot::Scratch },
+ )
+ .await
+ .map(|response| response.into_inner().base64_data),
+ }
+ .with_context(|| {
+ format!(
+ "MGS {:?}: SP {:?}: rot page {:?}",
+ client.baseurl(),
+ sp,
+ which
+ )
+ });
+
+ let page = match result {
+ Err(error) => {
+ in_progress.found_error(InventoryError::from(error));
+ continue;
}
+ Ok(data_base64) => RotPage { data_base64 },
+ };
+ if let Err(error) = in_progress.found_rot_page(
+ &baseboard_id,
+ which,
+ client.baseurl(),
+ page,
+ ) {
+ error!(
+ log,
+ "error reporting rot page: {:?} {:?} {:?}: {:#}",
+ baseboard_id,
+ which,
+ client.baseurl(),
+ error
+ );
}
}
}
@@ -455,8 +498,13 @@ impl<'a> Collector<'a> {
for (idx, url) in urls.into_iter().enumerate() {
let log = self.log.new(o!("SledAgent" => url.clone()));
let task = async move {
- let result = collect_one_sled_agent(&url, log).await;
- (idx, url, result)
+ let (span, result) = timed(
+ url.clone(),
+ "sled_agent",
+ collect_one_sled_agent(&url, log),
+ )
+ .await;
+ (idx, url, result, span)
};
if let Some(result) = tasks.spawn(task).await {
results.push(result);
@@ -467,8 +515,9 @@ impl<'a> Collector<'a> {
// Merge results in the order we enumerated the sleds, so that the
// collection's contents (in particular the order of its errors) do
// not depend on request completion order.
- results.sort_by_key(|(idx, _, _)| *idx);
- for (_, url, result) in results {
+ results.sort_by_key(|(idx, _, _, _)| *idx);
+ for (_, url, result, span) in results {
+ self.spans.push(span);
match result {
Err(error) => self.in_progress.found_error(error),
Ok(inventory) => {
@@ -521,8 +570,13 @@ impl<'a> Collector<'a> {
{
let log = self.log.clone();
let task = async move {
- let result = collect_one_timesync(&log, zone_id, &client).await;
- (idx, zone_id, result)
+ let (span, result) = timed(
+ zone_id.to_string(),
+ "timesync",
+ collect_one_timesync(&log, zone_id, &client),
+ )
+ .await;
+ (idx, zone_id, result, span)
};
if let Some(result) = tasks.spawn(task).await {
results.push(result);
@@ -530,8 +584,9 @@ impl<'a> Collector<'a> {
}
results.extend(tasks.join_remaining().await);
- results.sort_by_key(|(idx, _, _)| *idx);
- for (_, zone_id, result) in results {
+ results.sort_by_key(|(idx, _, _, _)| *idx);
+ for (_, zone_id, result, span) in results {
+ self.spans.push(span);
match result {
Err(error) => self.in_progress.found_error(error),
Ok(timesync) => {
@@ -566,8 +621,13 @@ impl<'a> Collector<'a> {
let client = client.clone();
let log = self.log.clone();
let task = async move {
- let result = collect_one_keeper(&client, &log).await;
- (idx, result)
+ let (span, result) = timed(
+ client.baseurl().to_string(),
+ "keeper",
+ collect_one_keeper(&client, &log),
+ )
+ .await;
+ (idx, result, span)
};
if let Some(result) = tasks.spawn(task).await {
results.push(result);
@@ -575,8 +635,9 @@ impl<'a> Collector<'a> {
}
results.extend(tasks.join_remaining().await);
- results.sort_by_key(|(idx, _)| *idx);
- for (_, result) in results {
+ results.sort_by_key(|(idx, _, _)| *idx);
+ for (_, result, span) in results {
+ self.spans.push(span);
match result {
Err(error) => self.in_progress.found_error(error),
Ok(membership) => self
@@ -650,9 +711,13 @@ impl<'a> Collector<'a> {
{
let log = self.log.clone();
let task = async move {
- let result =
- collect_one_dns_generation(&log, zone_id, &client).await;
- (idx, zone_id, result)
+ let (span, result) = timed(
+ zone_id.to_string(),
+ "dns_generation",
+ collect_one_dns_generation(&log, zone_id, &client),
+ )
+ .await;
+ (idx, zone_id, result, span)
};
if let Some(result) = tasks.spawn(task).await {
results.push(result);
@@ -660,8 +725,9 @@ impl<'a> Collector<'a> {
}
results.extend(tasks.join_remaining().await);
- results.sort_by_key(|(idx, _, _)| *idx);
- for (_, zone_id, result) in results {
+ results.sort_by_key(|(idx, _, _, _)| *idx);
+ for (_, zone_id, result, span) in results {
+ self.spans.push(span);
let result = result.and_then(|generation_status| {
self.in_progress
.found_internal_dns_generation_status(generation_status)
@@ -1191,7 +1257,7 @@ mod test {
&sled_enum,
log.clone(),
);
- let collection = collector
+ let (collection, spans) = collector
.collect_all()
.await
.expect("failed to carry out collection");
@@ -1202,6 +1268,36 @@ mod test {
);
assert_eq!(collection.collector, "test-suite");
+ // The timing spans should cover every phase, in the order the phases
+ // run, without overlap.
+ let phases: Vec<_> =
+ spans.iter().filter(|s| s.category == "phase").collect();
+ assert_eq!(
+ phases.iter().map(|s| s.name.as_str()).collect::>(),
+ vec![
+ "mgs",
+ "sled_agents",
+ "keepers",
+ "cockroach",
+ "timesync",
+ "dns_generations"
+ ]
+ );
+ for pair in phases.windows(2) {
+ assert!(pair[1].start >= pair[0].end);
+ }
+ for span in &spans {
+ assert!(span.end >= span.start);
+ }
+
+ // One span per collected target: two sled agents, one MGS client,
+ // and at least one SP behind it.
+ let count =
+ |cat: &str| spans.iter().filter(|s| s.category == cat).count();
+ assert_eq!(count("sled_agent"), 2);
+ assert_eq!(count("mgs_client"), 1);
+ assert!(count("sp") >= 1);
+
let s = dump_collection(&collection);
expectorate::assert_contents("tests/output/collector_basic.txt", &s);
@@ -1271,7 +1367,7 @@ mod test {
&sled_enum,
log.clone(),
);
- let collection = collector
+ let (collection, _spans) = collector
.collect_all()
.await
.expect("failed to carry out collection");
@@ -1321,7 +1417,7 @@ mod test {
&sled_enum,
log.clone(),
);
- let collection = collector
+ let (collection, _spans) = collector
.collect_all()
.await
.expect("failed to carry out collection");
@@ -1378,7 +1474,7 @@ mod test {
&sled_enum,
log.clone(),
);
- let collection = collector
+ let (collection, _spans) = collector
.collect_all()
.await
.expect("failed to carry out collection");
diff --git a/nexus/src/app/background/tasks/inventory_collection.rs b/nexus/src/app/background/tasks/inventory_collection.rs
index c6f8682996e..1784686b0c9 100644
--- a/nexus/src/app/background/tasks/inventory_collection.rs
+++ b/nexus/src/app/background/tasks/inventory_collection.rs
@@ -15,9 +15,11 @@ use nexus_db_queries::db::DataStore;
use nexus_inventory::InventoryError;
use nexus_networking::GatewayClient;
use nexus_types::deployment::SledFilter;
+use nexus_types::internal_api::background::InventoryCollectionStatus;
use nexus_types::inventory::Collection;
use omicron_cockroach_metrics::CockroachClusterAdminClient;
use omicron_uuid_kinds::CollectionUuid;
+use perfetto_trace::TraceSpan;
use serde_json::json;
use slog::{debug, o, warn};
use std::net::SocketAddr;
@@ -93,18 +95,21 @@ impl BackgroundTask for InventoryCollector {
"error" => message.clone());
json!({ "error": message })
}
- Ok(collection) => {
+ Ok((collection, spans)) => {
debug!(opctx.log, "inventory collection complete";
"collection_id" => collection.id.to_string(),
"time_started" => collection.time_started.to_string(),
);
- let json = json!({
- "collection_id": collection.id.to_string(),
- "time_started": collection.time_started.to_string(),
- "time_done": collection.time_done.to_string()
- });
+ let status = InventoryCollectionStatus {
+ collection_id: collection.id,
+ time_started: collection.time_started,
+ time_done: collection.time_done,
+ trace: Some(perfetto_trace::assemble(spans)),
+ };
self.tx.send_replace(Some(collection.id));
- json
+ serde_json::to_value(status).unwrap_or_else(
+ |error| json!({ "error": error.to_string() }),
+ )
}
}
}
@@ -120,7 +125,7 @@ async fn inventory_activate(
nkeep: u32,
disabled: bool,
cockroach_admin_client: &CockroachClusterAdminClient,
-) -> Result {
+) -> Result<(Collection, Vec), anyhow::Error> {
// If we're disabled, don't do anything. (This switch is only intended for
// unforeseen production emergencies.)
ensure!(!disabled, "disabled by explicit configuration");
@@ -211,7 +216,7 @@ async fn inventory_activate(
&sled_enum,
opctx.log.clone(),
);
- let collection =
+ let (collection, spans) =
inventory.collect_all().await.context("collecting inventory")?;
// Write it to the database.
@@ -220,7 +225,7 @@ async fn inventory_activate(
.await
.context("saving inventory to database")?;
- Ok(collection)
+ Ok((collection, spans))
}
/// Determine which sleds to inventory based on what's in the database
@@ -269,6 +274,7 @@ mod test {
use nexus_inventory::SledAgentEnumerator;
use nexus_test_utils_macros::nexus_test;
use nexus_types::identity::Asset;
+ use nexus_types::internal_api::background::InventoryCollectionStatus;
use omicron_common::api::external::ByteCount;
use omicron_common::api::external::LookupType;
use omicron_uuid_kinds::CollectionUuid;
@@ -313,7 +319,16 @@ mod test {
let nkeep = usize::try_from(nkeep).unwrap();
let mut all_our_collection_ids = Vec::new();
for i in 0..20 {
- let _ = task.activate(&opctx).await;
+ let value = task.activate(&opctx).await;
+
+ // The status should include a timing trace with at least one
+ // event per collection phase.
+ let status: InventoryCollectionStatus =
+ serde_json::from_value(value)
+ .expect("failed to parse activation status");
+ let trace = status.trace.expect("status should include a trace");
+ assert!(trace.trace_events.len() >= 6);
+
let collections = datastore.inventory_collections().await.unwrap();
// Nexus is creating inventory collections concurrently with us,
diff --git a/nexus/types/Cargo.toml b/nexus/types/Cargo.toml
index 331e64835c8..d5aaa403688 100644
--- a/nexus/types/Cargo.toml
+++ b/nexus/types/Cargo.toml
@@ -40,6 +40,7 @@ oximeter-db.workspace = true
oxnet.workspace = true
oxql-types.workspace = true
parse-display.workspace = true
+perfetto-trace.workspace = true
regex.workspace = true
schemars = { workspace = true, features = ["chrono", "uuid1", "url"] }
serde.workspace = true
diff --git a/nexus/types/src/internal_api/background.rs b/nexus/types/src/internal_api/background.rs
index 1e75a925929..73d95022161 100644
--- a/nexus/types/src/internal_api/background.rs
+++ b/nexus/types/src/internal_api/background.rs
@@ -725,6 +725,23 @@ impl slog::KV for DatasetsRendezvousStats {
}
}
+/// The status of a successful `inventory_collection` background task
+/// activation.
+#[derive(Clone, Debug, Deserialize, Serialize, PartialEq)]
+pub struct InventoryCollectionStatus {
+ pub collection_id: CollectionUuid,
+ pub time_started: DateTime,
+ pub time_done: DateTime,
+
+ /// Timing of the collection's phases and requests, in the Chrome Trace
+ /// Event format (load into ).
+ ///
+ /// Optional so that this type can also represent status reported by
+ /// versions of Nexus that did not record a trace.
+ #[serde(default, skip_serializing_if = "Option::is_none")]
+ pub trace: Option,
+}
+
/// The status of an `inventory_load` background task activation.
#[derive(Clone, Debug, Deserialize, Serialize, PartialEq, Eq)]
pub enum InventoryLoadStatus {
diff --git a/perfetto-trace/Cargo.toml b/perfetto-trace/Cargo.toml
new file mode 100644
index 00000000000..f63a5b3102a
--- /dev/null
+++ b/perfetto-trace/Cargo.toml
@@ -0,0 +1,14 @@
+[package]
+name = "perfetto-trace"
+version = "0.1.0"
+edition.workspace = true
+license = "MPL-2.0"
+
+[dependencies]
+chrono.workspace = true
+serde.workspace = true
+serde_json.workspace = true
+omicron-workspace-hack.workspace = true
+
+[lints]
+workspace = true
diff --git a/perfetto-trace/src/lib.rs b/perfetto-trace/src/lib.rs
new file mode 100644
index 00000000000..917af145992
--- /dev/null
+++ b/perfetto-trace/src/lib.rs
@@ -0,0 +1,254 @@
+// This Source Code Form is subject to the terms of the Mozilla Public
+// License, v. 2.0. If a copy of the MPL was not distributed with this
+// file, You can obtain one at https://mozilla.org/MPL/2.0/.
+
+//! Chrome Trace Event format support for visualizing operation timing
+//!
+//! Traces produced with this crate can be loaded into the Perfetto trace
+//! viewer () or `chrome://tracing`.
+
+use chrono::DateTime;
+use chrono::Utc;
+use serde::Deserialize;
+use serde::Serialize;
+
+/// Represents a Perfetto Trace Event format JSON file for visualization.
+///
+/// This format is used by the Perfetto trace viewer ()
+/// to visualize timing information for operations.
+#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
+pub struct Trace {
+ #[serde(rename = "traceEvents")]
+ pub trace_events: Vec,
+ /// Display unit for time values in the UI (e.g., "ms" for milliseconds)
+ #[serde(rename = "displayTimeUnit")]
+ pub display_time_unit: String,
+}
+
+/// A single event in the Perfetto Trace Event format.
+///
+/// This represents a complete event (duration event) showing when an operation
+/// started and how long it took.
+#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
+pub struct TraceEvent {
+ /// Human-readable name of the event
+ pub name: String,
+ /// Category name (abbreviated as "cat" in Perfetto format).
+ /// Used to group related events together in the trace viewer.
+ pub cat: String,
+ /// Phase type (abbreviated as "ph" in Perfetto format).
+ /// "X" means a "Complete" event with both timestamp and duration.
+ pub ph: String,
+ /// Timestamp in microseconds (abbreviated as "ts" in Perfetto format).
+ /// Represents when the event started, as microseconds since the epoch.
+ pub ts: i64,
+ /// Duration in microseconds (abbreviated as "dur" in Perfetto format).
+ /// How long the event took to complete.
+ pub dur: i64,
+ /// Process ID. Used to separate events into different process lanes
+ /// in the trace viewer.
+ pub pid: u32,
+ /// Thread ID. Used to separate events into different thread lanes
+ /// within a process in the trace viewer.
+ pub tid: usize,
+ /// Arbitrary key-value pairs with additional event metadata
+ pub args: serde_json::Value,
+}
+
+/// A completed timed operation, recorded with wall-clock timestamps.
+///
+/// Spans are collected while work runs and later converted into a [`Trace`]
+/// with [`assemble`].
+#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
+pub struct TraceSpan {
+ /// Name identifying the specific operation (e.g., a target URL)
+ pub name: String,
+ /// The kind of operation, used for grouping and filtering
+ pub category: String,
+ pub start: DateTime,
+ pub end: DateTime,
+ /// Arbitrary key-value pairs with additional metadata
+ pub args: serde_json::Value,
+}
+
+impl TraceSpan {
+ /// Creates a span lasting from `start` until now.
+ pub fn since(
+ name: impl Into,
+ category: impl Into,
+ start: DateTime,
+ ) -> Self {
+ TraceSpan {
+ name: name.into(),
+ category: category.into(),
+ start,
+ end: Utc::now(),
+ args: serde_json::Value::Null,
+ }
+ }
+}
+
+/// Runs `fut` to completion, returning its output along with a [`TraceSpan`]
+/// covering its execution.
+pub async fn timed(
+ name: impl Into,
+ category: impl Into,
+ fut: impl Future