fix: collect offline Top disk from the running service (#8174)

fix(connect): capture offline disk IO in the running service
This commit is contained in:
Chris
2026-09-28 01:35:40 +08:00
committed by GitHub
parent f5fbb5f4e1
commit e634df8611
8 changed files with 635 additions and 97 deletions
+4 -1
View File
@@ -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<String>,
/// Directory containing an enrolled Connect device identity
#[arg(long = "state-dir")]
pub state_dir: PathBuf,
+5 -2
View File
@@ -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;
+43 -10
View File
@@ -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<SavedTopExport, TopCaptureError> {
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<SavedTopExport, TopCaptureError> {
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();
+32 -1
View File
@@ -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<u8>,
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();
+477 -75
View File
@@ -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<super::top_disk::LocalTopDiskArchive, LocalTraceCaptureError> {
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::<CaptureResponse>(&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<tokio::net::unix::OwnedReadHalf>,
mut writer: tokio::net::unix::OwnedWriteHalf,
@@ -408,6 +505,182 @@ async fn handle_runtime_profile(
}
}
async fn handle_top_disk(
mut reader: BufReader<tokio::net::unix::OwnedReadHalf>,
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<super::top_api::SignedTopExport, RuntimeErrorCode> {
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<crate::connect::DeviceIdentity, RuntimeErrorCode> {
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::<super::super::top_disk::LocalTopDiskRequest>(value).is_err());
assert!(
serde_json::from_str::<super::CaptureRequest>(
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 {
@@ -63,3 +63,11 @@ pub(crate) async fn request_local_runtime_profile(
) -> Result<super::profile_cpu::SignedProfileExport, LocalTraceCaptureError> {
Err(LocalTraceCaptureError::RuntimeUnavailable)
}
pub(crate) async fn request_local_top_disk(
_state_root: &Path,
_request: super::top_disk::LocalTopDiskRequest,
_cancel: &CancellationToken,
) -> Result<super::top_disk::LocalTopDiskArchive, LocalTraceCaptureError> {
Err(LocalTraceCaptureError::RuntimeUnavailable)
}
+2 -2
View File
@@ -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,
+64 -6
View File
@@ -675,8 +675,8 @@ fn unix_now() -> Result<i64> {
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