feat(connect): execute signed disk diagnostics in the service (#8177)

This commit is contained in:
Chris
2026-09-28 08:36:49 +08:00
committed by GitHub
parent e634df8611
commit 1a251de4f3
2 changed files with 241 additions and 7 deletions
+31 -2
View File
@@ -24,7 +24,7 @@
"properties": {
"jobId": { "type": "string", "format": "uuid", "pattern": "^[0-9a-f]{8}-[0-9a-f]{4}-7[0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$" },
"protocolVersion": { "const": "v1" },
"jobType": { "enum": ["health.check", "profile.cpu", "profile.memory", "profile.threads", "performance.drive", "performance.network", "top.api", "top.locks", "top.rpc"] },
"jobType": { "enum": ["health.check", "profile.cpu", "profile.memory", "profile.threads", "performance.drive", "performance.network", "top.api", "top.disk", "top.locks", "top.rpc"] },
"schemaVersion": { "const": 1 },
"organizationName": { "type": "string", "pattern": "^organizations/[0-9a-f]{8}-[0-9a-f]{4}-7[0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$" },
"clusterName": { "type": "string", "pattern": "^organizations/[0-9a-f-]{36}/clusters/[0-9a-f]{8}-[0-9a-f]{4}-7[0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$" },
@@ -34,7 +34,7 @@
"nonce": { "type": "string", "pattern": "^[A-Za-z0-9_-]{43}$" },
"requiredCapabilities": {
"type": "array",
"prefixItems": [{ "enum": ["health.check.service@1", "profile.cpu@1", "profile.memory.service@1", "profile.threads@1", "performance.drive@1", "performance.network@1", "top.api@1", "top.locks@1", "top.rpc@1"] }],
"prefixItems": [{ "enum": ["health.check.service@1", "profile.cpu@1", "profile.memory.service@1", "profile.threads@1", "performance.drive@1", "performance.network@1", "top.api@1", "top.disk@1", "top.locks@1", "top.rpc@1"] }],
"items": false,
"minItems": 1,
"maxItems": 1,
@@ -308,6 +308,35 @@
}
}
},
{
"if": { "properties": { "jobType": { "const": "top.disk" } } },
"then": {
"properties": {
"requiredCapabilities": {
"type": "array",
"prefixItems": [{ "const": "top.disk@1" }],
"items": false,
"minItems": 1,
"maxItems": 1
},
"limits": {
"type": "object",
"properties": {
"maxCpuMillis": { "type": "integer", "maximum": 5000 },
"maxMemoryBytes": { "type": "integer", "minimum": 1048576 }
}
},
"parameters": {
"type": "object",
"not": {
"type": "object",
"properties": { "trafficBytes": true },
"required": ["trafficBytes"]
}
}
}
}
},
{
"if": { "properties": { "jobType": { "const": "top.rpc" } } },
"then": {
+210 -5
View File
@@ -35,11 +35,11 @@ use super::{
DriveProvenance, LocalDriveConsent, LocalNetworkConsent, LocalProfileConsent, LocalTopConsent, MAX_NETWORK_TRAFFIC_BYTES,
MAX_TOP_EXPORT_VALIDITY, NETWORK_CAPABILITY, NETWORK_SCHEMA_VERSION, NetworkOutcome, NetworkPerformanceError,
NetworkPerformanceRequest, NetworkProvenance, NetworkReasonCode, PROFILE_SCHEMA_VERSION, ProfileCaptureRequest,
ProfileOutcome, ProfileProvenance, THREAD_PROFILE_CAPABILITY, TOP_API_CAPABILITY, TOP_CLASSIFICATION, TOP_LOCKS_CAPABILITY,
TOP_RPC_CAPABILITY, TOP_SCHEMA_VERSION, ThreadProfileScope, TopApiOperation, TopCaptureLimits, TopCaptureRequest,
TopCaptureScope, TopOutcome, capture_cpu_profile, capture_thread_profile, capture_top_api, capture_top_locks,
capture_top_rpc, encode_signed_profile_export, measure_drive, measure_network, runtime_network_peer_aliases,
sign_drive_export, sign_network_export, sign_top_export_with_nonce,
ProfileOutcome, ProfileProvenance, THREAD_PROFILE_CAPABILITY, TOP_API_CAPABILITY, TOP_CLASSIFICATION, TOP_DISK_CAPABILITY,
TOP_LOCKS_CAPABILITY, TOP_RPC_CAPABILITY, TOP_SCHEMA_VERSION, ThreadProfileScope, TopApiOperation, TopCaptureLimits,
TopCaptureRequest, TopCaptureScope, TopOutcome, capture_cpu_profile, capture_thread_profile, capture_top_api,
capture_top_disk, capture_top_locks, capture_top_rpc, encode_signed_profile_export, measure_drive, measure_network,
runtime_network_peer_aliases, sign_drive_export, sign_network_export, sign_top_export_with_nonce,
};
use super::{
HEALTH_SCHEMA_VERSION, HEALTH_SERVICE_CAPABILITY, HEALTH_TIMEOUT_SECONDS, HealthError, HealthServiceRequest,
@@ -55,6 +55,7 @@ const PROFILE_MEMORY_JOB_TYPE: &str = "profile.memory";
const PERFORMANCE_DRIVE_JOB_TYPE: &str = "performance.drive";
const PERFORMANCE_NETWORK_JOB_TYPE: &str = "performance.network";
const TOP_API_JOB_TYPE: &str = "top.api";
const TOP_DISK_JOB_TYPE: &str = "top.disk";
const TOP_LOCKS_JOB_TYPE: &str = "top.locks";
const TOP_RPC_JOB_TYPE: &str = "top.rpc";
const HEALTH_JOB_TYPE: &str = "health.check";
@@ -73,6 +74,8 @@ const DRIVE_SCRATCH_BYTES: u64 = 524_288;
const DRIVE_BLOCK_BYTES: u64 = 65_536;
const DRIVE_SAMPLE_PERIOD_MICROS: u64 = 10_000;
const DRIVE_SCRATCH_DIRECTORY: &str = ".rustfs-connect-drive-scratch";
const MAX_TOP_DISK_CPU_MILLIS: u64 = 5_000;
const MIN_TOP_DISK_MEMORY_BYTES: u64 = 1_048_576;
const MAX_TOP_API_CPU_MILLIS: u64 = 5_000;
const MIN_TOP_API_MEMORY_BYTES: u64 = 1_048_576;
const MAX_TOP_LOCKS_CPU_MILLIS: u64 = 5_000;
@@ -89,6 +92,7 @@ enum DiagnosticJobKind {
PerformanceDrive,
PerformanceNetwork,
TopApi,
TopDisk,
TopLocks,
TopRpc,
}
@@ -468,6 +472,11 @@ impl DiagnosticJobEnvelope {
{
return Err(DiagnosticJobError::LimitExceeded);
}
if kind == DiagnosticJobKind::TopDisk
&& (self.limits.max_cpu_millis > MAX_TOP_DISK_CPU_MILLIS || self.limits.max_memory_bytes < MIN_TOP_DISK_MEMORY_BYTES)
{
return Err(DiagnosticJobError::LimitExceeded);
}
if kind == DiagnosticJobKind::ProfileThreads
&& (self.limits.max_memory_bytes != MAX_MEMORY_BYTES || self.limits.max_cpu_millis != MAX_CPU_MILLIS)
{
@@ -551,6 +560,11 @@ impl DiagnosticJobEnvelope {
{
Ok(DiagnosticJobKind::TopApi)
}
(TOP_DISK_JOB_TYPE, [capability], version)
if capability == TOP_DISK_CAPABILITY && version == u16::from(TOP_SCHEMA_VERSION) =>
{
Ok(DiagnosticJobKind::TopDisk)
}
(TOP_LOCKS_JOB_TYPE, [capability], version)
if capability == TOP_LOCKS_CAPABILITY && version == u16::from(TOP_SCHEMA_VERSION) =>
{
@@ -594,6 +608,7 @@ pub async fn execute_diagnostic_job(
DiagnosticJobKind::PerformanceNetwork => {
execute_performance_network_job(envelope, nonce, identity, provenance, cancel).await
}
DiagnosticJobKind::TopDisk => execute_top_disk_job(envelope, nonce, identity, provenance, cancel).await,
DiagnosticJobKind::TopApi => execute_top_api_job(envelope, nonce, identity, provenance, cancel).await,
DiagnosticJobKind::TopLocks => execute_top_locks_job(envelope, nonce, identity, provenance, cancel).await,
DiagnosticJobKind::TopRpc => execute_top_rpc_job(envelope, nonce, identity, provenance, cancel).await,
@@ -1117,6 +1132,70 @@ async fn execute_top_api_job(
})
}
async fn execute_top_disk_job(
envelope: DiagnosticJobEnvelope,
nonce: [u8; 32],
identity: &DeviceIdentity,
provenance: ProfileProvenance,
cancel: &CancellationToken,
) -> Result<DiagnosticJobExecution, DiagnosticJobError> {
let expire = parse_time(&envelope.expire_time)?;
let consent_expire = parse_time(&envelope.parameters.consent_expires_at)?;
let request = TopCaptureRequest {
scope: TopCaptureScope {
organization_name: envelope.organization_name,
cluster_name: envelope.cluster_name,
device_name: envelope.device_name,
run_uid: envelope.job_id.clone(),
artifact_uid: envelope.parameters.artifact_uid,
policy_revision: envelope.parameters.consent_policy_revision,
run_expires_at_unix: expire.timestamp(),
executable_sha256: provenance.executable_sha256().to_owned(),
build_features: provenance.build_features().to_vec(),
consent: LocalTopConsent {
uid: envelope.parameters.consent_uid,
tool_id: TOP_DISK_JOB_TYPE.to_owned(),
classification: TOP_CLASSIFICATION.to_owned(),
active: true,
expires_at_unix: consent_expire.timestamp(),
},
},
limits: TopCaptureLimits {
max_duration_millis: envelope.parameters.duration_millis,
max_working_memory_bytes: envelope.limits.max_memory_bytes,
max_cpu_millis: envelope.limits.max_cpu_millis,
..TopCaptureLimits::default()
},
window: Duration::from_millis(envelope.parameters.duration_millis),
export_validity: MAX_TOP_EXPORT_VALIDITY,
};
let result = capture_top_disk(&request, cancel).await.map_err(top_capture_failure)?;
let outcome = result.outcome.as_str().to_owned();
let reason = result.reason_code.as_str().to_owned();
if !matches!(result.outcome, TopOutcome::Succeeded | TopOutcome::Partial) {
return Ok(DiagnosticJobExecution {
job_id: envelope.job_id,
outcome,
reason,
artifact_uid: None,
artifact_sha256: None,
artifact_bytes: None,
});
}
let export = sign_top_export_with_nonce(&request, &result, identity, cancel, nonce).map_err(top_export_failure)?;
if export.archive_bytes.len() > usize::try_from(envelope.limits.max_output_bytes).unwrap_or(usize::MAX) {
return Err(DiagnosticJobError::LimitExceeded);
}
Ok(DiagnosticJobExecution {
job_id: envelope.job_id,
outcome,
reason,
artifact_uid: Some(export.artifact_uid),
artifact_sha256: Some(export.archive_sha256),
artifact_bytes: Some(export.archive_bytes),
})
}
async fn execute_top_locks_job(
envelope: DiagnosticJobEnvelope,
nonce: [u8; 32],
@@ -1646,6 +1725,132 @@ mod tests {
}
}
fn disk_envelope() -> DiagnosticJobEnvelope {
let mut top = envelope();
top.job_type = TOP_DISK_JOB_TYPE.to_owned();
top.required_capabilities = vec![TOP_DISK_CAPABILITY.to_owned()];
top.limits.max_cpu_millis = MAX_TOP_DISK_CPU_MILLIS;
top
}
#[test]
fn top_disk_job_rejects_mismatched_capabilities_parameters_consent_and_budgets() {
let now = "2030-01-01T00:00:10Z".parse().unwrap();
let (top, signer) = signed_envelope(disk_envelope());
assert!(signer.verify(&top, &target(&top), now).is_ok());
for case in 0..11 {
let mut invalid = disk_envelope();
let expected = match case {
0 => {
invalid.required_capabilities = vec![TOP_API_CAPABILITY.to_owned()];
DiagnosticJobError::Unsupported
}
1 => {
invalid.schema_version = 2;
DiagnosticJobError::Unsupported
}
2 => {
invalid.parameters.thread_scope = Some(ThreadProfileScope::TokioRuntime);
DiagnosticJobError::Invalid
}
3 => {
invalid.parameters.traffic_bytes = Some(1);
DiagnosticJobError::Invalid
}
4 => {
invalid.parameters.target_alias = Some("drive-1".to_owned());
DiagnosticJobError::Invalid
}
5 => {
invalid.parameters.consent_policy_revision = 0;
DiagnosticJobError::Invalid
}
6 => {
invalid.parameters.consent_expires_at = invalid.create_time.clone();
DiagnosticJobError::Expired
}
7 => {
invalid.limits.max_cpu_millis = MAX_TOP_DISK_CPU_MILLIS + 1;
DiagnosticJobError::LimitExceeded
}
8 => {
invalid.limits.max_memory_bytes = MIN_TOP_DISK_MEMORY_BYTES - 1;
DiagnosticJobError::LimitExceeded
}
9 => {
invalid.parameters.duration_millis = MAX_TOP_DISK_CPU_MILLIS + 1;
DiagnosticJobError::LimitExceeded
}
_ => {
invalid.parameters.sample_period_micros = 0;
DiagnosticJobError::LimitExceeded
}
};
let (invalid, signer) = signed_envelope(invalid);
assert_eq!(signer.verify(&invalid, &target(&invalid), now), Err(expected), "case {case}");
}
}
#[tokio::test]
#[serial_test::serial]
async fn top_disk_job_collects_service_io_and_preserves_signed_nonce() {
let now = Utc::now();
let mut top = disk_envelope();
top.parameters.duration_millis = 200;
top.create_time = now.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
top.expire_time = (now + chrono::Duration::seconds(30)).to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
top.parameters.consent_expires_at =
(now + chrono::Duration::seconds(60)).to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
let (top, signer) = signed_envelope(top);
let verified = signer.verify(&top, &target(&top), now).unwrap();
let identity = DeviceIdentity::generate();
let cancel = CancellationToken::new();
let collect = execute_diagnostic_job(
verified,
&identity,
ProfileProvenance::new("a".repeat(40), "b".repeat(64), "1.0.0", vec![]),
&cancel,
);
let workload = async {
#[cfg(target_os = "linux")]
{
use std::io::Write as _;
tokio::time::sleep(Duration::from_millis(50)).await;
let mut file = tempfile::tempfile().unwrap();
file.write_all(&[1u8; 65_536]).unwrap();
file.sync_all().unwrap();
}
};
let (result, ()) = tokio::join!(collect, workload);
let execution = result.unwrap();
#[cfg(target_os = "linux")]
{
assert_eq!(execution.outcome, "SUCCEEDED");
let mut archive = zip::ZipArchive::new(std::io::Cursor::new(execution.artifact_bytes.unwrap())).unwrap();
let result: serde_json::Value = serde_json::from_reader(archive.by_name("result.json").unwrap()).unwrap();
assert_eq!(result["toolId"], "top.disk");
assert_eq!(result["data"]["resourceAlias"], "resource-1");
assert!(result["data"]["writeBytes"].as_u64().unwrap() >= 65_536);
assert!(result["data"]["ioCount"].as_u64().unwrap() > 0);
use std::io::Read as _;
let mut bytes = Vec::new();
archive.by_name("envelope.json").unwrap().read_to_end(&mut bytes).unwrap();
let envelope: serde_json::Value = serde_json::from_slice(&bytes).unwrap();
assert_eq!(envelope["nonce"], top.nonce);
assert_eq!(envelope["classification"], "L3");
let signature: serde_json::Value = serde_json::from_reader(archive.by_name("envelope.sig").unwrap()).unwrap();
let mut signed = b"rustfs-diagnostic-envelope-v1\0".to_vec();
signed.extend_from_slice(&bytes);
assert!(identity.verifies_pending_registration_state(&signed, signature["value"].as_str().unwrap()));
}
#[cfg(not(target_os = "linux"))]
{
assert_eq!(execution.outcome, "UNSUPPORTED");
assert_eq!(execution.reason, "UNSUPPORTED_PLATFORM");
assert!(execution.artifact_bytes.is_none());
}
}
#[test]
fn accepts_only_the_bounded_top_api_capability_pair() {
let mut top = envelope();