diff --git a/rustfs/src/connect/diagnostics/profile_cpu.rs b/rustfs/src/connect/diagnostics/profile_cpu.rs index 7ed53f38c..9b75c3357 100644 --- a/rustfs/src/connect/diagnostics/profile_cpu.rs +++ b/rustfs/src/connect/diagnostics/profile_cpu.rs @@ -86,7 +86,7 @@ impl ProfileTool { } } -#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize)] +#[derive(Clone, Copy, Debug, Deserialize, PartialEq, Eq, Serialize)] #[serde(rename_all = "SCREAMING_SNAKE_CASE")] pub enum ProfileOutcome { Succeeded, @@ -108,7 +108,7 @@ impl ProfileOutcome { } } -#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize)] +#[derive(Clone, Copy, Debug, Deserialize, PartialEq, Eq, Serialize)] #[serde(rename_all = "SCREAMING_SNAKE_CASE")] pub enum ProfileReasonCode { Complete, diff --git a/rustfs/src/connect/diagnostics/trace_runtime.rs b/rustfs/src/connect/diagnostics/trace_runtime.rs index 920f38799..c2714f4c6 100644 --- a/rustfs/src/connect/diagnostics/trace_runtime.rs +++ b/rustfs/src/connect/diagnostics/trace_runtime.rs @@ -274,6 +274,8 @@ enum CaptureResponse { archive_base64: String, archive_sha256: String, artifact_uid: String, + outcome: ProfileOutcome, + reason_code: ProfileReasonCode, }, RuntimeError { code: RuntimeErrorCode, @@ -627,6 +629,7 @@ async fn request_local_profile_capture( let tool = match request.capability.as_str() { super::profile_cpu::THREAD_PROFILE_CAPABILITY => ProfileTool::Threads, super::profile_cpu::MEMORY_PROFILE_CAPABILITY if !native_threads => ProfileTool::Memory, + super::profile_cpu::CPU_PROFILE_CAPABILITY if !native_threads => ProfileTool::Cpu, _ => return Err(LocalTraceCaptureError::Protocol), }; let owner = private_state_owner(state_root)?; @@ -690,6 +693,8 @@ async fn request_local_profile_capture( archive_base64, archive_sha256, artifact_uid: returned_uid, + outcome, + reason_code, } => { let archive_bytes = URL_SAFE_NO_PAD .decode_to_vec(&archive_base64) @@ -705,8 +710,8 @@ async fn request_local_profile_capture( Ok(SignedProfileExport { artifact_uid, tool, - outcome: ProfileOutcome::Succeeded, - reason_code: ProfileReasonCode::Complete, + outcome, + reason_code, archive_bytes, archive_sha256, }) @@ -910,6 +915,8 @@ async fn handle_runtime_profile( archive_base64: URL_SAFE_NO_PAD.encode_to_string(&export.archive_bytes), archive_sha256: export.archive_sha256, artifact_uid: export.artifact_uid, + outcome: export.outcome, + reason_code: export.reason_code, }, Err(code) => CaptureResponse::RuntimeError { code }, }; @@ -1496,6 +1503,7 @@ async fn capture_local_runtime_profile( let tool = match input.capability.as_str() { super::profile_cpu::THREAD_PROFILE_CAPABILITY => ProfileTool::Threads, super::profile_cpu::MEMORY_PROFILE_CAPABILITY => ProfileTool::Memory, + super::profile_cpu::CPU_PROFILE_CAPABILITY => ProfileTool::Cpu, _ => return Err(RuntimeErrorCode::InvalidRequest), }; if input.schema_version != 1 { @@ -1549,7 +1557,15 @@ async fn capture_local_runtime_profile( .map_err(|_| RuntimeErrorCode::CollectionFailed)? .map_err(RuntimeErrorCode::from) } - ProfileTool::Cpu => Err(RuntimeErrorCode::InvalidRequest), + ProfileTool::Cpu => { + let export = super::profile_cpu::export_cpu_profile(&request, &key, cancel) + .await + .map_err(RuntimeErrorCode::from)?; + if export.outcome == ProfileOutcome::Unsupported { + return Err(RuntimeErrorCode::SourceUnavailable); + } + Ok(export) + } } } @@ -2995,6 +3011,69 @@ mod tests { assert!(!state.path().join("identity").exists()); } + #[tokio::test] + #[serial] + async fn local_cpu_profile_uses_service_capture_and_offline_identity() { + #[cfg(feature = "pyroscope")] + use std::io::Read as _; + let state = tempfile::tempdir().unwrap(); + std::fs::set_permissions(state.path(), std::fs::Permissions::from_mode(0o700)).unwrap(); + let (mut request, _key) = runtime_request(state.path()); + request.capability = super::super::profile_cpu::CPU_PROFILE_CAPABILITY.to_owned(); + request.duration_millis = 1_000; + request.sample_period_micros = 10_000; + let mut invalid = request.clone(); + invalid.acknowledge_l3 = false; + assert!(matches!( + super::capture_local_runtime_profile(state.path(), 1, invalid, &CancellationToken::new()).await, + Err(super::RuntimeErrorCode::ConsentRequired) + )); + let mut invalid = request.clone(); + invalid.offline_key_id = "0".repeat(64); + assert!(matches!( + super::capture_local_runtime_profile(state.path(), 1, invalid, &CancellationToken::new()).await, + Err(super::RuntimeErrorCode::IdentityUnavailable) + )); + let runtime = spawn_local_trace_capture_runtime(state.path(), &CancellationToken::new()).unwrap(); + let export = super::request_local_runtime_profile(state.path(), request.clone(), &CancellationToken::new()).await; + runtime.shutdown().await; + + #[cfg(not(feature = "pyroscope"))] + { + assert!(matches!( + export, + Err(LocalTraceCaptureError::RuntimeProfile(ref code)) if code == "SourceUnavailable" + )); + return; + } + + #[cfg(feature = "pyroscope")] + { + let export = export.unwrap(); + assert_eq!(export.tool, super::ProfileTool::Cpu); + assert!(matches!( + export.outcome, + super::ProfileOutcome::Succeeded | super::ProfileOutcome::Partial + )); + 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.cpu"); + assert_eq!(result["outcome"], export.outcome.as_str()); + assert_eq!(result["reasonCode"], export.reason_code.as_str()); + assert!(result["data"]["samples"].is_array()); + 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] #[serial] async fn local_runtime_profile_rejects_invalid_consent_identity_and_protocol() { diff --git a/rustfs/src/startup_entrypoint.rs b/rustfs/src/startup_entrypoint.rs index b378b5f9f..e2402b2a4 100644 --- a/rustfs/src/startup_entrypoint.rs +++ b/rustfs/src/startup_entrypoint.rs @@ -1517,8 +1517,8 @@ async fn execute_connect_profile(options: ConnectProfileOpts) -> Result<()> { use rand::{TryRng as _, rngs::SysRng}; 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() { + || (options.tool != ConnectProfileTool::Threads && options.offline_key_id.is_some()); + if options.tool != ConnectProfileTool::Threads && options.thread_scope.is_some() { return Err(Error::other("--thread-scope is valid only for the threads profile")); } if !service_profile && options.offline_key_id.is_some() {