From 4e539a0bf9ec2bdc66b66569ae01f10c4c17a2d5 Mon Sep 17 00:00:00 2001 From: Roopesh Tamma Date: Thu, 6 Aug 2026 01:36:40 -0700 Subject: [PATCH 1/5] perf(nvlink-manager): process partition monitor groups concurrently - Batch chassis NMX-C endpoint lookups instead of one query per chassis serial - Pre-split machine_nvlink_info into per-group disjoint shards before the concurrent fan-out; groups are mutually exclusive so no locking is needed - Add a configurable concurreny limit (default 16) Signed-off-by: Roopesh Tamma Signed-off-by: Roopesh Tamma --- Cargo.lock | 1 + crates/api-db/src/nvlink_nmxc_endpoints.rs | 73 ++++ crates/nvlink-manager/Cargo.toml | 1 + crates/nvlink-manager/src/config.rs | 14 + crates/nvlink-manager/src/lib.rs | 418 +++++++++++++------- crates/nvlink-manager/src/metrics.rs | 195 +++++++++ crates/nvlink-manager/src/nmx_c_endpoint.rs | 207 ---------- 7 files changed, 566 insertions(+), 343 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 8aec76435c..6b88530882 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2650,6 +2650,7 @@ dependencies = [ "config-version", "duration-str", "eyre", + "futures", "hex", "http", "libnmxc", 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..1a6b8bf856 100644 --- a/crates/nvlink-manager/src/config.rs +++ b/crates/nvlink-manager/src/config.rs @@ -54,12 +54,22 @@ 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. + #[serde(default = "NvLinkConfig::default_partition_monitor_max_concurrent_groups")] + pub partition_monitor_max_concurrent_groups: usize, } 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() -> usize { + 16 + } } #[derive(Clone, Debug, Deserialize, Serialize, PartialEq)] @@ -131,6 +141,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 +169,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(), } ); } diff --git a/crates/nvlink-manager/src/lib.rs b/crates/nvlink-manager/src/lib.rs index 2526447dc3..e9edd029f5 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); + 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 854851ca71..9d358f22d6 100644 --- a/crates/nvlink-manager/src/metrics.rs +++ b/crates/nvlink-manager/src/metrics.rs @@ -244,6 +244,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 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 { @@ -662,6 +728,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"; diff --git a/crates/nvlink-manager/src/nmx_c_endpoint.rs b/crates/nvlink-manager/src/nmx_c_endpoint.rs index 3312b4f49e..a253e0f3a2 100644 --- a/crates/nvlink-manager/src/nmx_c_endpoint.rs +++ b/crates/nvlink-manager/src/nmx_c_endpoint.rs @@ -17,10 +17,6 @@ use std::net::IpAddr; -use carbide_uuid::rack::RackId; -use db::db_read::DbReader; -use db::{self, DatabaseResult}; - use crate::config::NvLinkConfig; /// Whether an NMX-C monitor group is keyed by chassis serial or rack id. @@ -75,69 +71,10 @@ pub fn nmx_c_endpoint_url_from_nvos_ip( ) } -/// Resolves the NMX-C gRPC endpoint URL for a chassis- or rack-scoped machine group. -/// -/// - [`ManagedHostGroupType::Chassis`]: looks up `nvlink_nmxc_endpoints` by `chassis_serial`. -/// - [`ManagedHostGroupType::Rack`]: uses the first ready Fabric Manager control-plane switch's -/// NVOS IP in `rack_id`. Does not fall back to the chassis-serial mapping. -pub async fn resolve_nmx_c_endpoint_url( - db: &mut DB, - group_type: ManagedHostGroupType, - rack_id: Option<&RackId>, - chassis_serial: Option<&str>, - nvlink_config: &NvLinkConfig, -) -> DatabaseResult> -where - for<'db> &'db mut DB: DbReader<'db>, -{ - match group_type { - ManagedHostGroupType::Chassis => { - let Some(chassis_serial) = chassis_serial else { - return Ok(None); - }; - Ok( - db::nvlink_nmxc_endpoints::find_by_chassis_serial(&mut *db, chassis_serial.trim()) - .await? - .map(|row| row.endpoint), - ) - } - ManagedHostGroupType::Rack => { - let Some(rack_id) = rack_id else { - return Ok(None); - }; - - let switch_ids = db::switch::find_ready_control_plane_configured_switch_ids_in_rack( - &mut *db, rack_id, - ) - .await?; - - let Some(switch_id) = switch_ids.first() else { - return Ok(None); - }; - - let endpoint_rows = - db::switch::find_switch_endpoints_by_ids(&mut *db, &[*switch_id]).await?; - - Ok(endpoint_rows - .first() - .and_then(|row| row.nvos_ip.as_ref()) - .map(|nvos_ip| nmx_c_endpoint_url_from_nvos_ip(nvos_ip, None, nvlink_config))) - } - } -} - #[cfg(test)] mod tests { use std::net::Ipv4Addr; - use carbide_macros::sqlx_test; - use carbide_uuid::rack::{RackId, RackProfileId}; - use model::rack::RackConfig; - use model::switch::{ - CONTROL_PLANE_STATE_CONFIGURED, FabricManagerState, FabricManagerStatus, - SwitchControllerState, - }; - use super::*; #[test] @@ -172,148 +109,4 @@ mod tests { "https://10.0.0.1:9601" ); } - - #[sqlx_test] - async fn resolve_chassis_uses_nvlink_nmxc_endpoints_mapping( - pool: sqlx::PgPool, - ) -> Result<(), Box> { - let chassis_serial = "CHASSIS-ENDPOINT-1"; - let mapped_endpoint = "https://nmxc.example:9370"; - let config = NvLinkConfig::default(); - - let mut txn = pool.begin().await?; - db::nvlink_nmxc_endpoints::create(txn.as_mut(), chassis_serial, mapped_endpoint).await?; - - let resolved = resolve_nmx_c_endpoint_url( - txn.as_mut(), - ManagedHostGroupType::Chassis, - None, - Some(chassis_serial), - &config, - ) - .await?; - assert_eq!(resolved.as_deref(), Some(mapped_endpoint)); - txn.rollback().await?; - Ok(()) - } - - #[sqlx_test] - async fn resolve_rack_uses_ready_switch_nvos_ip( - pool: sqlx::PgPool, - ) -> Result<(), Box> { - let rack_id: RackId = "rack-nmxc-endpoint".parse()?; - let config = NvLinkConfig::default(); - - let mut txn = pool.begin().await?; - db::rack::create( - txn.as_mut(), - &rack_id, - Some(&RackProfileId::new("NVL72")), - &RackConfig::default(), - None, - ) - .await?; - txn.commit().await?; - - let mut txn = pool.begin().await?; - let switch = - db::test_support::switch::create_seeded_discovered(txn.as_mut(), 1, "Switch1").await?; - txn.commit().await?; - - let mut txn = pool.begin().await?; - sqlx::query("UPDATE switches SET rack_id = $1 WHERE id = $2") - .bind(&rack_id) - .bind(switch.id) - .execute(txn.as_mut()) - .await?; - - let switch = db::switch::find_by_id(txn.as_mut(), &switch.id) - .await? - .expect("switch should exist"); - assert!( - db::switch::try_update_controller_state( - txn.as_mut(), - switch.id, - switch.controller_state.version, - switch.controller_state.version.increment(), - &SwitchControllerState::Ready, - ) - .await? - ); - db::switch::update_fabric_manager_status( - txn.as_mut(), - switch.id, - Some(&FabricManagerStatus { - fabric_manager_state: FabricManagerState::Ok, - addition_info: Some(CONTROL_PLANE_STATE_CONFIGURED.to_string()), - reason: None, - error_message: None, - }), - ) - .await?; - - let expected_nvos_ip = db::switch::find_switch_endpoints_by_ids(txn.as_mut(), &[switch.id]) - .await? - .pop() - .expect("switch endpoint") - .nvos_ip - .expect("seeded switch should have an NVOS IP"); - let expected_url = nmx_c_endpoint_url_from_nvos_ip(&expected_nvos_ip, None, &config); - - let resolved = resolve_nmx_c_endpoint_url( - txn.as_mut(), - ManagedHostGroupType::Rack, - Some(&rack_id), - None, - &config, - ) - .await?; - assert_eq!(resolved.as_deref(), Some(expected_url.as_str())); - txn.rollback().await?; - Ok(()) - } - - #[sqlx_test] - async fn resolve_rack_does_not_fall_back_to_chassis_mapping( - pool: sqlx::PgPool, - ) -> Result<(), Box> { - let rack_id: RackId = "rack-nmxc-no-fallback".parse()?; - let chassis_serial = "CHASSIS-NO-FALLBACK"; - let mapped_endpoint = "https://chassis-fallback.example:9370"; - let config = NvLinkConfig::default(); - - let mut txn = pool.begin().await?; - db::rack::create( - txn.as_mut(), - &rack_id, - Some(&RackProfileId::new("NVL72")), - &RackConfig::default(), - None, - ) - .await?; - db::nvlink_nmxc_endpoints::create(txn.as_mut(), chassis_serial, mapped_endpoint).await?; - - // Rack has no ready control-plane switch. Chassis mapping exists and must not be used. - let rack_resolved = resolve_nmx_c_endpoint_url( - txn.as_mut(), - ManagedHostGroupType::Rack, - Some(&rack_id), - Some(chassis_serial), - &config, - ) - .await?; - assert_eq!(rack_resolved, None); - - let chassis_resolved = resolve_nmx_c_endpoint_url( - txn.as_mut(), - ManagedHostGroupType::Chassis, - None, - Some(chassis_serial), - &config, - ) - .await?; - assert_eq!(chassis_resolved.as_deref(), Some(mapped_endpoint)); - txn.rollback().await?; - Ok(()) - } } From 8324d4fbe196a0b2b2c660f27914868a7c265551 Mon Sep 17 00:00:00 2001 From: Roopesh Tamma Date: Thu, 6 Aug 2026 22:59:26 -0700 Subject: [PATCH 2/5] chore: code review fix, prevent config value of 0 for partiton_monitor_mzx_concurrent_groups Signed-off-by: Roopesh Tamma --- crates/nvlink-manager/src/config.rs | 16 ++++++++++++---- crates/nvlink-manager/src/lib.rs | 2 +- 2 files changed, 13 insertions(+), 5 deletions(-) diff --git a/crates/nvlink-manager/src/config.rs b/crates/nvlink-manager/src/config.rs index 1a6b8bf856..d0ad0791d1 100644 --- a/crates/nvlink-manager/src/config.rs +++ b/crates/nvlink-manager/src/config.rs @@ -57,9 +57,9 @@ pub struct NvLinkConfig { /// 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. + /// 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: usize, + pub partition_monitor_max_concurrent_groups: std::num::NonZeroUsize, } impl NvLinkConfig { @@ -67,8 +67,9 @@ impl NvLinkConfig { std::time::Duration::from_secs(60) } - pub const fn default_partition_monitor_max_concurrent_groups() -> usize { - 16 + pub const fn default_partition_monitor_max_concurrent_groups() -> std::num::NonZeroUsize { + // SAFETY: 16 is non-zero. + unsafe { std::num::NonZeroUsize::new_unchecked(16) } } } @@ -186,6 +187,13 @@ 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 e9edd029f5..c0429e5569 100644 --- a/crates/nvlink-manager/src/lib.rs +++ b/crates/nvlink-manager/src/lib.rs @@ -1403,7 +1403,7 @@ impl NvlPartitionMonitor { // 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); + 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. From d32eea5c000286ceffddbee81a06b4e2213a013c Mon Sep 17 00:00:00 2001 From: Roopesh Tamma Date: Thu, 6 Aug 2026 23:53:28 -0700 Subject: [PATCH 3/5] chore: Update README.md, add partition_monitor_max_concurrent_groups to NvLinkConfig table Signed-off-by: Roopesh Tamma --- crates/api-core/src/cfg/README.md | 1 + 1 file changed, 1 insertion(+) 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` From 2f61d4f71ed2584dc25e45e2fbd9e9fbd0082ff3 Mon Sep 17 00:00:00 2001 From: Roopesh Tamma Date: Sun, 9 Aug 2026 22:07:44 -0700 Subject: [PATCH 4/5] chore: clippy fixes Signed-off-by: Roopesh Tamma --- crates/nvlink-manager/src/metrics.rs | 2 +- crates/nvlink-manager/src/nmx_c_endpoint.rs | 207 ++++++++++++++++++++ 2 files changed, 208 insertions(+), 1 deletion(-) diff --git a/crates/nvlink-manager/src/metrics.rs b/crates/nvlink-manager/src/metrics.rs index 92412bae3f..6612fabf8b 100644 --- a/crates/nvlink-manager/src/metrics.rs +++ b/crates/nvlink-manager/src/metrics.rs @@ -258,7 +258,7 @@ impl NvlPartitionMonitorMetrics { /// /// Exhaustive destructuring is intentional: adding a field to /// [`NvlPartitionMonitorMetrics`] or [`NmxcMetrics`] must force a merge decision here. - pub fn merge_from(&mut self, other: Self) { + pub(crate) fn merge_from(&mut self, other: Self) { let Self { nmxc: NmxcMetrics { diff --git a/crates/nvlink-manager/src/nmx_c_endpoint.rs b/crates/nvlink-manager/src/nmx_c_endpoint.rs index a253e0f3a2..3312b4f49e 100644 --- a/crates/nvlink-manager/src/nmx_c_endpoint.rs +++ b/crates/nvlink-manager/src/nmx_c_endpoint.rs @@ -17,6 +17,10 @@ use std::net::IpAddr; +use carbide_uuid::rack::RackId; +use db::db_read::DbReader; +use db::{self, DatabaseResult}; + use crate::config::NvLinkConfig; /// Whether an NMX-C monitor group is keyed by chassis serial or rack id. @@ -71,10 +75,69 @@ pub fn nmx_c_endpoint_url_from_nvos_ip( ) } +/// Resolves the NMX-C gRPC endpoint URL for a chassis- or rack-scoped machine group. +/// +/// - [`ManagedHostGroupType::Chassis`]: looks up `nvlink_nmxc_endpoints` by `chassis_serial`. +/// - [`ManagedHostGroupType::Rack`]: uses the first ready Fabric Manager control-plane switch's +/// NVOS IP in `rack_id`. Does not fall back to the chassis-serial mapping. +pub async fn resolve_nmx_c_endpoint_url( + db: &mut DB, + group_type: ManagedHostGroupType, + rack_id: Option<&RackId>, + chassis_serial: Option<&str>, + nvlink_config: &NvLinkConfig, +) -> DatabaseResult> +where + for<'db> &'db mut DB: DbReader<'db>, +{ + match group_type { + ManagedHostGroupType::Chassis => { + let Some(chassis_serial) = chassis_serial else { + return Ok(None); + }; + Ok( + db::nvlink_nmxc_endpoints::find_by_chassis_serial(&mut *db, chassis_serial.trim()) + .await? + .map(|row| row.endpoint), + ) + } + ManagedHostGroupType::Rack => { + let Some(rack_id) = rack_id else { + return Ok(None); + }; + + let switch_ids = db::switch::find_ready_control_plane_configured_switch_ids_in_rack( + &mut *db, rack_id, + ) + .await?; + + let Some(switch_id) = switch_ids.first() else { + return Ok(None); + }; + + let endpoint_rows = + db::switch::find_switch_endpoints_by_ids(&mut *db, &[*switch_id]).await?; + + Ok(endpoint_rows + .first() + .and_then(|row| row.nvos_ip.as_ref()) + .map(|nvos_ip| nmx_c_endpoint_url_from_nvos_ip(nvos_ip, None, nvlink_config))) + } + } +} + #[cfg(test)] mod tests { use std::net::Ipv4Addr; + use carbide_macros::sqlx_test; + use carbide_uuid::rack::{RackId, RackProfileId}; + use model::rack::RackConfig; + use model::switch::{ + CONTROL_PLANE_STATE_CONFIGURED, FabricManagerState, FabricManagerStatus, + SwitchControllerState, + }; + use super::*; #[test] @@ -109,4 +172,148 @@ mod tests { "https://10.0.0.1:9601" ); } + + #[sqlx_test] + async fn resolve_chassis_uses_nvlink_nmxc_endpoints_mapping( + pool: sqlx::PgPool, + ) -> Result<(), Box> { + let chassis_serial = "CHASSIS-ENDPOINT-1"; + let mapped_endpoint = "https://nmxc.example:9370"; + let config = NvLinkConfig::default(); + + let mut txn = pool.begin().await?; + db::nvlink_nmxc_endpoints::create(txn.as_mut(), chassis_serial, mapped_endpoint).await?; + + let resolved = resolve_nmx_c_endpoint_url( + txn.as_mut(), + ManagedHostGroupType::Chassis, + None, + Some(chassis_serial), + &config, + ) + .await?; + assert_eq!(resolved.as_deref(), Some(mapped_endpoint)); + txn.rollback().await?; + Ok(()) + } + + #[sqlx_test] + async fn resolve_rack_uses_ready_switch_nvos_ip( + pool: sqlx::PgPool, + ) -> Result<(), Box> { + let rack_id: RackId = "rack-nmxc-endpoint".parse()?; + let config = NvLinkConfig::default(); + + let mut txn = pool.begin().await?; + db::rack::create( + txn.as_mut(), + &rack_id, + Some(&RackProfileId::new("NVL72")), + &RackConfig::default(), + None, + ) + .await?; + txn.commit().await?; + + let mut txn = pool.begin().await?; + let switch = + db::test_support::switch::create_seeded_discovered(txn.as_mut(), 1, "Switch1").await?; + txn.commit().await?; + + let mut txn = pool.begin().await?; + sqlx::query("UPDATE switches SET rack_id = $1 WHERE id = $2") + .bind(&rack_id) + .bind(switch.id) + .execute(txn.as_mut()) + .await?; + + let switch = db::switch::find_by_id(txn.as_mut(), &switch.id) + .await? + .expect("switch should exist"); + assert!( + db::switch::try_update_controller_state( + txn.as_mut(), + switch.id, + switch.controller_state.version, + switch.controller_state.version.increment(), + &SwitchControllerState::Ready, + ) + .await? + ); + db::switch::update_fabric_manager_status( + txn.as_mut(), + switch.id, + Some(&FabricManagerStatus { + fabric_manager_state: FabricManagerState::Ok, + addition_info: Some(CONTROL_PLANE_STATE_CONFIGURED.to_string()), + reason: None, + error_message: None, + }), + ) + .await?; + + let expected_nvos_ip = db::switch::find_switch_endpoints_by_ids(txn.as_mut(), &[switch.id]) + .await? + .pop() + .expect("switch endpoint") + .nvos_ip + .expect("seeded switch should have an NVOS IP"); + let expected_url = nmx_c_endpoint_url_from_nvos_ip(&expected_nvos_ip, None, &config); + + let resolved = resolve_nmx_c_endpoint_url( + txn.as_mut(), + ManagedHostGroupType::Rack, + Some(&rack_id), + None, + &config, + ) + .await?; + assert_eq!(resolved.as_deref(), Some(expected_url.as_str())); + txn.rollback().await?; + Ok(()) + } + + #[sqlx_test] + async fn resolve_rack_does_not_fall_back_to_chassis_mapping( + pool: sqlx::PgPool, + ) -> Result<(), Box> { + let rack_id: RackId = "rack-nmxc-no-fallback".parse()?; + let chassis_serial = "CHASSIS-NO-FALLBACK"; + let mapped_endpoint = "https://chassis-fallback.example:9370"; + let config = NvLinkConfig::default(); + + let mut txn = pool.begin().await?; + db::rack::create( + txn.as_mut(), + &rack_id, + Some(&RackProfileId::new("NVL72")), + &RackConfig::default(), + None, + ) + .await?; + db::nvlink_nmxc_endpoints::create(txn.as_mut(), chassis_serial, mapped_endpoint).await?; + + // Rack has no ready control-plane switch. Chassis mapping exists and must not be used. + let rack_resolved = resolve_nmx_c_endpoint_url( + txn.as_mut(), + ManagedHostGroupType::Rack, + Some(&rack_id), + Some(chassis_serial), + &config, + ) + .await?; + assert_eq!(rack_resolved, None); + + let chassis_resolved = resolve_nmx_c_endpoint_url( + txn.as_mut(), + ManagedHostGroupType::Chassis, + None, + Some(chassis_serial), + &config, + ) + .await?; + assert_eq!(chassis_resolved.as_deref(), Some(mapped_endpoint)); + txn.rollback().await?; + Ok(()) + } } From ad7e00bc1c86efee9e497fb1950a663c7bdfb20b Mon Sep 17 00:00:00 2001 From: Roopesh Tamma Date: Sun, 9 Aug 2026 22:25:48 -0700 Subject: [PATCH 5/5] chore: format nightly cleanup Signed-off-by: Roopesh Tamma --- crates/nvlink-manager/src/config.rs | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/crates/nvlink-manager/src/config.rs b/crates/nvlink-manager/src/config.rs index d0ad0791d1..2019f2115b 100644 --- a/crates/nvlink-manager/src/config.rs +++ b/crates/nvlink-manager/src/config.rs @@ -189,8 +189,9 @@ 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}"#); + 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"); }