diff --git a/Cargo.lock b/Cargo.lock index 2ef0ba5fa9..aaed27f196 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2651,6 +2651,7 @@ dependencies = [ "config-version", "duration-str", "eyre", + "futures", "hex", "http", "libnmxc", diff --git a/crates/api-core/src/cfg/README.md b/crates/api-core/src/cfg/README.md index 6ed6e3b505..66ff9f2cde 100644 --- a/crates/api-core/src/cfg/README.md +++ b/crates/api-core/src/cfg/README.md @@ -350,6 +350,7 @@ extracted identifier contains characters such as `/` or `.`. | `allow_insecure` | `bool` | `false` | Skip TLS verification for NMX-C. | | `nmx_c_endpoint_port` | `Option` | — | TCP port for NMX-C endpoints derived from switch NVOS IP. Unset uses the production NMX-C port. | | `nmx_c_certificate_rotation` | `NmxCCertificateRotationConfig` | *(default)* | Optional expiry-driven rotation for NMX-C server certificates. | +| `partition_monitor_max_concurrent_groups` | `NonZeroUsize` | `16` | Maximum number of NMX-C machine groups (chassis or rack) processed concurrently per monitor iteration. Bounds DB pool usage and gRPC fan-out. Must be ≥ 1. | ### `NmxCCertificateRotationConfig` diff --git a/crates/api-db/src/nvlink_nmxc_endpoints.rs b/crates/api-db/src/nvlink_nmxc_endpoints.rs index 97bd8ecd85..1e1a5fcc20 100644 --- a/crates/api-db/src/nvlink_nmxc_endpoints.rs +++ b/crates/api-db/src/nvlink_nmxc_endpoints.rs @@ -34,6 +34,19 @@ pub async fn find_by_chassis_serial( .map_err(|e| DatabaseError::new(Q, e)) } +pub async fn find_by_chassis_serials( + txn: impl DbReader<'_>, + chassis_serials: &[&str], +) -> DatabaseResult> { + const Q: &str = + "SELECT chassis_serial, endpoint FROM nvlink_nmxc_endpoints WHERE chassis_serial = ANY($1)"; + sqlx::query_as(Q) + .bind(chassis_serials) + .fetch_all(txn) + .await + .map_err(|e| DatabaseError::new(Q, e)) +} + pub async fn find_all(txn: impl DbReader<'_>) -> DatabaseResult> { const Q: &str = "SELECT chassis_serial, endpoint FROM nvlink_nmxc_endpoints ORDER BY chassis_serial"; @@ -93,3 +106,63 @@ pub async fn update( .await .map_err(|e| DatabaseError::new(Q, e)) } + +#[cfg(test)] +mod tests { + use super::*; + + #[crate::sqlx_test] + async fn find_by_chassis_serials_returns_matching_rows(pool: sqlx::PgPool) { + let mut txn = pool.begin().await.unwrap(); + create(txn.as_mut(), "SN-A", "https://a.example:9370") + .await + .unwrap(); + create(txn.as_mut(), "SN-B", "https://b.example:9370") + .await + .unwrap(); + create(txn.as_mut(), "SN-C", "https://c.example:9370") + .await + .unwrap(); + + let mut rows = find_by_chassis_serials(txn.as_mut(), &["SN-A", "SN-C"]) + .await + .unwrap(); + rows.sort_by(|a, b| a.chassis_serial.cmp(&b.chassis_serial)); + assert_eq!(rows.len(), 2); + assert_eq!(rows[0].chassis_serial, "SN-A"); + assert_eq!(rows[0].endpoint, "https://a.example:9370"); + assert_eq!(rows[1].chassis_serial, "SN-C"); + assert_eq!(rows[1].endpoint, "https://c.example:9370"); + + txn.rollback().await.unwrap(); + } + + #[crate::sqlx_test] + async fn find_by_chassis_serials_unknown_serials_are_excluded(pool: sqlx::PgPool) { + let mut txn = pool.begin().await.unwrap(); + create(txn.as_mut(), "SN-KNOWN", "https://known.example:9370") + .await + .unwrap(); + + let rows = find_by_chassis_serials(txn.as_mut(), &["SN-KNOWN", "SN-MISSING"]) + .await + .unwrap(); + assert_eq!(rows.len(), 1); + assert_eq!(rows[0].chassis_serial, "SN-KNOWN"); + + txn.rollback().await.unwrap(); + } + + #[crate::sqlx_test] + async fn find_by_chassis_serials_empty_slice_returns_empty(pool: sqlx::PgPool) { + let mut txn = pool.begin().await.unwrap(); + create(txn.as_mut(), "SN-X", "https://x.example:9370") + .await + .unwrap(); + + let rows = find_by_chassis_serials(txn.as_mut(), &[]).await.unwrap(); + assert!(rows.is_empty()); + + txn.rollback().await.unwrap(); + } +} diff --git a/crates/nvlink-manager/Cargo.toml b/crates/nvlink-manager/Cargo.toml index 200a23a908..f428ef12f3 100644 --- a/crates/nvlink-manager/Cargo.toml +++ b/crates/nvlink-manager/Cargo.toml @@ -32,6 +32,7 @@ async-trait = { workspace = true } chrono = { workspace = true } duration-str = { workspace = true } eyre = { workspace = true } +futures = { workspace = true, features = ["std"] } hex = { workspace = true } http = { workspace = true } librms = { workspace = true } diff --git a/crates/nvlink-manager/src/config.rs b/crates/nvlink-manager/src/config.rs index c6fca9b06b..2019f2115b 100644 --- a/crates/nvlink-manager/src/config.rs +++ b/crates/nvlink-manager/src/config.rs @@ -54,12 +54,23 @@ pub struct NvLinkConfig { /// Optional expiry-driven rotation for NMX-C server certificates. #[serde(default)] pub nmx_c_certificate_rotation: NmxCCertificateRotationConfig, + + /// Maximum number of NMX-C machine groups (chassis or rack) processed concurrently + /// during a partition monitor iteration. Bounds DB pool usage and gRPC fan-out. + /// Defaults to 16. Must be non-zero; deserialization rejects 0. + #[serde(default = "NvLinkConfig::default_partition_monitor_max_concurrent_groups")] + pub partition_monitor_max_concurrent_groups: std::num::NonZeroUsize, } impl NvLinkConfig { pub const fn default_monitor_run_interval() -> std::time::Duration { std::time::Duration::from_secs(60) } + + pub const fn default_partition_monitor_max_concurrent_groups() -> std::num::NonZeroUsize { + // SAFETY: 16 is non-zero. + unsafe { std::num::NonZeroUsize::new_unchecked(16) } + } } #[derive(Clone, Debug, Deserialize, Serialize, PartialEq)] @@ -131,6 +142,8 @@ impl Default for NvLinkConfig { nmx_c_endpoint_port: None, allow_insecure: false, nmx_c_certificate_rotation: NmxCCertificateRotationConfig::default(), + partition_monitor_max_concurrent_groups: + Self::default_partition_monitor_max_concurrent_groups(), } } } @@ -157,6 +170,8 @@ mod test { nmx_c_endpoint_port: None, allow_insecure: true, nmx_c_certificate_rotation: NmxCCertificateRotationConfig::default(), + partition_monitor_max_concurrent_groups: + NvLinkConfig::default_partition_monitor_max_concurrent_groups(), } ); } @@ -172,6 +187,14 @@ mod test { ); } + #[test] + fn deserialize_zero_concurrent_groups_is_rejected() { + let err = serde_json::from_str::( + r#"{"allow_insecure":false,"partition_monitor_max_concurrent_groups":0}"#, + ); + assert!(err.is_err(), "zero must be rejected by NonZeroUsize"); + } + #[test] fn deserialize_legacy_expiry_warning_window_as_rotation_window() { let config: NmxCCertificateRotationConfig = diff --git a/crates/nvlink-manager/src/lib.rs b/crates/nvlink-manager/src/lib.rs index 2526447dc3..c0429e5569 100644 --- a/crates/nvlink-manager/src/lib.rs +++ b/crates/nvlink-manager/src/lib.rs @@ -42,6 +42,7 @@ use db::nvl_partition::IdColumn; use db::work_lock_manager::WorkLockManagerHandle; use db::{self, ObjectColumnFilter, TransactionVending, machine}; use errors::{NvLinkManagerError, NvLinkManagerResult}; +use futures::future; use libnmxc::nmxc_model::{ GetComputeNodeInfoListRequest, GetGpuInfoListRequest, GetPartitionInfoListRequest, PartitionInfo, @@ -63,6 +64,7 @@ use model::nvl_partition::{NvlPartition, NvlPartitionName}; use sqlx::PgPool; #[cfg(feature = "test-support")] pub use switch_cert_monitor::{SwitchCertificateMonitor, SwitchCertificateMonitorIterationResult}; +use tokio::sync::Semaphore; use tokio::task::JoinSet; use tokio_util::sync::CancellationToken; use tracing::Instrument; @@ -1124,7 +1126,7 @@ struct PendingNullNvlinkObservation { /// Shared inputs for processing one chassis- or rack-scoped NMX-C monitor group. struct ProcessMachineGroupInput<'a> { - group_id: &'a str, + group_id: String, group_type: nmx_c_endpoint::ManagedHostGroupType, snapshots: &'a [&'a ManagedHostStateSnapshot], endpoint_url: Option<&'a str>, @@ -1132,11 +1134,19 @@ struct ProcessMachineGroupInput<'a> { /// Rack associated with the selected switch endpoint; absent for chassis mappings. rack_id: Option<&'a RackId>, all_managed_host_snapshots: &'a HashMap, - machine_nvlink_info: &'a mut HashMap>, + /// Pre-split shard of `machine_nvlink_info` containing only this group's machines. + machine_nvlink_info: HashMap>, db_nvl_partitions: &'a [NvlPartition], db_nvl_logical_partitions: &'a [LogicalPartition], } +/// Output of processing one NMX-C monitor group, collected and merged by the caller. +struct GroupResult { + completed_operations: usize, + null_observations: Vec, + partial_metrics: NvlPartitionMonitorMetrics, +} + impl NvlPartitionMonitor { const ITERATION_WORK_KEY: &'static str = "NvlPartitionMonitor::run_single_iteration"; @@ -1285,21 +1295,17 @@ impl NvlPartitionMonitor { db::nvl_logical_partition::find_by(&mut txn, ObjectColumnFilter::::All) .await?; - let mut chassis_serial_to_resolved_endpoint = HashMap::new(); - for chassis_serial in managed_host_snapshots_by_chassis_serial.keys() { - if let Some(endpoint_url) = nmx_c_endpoint::resolve_nmx_c_endpoint_url( - &mut txn, - nmx_c_endpoint::ManagedHostGroupType::Chassis, - None, - Some(chassis_serial), - &self.config, - ) - .await - .map_err(NvLinkManagerError::from)? - { - chassis_serial_to_resolved_endpoint.insert(chassis_serial.clone(), endpoint_url); - } - } + let chassis_serials: Vec<&str> = managed_host_snapshots_by_chassis_serial + .keys() + .map(String::as_str) + .collect(); + let chassis_serial_to_resolved_endpoint: HashMap = + db::nvlink_nmxc_endpoints::find_by_chassis_serials(&mut txn, &chassis_serials) + .await + .map_err(NvLinkManagerError::from)? + .into_iter() + .map(|row| (row.chassis_serial, row.endpoint)) + .collect(); // Close the inventory transaction before contacting NMX-C. Endpoint // lookup is best effort only when no rack partition work depends on it. @@ -1340,54 +1346,84 @@ impl NvlPartitionMonitor { metrics.num_logical_partitions = db_nvl_logical_partitions.len(); metrics.num_physical_partitions = db_nvl_partitions.len(); - let mut total_completed_operations = 0; - let mut pending_null_nvlink_observations = Vec::new(); + // Pre-split machine_nvlink_info into per-group shards before concurrent execution. + // Groups are disjoint (each MachineId belongs to exactly one group), so the split + // is lossless: each entry is moved into exactly one shard via remove(). + let mut all_group_inputs: Vec> = Vec::new(); - for (chassis_serial, chassis_snapshots) in &managed_host_snapshots_by_chassis_serial { - total_completed_operations += self - .process_nmx_c_partition_monitor_group( - ProcessMachineGroupInput { - group_id: chassis_serial, - group_type: nmx_c_endpoint::ManagedHostGroupType::Chassis, - snapshots: chassis_snapshots, - endpoint_url: chassis_serial_to_resolved_endpoint - .get(chassis_serial) - .map(String::as_str), - rack_id: None, - all_managed_host_snapshots: &managed_host_snapshots, - machine_nvlink_info: &mut machine_nvlink_info, - db_nvl_partitions: &db_nvl_partitions, - db_nvl_logical_partitions: &db_nvl_logical_partitions, - }, - metrics, - &mut pending_null_nvlink_observations, - ) - .await; + for (serial, snapshots) in &managed_host_snapshots_by_chassis_serial { + let shard = snapshots + .iter() + .filter_map(|s| { + machine_nvlink_info + .remove(&s.host_snapshot.id) + .map(|info| (s.host_snapshot.id, info)) + }) + .collect(); + all_group_inputs.push(ProcessMachineGroupInput { + group_id: serial.clone(), + group_type: nmx_c_endpoint::ManagedHostGroupType::Chassis, + snapshots, + endpoint_url: chassis_serial_to_resolved_endpoint + .get(serial) + .map(String::as_str), + rack_id: None, + all_managed_host_snapshots: &managed_host_snapshots, + machine_nvlink_info: shard, + db_nvl_partitions: &db_nvl_partitions, + db_nvl_logical_partitions: &db_nvl_logical_partitions, + }); } // A rack is one NVLink domain, so all hosts in the rack are reconciled // through the same NMX-C endpoint and hello response. - for (rack_id, rack_snapshots) in &managed_host_snapshots_by_rack_id { - let rack_id_str = rack_id.to_string(); - total_completed_operations += self - .process_nmx_c_partition_monitor_group( - ProcessMachineGroupInput { - group_id: &rack_id_str, - group_type: nmx_c_endpoint::ManagedHostGroupType::Rack, - snapshots: rack_snapshots, - endpoint_url: rack_id_to_resolved_endpoint - .get(rack_id) - .map(String::as_str), - rack_id: Some(rack_id), - all_managed_host_snapshots: &managed_host_snapshots, - machine_nvlink_info: &mut machine_nvlink_info, - db_nvl_partitions: &db_nvl_partitions, - db_nvl_logical_partitions: &db_nvl_logical_partitions, - }, - metrics, - &mut pending_null_nvlink_observations, - ) - .await; + for (rack_id, snapshots) in &managed_host_snapshots_by_rack_id { + let shard = snapshots + .iter() + .filter_map(|s| { + machine_nvlink_info + .remove(&s.host_snapshot.id) + .map(|info| (s.host_snapshot.id, info)) + }) + .collect(); + all_group_inputs.push(ProcessMachineGroupInput { + group_id: rack_id.to_string(), + group_type: nmx_c_endpoint::ManagedHostGroupType::Rack, + snapshots, + endpoint_url: rack_id_to_resolved_endpoint + .get(rack_id) + .map(String::as_str), + rack_id: Some(rack_id), + all_managed_host_snapshots: &managed_host_snapshots, + machine_nvlink_info: shard, + db_nvl_partitions: &db_nvl_partitions, + db_nvl_logical_partitions: &db_nvl_logical_partitions, + }); + } + + // Bound concurrency so group processing cannot exhaust the shared DB pool + // or open an unbounded number of NMX-C gRPC clients at once. + let concurrency = Semaphore::new(self.config.partition_monitor_max_concurrent_groups.get()); + let all_group_results = future::join_all(all_group_inputs.into_iter().map(|input| { + // Borrow outside `async move` so the closure copies the &Semaphore reference + // (which is Copy) rather than trying to move the Semaphore itself. + let concurrency = &concurrency; + async move { + let _permit = concurrency + .acquire() + .await + .expect("NMX-C group concurrency semaphore is never closed"); + self.process_nmx_c_partition_monitor_group(input).await + } + })) + .await; + + let mut total_completed_operations = 0; + let mut pending_null_nvlink_observations = Vec::new(); + for result in all_group_results { + total_completed_operations += result.completed_operations; + pending_null_nvlink_observations.extend(result.null_observations); + metrics.merge_from(result.partial_metrics); } // Rack groups already observe Hello while reconciling partitions. Racks @@ -1416,9 +1452,7 @@ impl NvlPartitionMonitor { async fn process_nmx_c_partition_monitor_group( &self, input: ProcessMachineGroupInput<'_>, - metrics: &mut NvlPartitionMonitorMetrics, - pending_null_nvlink_observations: &mut Vec, - ) -> usize { + ) -> GroupResult { let ProcessMachineGroupInput { group_id, group_type, @@ -1426,11 +1460,30 @@ impl NvlPartitionMonitor { endpoint_url, rack_id, all_managed_host_snapshots, - machine_nvlink_info, + mut machine_nvlink_info, db_nvl_partitions, db_nvl_logical_partitions, } = input; let group_type_label = group_type.as_str(); + let mut group_metrics = NvlPartitionMonitorMetrics::new(); + let mut null_observations: Vec = Vec::new(); + + macro_rules! early_return { + ($reason:expr) => {{ + Self::queue_null_nvlink_status_observation( + &mut null_observations, + &group_id, + group_type, + snapshots, + $reason, + ); + return GroupResult { + completed_operations: 0, + null_observations, + partial_metrics: group_metrics, + }; + }}; + } let Some(endpoint_url) = endpoint_url else { tracing::warn!( @@ -1438,14 +1491,7 @@ impl NvlPartitionMonitor { group_type = group_type_label, "No NMX-C endpoint (switch NVOS IP or nvlink_nmxc_endpoints mapping); skipping partition monitor work" ); - Self::queue_null_nvlink_status_observation( - pending_null_nvlink_observations, - group_id, - group_type, - snapshots, - ChassisNmxCUnreachableReason::NoEndpoint, - ); - return 0; + early_return!(ChassisNmxCUnreachableReason::NoEndpoint); }; let nmxc_endpoint = match Endpoint::new(endpoint_url) { @@ -1458,14 +1504,7 @@ impl NvlPartitionMonitor { error = %e, "Invalid NMX-C endpoint URI; skipping partition monitor work" ); - Self::queue_null_nvlink_status_observation( - pending_null_nvlink_observations, - group_id, - group_type, - snapshots, - ChassisNmxCUnreachableReason::InvalidEndpointUri, - ); - return 0; + early_return!(ChassisNmxCUnreachableReason::InvalidEndpointUri); } }; @@ -1479,14 +1518,7 @@ impl NvlPartitionMonitor { error = %e, "Failed to create NMX-C client; skipping partition monitor work" ); - Self::queue_null_nvlink_status_observation( - pending_null_nvlink_observations, - group_id, - group_type, - snapshots, - ChassisNmxCUnreachableReason::ClientCreateFailed, - ); - return 0; + early_return!(ChassisNmxCUnreachableReason::ClientCreateFailed); } }; let hello = match nmxc_client.hello(NMX_C_GATEWAY_ID).await { @@ -1499,14 +1531,7 @@ impl NvlPartitionMonitor { error = %e, "NMX-C hello failed; skipping partition monitor work" ); - Self::queue_null_nvlink_status_observation( - pending_null_nvlink_observations, - group_id, - group_type, - snapshots, - ChassisNmxCUnreachableReason::HelloFailed, - ); - return 0; + early_return!(ChassisNmxCUnreachableReason::HelloFailed); } }; let domain_uuid = match domain_uuid_from_nmx_c_hello(&hello) { @@ -1519,14 +1544,7 @@ impl NvlPartitionMonitor { error = %e, "Failed to parse domain UUID from NMX-C hello; skipping partition monitor work" ); - Self::queue_null_nvlink_status_observation( - pending_null_nvlink_observations, - group_id, - group_type, - snapshots, - ChassisNmxCUnreachableReason::DomainUuidParseFailed, - ); - return 0; + early_return!(ChassisNmxCUnreachableReason::DomainUuidParseFailed); } }; @@ -1538,8 +1556,8 @@ impl NvlPartitionMonitor { } // Endpoint + component versions for this group, so the metric-read failures below log them. - metrics.nmxc.endpoint = endpoint_url.to_string(); - metrics.nmxc.version = hello + group_metrics.nmxc.endpoint = endpoint_url.to_string(); + group_metrics.nmxc.version = hello .components_ver .iter() .map(|kv| format!("{}={}", kv.key, kv.value)) @@ -1572,7 +1590,7 @@ impl NvlPartitionMonitor { continue; } nvlink_info_db_updates.extend(populate_machine_nvlink_info_if_needed( - machine_nvlink_info, + &mut machine_nvlink_info, all_managed_host_snapshots, snapshot_chassis_serial.as_deref(), &[machine_id], @@ -1609,10 +1627,10 @@ impl NvlPartitionMonitor { .cloned() .collect(); - match self + let completed_operations = match self .check_partitions_and_apply_nmx_c_operations( nmxc_client.as_mut(), - metrics, + &mut group_metrics, domain_uuid, CheckPartitionsInput { db_nvl_logical_partitions: db_nvl_logical_partitions.to_vec(), @@ -1633,14 +1651,20 @@ impl NvlPartitionMonitor { "Partition monitor work failed; queuing null nvlink status observations" ); Self::queue_null_nvlink_status_observation( - pending_null_nvlink_observations, - group_id, + &mut null_observations, + &group_id, group_type, snapshots, ChassisNmxCUnreachableReason::PartitionMonitorWorkFailed, ); 0 } + }; + + GroupResult { + completed_operations, + null_observations, + partial_metrics: group_metrics, } } @@ -3365,9 +3389,9 @@ mod machine_group_tests { use tokio::task::JoinSet; use super::{ - ChassisNmxCUnreachableReason, NvLinkConfig, NvlPartitionMonitor, - NvlPartitionMonitorMetrics, PendingNullNvlinkObservation, ProcessMachineGroupInput, - group_managed_hosts_by_group_type, nmx_c_endpoint, + ChassisNmxCUnreachableReason, GroupResult, NvLinkConfig, NvlPartitionMonitor, + NvlPartitionMonitorMetrics, ProcessMachineGroupInput, group_managed_hosts_by_group_type, + nmx_c_endpoint, }; use crate::nvlink::test_support::NmxcSimClient; @@ -3652,30 +3676,27 @@ mod machine_group_tests { ]; for (scenario, endpoint_url, expected_reason, expected_machines_scanned) in cases { - let mut metrics = NvlPartitionMonitorMetrics::new(); - let mut pending: Vec = Vec::new(); - let mut machine_nvlink_info = - HashMap::from([(machine_id, machine_nvlink_info.clone())]); - - let completed = monitor - .process_nmx_c_partition_monitor_group( - ProcessMachineGroupInput { - group_id: "rack-1", - group_type: nmx_c_endpoint::ManagedHostGroupType::Rack, - snapshots: &rack_snapshots, - endpoint_url, - rack_id: Some(&rack_id), - all_managed_host_snapshots: &all_snapshots, - machine_nvlink_info: &mut machine_nvlink_info, - db_nvl_partitions: &[], - db_nvl_logical_partitions: &[], - }, - &mut metrics, - &mut pending, - ) + let machine_nvlink_info = HashMap::from([(machine_id, machine_nvlink_info.clone())]); + + let GroupResult { + completed_operations, + null_observations: pending, + partial_metrics: metrics, + } = monitor + .process_nmx_c_partition_monitor_group(ProcessMachineGroupInput { + group_id: "rack-1".to_string(), + group_type: nmx_c_endpoint::ManagedHostGroupType::Rack, + snapshots: &rack_snapshots, + endpoint_url, + rack_id: Some(&rack_id), + all_managed_host_snapshots: &all_snapshots, + machine_nvlink_info, + db_nvl_partitions: &[], + db_nvl_logical_partitions: &[], + }) .await; - assert_eq!(completed, 0, "{scenario}"); + assert_eq!(completed_operations, 0, "{scenario}"); assert_eq!( metrics.num_machines_scanned, expected_machines_scanned, "{scenario}" @@ -3740,4 +3761,129 @@ mod machine_group_tests { Ok(()) } + + #[tokio::test] + async fn concurrent_group_results_merge_into_iteration_metrics() { + use futures::future; + + use super::{AppliedChange, NmxcMetricOperation, NmxcMetricOperationStatus}; + + let applied_create = AppliedChange { + operation: NmxcMetricOperation::Create, + status: NmxcMetricOperationStatus::Completed, + }; + let domain_a = "domain-a".to_string(); + let domain_b = "domain-b".to_string(); + + let mut metrics_a = NvlPartitionMonitorMetrics::new(); + metrics_a.num_machines_scanned = 2; + metrics_a.num_instances_scanned = 1; + metrics_a.applied_changes.insert(applied_create.clone(), 1); + metrics_a + .nmxc + .partition_health + .insert((domain_a.clone(), "healthy"), 3); + metrics_a + .nmxc + .gpu_health + .insert((domain_a.clone(), "healthy"), 4); + metrics_a + .nmxc + .compute_node_health + .insert((domain_a.clone(), "healthy"), 1); + metrics_a.nmxc.endpoint = "http://nmxc-a.example:9370".to_string(); + + let mut metrics_b = NvlPartitionMonitorMetrics::new(); + metrics_b.num_machines_scanned = 3; + metrics_b.num_instances_scanned = 2; + metrics_b.applied_changes.insert(applied_create.clone(), 2); + metrics_b + .nmxc + .partition_health + .insert((domain_b.clone(), "healthy"), 5); + metrics_b + .nmxc + .gpu_health + .insert((domain_b.clone(), "healthy"), 6); + metrics_b + .nmxc + .compute_node_health + .insert((domain_b.clone(), "healthy"), 2); + metrics_b.nmxc.endpoint = "http://nmxc-b.example:9370".to_string(); + + let machine_c = machine_id(23); + let prepared = vec![ + GroupResult { + completed_operations: 2, + null_observations: vec![], + partial_metrics: metrics_a, + }, + GroupResult { + completed_operations: 1, + null_observations: vec![], + partial_metrics: metrics_b, + }, + GroupResult { + completed_operations: 0, + null_observations: vec![super::PendingNullNvlinkObservation { + group_id: "CHASSIS-C".to_string(), + group_type: nmx_c_endpoint::ManagedHostGroupType::Chassis, + reason: ChassisNmxCUnreachableReason::NoEndpoint, + machine_ids: vec![machine_c], + }], + partial_metrics: NvlPartitionMonitorMetrics::new(), + }, + ]; + let group_results = + future::join_all(prepared.into_iter().map(|result| async move { result })).await; + + // Mirror the fan-in fold in run_single_iteration_inner. + let mut metrics = NvlPartitionMonitorMetrics::new(); + metrics.num_logical_partitions = 3; + metrics.num_physical_partitions = 1; + let mut total_completed_operations = 0; + let mut pending_null_observations = Vec::new(); + for result in group_results { + total_completed_operations += result.completed_operations; + pending_null_observations.extend(result.null_observations); + metrics.merge_from(result.partial_metrics); + } + metrics.num_completed_operations = total_completed_operations; + + assert_eq!(total_completed_operations, 3); + assert_eq!(metrics.num_completed_operations, 3); + assert_eq!(metrics.num_machines_scanned, 5); + assert_eq!(metrics.num_instances_scanned, 3); + assert_eq!(metrics.applied_changes[&applied_create], 3); + assert_eq!( + metrics.nmxc.partition_health, + HashMap::from([ + ((domain_a.clone(), "healthy"), 3), + ((domain_b.clone(), "healthy"), 5), + ]) + ); + assert_eq!( + metrics.nmxc.gpu_health, + HashMap::from([ + ((domain_a.clone(), "healthy"), 4), + ((domain_b.clone(), "healthy"), 6), + ]) + ); + assert_eq!( + metrics.nmxc.compute_node_health, + HashMap::from([((domain_a, "healthy"), 1), ((domain_b, "healthy"), 2),]) + ); + assert_eq!(metrics.nmxc.endpoint, "http://nmxc-b.example:9370"); + assert_eq!(metrics.num_logical_partitions, 3); + assert_eq!(metrics.num_physical_partitions, 1); + assert!(metrics.num_nmx_c_unreachable_chassis.is_empty()); + + assert_eq!(pending_null_observations.len(), 1); + assert_eq!(pending_null_observations[0].group_id, "CHASSIS-C"); + assert_eq!( + pending_null_observations[0].reason, + ChassisNmxCUnreachableReason::NoEndpoint + ); + assert_eq!(pending_null_observations[0].machine_ids, vec![machine_c]); + } } diff --git a/crates/nvlink-manager/src/metrics.rs b/crates/nvlink-manager/src/metrics.rs index 1825315e16..6612fabf8b 100644 --- a/crates/nvlink-manager/src/metrics.rs +++ b/crates/nvlink-manager/src/metrics.rs @@ -243,6 +243,72 @@ impl NvlPartitionMonitorMetrics { }, } } + + /// Accumulates per-group metrics collected during concurrent group processing into `self`. + /// + /// Fields that are counters or collections are summed/extended. Single-valued NMX-C metadata + /// (endpoint, version, connect_error) use last-non-empty-wins, matching the previous + /// sequential behaviour where the last processed group determined those values. Health maps + /// are keyed by `(domain_uuid, state)` so entries from different groups have distinct keys + /// and can be merged with `extend`. + /// + /// Fields managed by the caller (`recording_started_at`, `num_logical_partitions`, + /// `num_physical_partitions`, `num_completed_operations`, `num_nmx_c_unreachable_chassis`) + /// are not touched here. + /// + /// Exhaustive destructuring is intentional: adding a field to + /// [`NvlPartitionMonitorMetrics`] or [`NmxcMetrics`] must force a merge decision here. + pub(crate) fn merge_from(&mut self, other: Self) { + let Self { + nmxc: + NmxcMetrics { + endpoint, + connect_error, + version, + partition_health, + gpu_health, + compute_node_health, + }, + num_machines_scanned, + num_instances_scanned, + num_gpus_scanned, + num_machine_nvl_status_updates, + num_nvlink_info_mismatches, + num_stale_partitions_deleted, + applied_changes, + nvlink_config_apply_durations_ms, + // Caller-owned — intentionally not merged. + recording_started_at: _, + num_logical_partitions: _, + num_physical_partitions: _, + num_completed_operations: _, + num_nmx_c_unreachable_chassis: _, + } = other; + + if !endpoint.is_empty() { + self.nmxc.endpoint = endpoint; + } + if !version.is_empty() { + self.nmxc.version = version; + } + if !connect_error.is_empty() { + self.nmxc.connect_error = connect_error; + } + self.nmxc.partition_health.extend(partition_health); + self.nmxc.gpu_health.extend(gpu_health); + self.nmxc.compute_node_health.extend(compute_node_health); + self.num_machines_scanned += num_machines_scanned; + self.num_instances_scanned += num_instances_scanned; + self.num_gpus_scanned += num_gpus_scanned; + self.num_machine_nvl_status_updates += num_machine_nvl_status_updates; + self.num_nvlink_info_mismatches += num_nvlink_info_mismatches; + self.num_stale_partitions_deleted += num_stale_partitions_deleted; + for (k, v) in applied_changes { + *self.applied_changes.entry(k).or_default() += v; + } + self.nvlink_config_apply_durations_ms + .extend(nvlink_config_apply_durations_ms); + } } impl Display for NvlPartitionMonitorMetrics { @@ -658,6 +724,135 @@ mod tests { use super::*; + #[test] + fn merge_from_sums_group_fields_and_preserves_caller_owned() { + let applied_create = AppliedChange { + operation: NmxcMetricOperation::Create, + status: NmxcMetricOperationStatus::Completed, + }; + let applied_remove = AppliedChange { + operation: NmxcMetricOperation::Remove, + status: NmxcMetricOperationStatus::Failed, + }; + + let mut base = NvlPartitionMonitorMetrics::new(); + base.recording_started_at = std::time::Instant::now(); + let started_at = base.recording_started_at; + base.num_logical_partitions = 4; + base.num_physical_partitions = 2; + base.num_completed_operations = 7; + base.num_nmx_c_unreachable_chassis + .insert(ChassisNmxCUnreachableReason::NoEndpoint, 1); + base.num_machines_scanned = 1; + base.num_instances_scanned = 2; + base.num_gpus_scanned = 3; + base.num_machine_nvl_status_updates = 1; + base.num_nvlink_info_mismatches = 1; + base.num_stale_partitions_deleted = 1; + base.applied_changes.insert(applied_create.clone(), 2); + base.nvlink_config_apply_durations_ms.push(10.0); + base.nmxc.endpoint = "https://first.example:9370".to_string(); + base.nmxc.version = "first=1".to_string(); + base.nmxc.connect_error = "first-error".to_string(); + base.nmxc + .partition_health + .insert(("domain-a".to_string(), "healthy"), 1); + base.nmxc + .gpu_health + .insert(("domain-a".to_string(), "healthy"), 2); + base.nmxc + .compute_node_health + .insert(("domain-a".to_string(), "healthy"), 3); + + let mut other = NvlPartitionMonitorMetrics::new(); + other.num_logical_partitions = 99; + other.num_physical_partitions = 99; + other.num_completed_operations = 99; + other + .num_nmx_c_unreachable_chassis + .insert(ChassisNmxCUnreachableReason::HelloFailed, 5); + other.num_machines_scanned = 10; + other.num_instances_scanned = 20; + other.num_gpus_scanned = 30; + other.num_machine_nvl_status_updates = 4; + other.num_nvlink_info_mismatches = 5; + other.num_stale_partitions_deleted = 6; + other.applied_changes.insert(applied_create.clone(), 3); + other.applied_changes.insert(applied_remove.clone(), 1); + other.nvlink_config_apply_durations_ms.push(20.0); + other.nmxc.endpoint = "https://second.example:9370".to_string(); + other.nmxc.version = "second=2".to_string(); + other.nmxc.connect_error = "second-error".to_string(); + other + .nmxc + .partition_health + .insert(("domain-b".to_string(), "healthy"), 4); + other + .nmxc + .gpu_health + .insert(("domain-b".to_string(), "healthy"), 5); + other + .nmxc + .compute_node_health + .insert(("domain-b".to_string(), "healthy"), 6); + + base.merge_from(other); + + assert_eq!(base.num_machines_scanned, 11); + assert_eq!(base.num_instances_scanned, 22); + assert_eq!(base.num_gpus_scanned, 33); + assert_eq!(base.num_machine_nvl_status_updates, 5); + assert_eq!(base.num_nvlink_info_mismatches, 6); + assert_eq!(base.num_stale_partitions_deleted, 7); + assert_eq!(base.applied_changes[&applied_create], 5); + assert_eq!(base.applied_changes[&applied_remove], 1); + assert_eq!(base.nvlink_config_apply_durations_ms, vec![10.0, 20.0]); + assert_eq!(base.nmxc.endpoint, "https://second.example:9370"); + assert_eq!(base.nmxc.version, "second=2"); + assert_eq!(base.nmxc.connect_error, "second-error"); + assert_eq!( + base.nmxc.partition_health, + HashMap::from([ + (("domain-a".to_string(), "healthy"), 1), + (("domain-b".to_string(), "healthy"), 4), + ]) + ); + assert_eq!( + base.nmxc.gpu_health, + HashMap::from([ + (("domain-a".to_string(), "healthy"), 2), + (("domain-b".to_string(), "healthy"), 5), + ]) + ); + assert_eq!( + base.nmxc.compute_node_health, + HashMap::from([ + (("domain-a".to_string(), "healthy"), 3), + (("domain-b".to_string(), "healthy"), 6), + ]) + ); + + // Caller-owned fields must not change during merge. + assert_eq!(base.recording_started_at, started_at); + assert_eq!(base.num_logical_partitions, 4); + assert_eq!(base.num_physical_partitions, 2); + assert_eq!(base.num_completed_operations, 7); + assert_eq!( + base.num_nmx_c_unreachable_chassis, + HashMap::from([(ChassisNmxCUnreachableReason::NoEndpoint, 1)]) + ); + + // Empty metadata from `other` must not overwrite existing values. + let mut keep = NvlPartitionMonitorMetrics::new(); + keep.nmxc.endpoint = "https://keep.example:9370".to_string(); + keep.nmxc.version = "keep=1".to_string(); + keep.nmxc.connect_error = "keep-error".to_string(); + keep.merge_from(NvlPartitionMonitorMetrics::new()); + assert_eq!(keep.nmxc.endpoint, "https://keep.example:9370"); + assert_eq!(keep.nmxc.version, "keep=1"); + assert_eq!(keep.nmxc.connect_error, "keep-error"); + } + #[test] fn partition_monitor_iteration_records_latency_and_warns_only_on_failure() { const METRIC_NAME: &str = "carbide_nvlink_partition_monitor_iteration_latency_milliseconds";