feat: Capture offline CPU profiles in serving process (#8245)

feat: capture offline CPU profiles in serving process
This commit is contained in:
Chris
2026-09-29 22:37:10 +08:00
committed by GitHub
parent 380e98a42c
commit 1e7065101d
3 changed files with 86 additions and 7 deletions
@@ -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,
@@ -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() {
+2 -2
View File
@@ -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() {