From f0ce628a7fecd3f6b8468ec30ff7eee706d6393b Mon Sep 17 00:00:00 2001 From: Chris Date: Tue, 29 Sep 2026 08:48:56 +0800 Subject: [PATCH] feat(connect): capture offline native threads in service process (#8215) --- rustfs/src/connect/diagnostics/mod.rs | 4 +- .../src/connect/diagnostics/trace_runtime.rs | 278 +++++++++++++++++- .../diagnostics/trace_runtime_unsupported.rs | 8 + rustfs/src/connect/mod.rs | 4 +- rustfs/src/startup_entrypoint.rs | 18 +- 5 files changed, 288 insertions(+), 24 deletions(-) diff --git a/rustfs/src/connect/diagnostics/mod.rs b/rustfs/src/connect/diagnostics/mod.rs index f55370bea..74ce71486 100644 --- a/rustfs/src/connect/diagnostics/mod.rs +++ b/rustfs/src/connect/diagnostics/mod.rs @@ -172,8 +172,8 @@ 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_top_disk, - request_local_trace_capture, spawn_local_trace_capture_runtime, + LocalTraceCaptureError, LocalTraceCaptureRuntime, request_local_native_threads_profile, 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; diff --git a/rustfs/src/connect/diagnostics/trace_runtime.rs b/rustfs/src/connect/diagnostics/trace_runtime.rs index 2428c7e08..59ff0a041 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, RUNTIME_PROFILE and TOP_DISK. Signed requests select +//! Version 1 accepts TRACE_RECORD, RUNTIME_PROFILE, NATIVE_THREADS_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 @@ -110,6 +110,10 @@ enum CaptureRequest { protocol_version: u16, request: LocalRuntimeProfileRequest, }, + NativeThreadsProfile { + protocol_version: u16, + request: LocalRuntimeProfileRequest, + }, TopDisk { protocol_version: u16, request: super::top_disk::LocalTopDiskRequest, @@ -305,10 +309,27 @@ pub(crate) async fn request_local_runtime_profile( state_root: &Path, request: LocalRuntimeProfileRequest, cancel: &CancellationToken, +) -> Result { + request_local_profile_capture(state_root, request, cancel, false).await +} + +pub(crate) async fn request_local_native_threads_profile( + state_root: &Path, + request: LocalRuntimeProfileRequest, + cancel: &CancellationToken, +) -> Result { + request_local_profile_capture(state_root, request, cancel, true).await +} + +async fn request_local_profile_capture( + state_root: &Path, + request: LocalRuntimeProfileRequest, + cancel: &CancellationToken, + native_threads: bool, ) -> Result { let tool = match request.capability.as_str() { super::profile_cpu::THREAD_PROFILE_CAPABILITY => ProfileTool::Threads, - super::profile_cpu::MEMORY_PROFILE_CAPABILITY => ProfileTool::Memory, + super::profile_cpu::MEMORY_PROFILE_CAPABILITY if !native_threads => ProfileTool::Memory, _ => return Err(LocalTraceCaptureError::Protocol), }; let owner = private_state_owner(state_root)?; @@ -319,9 +340,16 @@ pub(crate) async fn request_local_runtime_profile( return Err(LocalTraceCaptureError::StateSecurity); } let artifact_uid = request.artifact_uid.clone(); - let message = CaptureRequest::RuntimeProfile { - protocol_version: PROTOCOL_VERSION, - request, + let message = if native_threads { + CaptureRequest::NativeThreadsProfile { + protocol_version: PROTOCOL_VERSION, + request, + } + } else { + CaptureRequest::RuntimeProfile { + protocol_version: PROTOCOL_VERSION, + request, + } }; let mut bytes = serde_json::to_vec(&message).map_err(|_| LocalTraceCaptureError::Protocol)?; bytes.push(b'\n'); @@ -480,9 +508,16 @@ async fn handle_runtime_profile( protocol_version: u16, request: LocalRuntimeProfileRequest, shutdown: CancellationToken, + native_threads: bool, ) { let cancel = shutdown.child_token(); - let capture = capture_local_runtime_profile(state_root, protocol_version, request, &cancel); + let capture = async { + if native_threads { + capture_local_native_threads_profile(state_root, protocol_version, request, &cancel).await + } else { + capture_local_runtime_profile(state_root, protocol_version, request, &cancel).await + } + }; tokio::pin!(capture); let mut unexpected = [0_u8; 1]; let result = tokio::select! { @@ -782,6 +817,91 @@ async fn capture_local_runtime_profile( } } +async fn capture_local_native_threads_profile( + state_root: &Path, + protocol_version: u16, + input: LocalRuntimeProfileRequest, + cancel: &CancellationToken, +) -> Result { + use rand::{TryRng as _, rngs::SysRng}; + 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)) + || input.schema_version != 1 + || input.capability != super::profile_cpu::THREAD_PROFILE_CAPABILITY + { + return Err(RuntimeErrorCode::InvalidRequest); + } + if cancel.is_cancelled() { + return Err(RuntimeErrorCode::Cancelled); + } + 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 || input.expires_at_unix > input.consent_expires_at_unix { + return Err(RuntimeErrorCode::ConsentExpired); + } + if input.expires_at_unix <= now || input.expires_at_unix.saturating_sub(now) > super::profile_cpu::MAX_VALIDITY_SECONDS { + return Err(RuntimeErrorCode::Expired); + } + if input.duration_millis == 0 + || input.duration_millis > 30_000 + || input.sample_period_micros == 0 + || input.sample_period_micros > input.duration_millis * 1_000 + { + 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 mut nonce = [0_u8; 32]; + SysRng + .try_fill_bytes(&mut nonce) + .map_err(|_| RuntimeErrorCode::CollectionFailed)?; + let request = ProfileCaptureRequest { + 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, + schema_version: input.schema_version, + capability: input.capability, + consent: LocalProfileConsent { + consent_uid: input.consent_uid, + policy_revision: input.policy_revision, + expires_at_unix: input.consent_expires_at_unix, + confirmed: input.acknowledge_l3, + }, + produced_at_unix: unix_now().map_err(|_| RuntimeErrorCode::CollectionFailed)?, + expires_at_unix: input.expires_at_unix, + nonce, + duration: Duration::from_millis(input.duration_millis), + sample_period: Duration::from_micros(input.sample_period_micros), + provenance, + }; + let cancel = cancel.clone(); + // The blocking collector runs in this serving process; joining it also releases its lease. + tokio::task::spawn_blocking(move || { + let result = super::profile_threads::capture_thread_profile( + &request, + super::profile_cpu::ThreadProfileScope::NativeThreads, + &cancel, + )?; + if result.outcome() != ProfileOutcome::Succeeded { + return Err(ProfileError::CollectionFailed); + } + encode_signed_profile_export(&request, &result, &key, &cancel) + }) + .await + .map_err(|_| RuntimeErrorCode::CollectionFailed)? + .map_err(RuntimeErrorCode::from) +} + async fn run_listener(listener: UnixListener, owner: u32, state_root: PathBuf, shutdown: CancellationToken) { let mut connections = JoinSet::new(); loop { @@ -820,7 +940,14 @@ async fn handle_connection(stream: UnixStream, state_root: PathBuf, shutdown: Ca protocol_version, request, } => { - handle_runtime_profile(reader, writer, &state_root, protocol_version, request, shutdown).await; + handle_runtime_profile(reader, writer, &state_root, protocol_version, request, shutdown, false).await; + return; + } + CaptureRequest::NativeThreadsProfile { + protocol_version, + request, + } => { + handle_runtime_profile(reader, writer, &state_root, protocol_version, request, shutdown, true).await; return; } CaptureRequest::TopDisk { @@ -1356,6 +1483,127 @@ mod tests { assert!(!state.path().join("identity").exists()); } + #[cfg(target_os = "linux")] + #[test] + fn local_native_threads_service_child() { + let Some(state) = std::env::var_os("RUSTFS_TEST_NATIVE_THREADS_STATE") else { + return; + }; + let stop = CancellationToken::new(); + let input_stop = stop.clone(); + std::thread::spawn(move || { + use std::io::Read as _; + let _ = std::io::stdin().read(&mut [0u8]); + input_stop.cancel(); + }); + let sleepers: Vec<_> = (0..16) + .map(|_| { + let stop = stop.clone(); + std::thread::spawn(move || { + while !stop.is_cancelled() { + std::thread::park_timeout(Duration::from_millis(100)); + } + }) + }) + .collect(); + tokio::runtime::Runtime::new().unwrap().block_on(async { + let runtime = spawn_local_trace_capture_runtime(std::path::Path::new(&state), &stop).unwrap(); + stop.cancelled().await; + runtime.shutdown().await; + }); + for sleeper in sleepers { + sleeper.join().unwrap(); + } + } + + #[cfg(target_os = "linux")] + #[tokio::test] + #[serial] + async fn local_native_threads_profiles_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, key) = runtime_request(state.path()); + let mut child = std::process::Command::new(std::env::current_exe().unwrap()) + .args([ + "connect::diagnostics::trace_runtime::tests::local_native_threads_service_child", + "--exact", + "--nocapture", + ]) + .env("RUSTFS_TEST_NATIVE_THREADS_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 export = if ready.is_ok() { + super::request_local_native_threads_profile(state.path(), request.clone(), &CancellationToken::new()).await + } else { + Err(LocalTraceCaptureError::Protocol) + }; + drop(child.stdin.take()); + assert!(child.wait().unwrap().success()); + let export = export.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["toolId"], "profile.threads"); + assert_eq!(result["outcome"], "SUCCEEDED"); + assert_eq!(result["data"]["scope"], "NATIVE_THREADS"); + let states = result["data"]["states"].as_array().unwrap(); + assert_eq!(states.len(), 4); + assert!(states.iter().map(|state| state["threadCount"].as_u64().unwrap()).sum::() >= 16); + assert!(result["data"].get("threadNames").is_none()); + let mut envelope = Vec::new(); + zip.by_name("envelope.json").unwrap().read_to_end(&mut envelope).unwrap(); + let metadata: serde_json::Value = serde_json::from_slice(&envelope).unwrap(); + assert_eq!(metadata["classification"], "L3"); + assert_eq!(metadata["deviceKeyId"], request.offline_key_id); + 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())); + assert!(!state.path().join("identity").exists()); + } + + #[tokio::test] + async fn local_native_threads_rejects_consent_scope_and_missing_identity() { + 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 cancel = CancellationToken::new(); + assert!(matches!( + super::capture_local_native_threads_profile(state.path(), 2, request.clone(), &cancel).await, + Err(super::RuntimeErrorCode::InvalidRequest) + )); + let mut invalid = request.clone(); + invalid.acknowledge_l3 = false; + assert!(matches!( + super::capture_local_native_threads_profile(state.path(), 1, invalid, &cancel).await, + Err(super::RuntimeErrorCode::ConsentRequired) + )); + let mut invalid = request.clone(); + invalid.capability = "profile.memory@1".to_owned(); + assert!(matches!( + super::capture_local_native_threads_profile(state.path(), 1, invalid, &cancel).await, + Err(super::RuntimeErrorCode::InvalidRequest) + )); + let mut invalid = request.clone(); + invalid.expires_at_unix = 1; + assert!(matches!( + super::capture_local_native_threads_profile(state.path(), 1, invalid, &cancel).await, + Err(super::RuntimeErrorCode::Expired) + )); + std::fs::remove_file(crate::connect::OfflineKeyStore::new(state.path()).key_path()).unwrap(); + assert!(matches!( + super::capture_local_native_threads_profile(state.path(), 1, request, &cancel).await, + Err(super::RuntimeErrorCode::IdentityUnavailable) + )); + } + #[tokio::test] #[serial] async fn local_memory_profile_signs_service_capture_with_offline_identity() { @@ -1484,7 +1732,7 @@ mod tests { #[tokio::test] #[serial] async fn local_signed_capture_does_not_claim_unacknowledged_cancellation() { - for disk in [false, true] { + for kind in 0..3 { let state = tempfile::tempdir().unwrap(); std::fs::set_permissions(state.path(), std::fs::Permissions::from_mode(0o700)).unwrap(); let (request, _) = runtime_request(state.path()); @@ -1506,14 +1754,16 @@ mod tests { 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) + match kind { + 0 => super::request_local_runtime_profile(&state_root, request, &task_cancel) .await - .map(|_| ()) - } else { - super::request_local_runtime_profile(&state_root, request, &task_cancel) + .map(|_| ()), + 1 => super::request_local_native_threads_profile(&state_root, request, &task_cancel) .await - .map(|_| ()) + .map(|_| ()), + _ => super::request_local_top_disk(&state_root, disk_input, &task_cancel) + .await + .map(|_| ()), } }); received.await.unwrap(); diff --git a/rustfs/src/connect/diagnostics/trace_runtime_unsupported.rs b/rustfs/src/connect/diagnostics/trace_runtime_unsupported.rs index e13c3f703..deca5a997 100644 --- a/rustfs/src/connect/diagnostics/trace_runtime_unsupported.rs +++ b/rustfs/src/connect/diagnostics/trace_runtime_unsupported.rs @@ -64,6 +64,14 @@ pub(crate) async fn request_local_runtime_profile( Err(LocalTraceCaptureError::RuntimeUnavailable) } +pub(crate) async fn request_local_native_threads_profile( + _state_root: &Path, + _request: super::profile_cpu::LocalRuntimeProfileRequest, + _cancel: &CancellationToken, +) -> Result { + Err(LocalTraceCaptureError::RuntimeUnavailable) +} + pub(crate) async fn request_local_top_disk( _state_root: &Path, _request: super::top_disk::LocalTopDiskRequest, diff --git a/rustfs/src/connect/mod.rs b/rustfs/src/connect/mod.rs index 2faaf6254..bbe1a3763 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_top_disk, - request_local_trace_capture, spawn_local_trace_capture_runtime, + LocalTraceCaptureError, LocalTraceCaptureRuntime, request_local_native_threads_profile, 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 0ec730984..9fe5a250b 100644 --- a/rustfs/src/startup_entrypoint.rs +++ b/rustfs/src/startup_entrypoint.rs @@ -1401,17 +1401,17 @@ async fn execute_connect_profile(options: ConnectProfileOpts) -> Result<()> { }; use rand::{TryRng as _, rngs::SysRng}; - let runtime_scope = (options.tool == ConnectProfileTool::Threads - && options.thread_scope == Some(ConnectThreadProfileScope::TokioRuntime)) + let service_profile = (options.tool == ConnectProfileTool::Threads && options.thread_scope.is_some()) || (options.tool == ConnectProfileTool::Memory && options.offline_key_id.is_some()); if options.tool == ConnectProfileTool::Memory && options.thread_scope.is_some() { return Err(Error::other("--thread-scope is valid only for the threads profile")); } - if !runtime_scope && options.offline_key_id.is_some() { - return Err(Error::other("--offline-key-id is valid only for a service runtime profile")); + if !service_profile && options.offline_key_id.is_some() { + return Err(Error::other("--offline-key-id is valid only for a service profile")); } let cancel = tokio_util::sync::CancellationToken::new(); - let export = if runtime_scope { + let export = if service_profile { + let native_threads = options.thread_scope == Some(ConnectThreadProfileScope::NativeThreads); let offline_key_id = options .offline_key_id .ok_or_else(|| Error::other("--offline-key-id must explicitly select an existing offline identity"))?; @@ -1432,7 +1432,13 @@ async fn execute_connect_profile(options: ConnectProfileOpts) -> Result<()> { duration_millis: options.duration_millis, sample_period_micros: options.sample_period_micros, }; - let capture = crate::connect::request_local_runtime_profile(&options.state_dir, input, &cancel); + let capture = async { + if native_threads { + crate::connect::request_local_native_threads_profile(&options.state_dir, input, &cancel).await + } else { + crate::connect::request_local_runtime_profile(&options.state_dir, input, &cancel).await + } + }; tokio::pin!(capture); tokio::select! { biased;