diff --git a/rustfs/src/connect/environment.rs b/rustfs/src/connect/environment.rs new file mode 100644 index 000000000..b8f7f21c9 --- /dev/null +++ b/rustfs/src/connect/environment.rs @@ -0,0 +1,352 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! Bounded, identifier-free deployment environment inventory. + +use std::collections::BTreeSet; +use std::sync::{Arc, LazyLock}; +use std::time::Duration; + +use serde::Serialize; +use sysinfo::{Disks, Networks, RefreshKind, System}; +use thiserror::Error; +use tokio::sync::Semaphore; +use tokio_util::sync::CancellationToken; + +use super::inventory::InventorySnapshot; + +pub const ENVIRONMENT_CAPABILITY: &str = "inventory.environment@1"; +pub const ENVIRONMENT_SCHEMA_VERSION: u16 = 1; +pub const MAX_ENVIRONMENT_DURATION: Duration = Duration::from_secs(30); + +static SYSTEM_SCAN_PERMIT: LazyLock> = LazyLock::new(|| Arc::new(Semaphore::new(1))); + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct EnvironmentCollectionRequest { + timeout: Duration, +} + +impl EnvironmentCollectionRequest { + pub fn negotiate(schema_version: u16, capability: &str, timeout: Duration) -> Result { + if schema_version != ENVIRONMENT_SCHEMA_VERSION { + return Err(EnvironmentError::UnsupportedVersion); + } + if capability != ENVIRONMENT_CAPABILITY { + return Err(EnvironmentError::UnsupportedCapability); + } + if timeout.is_zero() || timeout > MAX_ENVIRONMENT_DURATION { + return Err(EnvironmentError::InvalidTimeout); + } + Ok(Self { timeout }) + } +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Serialize)] +#[serde(rename_all = "UPPERCASE")] +pub enum EnvironmentOsFamily { + Linux, + Darwin, + Windows, + Freebsd, + Other, +} + +impl EnvironmentOsFamily { + fn current() -> Self { + match std::env::consts::OS { + "linux" => Self::Linux, + "macos" => Self::Darwin, + "windows" => Self::Windows, + "freebsd" => Self::Freebsd, + _ => Self::Other, + } + } +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Serialize)] +#[serde(rename_all = "lowercase")] +pub enum EnvironmentFilesystemType { + Ext4, + Xfs, + Zfs, + Apfs, + Other, +} + +impl EnvironmentFilesystemType { + fn from_reported(value: &str) -> Self { + match value.to_ascii_lowercase().as_str() { + "ext4" => Self::Ext4, + "xfs" => Self::Xfs, + "zfs" => Self::Zfs, + "apfs" => Self::Apfs, + _ => Self::Other, + } + } +} + +fn filesystem_types<'a>(reported: impl IntoIterator) -> Result, EnvironmentError> { + let values = reported + .into_iter() + .map(EnvironmentFilesystemType::from_reported) + .collect::>() + .into_iter() + .collect::>(); + if values.is_empty() { + return Err(EnvironmentError::SourceUnavailable(EnvironmentSource::Filesystem)); + } + Ok(values) +} + +#[derive(Clone, Debug, PartialEq, Eq, Serialize)] +#[serde(deny_unknown_fields, rename_all = "camelCase")] +pub struct EnvironmentInventory { + node_count: u16, + drive_count: u32, + os_family: EnvironmentOsFamily, + filesystem_types: Vec, +} + +impl EnvironmentInventory { + pub fn node_count(&self) -> u16 { + self.node_count + } + + pub fn drive_count(&self) -> u32 { + self.drive_count + } + + pub fn os_family(&self) -> EnvironmentOsFamily { + self.os_family + } + + pub fn filesystem_types(&self) -> &[EnvironmentFilesystemType] { + &self.filesystem_types + } +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum EnvironmentSource { + Filesystem, + Cpu, + Memory, + Network, +} + +#[derive(Debug, Error, PartialEq, Eq)] +pub enum EnvironmentError { + #[error("inventory_environment_unsupported_version")] + UnsupportedVersion, + #[error("inventory_environment_unsupported_capability")] + UnsupportedCapability, + #[error("inventory_environment_invalid_timeout")] + InvalidTimeout, + #[error("inventory_environment_source_unavailable")] + SourceUnavailable(EnvironmentSource), + #[error("inventory_environment_cancelled")] + Cancelled, + #[error("inventory_environment_timed_out")] + TimedOut, + #[error("inventory_environment_task_failed")] + TaskFailed, +} + +#[derive(Debug, PartialEq, Eq)] +pub(crate) struct HostEnvironment { + pub(crate) os_summary: String, + pub(crate) kernel_summary: String, + pub(crate) architecture: &'static str, + pub(crate) cores: usize, + pub(crate) total_memory_bytes: u64, + pub(crate) under_memory_pressure: bool, + pub(crate) filesystem_types: Vec, + pub(crate) interface_count: usize, + pub(crate) bond_count: usize, +} + +impl HostEnvironment { + fn collect() -> Result { + #[cfg(test)] + let _scan = test_support::ScanGuard::start(); + + // Processes, names, addresses, paths, mount options and device labels are outside this schema. + let system = System::new_with_specifics(RefreshKind::everything().without_processes()); + let cores = system.cpus().len(); + if cores == 0 { + return Err(EnvironmentError::SourceUnavailable(EnvironmentSource::Cpu)); + } + let total_memory_bytes = system.total_memory(); + if total_memory_bytes == 0 { + return Err(EnvironmentError::SourceUnavailable(EnvironmentSource::Memory)); + } + let available_memory = system.available_memory(); + let disks = Disks::new_with_refreshed_list(); + let reported_filesystems = disks + .iter() + .map(|disk| disk.file_system().to_string_lossy()) + .collect::>(); + let filesystem_types = filesystem_types(reported_filesystems.iter().map(AsRef::as_ref))?; + let networks = Networks::new_with_refreshed_list(); + if networks.is_empty() { + return Err(EnvironmentError::SourceUnavailable(EnvironmentSource::Network)); + } + + Ok(Self { + os_summary: System::long_os_version().unwrap_or_else(|| "unknown".to_owned()), + kernel_summary: System::kernel_long_version(), + architecture: std::env::consts::ARCH, + cores, + total_memory_bytes, + under_memory_pressure: available_memory.saturating_mul(10) < total_memory_bytes, + filesystem_types, + interface_count: networks.len(), + bond_count: networks.keys().filter(|name| name.starts_with("bond")).count(), + }) + } +} + +pub async fn collect_environment( + inventory: &InventorySnapshot, + request: EnvironmentCollectionRequest, + cancel: &CancellationToken, +) -> Result { + let host = collect_host_environment(request.timeout, cancel).await?; + Ok(EnvironmentInventory { + node_count: inventory.node_count(), + drive_count: inventory.drive_count(), + os_family: EnvironmentOsFamily::current(), + filesystem_types: host.filesystem_types, + }) +} + +pub(crate) async fn collect_host_environment( + timeout: Duration, + cancel: &CancellationToken, +) -> Result { + let deadline = tokio::time::Instant::now() + timeout; + let permit = tokio::select! { + biased; + () = cancel.cancelled() => return Err(EnvironmentError::Cancelled), + result = tokio::time::timeout_at(deadline, SYSTEM_SCAN_PERMIT.clone().acquire_owned()) => { + match result { + Ok(Ok(permit)) => permit, + Ok(Err(_)) => return Err(EnvironmentError::TaskFailed), + Err(_) => return Err(EnvironmentError::TimedOut), + } + } + }; + let task = tokio::task::spawn_blocking(move || { + let _permit = permit; + HostEnvironment::collect() + }); + tokio::select! { + biased; + () = cancel.cancelled() => Err(EnvironmentError::Cancelled), + result = tokio::time::timeout_at(deadline, task) => { + match result { + Ok(Ok(environment)) => environment, + Ok(Err(_)) => Err(EnvironmentError::TaskFailed), + Err(_) => Err(EnvironmentError::TimedOut), + } + } + } +} + +#[cfg(test)] +mod test_support { + use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering}; + use std::time::Duration; + + pub(super) static DELAY_MILLIS: AtomicU64 = AtomicU64::new(0); + pub(super) static ACTIVE: AtomicUsize = AtomicUsize::new(0); + pub(super) static MAX_ACTIVE: AtomicUsize = AtomicUsize::new(0); + + pub(super) struct ScanGuard; + + impl ScanGuard { + pub(super) fn start() -> Self { + let active = ACTIVE.fetch_add(1, Ordering::SeqCst) + 1; + MAX_ACTIVE.fetch_max(active, Ordering::SeqCst); + let delay = DELAY_MILLIS.load(Ordering::SeqCst); + if delay != 0 { + std::thread::sleep(Duration::from_millis(delay)); + } + Self + } + } + + impl Drop for ScanGuard { + fn drop(&mut self) { + ACTIVE.fetch_sub(1, Ordering::SeqCst); + } + } +} + +#[cfg(test)] +mod tests { + use std::sync::atomic::Ordering; + + use super::*; + + async fn wait_for_active(expected: usize) { + tokio::time::timeout(Duration::from_secs(1), async { + while test_support::ACTIVE.load(Ordering::SeqCst) != expected { + tokio::time::sleep(Duration::from_millis(5)).await; + } + }) + .await + .expect("system scan reaches expected state"); + } + + #[tokio::test] + async fn timed_out_and_cancelled_scans_remain_single_flight() { + test_support::MAX_ACTIVE.store(0, Ordering::SeqCst); + test_support::DELAY_MILLIS.store(150, Ordering::SeqCst); + + let cancel = CancellationToken::new(); + assert_eq!( + collect_host_environment(Duration::from_millis(20), &cancel).await, + Err(EnvironmentError::TimedOut) + ); + assert_eq!(test_support::ACTIVE.load(Ordering::SeqCst), 1); + + let second_cancel = CancellationToken::new(); + let second = tokio::spawn({ + let second_cancel = second_cancel.clone(); + async move { collect_host_environment(Duration::from_secs(1), &second_cancel).await } + }); + tokio::time::sleep(Duration::from_millis(20)).await; + assert_eq!(test_support::MAX_ACTIVE.load(Ordering::SeqCst), 1); + second_cancel.cancel(); + assert_eq!(second.await.expect("second scan"), Err(EnvironmentError::Cancelled)); + wait_for_active(0).await; + test_support::DELAY_MILLIS.store(0, Ordering::SeqCst); + } + + #[test] + fn filesystem_projection_drops_paths_options_and_secret_like_values() { + assert_eq!( + filesystem_types(["xfs", "ext4", "/srv/customer-a", "rw,password=SYNTHETIC_SECRET_123"]), + Ok(vec![ + EnvironmentFilesystemType::Ext4, + EnvironmentFilesystemType::Xfs, + EnvironmentFilesystemType::Other, + ]) + ); + assert_eq!( + filesystem_types(std::iter::empty()), + Err(EnvironmentError::SourceUnavailable(EnvironmentSource::Filesystem)) + ); + } +} diff --git a/rustfs/src/connect/mod.rs b/rustfs/src/connect/mod.rs index ec3a1f999..774003d6b 100644 --- a/rustfs/src/connect/mod.rs +++ b/rustfs/src/connect/mod.rs @@ -28,6 +28,7 @@ pub mod client; pub mod config; pub mod credential_store; +pub mod environment; pub mod heartbeat; pub mod identity; pub mod identity_store; @@ -41,6 +42,10 @@ mod telemetry; pub use client::{ClientError, ConnectClient, ConnectConfig}; pub use config::{HeartbeatConfig, HeartbeatConfigError, HeartbeatSchedule}; pub use credential_store::{CredentialStore, DeviceCredential}; +pub use environment::{ + ENVIRONMENT_CAPABILITY, ENVIRONMENT_SCHEMA_VERSION, EnvironmentCollectionRequest, EnvironmentError, + EnvironmentFilesystemType, EnvironmentInventory, EnvironmentOsFamily, MAX_ENVIRONMENT_DURATION, collect_environment, +}; pub use heartbeat::{CoarseNodeSummary, HeartbeatError, HeartbeatStatus}; pub use identity::{DeviceIdentity, IdentityError, RegistrationProof, RegistrationTranscript}; pub use identity_store::{IdentityStore, StoreError}; diff --git a/rustfs/src/connect/offline/collectors.rs b/rustfs/src/connect/offline/collectors.rs index 4cfdce16a..3b51c5539 100644 --- a/rustfs/src/connect/offline/collectors.rs +++ b/rustfs/src/connect/offline/collectors.rs @@ -14,25 +14,21 @@ //! Fixed Q07 L0/L1 collectors for an operator-triggered offline diagnostic. -use std::collections::BTreeSet; use std::path::Path; -use std::sync::{Arc, LazyLock}; use std::time::Duration; use serde::Serialize; use serde_json::{Value, json}; -use sysinfo::{Disks, Networks, RefreshKind, System}; use thiserror::Error; -use tokio::sync::Semaphore; use tokio_util::sync::CancellationToken; +use super::super::environment::{EnvironmentError, HostEnvironment, collect_host_environment}; use super::super::inventory::{InventoryError, InventorySnapshot, InventoryStateStore}; use super::manifest_entry::ManifestEntry; use super::redaction::RedactionError; const COLLECT_TIMEOUT: Duration = Duration::from_secs(2); const MAX_ENTRY_BYTES: usize = 16 * 1024; -static SYSTEM_SCAN_PERMIT: LazyLock> = LazyLock::new(|| Arc::new(Semaphore::new(1))); #[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize)] pub enum DataClassification { @@ -118,7 +114,7 @@ impl OfflineCollector { self.field_id().split_once('.').expect("collector field ids are frozen").1 } - fn value(self, inventory: &InventorySnapshot, system: &SystemSnapshot) -> Value { + fn value(self, inventory: &InventorySnapshot, system: &HostEnvironment) -> Value { match self { Self::RustfsVersion => json!(inventory.rustfs_version()), Self::NodeCount => json!(inventory.node_count()), @@ -150,6 +146,8 @@ pub enum CollectorError { TimedOut, #[error("offline diagnostic collector task failed")] TaskFailed, + #[error("offline diagnostic source is unavailable")] + SourceUnavailable, #[error("offline diagnostic field {field_id} exceeds its {limit} byte entry budget")] EntryTooLarge { field_id: &'static str, limit: usize }, #[error("offline diagnostic entry is not representable as JSON")] @@ -160,6 +158,20 @@ pub enum CollectorError { Redaction(#[from] RedactionError), } +impl From for CollectorError { + fn from(error: EnvironmentError) -> Self { + match error { + EnvironmentError::Cancelled => Self::Cancelled, + EnvironmentError::TimedOut => Self::TimedOut, + EnvironmentError::TaskFailed => Self::TaskFailed, + EnvironmentError::SourceUnavailable(_) => Self::SourceUnavailable, + EnvironmentError::UnsupportedVersion | EnvironmentError::UnsupportedCapability | EnvironmentError::InvalidTimeout => { + Self::TaskFailed + } + } + } +} + /// The bounded entries plus the capture time of their persisted L0 source. #[derive(Debug, PartialEq)] pub struct OfflineDiagnostics { @@ -168,79 +180,6 @@ pub struct OfflineDiagnostics { pub inventory_age: Duration, } -#[derive(Debug)] -struct SystemSnapshot { - os_summary: String, - kernel_summary: String, - architecture: &'static str, - cores: usize, - total_memory_bytes: u64, - under_memory_pressure: bool, - filesystem_types: Vec, - interface_count: usize, - bond_count: usize, -} - -impl SystemSnapshot { - fn collect() -> Self { - #[cfg(test)] - let _scan = test_support::ScanGuard::start(); - - // Do not enumerate processes: Q07 allows only CPU and memory summaries. - let system = System::new_with_specifics(RefreshKind::everything().without_processes()); - let total_memory_bytes = system.total_memory(); - let available_memory = system.available_memory(); - let filesystem_types = Disks::new_with_refreshed_list() - .iter() - .map(|disk| disk.file_system().to_string_lossy().into_owned()) - .collect::>() - .into_iter() - .collect(); - let networks = Networks::new_with_refreshed_list(); - Self { - os_summary: System::long_os_version().unwrap_or_else(|| "unknown".to_owned()), - kernel_summary: System::kernel_long_version(), - architecture: std::env::consts::ARCH, - cores: system.cpus().len(), - total_memory_bytes, - under_memory_pressure: total_memory_bytes != 0 && available_memory.saturating_mul(10) < total_memory_bytes, - filesystem_types, - interface_count: networks.len(), - bond_count: networks.keys().filter(|name| name.starts_with("bond")).count(), - } - } -} - -async fn collect_system_snapshot(cancel: &CancellationToken) -> Result { - let deadline = tokio::time::Instant::now() + COLLECT_TIMEOUT; - let permit = tokio::select! { - biased; - () = cancel.cancelled() => return Err(CollectorError::Cancelled), - result = tokio::time::timeout_at(deadline, SYSTEM_SCAN_PERMIT.clone().acquire_owned()) => { - match result { - Ok(Ok(permit)) => permit, - Ok(Err(_)) => return Err(CollectorError::TaskFailed), - Err(_) => return Err(CollectorError::TimedOut), - } - } - }; - let task = tokio::task::spawn_blocking(move || { - let _permit = permit; - SystemSnapshot::collect() - }); - tokio::select! { - biased; - () = cancel.cancelled() => Err(CollectorError::Cancelled), - result = tokio::time::timeout_at(deadline, task) => { - match result { - Ok(Ok(snapshot)) => Ok(snapshot), - Ok(Err(_)) => Err(CollectorError::TaskFailed), - Err(_) => Err(CollectorError::TimedOut), - } - } - } -} - /// Collect all and only the Q07 offline L0/L1 fields after acquiring the /// stopped-runtime inventory lock. pub async fn collect_offline_diagnostics( @@ -253,7 +192,7 @@ pub async fn collect_offline_diagnostics( let store = InventoryStateStore::from_state_root(state_root)?; let _lock = store.try_runtime_lock()?; let persisted = store.read_latest(chrono::Utc::now())?; - let system = collect_system_snapshot(cancel).await?; + let system = collect_host_environment(COLLECT_TIMEOUT, cancel).await?; let mut entries = Vec::with_capacity(COLLECTORS.len()); for collector in COLLECTORS { @@ -269,102 +208,3 @@ pub async fn collect_offline_diagnostics( inventory_age: persisted.age, }) } - -#[cfg(test)] -mod test_support { - use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering}; - use std::time::Duration; - - pub(super) static DELAY_MILLIS: AtomicU64 = AtomicU64::new(0); - pub(super) static ACTIVE: AtomicUsize = AtomicUsize::new(0); - pub(super) static MAX_ACTIVE: AtomicUsize = AtomicUsize::new(0); - - pub(super) struct ScanGuard; - - impl ScanGuard { - pub(super) fn start() -> Self { - let active = ACTIVE.fetch_add(1, Ordering::SeqCst) + 1; - MAX_ACTIVE.fetch_max(active, Ordering::SeqCst); - let delay = DELAY_MILLIS.load(Ordering::SeqCst); - if delay != 0 { - std::thread::sleep(Duration::from_millis(delay)); - } - Self - } - } - - impl Drop for ScanGuard { - fn drop(&mut self) { - ACTIVE.fetch_sub(1, Ordering::SeqCst); - } - } -} - -#[cfg(test)] -mod tests { - use std::sync::atomic::Ordering; - - use super::*; - - async fn wait_for_active(expected: usize) { - tokio::time::timeout(Duration::from_secs(1), async { - while test_support::ACTIVE.load(Ordering::SeqCst) != expected { - tokio::time::sleep(Duration::from_millis(5)).await; - } - }) - .await - .expect("system scan reaches expected state"); - } - - #[tokio::test] - async fn connect_offline_collectors_timeout_and_cancel_never_overlap_system_scans() { - test_support::MAX_ACTIVE.store(0, Ordering::SeqCst); - test_support::DELAY_MILLIS.store((COLLECT_TIMEOUT + Duration::from_millis(200)).as_millis() as u64, Ordering::SeqCst); - - let first_cancel = CancellationToken::new(); - assert!(matches!(collect_system_snapshot(&first_cancel).await, Err(CollectorError::TimedOut))); - assert_eq!(test_support::ACTIVE.load(Ordering::SeqCst), 1, "timed-out blocking scan remains active"); - - let second_cancel = CancellationToken::new(); - let second = tokio::spawn({ - let second_cancel = second_cancel.clone(); - async move { collect_system_snapshot(&second_cancel).await } - }); - tokio::time::sleep(Duration::from_millis(50)).await; - assert_eq!( - test_support::MAX_ACTIVE.load(Ordering::SeqCst), - 1, - "a timed-out scan keeps the single-flight permit" - ); - second_cancel.cancel(); - assert!(matches!(second.await.expect("second collector task"), Err(CollectorError::Cancelled))); - wait_for_active(0).await; - - test_support::MAX_ACTIVE.store(0, Ordering::SeqCst); - test_support::DELAY_MILLIS.store(250, Ordering::SeqCst); - let third_cancel = CancellationToken::new(); - let third = tokio::spawn({ - let third_cancel = third_cancel.clone(); - async move { collect_system_snapshot(&third_cancel).await } - }); - wait_for_active(1).await; - third_cancel.cancel(); - assert!(matches!(third.await.expect("third collector task"), Err(CollectorError::Cancelled))); - - let fourth_cancel = CancellationToken::new(); - let fourth = tokio::spawn({ - let fourth_cancel = fourth_cancel.clone(); - async move { collect_system_snapshot(&fourth_cancel).await } - }); - tokio::time::sleep(Duration::from_millis(50)).await; - assert_eq!( - test_support::MAX_ACTIVE.load(Ordering::SeqCst), - 1, - "a cancelled scan keeps the single-flight permit" - ); - fourth_cancel.cancel(); - assert!(matches!(fourth.await.expect("fourth collector task"), Err(CollectorError::Cancelled))); - wait_for_active(0).await; - test_support::DELAY_MILLIS.store(0, Ordering::SeqCst); - } -} diff --git a/rustfs/tests/connect_environment.rs b/rustfs/tests/connect_environment.rs new file mode 100644 index 000000000..2a2751f1f --- /dev/null +++ b/rustfs/tests/connect_environment.rs @@ -0,0 +1,101 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::time::Duration; + +use rustfs::connect::{ + ENVIRONMENT_CAPABILITY, ENVIRONMENT_SCHEMA_VERSION, EnvironmentCollectionRequest, EnvironmentError, + EnvironmentFilesystemType, InventorySnapshot, MAX_ENVIRONMENT_DURATION, collect_environment, +}; +use serde_json::Value; +use tokio_util::sync::CancellationToken; + +fn inventory() -> InventorySnapshot { + InventorySnapshot::current(2, 8, 8_000_000, 2_000_000, []).expect("known small deployment inventory") +} + +fn request() -> EnvironmentCollectionRequest { + EnvironmentCollectionRequest::negotiate(ENVIRONMENT_SCHEMA_VERSION, ENVIRONMENT_CAPABILITY, Duration::from_secs(2)) + .expect("supported inventory.environment request") +} + +#[test] +fn environment_negotiation_rejects_old_versions_unknown_capabilities_and_invalid_budgets() { + assert_eq!( + EnvironmentCollectionRequest::negotiate(0, ENVIRONMENT_CAPABILITY, Duration::from_secs(1)), + Err(EnvironmentError::UnsupportedVersion) + ); + assert_eq!( + EnvironmentCollectionRequest::negotiate(ENVIRONMENT_SCHEMA_VERSION, "inventory.environment@2", Duration::from_secs(1)), + Err(EnvironmentError::UnsupportedCapability) + ); + for timeout in [Duration::ZERO, MAX_ENVIRONMENT_DURATION + Duration::from_millis(1)] { + assert_eq!( + EnvironmentCollectionRequest::negotiate(ENVIRONMENT_SCHEMA_VERSION, ENVIRONMENT_CAPABILITY, timeout), + Err(EnvironmentError::InvalidTimeout) + ); + } +} + +#[tokio::test] +async fn environment_collection_acknowledges_cancellation_before_sampling() { + let cancel = CancellationToken::new(); + cancel.cancel(); + let result = collect_environment(&inventory(), request(), &cancel).await; + assert_eq!(result, Err(EnvironmentError::Cancelled)); +} + +#[tokio::test] +async fn environment_collection_emits_only_the_closed_identifier_free_schema() { + let result = collect_environment(&inventory(), request(), &CancellationToken::new()) + .await + .expect("this host exposes bounded environment inventory"); + assert_eq!(result.node_count(), 2); + assert_eq!(result.drive_count(), 8); + assert!(!result.filesystem_types().is_empty()); + assert!(result.filesystem_types().len() <= 5); + assert!(result.filesystem_types().iter().all(|value| matches!( + value, + EnvironmentFilesystemType::Ext4 + | EnvironmentFilesystemType::Xfs + | EnvironmentFilesystemType::Zfs + | EnvironmentFilesystemType::Apfs + | EnvironmentFilesystemType::Other + ))); + + let serialized = serde_json::to_value(&result).expect("environment JSON"); + let object = serialized.as_object().expect("environment object"); + assert_eq!( + object.keys().map(String::as_str).collect::>(), + ["driveCount", "filesystemTypes", "nodeCount", "osFamily"] + .into_iter() + .collect() + ); + let encoded = serde_json::to_string(&serialized).expect("environment JSON text"); + for forbidden in [ + "hostname", + "mountPath", + "mountOptions", + "device", + "serial", + "credential", + "AWS_SECRET_ACCESS_KEY", + "/srv/customer-a", + ] { + assert!(!encoded.contains(forbidden), "environment output exposed {forbidden}"); + } + + println!("inventory.environment actual output: {encoded}"); + assert!(matches!(serialized, Value::Object(_))); +}