|
|
|
@@ -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<SignedProfileExport, LocalTraceCaptureError> {
|
|
|
|
|
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<SignedProfileExport, LocalTraceCaptureError> {
|
|
|
|
|
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<SignedProfileExport, LocalTraceCaptureError> {
|
|
|
|
|
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<SignedProfileExport, RuntimeErrorCode> {
|
|
|
|
|
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::<u64>() >= 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();
|
|
|
|
|