diff --git a/Cargo.lock b/Cargo.lock index 91f6ed0d8..daf4634ac 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -9589,6 +9589,7 @@ dependencies = [ "rustfs-trusted-proxies", "rustfs-utils", "rustfs-zip", + "rustix", "rustls", "rustls-pki-types", "s3s", diff --git a/rustfs/Cargo.toml b/rustfs/Cargo.toml index cf985363c..abfee7bde 100644 --- a/rustfs/Cargo.toml +++ b/rustfs/Cargo.toml @@ -321,6 +321,7 @@ zstd.workspace = true ed25519-dalek.workspace = true rustls = { workspace = true, default-features = false, features = ["aws-lc-rs", "logging", "tls12", "prefer-post-quantum", "std"] } rustls-pki-types = { workspace = true } +rustix.workspace = true x509-parser = { workspace = true } subtle = { workspace = true } jiff = { workspace = true, features = ["serde"] } diff --git a/rustfs/src/connect/diagnostics/job.rs b/rustfs/src/connect/diagnostics/job.rs new file mode 100644 index 000000000..31f235d2b --- /dev/null +++ b/rustfs/src/connect/diagnostics/job.rs @@ -0,0 +1,517 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! Verification and execution of the small allow-list of Connect diagnostic jobs. + +use std::time::Duration; +use std::{fs, io::Read as _, path::Path}; + +#[cfg(unix)] +use std::os::unix::fs::{MetadataExt as _, OpenOptionsExt as _, PermissionsExt as _}; + +use base64_simd::URL_SAFE_NO_PAD; +use chrono::{DateTime, Utc}; +use ed25519_dalek::{Signature, Verifier as _, VerifyingKey}; +use serde::{Deserialize, Serialize}; +use sha2::{Digest as _, Sha256}; +use thiserror::Error; +use tokio_util::sync::CancellationToken; +use uuid::{Uuid, Variant, Version}; + +use super::{ + CPU_PROFILE_CAPABILITY, LocalProfileConsent, PROFILE_SCHEMA_VERSION, ProfileCaptureRequest, ProfileProvenance, + capture_cpu_profile, encode_signed_profile_export, +}; +use crate::connect::DeviceIdentity; + +const PROTOCOL_VERSION: &str = "v1"; +const JOB_TYPE: &str = "profile.cpu"; +pub const DIAGNOSTIC_JOB_SIGNATURE_DOMAIN: &[u8] = b"rustfs-connect-agent-job-v1\0"; +const MAX_JOB_LIFETIME_SECONDS: i64 = 1_800; +const MAX_FUTURE_SKEW_SECONDS: i64 = 300; +const MAX_OUTPUT_BYTES: u64 = 524_288; +const MAX_MEMORY_BYTES: u64 = 64 * 1024 * 1024; +const MAX_CPU_MILLIS: u64 = 30_000; + +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct DiagnosticJobTarget { + pub organization_name: String, + pub cluster_name: String, + pub device_name: String, +} + +#[derive(Clone, Debug, PartialEq, Eq, Deserialize, Serialize)] +#[serde(deny_unknown_fields, rename_all = "camelCase")] +pub struct DiagnosticJobLimits { + pub timeout_seconds: u64, + pub max_output_bytes: u64, + pub max_memory_bytes: u64, + pub max_cpu_millis: u64, +} + +#[derive(Clone, Debug, PartialEq, Eq, Deserialize, Serialize)] +#[serde(deny_unknown_fields, rename_all = "camelCase")] +pub struct DiagnosticJobParameters { + pub artifact_uid: String, + pub consent_uid: String, + pub consent_policy_revision: u64, + pub consent_expires_at: String, + pub duration_millis: u64, + pub sample_period_micros: u64, +} + +#[derive(Clone, Debug, PartialEq, Eq, Deserialize, Serialize)] +#[serde(deny_unknown_fields, rename_all = "camelCase")] +pub struct DiagnosticJobSignature { + algorithm: String, + key_id: String, + value: String, +} + +#[derive(Clone, Debug, PartialEq, Eq, Deserialize, Serialize)] +#[serde(deny_unknown_fields, rename_all = "camelCase")] +pub struct DiagnosticJobEnvelope { + pub job_id: String, + pub protocol_version: String, + pub job_type: String, + pub schema_version: u16, + pub organization_name: String, + pub cluster_name: String, + pub device_name: String, + pub create_time: String, + pub expire_time: String, + pub nonce: String, + pub required_capabilities: Vec, + pub limits: DiagnosticJobLimits, + pub parameters: DiagnosticJobParameters, + pub signature: DiagnosticJobSignature, +} + +#[derive(Serialize)] +#[serde(rename_all = "camelCase")] +struct UnsignedDiagnosticJob<'a> { + job_id: &'a str, + protocol_version: &'a str, + job_type: &'a str, + schema_version: u16, + organization_name: &'a str, + cluster_name: &'a str, + device_name: &'a str, + create_time: &'a str, + expire_time: &'a str, + nonce: &'a str, + required_capabilities: &'a [String], + limits: &'a DiagnosticJobLimits, + parameters: &'a DiagnosticJobParameters, +} + +#[derive(Clone, Debug)] +pub struct TrustedDiagnosticJobSigner { + key_id: String, + key: VerifyingKey, +} + +impl TrustedDiagnosticJobSigner { + pub fn new(key_id: String, public_key: [u8; 32]) -> Result { + if !lower_hex(&key_id, 64) || hex_lower(&Sha256::digest(public_key)) != key_id { + return Err(DiagnosticJobError::TrustInvalid); + } + let key = VerifyingKey::from_bytes(&public_key).map_err(|_| DiagnosticJobError::TrustInvalid)?; + Ok(Self { key_id, key }) + } + + pub fn from_public_key_file(path: &Path, key_id: String) -> Result { + #[cfg(not(unix))] + { + let _ = (path, key_id); + return Err(DiagnosticJobError::TrustInvalid); + } + #[cfg(unix)] + { + let initial = fs::symlink_metadata(path).map_err(|_| DiagnosticJobError::TrustInvalid)?; + if !initial.file_type().is_file() || initial.permissions().mode() & 0o077 != 0 || initial.len() > 256 { + return Err(DiagnosticJobError::TrustInvalid); + } + let mut options = fs::OpenOptions::new(); + options.read(true).custom_flags(libc::O_NOFOLLOW | libc::O_CLOEXEC); + let mut file = options.open(path).map_err(|_| DiagnosticJobError::TrustInvalid)?; + let opened = file.metadata().map_err(|_| DiagnosticJobError::TrustInvalid)?; + if !opened.is_file() + || opened.uid() != rustix::process::geteuid().as_raw() + || opened.dev() != initial.dev() + || opened.ino() != initial.ino() + { + return Err(DiagnosticJobError::TrustInvalid); + } + let mut encoded = Vec::new(); + file.take(257) + .read_to_end(&mut encoded) + .map_err(|_| DiagnosticJobError::TrustInvalid)?; + if encoded.len() > 256 { + return Err(DiagnosticJobError::TrustInvalid); + } + let encoded = std::str::from_utf8(&encoded) + .map_err(|_| DiagnosticJobError::TrustInvalid)? + .trim(); + let public_key = URL_SAFE_NO_PAD + .decode_to_vec(encoded.as_bytes()) + .map_err(|_| DiagnosticJobError::TrustInvalid)?; + if URL_SAFE_NO_PAD.encode_to_string(&public_key) != encoded { + return Err(DiagnosticJobError::TrustInvalid); + } + Self::new(key_id, public_key.try_into().map_err(|_| DiagnosticJobError::TrustInvalid)?) + } + } + + pub fn verify( + &self, + envelope: &DiagnosticJobEnvelope, + target: &DiagnosticJobTarget, + now: DateTime, + ) -> Result { + envelope.validate(target, now)?; + if envelope.signature.algorithm != "Ed25519" || envelope.signature.key_id != self.key_id { + return Err(DiagnosticJobError::SignerUntrusted); + } + let encoded = URL_SAFE_NO_PAD + .decode_to_vec(envelope.signature.value.as_bytes()) + .map_err(|_| DiagnosticJobError::SignatureInvalid)?; + let signature = Signature::from_slice(&encoded).map_err(|_| DiagnosticJobError::SignatureInvalid)?; + let payload = envelope.signing_payload()?; + self.key + .verify(&payload, &signature) + .map_err(|_| DiagnosticJobError::SignatureInvalid)?; + let nonce = URL_SAFE_NO_PAD + .decode_to_vec(envelope.nonce.as_bytes()) + .map_err(|_| DiagnosticJobError::Invalid)? + .try_into() + .map_err(|_| DiagnosticJobError::Invalid)?; + Ok(VerifiedDiagnosticJob { + envelope: envelope.clone(), + nonce, + }) + } +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct VerifiedDiagnosticJob { + envelope: DiagnosticJobEnvelope, + nonce: [u8; 32], +} + +#[derive(Clone, Debug, PartialEq, Eq, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct DiagnosticJobExecution { + pub job_id: String, + pub outcome: &'static str, + pub reason: &'static str, + pub artifact_uid: Option, + pub artifact_sha256: Option, + pub artifact_bytes: Option>, +} + +#[derive(Debug, Error, PartialEq, Eq)] +pub enum DiagnosticJobError { + #[error("connect_diagnostic_job_invalid")] + Invalid, + #[error("connect_diagnostic_job_target_mismatch")] + TargetMismatch, + #[error("connect_diagnostic_job_expired")] + Expired, + #[error("connect_diagnostic_job_unsupported")] + Unsupported, + #[error("connect_diagnostic_job_limit_exceeded")] + LimitExceeded, + #[error("connect_diagnostic_job_trust_invalid")] + TrustInvalid, + #[error("connect_diagnostic_job_signer_untrusted")] + SignerUntrusted, + #[error("connect_diagnostic_job_signature_invalid")] + SignatureInvalid, + #[error("connect_diagnostic_job_encoding_failed")] + Encoding, + #[error("connect_diagnostic_job_cancelled")] + Cancelled, + #[error("connect_diagnostic_job_collection_failed")] + CollectionFailed, +} + +impl DiagnosticJobEnvelope { + fn unsigned(&self) -> UnsignedDiagnosticJob<'_> { + UnsignedDiagnosticJob { + job_id: &self.job_id, + protocol_version: &self.protocol_version, + job_type: &self.job_type, + schema_version: self.schema_version, + organization_name: &self.organization_name, + cluster_name: &self.cluster_name, + device_name: &self.device_name, + create_time: &self.create_time, + expire_time: &self.expire_time, + nonce: &self.nonce, + required_capabilities: &self.required_capabilities, + limits: &self.limits, + parameters: &self.parameters, + } + } + + pub fn signing_payload(&self) -> Result, DiagnosticJobError> { + let payload = serde_json::to_vec(&self.unsigned()).map_err(|_| DiagnosticJobError::Encoding)?; + let mut signed = Vec::with_capacity(DIAGNOSTIC_JOB_SIGNATURE_DOMAIN.len() + payload.len()); + signed.extend_from_slice(DIAGNOSTIC_JOB_SIGNATURE_DOMAIN); + signed.extend_from_slice(&payload); + Ok(signed) + } + + fn validate(&self, target: &DiagnosticJobTarget, now: DateTime) -> Result<(), DiagnosticJobError> { + if self.protocol_version != PROTOCOL_VERSION + || self.job_type != JOB_TYPE + || self.schema_version != PROFILE_SCHEMA_VERSION + || self.required_capabilities != [CPU_PROFILE_CAPABILITY] + { + return Err(DiagnosticJobError::Unsupported); + } + if self.organization_name != target.organization_name + || self.cluster_name != target.cluster_name + || self.device_name != target.device_name + { + return Err(DiagnosticJobError::TargetMismatch); + } + if !uuid7(&self.job_id) + || !uuid7(&self.parameters.artifact_uid) + || !uuid7(&self.parameters.consent_uid) + || self.parameters.consent_policy_revision == 0 + { + return Err(DiagnosticJobError::Invalid); + } + let create = parse_time(&self.create_time)?; + let expire = parse_time(&self.expire_time)?; + let consent_expire = parse_time(&self.parameters.consent_expires_at)?; + if create > now + chrono::Duration::seconds(MAX_FUTURE_SKEW_SECONDS) + || expire <= now + || expire <= create + || expire > create + chrono::Duration::seconds(MAX_JOB_LIFETIME_SECONDS) + || consent_expire < expire + { + return Err(DiagnosticJobError::Expired); + } + let nonce = URL_SAFE_NO_PAD + .decode_to_vec(self.nonce.as_bytes()) + .map_err(|_| DiagnosticJobError::Invalid)?; + if nonce.len() != 32 + || URL_SAFE_NO_PAD.encode_to_string(&nonce) != self.nonce + || self.limits.timeout_seconds == 0 + || self.limits.timeout_seconds > 30 + || self.limits.max_output_bytes == 0 + || self.limits.max_output_bytes > MAX_OUTPUT_BYTES + || self.limits.max_memory_bytes == 0 + || self.limits.max_memory_bytes > MAX_MEMORY_BYTES + || self.limits.max_cpu_millis == 0 + || self.limits.max_cpu_millis > MAX_CPU_MILLIS + || self.parameters.duration_millis == 0 + || self.parameters.duration_millis > self.limits.timeout_seconds.saturating_mul(1_000) + || self.parameters.duration_millis > self.limits.max_cpu_millis + || self.parameters.sample_period_micros == 0 + || self.parameters.sample_period_micros > self.parameters.duration_millis.saturating_mul(1_000) + { + return Err(DiagnosticJobError::LimitExceeded); + } + Ok(()) + } +} + +pub async fn execute_diagnostic_job( + job: VerifiedDiagnosticJob, + identity: &DeviceIdentity, + provenance: ProfileProvenance, + cancel: &CancellationToken, +) -> Result { + if cancel.is_cancelled() { + return Err(DiagnosticJobError::Cancelled); + } + let envelope = job.envelope; + let expire = parse_time(&envelope.expire_time)?; + let consent_expire = parse_time(&envelope.parameters.consent_expires_at)?; + let request = ProfileCaptureRequest { + 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.clone(), + schema_version: envelope.schema_version, + capability: CPU_PROFILE_CAPABILITY.to_owned(), + consent: LocalProfileConsent { + consent_uid: envelope.parameters.consent_uid, + policy_revision: envelope.parameters.consent_policy_revision, + expires_at_unix: consent_expire.timestamp(), + confirmed: true, + }, + produced_at_unix: Utc::now().timestamp(), + expires_at_unix: expire.timestamp(), + nonce: job.nonce, + duration: Duration::from_millis(envelope.parameters.duration_millis), + sample_period: Duration::from_micros(envelope.parameters.sample_period_micros), + provenance, + }; + let result = capture_cpu_profile(&request, cancel).await.map_err(|error| match error { + super::ProfileError::Cancelled => DiagnosticJobError::Cancelled, + _ => DiagnosticJobError::CollectionFailed, + })?; + let export = encode_signed_profile_export(&request, &result, identity, cancel).map_err(|error| match error { + super::ProfileError::Cancelled => DiagnosticJobError::Cancelled, + _ => DiagnosticJobError::CollectionFailed, + })?; + if export.archive_bytes.len() > usize::try_from(envelope.limits.max_output_bytes).unwrap_or(usize::MAX) { + return Err(DiagnosticJobError::LimitExceeded); + } + let outcome = result.outcome().as_str(); + let reason = result.reason_code().as_str(); + 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), + }) +} + +fn parse_time(value: &str) -> Result, DiagnosticJobError> { + DateTime::parse_from_rfc3339(value) + .map(|value| value.with_timezone(&Utc)) + .map_err(|_| DiagnosticJobError::Invalid) +} + +fn uuid7(value: &str) -> bool { + Uuid::parse_str(value).is_ok_and(|uid| { + uid.get_variant() == Variant::RFC4122 && uid.get_version() == Some(Version::SortRand) && uid.to_string() == value + }) +} + +fn lower_hex(value: &str, length: usize) -> bool { + value.len() == length + && value + .bytes() + .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte)) +} + +fn hex_lower(bytes: &[u8]) -> String { + hex_simd::encode_to_string(bytes, hex_simd::AsciiCase::Lower) +} + +#[cfg(test)] +mod tests { + use super::*; + use ed25519_dalek::{Signer as _, SigningKey}; + + fn envelope() -> DiagnosticJobEnvelope { + DiagnosticJobEnvelope { + job_id: "018cc251-f400-7abc-8def-0123456789ab".to_owned(), + protocol_version: "v1".to_owned(), + job_type: "profile.cpu".to_owned(), + schema_version: 1, + organization_name: "organizations/018cc251-f400-7abc-8def-0123456789ab".to_owned(), + cluster_name: "organizations/018cc251-f400-7abc-8def-0123456789ab/clusters/018cc251-f400-7abc-8def-0123456789ac".to_owned(), + device_name: "organizations/018cc251-f400-7abc-8def-0123456789ab/clusters/018cc251-f400-7abc-8def-0123456789ac/clusterDevices/018cc251-f400-7abc-8def-0123456789ad".to_owned(), + create_time: "2030-01-01T00:00:00Z".to_owned(), + expire_time: "2030-01-01T00:00:30Z".to_owned(), + nonce: URL_SAFE_NO_PAD.encode_to_string([7_u8; 32]), + required_capabilities: vec![CPU_PROFILE_CAPABILITY.to_owned()], + limits: DiagnosticJobLimits { + timeout_seconds: 30, + max_output_bytes: 524_288, + max_memory_bytes: 64 * 1024 * 1024, + max_cpu_millis: 30_000, + }, + parameters: DiagnosticJobParameters { + artifact_uid: "018cc251-f400-7abc-8def-0123456789ae".to_owned(), + consent_uid: "018cc251-f400-7abc-8def-0123456789af".to_owned(), + consent_policy_revision: 1, + consent_expires_at: "2030-01-01T00:01:00Z".to_owned(), + duration_millis: 1_000, + sample_period_micros: 10_000, + }, + signature: DiagnosticJobSignature { + algorithm: "Ed25519".to_owned(), + key_id: String::new(), + value: String::new(), + }, + } + } + + fn signed() -> (DiagnosticJobEnvelope, TrustedDiagnosticJobSigner) { + let signing = SigningKey::from_bytes(&[9_u8; 32]); + let key_id = hex_lower(&Sha256::digest(signing.verifying_key().as_bytes())); + let trusted = + TrustedDiagnosticJobSigner::new(key_id.clone(), *signing.verifying_key().as_bytes()).expect("trusted signer"); + let mut envelope = envelope(); + envelope.signature.key_id = key_id; + envelope.signature.value = + URL_SAFE_NO_PAD.encode_to_string(signing.sign(&envelope.signing_payload().expect("payload")).to_bytes()); + (envelope, trusted) + } + + fn target(envelope: &DiagnosticJobEnvelope) -> DiagnosticJobTarget { + DiagnosticJobTarget { + organization_name: envelope.organization_name.clone(), + cluster_name: envelope.cluster_name.clone(), + device_name: envelope.device_name.clone(), + } + } + + #[test] + fn accepts_a_bounded_signed_profile_job_for_the_exact_device() { + let (envelope, signer) = signed(); + signer + .verify(&envelope, &target(&envelope), "2030-01-01T00:00:10Z".parse().expect("time")) + .expect("valid job"); + } + + #[test] + fn rejects_tampering_cross_device_replay_and_expiry() { + let (envelope, signer) = signed(); + let mut tampered = envelope.clone(); + tampered.parameters.duration_millis = 2_000; + assert_eq!( + signer.verify(&tampered, &target(&tampered), "2030-01-01T00:00:10Z".parse().expect("time")), + Err(DiagnosticJobError::SignatureInvalid) + ); + let mut wrong_target = target(&envelope); + wrong_target.device_name.push('0'); + assert_eq!( + signer.verify(&envelope, &wrong_target, "2030-01-01T00:00:10Z".parse().expect("time")), + Err(DiagnosticJobError::TargetMismatch) + ); + assert_eq!( + signer.verify(&envelope, &target(&envelope), "2030-01-01T00:00:30Z".parse().expect("time")), + Err(DiagnosticJobError::Expired) + ); + } + + #[test] + fn rejects_unbounded_and_non_allow_listed_jobs_before_signature_use() { + let (mut envelope, signer) = signed(); + envelope.limits.max_output_bytes += 1; + assert_eq!( + signer.verify(&envelope, &target(&envelope), "2030-01-01T00:00:10Z".parse().expect("time")), + Err(DiagnosticJobError::LimitExceeded) + ); + envelope = signed().0; + envelope.job_type = "shell.exec".to_owned(); + assert_eq!( + signer.verify(&envelope, &target(&envelope), "2030-01-01T00:00:10Z".parse().expect("time")), + Err(DiagnosticJobError::Unsupported) + ); + } +} diff --git a/rustfs/src/connect/diagnostics/mod.rs b/rustfs/src/connect/diagnostics/mod.rs index 81c7c914c..03012ab51 100644 --- a/rustfs/src/connect/diagnostics/mod.rs +++ b/rustfs/src/connect/diagnostics/mod.rs @@ -13,6 +13,7 @@ // limitations under the License. mod inspect; +mod job; mod logs; mod perf_client; mod perf_drive; @@ -44,6 +45,10 @@ pub use inspect::{ InspectOutcome, InspectProvenance, InspectReason, InspectReasonCode, InspectRequest, InspectRule, InspectRuleOutcome, InspectRun, Reconstruction, SavedInspectExport, SignedInspectExport, export_inspect_summary, save_signed_inspect_export, }; +pub use job::{ + DIAGNOSTIC_JOB_SIGNATURE_DOMAIN, DiagnosticJobEnvelope, DiagnosticJobError, DiagnosticJobExecution, DiagnosticJobLimits, + DiagnosticJobParameters, DiagnosticJobTarget, TrustedDiagnosticJobSigner, VerifiedDiagnosticJob, execute_diagnostic_job, +}; pub use logs::{ CaptureMode, LOGS_CAPABILITY, LOGS_SCHEMA_VERSION, LocalLogConsent, LogCaptureError, LogCaptureRequest, LogProvenance, SavedLogExport, SignedLogExport, export_logs, save_signed_log_export, diff --git a/rustfs/src/connect/mod.rs b/rustfs/src/connect/mod.rs index 428267f22..4218bc43f 100644 --- a/rustfs/src/connect/mod.rs +++ b/rustfs/src/connect/mod.rs @@ -81,6 +81,10 @@ pub use diagnostics::{ save_signed_telemetry_export, sign_client_export, sign_drive_export, spawn_environment_schedule, validate_client_limits, validate_drive_limits, }; +pub use diagnostics::{ + DIAGNOSTIC_JOB_SIGNATURE_DOMAIN, DiagnosticJobEnvelope, DiagnosticJobError, DiagnosticJobExecution, DiagnosticJobLimits, + DiagnosticJobParameters, DiagnosticJobTarget, TrustedDiagnosticJobSigner, VerifiedDiagnosticJob, execute_diagnostic_job, +}; pub use diagnostics::{ LocalNetworkConsent, MAX_NETWORK_ARCHIVE_BYTES, MAX_NETWORK_BANDWIDTH_BYTES_PER_SECOND, MAX_NETWORK_DECOMPRESSED_BYTES, MAX_NETWORK_DURATION, MAX_NETWORK_ENVELOPE_BYTES, MAX_NETWORK_OPERATIONS, MAX_NETWORK_PEERS, MAX_NETWORK_RESULT_BYTES,