diff --git a/rustfs/src/config/cli.rs b/rustfs/src/config/cli.rs index d616f8196..742b08b23 100644 --- a/rustfs/src/config/cli.rs +++ b/rustfs/src/config/cli.rs @@ -379,7 +379,7 @@ pub struct ConnectTopOpts { pub enum ConnectTopCommands { /// Capture API activity when an approved typed source is available Api(ConnectTopCaptureOpts), - /// Capture process disk I/O on supported platforms + /// Capture running service process disk I/O on supported platforms Disk(ConnectTopCaptureOpts), /// Capture lock activity when an approved typed source is available Locks(ConnectTopCaptureOpts), @@ -391,6 +391,9 @@ pub enum ConnectTopCommands { #[derive(Args, Clone)] pub struct ConnectTopCaptureOpts { + /// Explicit existing offline identity pin, required for service disk capture + #[arg(long = "offline-key-id", value_parser = NonEmptyStringValueParser::new())] + pub offline_key_id: Option, /// Directory containing an enrolled Connect device identity #[arg(long = "state-dir")] pub state_dir: PathBuf, diff --git a/rustfs/src/connect/diagnostics/mod.rs b/rustfs/src/connect/diagnostics/mod.rs index 81ace7742..f55370bea 100644 --- a/rustfs/src/connect/diagnostics/mod.rs +++ b/rustfs/src/connect/diagnostics/mod.rs @@ -172,6 +172,9 @@ pub use trace_record::{ }; pub use trace_replay::{LocallyReviewedTraceArtifact, ReplayedTrace, TraceReplayError, replay_trace, replay_trace_result}; pub(crate) use trace_runtime::{ - LocalTraceCaptureError, LocalTraceCaptureRuntime, request_local_runtime_profile, request_local_trace_capture, - spawn_local_trace_capture_runtime, + LocalTraceCaptureError, LocalTraceCaptureRuntime, request_local_runtime_profile, request_local_top_disk, + request_local_trace_capture, spawn_local_trace_capture_runtime, }; + +pub(crate) use top_api::save_top_archive; +pub(crate) use top_disk::LocalTopDiskRequest; diff --git a/rustfs/src/connect/diagnostics/top_api.rs b/rustfs/src/connect/diagnostics/top_api.rs index c89017e39..b5b7ae82f 100644 --- a/rustfs/src/connect/diagnostics/top_api.rs +++ b/rustfs/src/connect/diagnostics/top_api.rs @@ -58,7 +58,7 @@ const PROTOCOL_VERSION: &str = "v1"; const SIGNATURE_ALGORITHM: &str = "ES256"; const SIGNATURE_DOMAIN: &[u8] = b"rustfs-diagnostic-envelope-v1"; const MAX_ENVELOPE_BYTES: usize = 16_384; -const MAX_ARCHIVE_BYTES: usize = 524_288; +pub(crate) const MAX_ARCHIVE_BYTES: usize = 524_288; const MAX_DECOMPRESSED_BYTES: usize = 278_528; pub const MAX_TOP_EXPORT_VALIDITY: Duration = Duration::from_secs(2_592_000); const TOP_CAPTURE_WORKING_SET_BYTES: u64 = 1_048_576; @@ -635,16 +635,26 @@ pub fn save_signed_top_export( output: &Path, export: &SignedTopExport, cancel: &CancellationToken, +) -> Result { + save_top_archive(output, &export.artifact_uid, &export.archive_bytes, &export.archive_sha256, cancel) +} + +pub(crate) fn save_top_archive( + output: &Path, + artifact_uid: &str, + archive_bytes: &[u8], + archive_sha256: &str, + cancel: &CancellationToken, ) -> Result { if cancel.is_cancelled() { return Err(TopCaptureError::Cancelled); } - if !is_uuid_v7(&export.artifact_uid) { + if !is_uuid_v7(artifact_uid) { return Err(TopCaptureError::Scope); } - if export.archive_bytes.is_empty() - || export.archive_bytes.len() > MAX_ARCHIVE_BYTES - || hex_lower(&Sha256::digest(&export.archive_bytes)) != export.archive_sha256 + if archive_bytes.is_empty() + || archive_bytes.len() > MAX_ARCHIVE_BYTES + || hex_lower(&Sha256::digest(archive_bytes)) != archive_sha256 { return Err(TopCaptureError::Result); } @@ -653,7 +663,7 @@ pub fn save_signed_top_export( .filter(|path| !path.as_os_str().is_empty()) .unwrap_or_else(|| Path::new(".")); let filename = output.file_name().ok_or(TopCaptureError::Scope)?.to_string_lossy(); - let temporary = parent.join(format!(".{filename}.{}.partial", export.artifact_uid)); + let temporary = parent.join(format!(".{filename}.{}.partial", artifact_uid)); let mut options = OpenOptions::new(); options.write(true).create_new(true); #[cfg(unix)] @@ -663,7 +673,7 @@ pub fn save_signed_top_export( } let mut file = options.open(&temporary).map_err(map_create_error)?; let result = (|| { - file.write_all(&export.archive_bytes).map_err(io_error)?; + file.write_all(archive_bytes).map_err(io_error)?; if cancel.is_cancelled() { return Err(TopCaptureError::Cancelled); } @@ -680,9 +690,9 @@ pub fn save_signed_top_export( return Err(TopCaptureError::DurabilityAfterCommit(error.kind())); } Ok(SavedTopExport { - artifact_uid: export.artifact_uid.clone(), - archive_size_bytes: export.archive_bytes.len() as u64, - archive_sha256: export.archive_sha256.clone(), + artifact_uid: artifact_uid.to_owned(), + archive_size_bytes: archive_bytes.len() as u64, + archive_sha256: archive_sha256.to_owned(), }) })(); if result.is_err() { @@ -1003,6 +1013,29 @@ mod tests { use rustfs_s3_ops::S3Operation; use serial_test::serial; + #[cfg(target_os = "linux")] + #[tokio::test] + #[serial] + async fn queued_disk_capture_rechecks_expiry_before_sampling() { + let mut request = capture_request(Duration::from_secs(2)); + request.scope.consent.tool_id = "top.disk".to_owned(); + request.scope.run_expires_at_unix = OffsetDateTime::now_utc().unix_timestamp() + 3; + let cancel = CancellationToken::new(); + let permit = request.acquire(&cancel).await.unwrap().unwrap(); + let capture = super::super::top_disk::capture_top_disk(&request, &cancel); + tokio::pin!(capture); + assert!(tokio::time::timeout(Duration::from_millis(20), &mut capture).await.is_err()); + while OffsetDateTime::now_utc().unix_timestamp() < request.scope.run_expires_at_unix { + tokio::time::sleep(Duration::from_millis(20)).await; + } + drop(permit); + // The old path sampled first and waited the entire two-second window. + let result = tokio::time::timeout(Duration::from_millis(500), capture) + .await + .expect("expired queued capture must not start its sampling window"); + assert!(matches!(result, Err(TopCaptureError::Expired))); + } + fn capture_request(window: Duration) -> TopCaptureRequest { let organization_uid = Uuid::now_v7(); let cluster_uid = Uuid::now_v7(); diff --git a/rustfs/src/connect/diagnostics/top_disk.rs b/rustfs/src/connect/diagnostics/top_disk.rs index a506f85f2..24faa17ed 100644 --- a/rustfs/src/connect/diagnostics/top_disk.rs +++ b/rustfs/src/connect/diagnostics/top_disk.rs @@ -14,7 +14,7 @@ //! Bounded process disk-I/O window backed by RustFS's existing process sampler. -use serde::Serialize; +use serde::{Deserialize, Serialize}; #[cfg(target_os = "linux")] use tokio::time::Instant; use tokio_util::sync::CancellationToken; @@ -24,6 +24,32 @@ use super::top_api::{MAX_SAFE_INTEGER, TopCaptureError, TopCaptureRequest, TopRe const TOOL_ID: &str = "top.disk"; pub const TOP_DISK_CAPABILITY: &str = "top.disk@1"; +/// Closed local-service request: no process selector, paths, or supplied provenance. +#[derive(Clone, Debug, Deserialize, Serialize)] +#[serde(deny_unknown_fields, rename_all = "camelCase")] +pub struct LocalTopDiskRequest { + pub offline_key_id: String, + pub organization_name: String, + pub cluster_name: String, + pub device_name: String, + pub run_uid: String, + pub artifact_uid: String, + pub consent_uid: String, + pub policy_revision: u64, + pub consent_expires_at_unix: i64, + pub acknowledge_l3: bool, + pub run_expires_at_unix: i64, + pub window_millis: u64, + pub export_validity_seconds: u64, +} + +/// Archive received from the owner-only service socket. +pub(crate) struct LocalTopDiskArchive { + pub artifact_uid: String, + pub archive_bytes: Vec, + pub archive_sha256: String, +} + #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub struct DiskCounterSnapshot { pub read_bytes: u64, @@ -58,6 +84,11 @@ pub async fn capture_top_disk( let Some(_permit) = request.acquire(cancel).await? else { return request.cancelled(TOOL_ID); }; + // A queued request may outlive its authorization while another capture owns the lease. + request.validate_capture(TOOL_ID)?; + if cancel.is_cancelled() { + return request.cancelled(TOOL_ID); + } let mut sampler = rustfs_io_metrics::ProcessSampler::new(); let before = process_snapshot(&mut sampler)?; let started = Instant::now(); diff --git a/rustfs/src/connect/diagnostics/trace_runtime.rs b/rustfs/src/connect/diagnostics/trace_runtime.rs index 664ce83f3..1f07a2162 100644 --- a/rustfs/src/connect/diagnostics/trace_runtime.rs +++ b/rustfs/src/connect/diagnostics/trace_runtime.rs @@ -14,7 +14,7 @@ //! Owner-only local transport between diagnostic CLI commands and the running server. //! -//! Version 1 accepts only TRACE_RECORD and RUNTIME_PROFILE. Runtime requests select +//! Version 1 accepts only TRACE_RECORD, RUNTIME_PROFILE and TOP_DISK. Signed requests select //! an existing offline key by SPKI digest; this is not proof of Connect enrollment. //! The receiver checks enrollment, target ownership and consent at import. The //! server owns provenance, nonce generation, capture and signing; the CLI receives @@ -66,7 +66,9 @@ pub(crate) enum LocalTraceCaptureError { Producer(#[from] TelemetryProducerError), #[error("runtime profile request rejected: {0}")] RuntimeProfile(String), - #[error("runtime profile cancellation was not acknowledged")] + #[error("top disk request rejected: {0}")] + TopDisk(String), + #[error("local diagnostic cancellation was not acknowledged")] CancellationUnconfirmed, } @@ -108,6 +110,10 @@ enum CaptureRequest { protocol_version: u16, request: LocalRuntimeProfileRequest, }, + TopDisk { + protocol_version: u16, + request: super::top_disk::LocalTopDiskRequest, + }, } #[derive(Clone, Copy, Debug, Deserialize, Serialize)] @@ -175,6 +181,14 @@ enum CaptureResponse { RuntimeError { code: RuntimeErrorCode, }, + TopDiskOk { + archive_base64: String, + archive_sha256: String, + artifact_uid: String, + }, + TopDiskError { + code: RuntimeErrorCode, + }, } pub(crate) fn spawn_local_trace_capture_runtime( @@ -265,6 +279,7 @@ enum RuntimeErrorCode { Busy, Cancelled, TimedOut, + UnsupportedPlatform, CollectionFailed, } @@ -371,6 +386,88 @@ pub(crate) async fn request_local_runtime_profile( } } +pub(crate) async fn request_local_top_disk( + state_root: &Path, + request: super::top_disk::LocalTopDiskRequest, + cancel: &CancellationToken, +) -> Result { + let owner = private_state_owner(state_root)?; + let socket_path = state_root.join(SOCKET_FILE); + socket_identity(&socket_path, owner)?; + let stream = UnixStream::connect(&socket_path).await.map_err(LocalTraceCaptureError::Io)?; + if !stream.peer_cred().is_ok_and(|credentials| credentials.uid() == owner) { + return Err(LocalTraceCaptureError::StateSecurity); + } + let artifact_uid = request.artifact_uid.clone(); + let message = CaptureRequest::TopDisk { + protocol_version: PROTOCOL_VERSION, + request, + }; + let mut bytes = serde_json::to_vec(&message).map_err(|_| LocalTraceCaptureError::Protocol)?; + bytes.push(b'\n'); + if bytes.len() as u64 > MAX_REQUEST_BYTES { + return Err(LocalTraceCaptureError::Protocol); + } + let (reader, mut writer) = stream.into_split(); + tokio::time::timeout(REQUEST_TIMEOUT, writer.write_all(&bytes)) + .await + .map_err(|_| LocalTraceCaptureError::Protocol)? + .map_err(LocalTraceCaptureError::Io)?; + // Keep the write half open: EOF is the server's cancellation signal. + let response = async { + let mut bytes = Vec::new(); + reader + .take(MAX_RUNTIME_RESPONSE_BYTES + 1) + .read_to_end(&mut bytes) + .await + .map_err(LocalTraceCaptureError::Io)?; + if bytes.is_empty() || bytes.len() as u64 > MAX_RUNTIME_RESPONSE_BYTES { + return Err(LocalTraceCaptureError::Protocol); + } + serde_json::from_slice::(&bytes).map_err(|_| LocalTraceCaptureError::Protocol) + }; + let response = tokio::time::timeout(MAX_PROFILE_DURATION + Duration::from_secs(5), response); + tokio::pin!(response); + let response = tokio::select! { + biased; + _ = cancel.cancelled() => { + writer.shutdown().await.map_err(|_| LocalTraceCaptureError::CancellationUnconfirmed)?; + // The service closes the response only after the blocking collector has joined. + let acknowledged = matches!(tokio::time::timeout(Duration::from_secs(2), &mut response).await, + Ok(Ok(Ok(CaptureResponse::TopDiskError { .. } | CaptureResponse::TopDiskOk { .. })))); + if !acknowledged { return Err(LocalTraceCaptureError::CancellationUnconfirmed); } + return Err(LocalTraceCaptureError::TopDisk("CANCELLED".to_owned())); + } + response = &mut response => response.map_err(|_| LocalTraceCaptureError::Protocol)??, + }; + match response { + CaptureResponse::TopDiskOk { + archive_base64, + archive_sha256, + artifact_uid: returned_uid, + } => { + let archive_bytes = URL_SAFE_NO_PAD + .decode_to_vec(&archive_base64) + .map_err(|_| LocalTraceCaptureError::Protocol)?; + if archive_bytes.is_empty() + || archive_bytes.len() > super::top_api::MAX_ARCHIVE_BYTES + || returned_uid != artifact_uid + || URL_SAFE_NO_PAD.encode_to_string(&archive_bytes) != archive_base64 + || hex_simd::encode_to_string(Sha256::digest(&archive_bytes), hex_simd::AsciiCase::Lower) != archive_sha256 + { + return Err(LocalTraceCaptureError::Protocol); + } + Ok(super::top_disk::LocalTopDiskArchive { + artifact_uid, + archive_bytes, + archive_sha256, + }) + } + CaptureResponse::TopDiskError { code } => Err(LocalTraceCaptureError::TopDisk(format!("{code:?}"))), + _ => Err(LocalTraceCaptureError::Protocol), + } +} + async fn handle_runtime_profile( mut reader: BufReader, mut writer: tokio::net::unix::OwnedWriteHalf, @@ -408,6 +505,182 @@ async fn handle_runtime_profile( } } +async fn handle_top_disk( + mut reader: BufReader, + mut writer: tokio::net::unix::OwnedWriteHalf, + state_root: &Path, + protocol_version: u16, + request: super::top_disk::LocalTopDiskRequest, + shutdown: CancellationToken, +) { + let cancel = shutdown.child_token(); + let capture = capture_local_top_disk(state_root, protocol_version, request, &cancel); + tokio::pin!(capture); + let mut unexpected = [0_u8; 1]; + let result = tokio::select! { + biased; + _ = shutdown.cancelled() => { cancel.cancel(); let _ = capture.await; return; } + _ = reader.read(&mut unexpected) => { cancel.cancel(); let _ = capture.await; Err(RuntimeErrorCode::Cancelled) } + result = &mut capture => result, + }; + let response = match result { + Ok(export) => CaptureResponse::TopDiskOk { + archive_base64: URL_SAFE_NO_PAD.encode_to_string(&export.archive_bytes), + archive_sha256: export.archive_sha256, + artifact_uid: export.artifact_uid, + }, + Err(code) => CaptureResponse::TopDiskError { code }, + }; + if let Ok(bytes) = serde_json::to_vec(&response) + && bytes.len() as u64 <= MAX_RUNTIME_RESPONSE_BYTES + { + let _ = tokio::time::timeout(REQUEST_TIMEOUT, async { + writer.write_all(&bytes).await?; + writer.shutdown().await + }) + .await; + } +} + +async fn capture_local_top_disk( + state_root: &Path, + protocol_version: u16, + input: super::top_disk::LocalTopDiskRequest, + cancel: &CancellationToken, +) -> Result { + use super::top_api::{LocalTopConsent, TOP_CLASSIFICATION, TopCaptureLimits, TopCaptureRequest, TopCaptureScope}; + if protocol_version != PROTOCOL_VERSION + || input.offline_key_id.len() != 64 + || !input + .offline_key_id + .bytes() + .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b)) + { + return Err(RuntimeErrorCode::InvalidRequest); + } + if cancel.is_cancelled() { + return Err(RuntimeErrorCode::Cancelled); + } + // Reject unconsented, expired and unbounded work before reading a key or executable. + if !input.acknowledge_l3 || input.policy_revision == 0 { + return Err(RuntimeErrorCode::ConsentRequired); + } + let now = unix_now().map_err(|_| RuntimeErrorCode::CollectionFailed)?; + if input.consent_expires_at_unix <= now { + return Err(RuntimeErrorCode::ConsentExpired); + } + if input.run_expires_at_unix <= now + || input.window_millis.div_ceil(1000) + > input + .run_expires_at_unix + .min(input.consent_expires_at_unix) + .saturating_sub(now) as u64 + { + return Err(RuntimeErrorCode::Expired); + } + if input.window_millis == 0 + || input.window_millis > 30_000 + || input.export_validity_seconds == 0 + || input.export_validity_seconds > super::top_api::MAX_TOP_EXPORT_VALIDITY.as_secs() + { + return Err(RuntimeErrorCode::LimitExceeded); + } + let key = load_offline_key(state_root, &input.offline_key_id)?; + let provenance = super::job_delivery::executable_provenance() + .await + .map_err(|_| RuntimeErrorCode::CollectionFailed)?; + let request = TopCaptureRequest { + scope: TopCaptureScope { + organization_name: input.organization_name, + cluster_name: input.cluster_name, + device_name: input.device_name, + run_uid: input.run_uid, + artifact_uid: input.artifact_uid, + policy_revision: input.policy_revision, + run_expires_at_unix: input.run_expires_at_unix, + executable_sha256: provenance.executable_sha256().to_owned(), + build_features: provenance.build_features().to_vec(), + consent: LocalTopConsent { + uid: input.consent_uid, + tool_id: "top.disk".to_owned(), + classification: TOP_CLASSIFICATION.to_owned(), + active: input.acknowledge_l3, + expires_at_unix: input.consent_expires_at_unix, + }, + }, + limits: TopCaptureLimits::default(), + window: Duration::from_millis(input.window_millis), + export_validity: Duration::from_secs(input.export_validity_seconds), + }; + let result = super::top_disk::capture_top_disk(&request, cancel).await.map_err(top_error)?; + if result.outcome == super::top_api::TopOutcome::Cancelled { + return Err(RuntimeErrorCode::Cancelled); + } + if result.outcome == super::top_api::TopOutcome::Unsupported { + return Err(RuntimeErrorCode::UnsupportedPlatform); + } + if result.outcome != super::top_api::TopOutcome::Succeeded { + return Err(RuntimeErrorCode::CollectionFailed); + } + super::top_api::sign_top_export(&request, &result, &key, cancel).map_err(top_error) +} + +fn top_error(error: super::top_api::TopCaptureError) -> RuntimeErrorCode { + use super::top_api::TopCaptureError; + match error { + TopCaptureError::ConsentRequired => RuntimeErrorCode::ConsentRequired, + TopCaptureError::ConsentExpired => RuntimeErrorCode::ConsentExpired, + TopCaptureError::Expired => RuntimeErrorCode::Expired, + TopCaptureError::Limits | TopCaptureError::ResultTooLarge => RuntimeErrorCode::LimitExceeded, + TopCaptureError::Cancelled => RuntimeErrorCode::Cancelled, + TopCaptureError::Scope | TopCaptureError::ConsentScope => RuntimeErrorCode::InvalidRequest, + _ => RuntimeErrorCode::CollectionFailed, + } +} + +fn load_offline_key(state_root: &Path, offline_key_id: &str) -> Result { + let owner = private_state_owner(state_root).map_err(|_| RuntimeErrorCode::IdentityUnavailable)?; + let key_root = state_root.join("offline"); + let key_directory = std::fs::symlink_metadata(&key_root).map_err(|_| RuntimeErrorCode::IdentityUnavailable)?; + // Existing offline enrollment creates this subdirectory with the process umask. + // The state root is 0700; require the child to remain owned and non-writable by others. + if !key_directory.is_dir() + || key_directory.file_type().is_symlink() + || key_directory.uid() != owner + || key_directory.permissions().mode() & 0o022 != 0 + { + return Err(RuntimeErrorCode::IdentityUnavailable); + } + // IdentityStore::load follows symlinks; this IPC boundary must reject them before reading. + use std::os::unix::fs::OpenOptionsExt as _; + let path = crate::connect::OfflineKeyStore::new(state_root).key_path(); + let file = std::fs::OpenOptions::new() + .read(true) + .custom_flags(libc::O_NOFOLLOW | libc::O_CLOEXEC | libc::O_NONBLOCK) + .open(path) + .map_err(|_| RuntimeErrorCode::IdentityUnavailable)?; + let metadata = file.metadata().map_err(|_| RuntimeErrorCode::IdentityUnavailable)?; + if !metadata.is_file() + || metadata.uid() != owner + || metadata.permissions().mode() & 0o7777 != 0o600 + || metadata.len() == 0 + || metadata.len() > 4_096 + { + return Err(RuntimeErrorCode::IdentityUnavailable); + } + let mut der = zeroize::Zeroizing::new(Vec::new()); + std::io::Read::read_to_end(&mut std::io::Read::take(file, 4_097), &mut der) + .map_err(|_| RuntimeErrorCode::IdentityUnavailable)?; + if der.len() > 4_096 { + return Err(RuntimeErrorCode::IdentityUnavailable); + } + let key = crate::connect::DeviceIdentity::from_pkcs8_der(&der).map_err(|_| RuntimeErrorCode::IdentityUnavailable)?; + if hex_simd::encode_to_string(Sha256::digest(key.public_key_der()), hex_simd::AsciiCase::Lower) != offline_key_id { + return Err(RuntimeErrorCode::IdentityUnavailable); + } + Ok(key) +} + async fn capture_local_runtime_profile( state_root: &Path, protocol_version: u16, @@ -447,45 +720,7 @@ async fn capture_local_runtime_profile( if input.schema_version != 1 || input.capability != super::profile_cpu::THREAD_PROFILE_CAPABILITY { return Err(RuntimeErrorCode::InvalidRequest); } - let owner = private_state_owner(state_root).map_err(|_| RuntimeErrorCode::IdentityUnavailable)?; - let key_root = state_root.join("offline"); - let key_directory = std::fs::symlink_metadata(&key_root).map_err(|_| RuntimeErrorCode::IdentityUnavailable)?; - // Existing offline enrollment creates this subdirectory with the process umask. - // The state root is 0700; require the child to remain owned and non-writable by others. - if !key_directory.is_dir() - || key_directory.file_type().is_symlink() - || key_directory.uid() != owner - || key_directory.permissions().mode() & 0o022 != 0 - { - return Err(RuntimeErrorCode::IdentityUnavailable); - } - // IdentityStore::load follows symlinks; this IPC boundary must reject them before reading. - use std::os::unix::fs::OpenOptionsExt as _; - let path = crate::connect::OfflineKeyStore::new(state_root).key_path(); - let file = std::fs::OpenOptions::new() - .read(true) - .custom_flags(libc::O_NOFOLLOW | libc::O_CLOEXEC | libc::O_NONBLOCK) - .open(path) - .map_err(|_| RuntimeErrorCode::IdentityUnavailable)?; - let metadata = file.metadata().map_err(|_| RuntimeErrorCode::IdentityUnavailable)?; - if !metadata.is_file() - || metadata.uid() != owner - || metadata.permissions().mode() & 0o7777 != 0o600 - || metadata.len() == 0 - || metadata.len() > 4_096 - { - return Err(RuntimeErrorCode::IdentityUnavailable); - } - let mut der = zeroize::Zeroizing::new(Vec::new()); - std::io::Read::read_to_end(&mut std::io::Read::take(file, 4_097), &mut der) - .map_err(|_| RuntimeErrorCode::IdentityUnavailable)?; - if der.len() > 4_096 { - return Err(RuntimeErrorCode::IdentityUnavailable); - } - let key = crate::connect::DeviceIdentity::from_pkcs8_der(&der).map_err(|_| RuntimeErrorCode::IdentityUnavailable)?; - if hex_simd::encode_to_string(Sha256::digest(key.public_key_der()), hex_simd::AsciiCase::Lower) != input.offline_key_id { - return Err(RuntimeErrorCode::IdentityUnavailable); - } + let key = load_offline_key(state_root, &input.offline_key_id)?; let provenance = super::job_delivery::executable_provenance() .await .map_err(|_| RuntimeErrorCode::CollectionFailed)?; @@ -570,6 +805,13 @@ async fn handle_connection(stream: UnixStream, state_root: PathBuf, shutdown: Ca handle_runtime_profile(reader, writer, &state_root, protocol_version, request, shutdown).await; return; } + CaptureRequest::TopDisk { + protocol_version, + request, + } => { + handle_top_disk(reader, writer, &state_root, protocol_version, request, shutdown).await; + return; + } CaptureRequest::TraceRecord { protocol_version, consent_expires_at_unix, @@ -901,6 +1143,138 @@ mod tests { ) } + fn disk_request(state: &std::path::Path) -> super::super::top_disk::LocalTopDiskRequest { + let (request, _) = runtime_request(state); + super::super::top_disk::LocalTopDiskRequest { + offline_key_id: request.offline_key_id, + organization_name: request.organization_name, + cluster_name: request.cluster_name, + device_name: request.device_name, + run_uid: request.run_uid, + artifact_uid: request.artifact_uid, + consent_uid: request.consent_uid, + policy_revision: request.policy_revision, + consent_expires_at_unix: request.consent_expires_at_unix, + acknowledge_l3: true, + run_expires_at_unix: request.expires_at_unix, + window_millis: 200, + export_validity_seconds: 30, + } + } + + #[tokio::test] + async fn local_top_disk_rejects_invalid_requests_without_identity_or_socket_fallback() { + let state = tempfile::tempdir().unwrap(); + std::fs::set_permissions(state.path(), std::fs::Permissions::from_mode(0o700)).unwrap(); + let mut request = disk_request(state.path()); + std::fs::remove_file(crate::connect::OfflineKeyStore::new(state.path()).key_path()).unwrap(); + request.acknowledge_l3 = false; + assert!(matches!( + super::capture_local_top_disk(state.path(), 1, request.clone(), &CancellationToken::new()).await, + Err(super::RuntimeErrorCode::ConsentRequired) + )); + request.acknowledge_l3 = true; + request.window_millis = 30_001; + assert!(matches!( + super::capture_local_top_disk(state.path(), 1, request.clone(), &CancellationToken::new()).await, + Err(super::RuntimeErrorCode::LimitExceeded) + )); + request.window_millis = 200; + assert!(matches!( + super::capture_local_top_disk(state.path(), 1, request.clone(), &CancellationToken::new()).await, + Err(super::RuntimeErrorCode::IdentityUnavailable) + )); + assert!( + super::request_local_top_disk(state.path(), request.clone(), &CancellationToken::new()) + .await + .is_err() + ); + let mut value = serde_json::to_value(&request).unwrap(); + value["pid"] = serde_json::json!(1); + assert!(serde_json::from_value::(value).is_err()); + assert!( + serde_json::from_str::( + r#"{"operation":"TOP_DISK","protocolVersion":1,"protocolVersion":1,"request":{}}"# + ) + .is_err() + ); + } + + #[cfg(target_os = "linux")] + #[test] + fn local_top_disk_service_child() { + use std::io::{Read as _, Write as _}; + let Some(state) = std::env::var_os("RUSTFS_TEST_TOP_DISK_STATE") else { + return; + }; + let state = std::path::PathBuf::from(state); + let stop = CancellationToken::new(); + let input_stop = stop.clone(); + std::thread::spawn(move || { + let _ = std::io::stdin().read(&mut [0u8]); + input_stop.cancel(); + }); + tokio::runtime::Runtime::new().unwrap().block_on(async { + let runtime = spawn_local_trace_capture_runtime(&state, &stop).unwrap(); + let mut file = tempfile::tempfile().unwrap(); + while !stop.is_cancelled() { + file.write_all(&[1u8; 4096]).unwrap(); + file.sync_all().unwrap(); + tokio::time::sleep(Duration::from_millis(5)).await; + } + runtime.shutdown().await; + }); + } + + #[cfg(target_os = "linux")] + #[tokio::test] + #[serial] + async fn local_top_disk_measures_service_process_and_signs_offline() { + use std::io::Read as _; + let state = tempfile::tempdir().unwrap(); + std::fs::set_permissions(state.path(), std::fs::Permissions::from_mode(0o700)).unwrap(); + let request = disk_request(state.path()); + let mut child = std::process::Command::new(std::env::current_exe().unwrap()) + .args([ + "connect::diagnostics::trace_runtime::tests::local_top_disk_service_child", + "--exact", + "--nocapture", + ]) + .env("RUSTFS_TEST_TOP_DISK_STATE", state.path()) + .stdin(std::process::Stdio::piped()) + .spawn() + .unwrap(); + let ready = tokio::time::timeout(Duration::from_secs(5), async { + while !state.path().join(super::SOCKET_FILE).exists() { + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await; + let result = if ready.is_ok() { + super::request_local_top_disk(state.path(), request.clone(), &CancellationToken::new()).await + } else { + Err(LocalTraceCaptureError::Protocol) + }; + drop(child.stdin.take()); + assert!(child.wait().unwrap().success()); + let export = result.unwrap(); + let mut zip = zip::ZipArchive::new(std::io::Cursor::new(&export.archive_bytes)).unwrap(); + let result: serde_json::Value = serde_json::from_reader(zip.by_name("result.json").unwrap()).unwrap(); + assert_eq!(result["outcome"], "SUCCEEDED"); + assert!(result["data"]["writeBytes"].as_u64().unwrap() >= 4096); + assert!(result["data"]["ioCount"].as_u64().unwrap() >= 20); + let mut envelope = Vec::new(); + zip.by_name("envelope.json").unwrap().read_to_end(&mut envelope).unwrap(); + let envelope_value: serde_json::Value = serde_json::from_slice(&envelope).unwrap(); + assert_eq!(envelope_value["deviceKeyId"], request.offline_key_id); + assert_eq!(envelope_value["classification"], "L3"); + let key = super::load_offline_key(state.path(), &request.offline_key_id).unwrap(); + let signature: serde_json::Value = serde_json::from_reader(zip.by_name("envelope.sig").unwrap()).unwrap(); + let mut signed = b"rustfs-diagnostic-envelope-v1\0".to_vec(); + signed.extend_from_slice(&envelope); + assert!(key.verifies_pending_registration_state(&signed, signature["value"].as_str().unwrap())); + } + #[tokio::test] #[serial] async fn local_runtime_profile_uses_server_workers_and_offline_signature() { @@ -1052,45 +1426,73 @@ mod tests { } #[tokio::test] #[serial] - async fn local_runtime_profile_does_not_claim_unacknowledged_cancellation() { - let state = tempfile::tempdir().unwrap(); - std::fs::set_permissions(state.path(), std::fs::Permissions::from_mode(0o700)).unwrap(); - let (request, _) = runtime_request(state.path()); - let owner = super::private_state_owner(state.path()).unwrap(); - let listener = super::bind_listener(&state.path().join(super::SOCKET_FILE), owner).unwrap(); - let (ready, received) = tokio::sync::oneshot::channel(); - let stop = CancellationToken::new(); - let peer_stop = stop.clone(); - let peer = tokio::spawn(async move { - let (stream, _) = listener.accept().await.unwrap(); - let (reader, _writer) = stream.into_split(); - let mut reader = tokio::io::BufReader::new(reader); - super::read_request(&mut reader).await.unwrap(); - ready.send(()).unwrap(); - peer_stop.cancelled().await; - }); - let cancel = CancellationToken::new(); - let task_cancel = cancel.clone(); - let state_root = state.path().to_path_buf(); - let task = tokio::spawn(async move { super::request_local_runtime_profile(&state_root, request, &task_cancel).await }); - received.await.unwrap(); - cancel.cancel(); - assert!(matches!(task.await.unwrap(), Err(LocalTraceCaptureError::CancellationUnconfirmed))); - stop.cancel(); - peer.await.unwrap(); + async fn local_signed_capture_does_not_claim_unacknowledged_cancellation() { + for disk in [false, true] { + let state = tempfile::tempdir().unwrap(); + std::fs::set_permissions(state.path(), std::fs::Permissions::from_mode(0o700)).unwrap(); + let (request, _) = runtime_request(state.path()); + let owner = super::private_state_owner(state.path()).unwrap(); + let listener = super::bind_listener(&state.path().join(super::SOCKET_FILE), owner).unwrap(); + let (ready, received) = tokio::sync::oneshot::channel(); + let stop = CancellationToken::new(); + let peer_stop = stop.clone(); + let peer = tokio::spawn(async move { + let (stream, _) = listener.accept().await.unwrap(); + let (reader, _writer) = stream.into_split(); + let mut reader = tokio::io::BufReader::new(reader); + super::read_request(&mut reader).await.unwrap(); + ready.send(()).unwrap(); + peer_stop.cancelled().await; + }); + let cancel = CancellationToken::new(); + let task_cancel = cancel.clone(); + let state_root = state.path().to_path_buf(); + let disk_input = disk_request(state.path()); + let task = tokio::spawn(async move { + if disk { + super::request_local_top_disk(&state_root, disk_input, &task_cancel) + .await + .map(|_| ()) + } else { + super::request_local_runtime_profile(&state_root, request, &task_cancel) + .await + .map(|_| ()) + } + }); + received.await.unwrap(); + cancel.cancel(); + assert!(matches!(task.await.unwrap(), Err(LocalTraceCaptureError::CancellationUnconfirmed))); + stop.cancel(); + peer.await.unwrap(); + } } #[tokio::test] async fn local_capture_wire_keeps_separate_request_limits() { use tokio::io::AsyncWriteExt as _; - for (size, accepted) in [(1_024, true), (1_025, false), (8_193, false)] { + let state = tempfile::tempdir().unwrap(); + std::fs::set_permissions(state.path(), std::fs::Permissions::from_mode(0o700)).unwrap(); + for (disk, size, accepted) in [ + (false, 1_024, true), + (false, 1_025, false), + (false, 8_193, false), + (true, 8_192, true), + (true, 8_193, false), + ] { let (mut sender, receiver) = tokio::net::UnixStream::pair().unwrap(); - let mut bytes = serde_json::to_vec(&super::CaptureRequest::TraceRecord { - protocol_version: 1, - consent_expires_at_unix: 1, - duration_millis: 1, - max_spans: 1, - }) - .unwrap(); + let request = if disk { + super::CaptureRequest::TopDisk { + protocol_version: 1, + request: disk_request(state.path()), + } + } else { + super::CaptureRequest::TraceRecord { + protocol_version: 1, + consent_expires_at_unix: 1, + duration_millis: 1, + max_spans: 1, + } + }; + let mut bytes = serde_json::to_vec(&request).unwrap(); bytes.resize(size - 1, b' '); bytes.push(b'\n'); let writer = tokio::spawn(async move { diff --git a/rustfs/src/connect/diagnostics/trace_runtime_unsupported.rs b/rustfs/src/connect/diagnostics/trace_runtime_unsupported.rs index 50a6ebd5f..e13c3f703 100644 --- a/rustfs/src/connect/diagnostics/trace_runtime_unsupported.rs +++ b/rustfs/src/connect/diagnostics/trace_runtime_unsupported.rs @@ -63,3 +63,11 @@ pub(crate) async fn request_local_runtime_profile( ) -> Result { Err(LocalTraceCaptureError::RuntimeUnavailable) } + +pub(crate) async fn request_local_top_disk( + _state_root: &Path, + _request: super::top_disk::LocalTopDiskRequest, + _cancel: &CancellationToken, +) -> Result { + Err(LocalTraceCaptureError::RuntimeUnavailable) +} diff --git a/rustfs/src/connect/mod.rs b/rustfs/src/connect/mod.rs index 712a62920..2faaf6254 100644 --- a/rustfs/src/connect/mod.rs +++ b/rustfs/src/connect/mod.rs @@ -129,8 +129,8 @@ pub use diagnostics::{ evaluate_network_window, save_signed_top_export, sign_top_export, }; pub(crate) use diagnostics::{ - LocalTraceCaptureError, LocalTraceCaptureRuntime, request_local_runtime_profile, request_local_trace_capture, - spawn_local_trace_capture_runtime, + LocalTraceCaptureError, LocalTraceCaptureRuntime, request_local_runtime_profile, request_local_top_disk, + request_local_trace_capture, spawn_local_trace_capture_runtime, }; pub use environment::{ ENVIRONMENT_CAPABILITY, ENVIRONMENT_SCHEMA_VERSION, EnvironmentCollectionRequest, EnvironmentError, EnvironmentExportRequest, diff --git a/rustfs/src/startup_entrypoint.rs b/rustfs/src/startup_entrypoint.rs index a6774d7d5..d2ce05da9 100644 --- a/rustfs/src/startup_entrypoint.rs +++ b/rustfs/src/startup_entrypoint.rs @@ -675,8 +675,8 @@ fn unix_now() -> Result { async fn execute_connect_top(command: ConnectTopCommands) -> Result<()> { use crate::connect::{ IdentityStore, LocalTopConsent, MAX_TOP_DURATION, MAX_TOP_EXPORT_VALIDITY, TOP_CLASSIFICATION, TopApiOperation, - TopCaptureLimits, TopCaptureRequest, TopCaptureScope, capture_top_api, capture_top_disk, capture_top_locks, - capture_top_net, capture_top_rpc, + TopCaptureLimits, TopCaptureRequest, TopCaptureScope, capture_top_api, capture_top_locks, capture_top_net, + capture_top_rpc, }; let (tool_id, options) = match command { @@ -692,6 +692,68 @@ async fn execute_connect_top(command: ConnectTopCommands) -> Result<()> { return Err(Error::other("connect_top_limits_invalid")); } + if tool_id == "top.disk" { + let offline_key_id = options + .offline_key_id + .ok_or_else(|| Error::other("--offline-key-id must select an existing offline identity"))?; + let input = crate::connect::diagnostics::LocalTopDiskRequest { + offline_key_id, + organization_name: options.organization, + cluster_name: options.cluster, + device_name: options.device, + run_uid: options.run_uid, + artifact_uid: options.artifact_uid, + consent_uid: options.consent_uid, + policy_revision: options.policy_revision, + consent_expires_at_unix: options.consent_expires_at_unix, + acknowledge_l3: options.acknowledge_l3, + run_expires_at_unix: options.run_expires_at_unix, + window_millis: options.window_millis, + export_validity_seconds: options.export_validity_seconds, + }; + let cancel = CancellationToken::new(); + let capture = crate::connect::request_local_top_disk(&options.state_dir, input, &cancel); + tokio::pin!(capture); + let export = tokio::select! { + biased; + signal = tokio::signal::ctrl_c() => { + signal.map_err(Error::other)?; + cancel.cancel(); + return match capture.await { + Err(error) => Err(Error::other(error)), + Ok(_) => Err(Error::other("top disk collection cancelled")), + }; + } + result = &mut capture => result.map_err(Error::other)?, + }; + let writer_cancel = cancel.clone(); + let mut writer = tokio::task::spawn_blocking(move || { + crate::connect::diagnostics::save_top_archive( + &options.output, + &export.artifact_uid, + &export.archive_bytes, + &export.archive_sha256, + &writer_cancel, + ) + }); + let receipt = tokio::select! { + biased; + signal = tokio::signal::ctrl_c() => { + signal.map_err(Error::other)?; cancel.cancel(); + writer.await.map_err(Error::other)?.map_err(Error::other)? + } + result = &mut writer => result.map_err(Error::other)?.map_err(Error::other)?, + }; + println!( + "artifact={} bytes={} sha256={}", + receipt.artifact_uid, receipt.archive_size_bytes, receipt.archive_sha256 + ); + println!("upload=not-performed"); + return Ok(()); + } + if options.offline_key_id.is_some() { + return Err(Error::other("--offline-key-id is only supported for top disk")); + } let identity = IdentityStore::new(options.state_dir.join("identity")) .load() .map_err(Error::other)? @@ -725,10 +787,6 @@ async fn execute_connect_top(command: ConnectTopCommands) -> Result<()> { let result = await_top_capture(capture_top_api(&request, TopApiOperation::GetObject, &cancel), &cancel).await?; finish_top_capture(&request, result, &identity, options.output, &cancel).await } - "top.disk" => { - let result = await_top_capture(capture_top_disk(&request, &cancel), &cancel).await?; - finish_top_capture(&request, result, &identity, options.output, &cancel).await - } "top.locks" => { let result = await_top_capture(capture_top_locks(&request, &cancel), &cancel).await?; finish_top_capture(&request, result, &identity, options.output, &cancel).await