feat(connect): export signed environment inventory (#7887)

feat: export signed environment inventory
This commit is contained in:
Chris
2026-09-15 03:02:22 +08:00
committed by GitHub
parent a015d8a181
commit f93237aa95
5 changed files with 495 additions and 7 deletions
+19
View File
@@ -348,6 +348,25 @@ pub struct ConnectEnvironmentInventoryOpts {
/// Confirm this explicit local L1 inventory operation
#[arg(long = "acknowledge-l1", required = true, action = clap::ArgAction::SetTrue)]
pub acknowledge_l1: bool,
/// New local signed archive path; omit to print the four-field JSON inventory
#[arg(long)]
pub output: Option<PathBuf>,
#[arg(long)]
pub organization: Option<String>,
#[arg(long)]
pub cluster: Option<String>,
#[arg(long)]
pub device: Option<String>,
#[arg(long = "run-uid")]
pub run_uid: Option<String>,
#[arg(long = "artifact-uid")]
pub artifact_uid: Option<String>,
#[arg(long = "consent-uid")]
pub consent_uid: Option<String>,
#[arg(long = "policy-revision")]
pub policy_revision: Option<u64>,
#[arg(long = "expires-at")]
pub expires_at_unix: Option<i64>,
}
#[derive(Args, Clone)]
+414 -1
View File
@@ -15,16 +15,26 @@
//! Bounded, identifier-free deployment environment inventory.
use std::collections::BTreeSet;
use std::fs::{self, File, OpenOptions};
use std::io::{Cursor, Write as _};
use std::path::Path;
use std::sync::{Arc, LazyLock};
use std::time::Duration;
use base64_simd::URL_SAFE_NO_PAD;
use p256::ecdsa::{Signature, SigningKey, signature::Signer as _};
use p256::pkcs8::DecodePrivateKey as _;
use serde::Serialize;
use sha2::{Digest as _, Sha256};
use sysinfo::{Disks, Networks, RefreshKind, System};
use thiserror::Error;
use time::{OffsetDateTime, format_description::well_known::Rfc3339};
use tokio::sync::Semaphore;
use tokio_util::sync::CancellationToken;
use uuid::{Uuid, Variant, Version};
use zip::{CompressionMethod, ZipWriter, write::SimpleFileOptions};
use super::inventory::InventorySnapshot;
use super::{DeviceIdentity, inventory::InventorySnapshot};
pub const ENVIRONMENT_CAPABILITY: &str = "inventory.environment@1";
pub const ENVIRONMENT_SCHEMA_VERSION: u16 = 1;
@@ -136,6 +146,347 @@ impl EnvironmentInventory {
}
}
const SIGNATURE_DOMAIN: &[u8] = b"rustfs-diagnostic-envelope-v1\0";
const OUTPUT_MODE: u32 = 0o600;
#[derive(Clone, Debug)]
pub struct EnvironmentExportRequest {
pub confirmed: bool,
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 produced_at_unix: i64,
pub expires_at_unix: i64,
pub nonce: [u8; 32],
pub source_commit: String,
pub executable_sha256: String,
pub rustfs_version: String,
pub build_features: Vec<String>,
}
#[derive(Clone, Debug)]
pub struct SignedEnvironmentExport {
pub artifact_uid: String,
pub archive_bytes: Vec<u8>,
pub archive_sha256: String,
}
#[derive(Clone, Debug)]
pub struct SavedEnvironmentExport {
pub artifact_uid: String,
pub archive_size_bytes: u64,
pub archive_sha256: String,
}
#[derive(Serialize)]
#[serde(rename_all = "camelCase")]
struct EnvironmentResult<'a> {
schema_version: u16,
run_uid: &'a str,
tool_id: &'static str,
capability: &'static str,
outcome: &'static str,
reason_code: &'static str,
duration_millis: u64,
provenance: EnvironmentProvenance<'a>,
coverage: EnvironmentCoverage,
data: &'a EnvironmentInventory,
}
#[derive(Serialize)]
#[serde(rename_all = "camelCase")]
struct EnvironmentProvenance<'a> {
repository: &'static str,
source_commit: &'a str,
executable_sha256: &'a str,
rustfs_version: &'a str,
os_family: EnvironmentOsFamily,
architecture: &'static str,
build_features: &'a [String],
}
#[derive(Serialize)]
#[serde(rename_all = "camelCase")]
struct EnvironmentCoverage {
requested_units: u8,
completed_units: u8,
unit: &'static str,
}
#[derive(Serialize)]
#[serde(rename_all = "camelCase")]
struct EnvironmentEnvelope<'a> {
format_version: &'static str,
protocol_version: &'static str,
organization_name: &'a str,
cluster_name: &'a str,
device_name: &'a str,
run_uid: &'a str,
artifact_uid: &'a str,
tool_id: &'static str,
schema_version: u16,
classification: &'static str,
consent_uid: &'a str,
policy_revision: u64,
produced_at: String,
expires_at: String,
nonce: String,
device_key_id: &'a str,
payload: EnvironmentPayload,
}
#[derive(Serialize)]
#[serde(rename_all = "camelCase")]
struct EnvironmentPayload {
path: &'static str,
media_type: &'static str,
size_bytes: u64,
sha256: String,
}
#[derive(Serialize)]
#[serde(rename_all = "camelCase")]
struct EnvironmentSignature<'a> {
algorithm: &'static str,
key_id: &'a str,
value: String,
}
pub fn sign_environment_inventory(
inventory: &EnvironmentInventory,
request: &EnvironmentExportRequest,
key: &DeviceIdentity,
duration: Duration,
cancel: &CancellationToken,
) -> Result<SignedEnvironmentExport, EnvironmentError> {
if cancel.is_cancelled() {
return Err(EnvironmentError::Cancelled);
}
let now = OffsetDateTime::now_utc().unix_timestamp();
if !request.confirmed
|| duration.is_zero()
|| duration > MAX_ENVIRONMENT_DURATION
|| request.policy_revision == 0
|| request.produced_at_unix > now.saturating_add(300)
|| request.expires_at_unix <= now
|| request.expires_at_unix <= request.produced_at_unix
|| request.expires_at_unix - request.produced_at_unix > 2_592_000
|| !uuid7(&request.run_uid)
|| !uuid7(&request.artifact_uid)
|| !uuid7(&request.consent_uid)
|| !request.organization_name.starts_with("organizations/")
|| !request
.cluster_name
.starts_with(&(request.organization_name.clone() + "/clusters/"))
|| !request
.device_name
.starts_with(&(request.cluster_name.clone() + "/clusterDevices/"))
|| !lower_hex(&request.source_commit, 40)
|| !lower_hex(&request.executable_sha256, 64)
|| !version(&request.rustfs_version)
|| request.build_features.len() > 64
|| !request.build_features.iter().all(|feature| build_feature(feature))
{
return Err(EnvironmentError::InvalidExport);
}
let result = EnvironmentResult {
schema_version: 1,
run_uid: &request.run_uid,
tool_id: "inventory.environment",
capability: ENVIRONMENT_CAPABILITY,
outcome: "SUCCEEDED",
reason_code: "COMPLETE",
duration_millis: u64::try_from(duration.as_millis()).unwrap_or(30_000).clamp(1, 30_000),
provenance: EnvironmentProvenance {
repository: "rustfs/rustfs",
source_commit: &request.source_commit,
executable_sha256: &request.executable_sha256,
rustfs_version: &request.rustfs_version,
os_family: EnvironmentOsFamily::current(),
architecture: match std::env::consts::ARCH {
"x86_64" => "x86_64",
"aarch64" => "aarch64",
_ => "other",
},
build_features: &request.build_features,
},
coverage: EnvironmentCoverage {
requested_units: 1,
completed_units: 1,
unit: "RESOURCE",
},
data: inventory,
};
let result_bytes = serde_json::to_vec(&result).map_err(|_| EnvironmentError::ExportEncoding)?;
let key_id = hex_lower(&Sha256::digest(key.public_key_der()));
let envelope = EnvironmentEnvelope {
format_version: "rustfs.connect.diagnosticEnvelope/1",
protocol_version: "v1",
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,
tool_id: "inventory.environment",
schema_version: 1,
classification: "L1",
consent_uid: &request.consent_uid,
policy_revision: request.policy_revision,
produced_at: timestamp(request.produced_at_unix)?,
expires_at: timestamp(request.expires_at_unix)?,
nonce: URL_SAFE_NO_PAD.encode_to_string(request.nonce),
device_key_id: &key_id,
payload: EnvironmentPayload {
path: "result.json",
media_type: "application/json",
size_bytes: result_bytes.len() as u64,
sha256: hex_lower(&Sha256::digest(&result_bytes)),
},
};
let envelope_bytes = serde_json::to_vec(&envelope).map_err(|_| EnvironmentError::ExportEncoding)?;
let pkcs8 = key.to_pkcs8_der().map_err(|_| EnvironmentError::ExportSigning)?;
let signing_key = SigningKey::from_pkcs8_der(pkcs8.as_slice()).map_err(|_| EnvironmentError::ExportSigning)?;
let mut input = Vec::with_capacity(SIGNATURE_DOMAIN.len() + envelope_bytes.len());
input.extend_from_slice(SIGNATURE_DOMAIN);
input.extend_from_slice(&envelope_bytes);
let signature: Signature = signing_key.sign(&input);
let signature_bytes = serde_json::to_vec(&EnvironmentSignature {
algorithm: "ES256",
key_id: &key_id,
value: URL_SAFE_NO_PAD.encode_to_string(signature.normalize_s().to_bytes()),
})
.map_err(|_| EnvironmentError::ExportEncoding)?;
if cancel.is_cancelled() {
return Err(EnvironmentError::Cancelled);
}
let cursor = Cursor::new(Vec::new());
let mut zip = ZipWriter::new(cursor);
let options = SimpleFileOptions::DEFAULT
.compression_method(CompressionMethod::Stored)
.unix_permissions(OUTPUT_MODE);
for (name, bytes) in [
("envelope.json", envelope_bytes.as_slice()),
("envelope.sig", signature_bytes.as_slice()),
("result.json", result_bytes.as_slice()),
] {
zip.start_file(name, options).map_err(|_| EnvironmentError::ExportEncoding)?;
zip.write_all(bytes).map_err(|_| EnvironmentError::ExportIo)?;
}
let archive_bytes = zip.finish().map_err(|_| EnvironmentError::ExportEncoding)?.into_inner();
if archive_bytes.len() > 65_536 {
return Err(EnvironmentError::InvalidExport);
}
Ok(SignedEnvironmentExport {
artifact_uid: request.artifact_uid.clone(),
archive_sha256: hex_lower(&Sha256::digest(&archive_bytes)),
archive_bytes,
})
}
pub fn save_signed_environment_export(
output: &Path,
export: &SignedEnvironmentExport,
cancel: &CancellationToken,
) -> Result<SavedEnvironmentExport, EnvironmentError> {
if cancel.is_cancelled() {
return Err(EnvironmentError::Cancelled);
}
let parent = output
.parent()
.filter(|p| !p.as_os_str().is_empty())
.unwrap_or_else(|| Path::new("."));
let temporary = parent.join(format!(
".{}.{}.partial",
output.file_name().ok_or(EnvironmentError::InvalidExport)?.to_string_lossy(),
export.artifact_uid
));
let mut options = OpenOptions::new();
options.write(true).create_new(true);
#[cfg(unix)]
{
use std::os::unix::fs::OpenOptionsExt as _;
options.mode(OUTPUT_MODE);
}
let mut file = options.open(&temporary).map_err(|_| EnvironmentError::ExportIo)?;
let saved = (|| {
file.write_all(&export.archive_bytes)
.map_err(|_| EnvironmentError::ExportIo)?;
if cancel.is_cancelled() {
return Err(EnvironmentError::Cancelled);
}
file.sync_all().map_err(|_| EnvironmentError::ExportIo)?;
fs::hard_link(&temporary, output).map_err(|_| EnvironmentError::ExportIo)?;
fs::remove_file(&temporary).map_err(|_| EnvironmentError::ExportIo)?;
#[cfg(unix)]
File::open(parent)
.and_then(|d| d.sync_all())
.map_err(|_| EnvironmentError::ExportIo)?;
Ok(SavedEnvironmentExport {
artifact_uid: export.artifact_uid.clone(),
archive_size_bytes: export.archive_bytes.len() as u64,
archive_sha256: export.archive_sha256.clone(),
})
})();
if saved.is_err() {
let _ = fs::remove_file(&temporary);
}
saved
}
fn timestamp(unix: i64) -> Result<String, EnvironmentError> {
OffsetDateTime::from_unix_timestamp(unix)
.map_err(|_| EnvironmentError::InvalidExport)?
.format(&Rfc3339)
.map_err(|_| EnvironmentError::ExportEncoding)
}
fn uuid7(value: &str) -> bool {
Uuid::parse_str(value).is_ok_and(|u| u.get_version() == Some(Version::SortRand) && u.get_variant() == Variant::RFC4122)
}
fn lower_hex(value: &str, len: usize) -> bool {
value.len() == len && value.bytes().all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
}
fn version(value: &str) -> bool {
if value.is_empty()
|| value.len() > 64
|| !value
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'-'))
{
return false;
}
let (core, suffix) = value
.split_once('-')
.map_or((value, None), |(core, suffix)| (core, Some(suffix)));
if suffix.is_some_and(str::is_empty) {
return false;
}
let mut parts = core.split('.');
parts.clone().count() == 3 && parts.all(|part| !part.is_empty() && part.bytes().all(|byte| byte.is_ascii_digit()))
}
fn build_feature(value: &str) -> bool {
!value.is_empty()
&& value.len() <= 64
&& value.as_bytes()[0].is_ascii_lowercase()
&& value
.bytes()
.all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || matches!(byte, b'_' | b'-'))
}
fn hex_lower(bytes: &[u8]) -> String {
const H: &[u8; 16] = b"0123456789abcdef";
let mut out = String::with_capacity(bytes.len() * 2);
for b in bytes {
out.push(H[(b >> 4) as usize] as char);
out.push(H[(b & 15) as usize] as char);
}
out
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum EnvironmentSource {
Filesystem,
@@ -160,6 +511,14 @@ pub enum EnvironmentError {
TimedOut,
#[error("inventory_environment_task_failed")]
TaskFailed,
#[error("inventory_environment_invalid_export")]
InvalidExport,
#[error("inventory_environment_export_encoding_failed")]
ExportEncoding,
#[error("inventory_environment_export_signing_failed")]
ExportSigning,
#[error("inventory_environment_export_io_failed")]
ExportIo,
}
#[derive(Debug, PartialEq, Eq)]
@@ -295,6 +654,7 @@ mod test_support {
#[cfg(test)]
mod tests {
use std::io::Read as _;
use std::sync::atomic::Ordering;
use super::*;
@@ -349,4 +709,57 @@ mod tests {
Err(EnvironmentError::SourceUnavailable(EnvironmentSource::Filesystem))
);
}
#[test]
fn signed_export_contains_only_the_allow_list_and_never_clobbers() {
let inventory = EnvironmentInventory {
node_count: 1,
drive_count: 0,
os_family: EnvironmentOsFamily::Linux,
filesystem_types: vec![EnvironmentFilesystemType::Xfs],
};
let now = OffsetDateTime::now_utc().unix_timestamp();
let request = EnvironmentExportRequest {
confirmed: true,
organization_name: "organizations/019e3ae0-0000-7000-8000-000000000001".into(),
cluster_name: "organizations/019e3ae0-0000-7000-8000-000000000001/clusters/019e3ae0-0000-7000-8000-000000000002".into(),
device_name: "organizations/019e3ae0-0000-7000-8000-000000000001/clusters/019e3ae0-0000-7000-8000-000000000002/clusterDevices/019e3ae0-0000-7000-8000-000000000003".into(),
run_uid: "019e3ae0-0000-7000-8000-000000000004".into(),
artifact_uid: "019e3ae0-0000-7000-8000-000000000005".into(),
consent_uid: "019e3ae0-0000-7000-8000-000000000006".into(),
policy_revision: 1,
produced_at_unix: now,
expires_at_unix: now + 60,
nonce: [7; 32],
source_commit: "a".repeat(40),
executable_sha256: "b".repeat(64),
rustfs_version: "1.0.0-rc.6".into(),
build_features: vec![],
};
let cancel = CancellationToken::new();
let mut unconfirmed = request.clone();
unconfirmed.confirmed = false;
assert!(matches!(
sign_environment_inventory(&inventory, &unconfirmed, &DeviceIdentity::generate(), Duration::from_millis(1), &cancel,),
Err(EnvironmentError::InvalidExport)
));
let export =
sign_environment_inventory(&inventory, &request, &DeviceIdentity::generate(), Duration::from_millis(1), &cancel)
.expect("signed export");
let mut zip = zip::ZipArchive::new(Cursor::new(export.archive_bytes.as_slice())).expect("zip");
let mut result = String::new();
zip.by_name("result.json")
.expect("result")
.read_to_string(&mut result)
.expect("read");
let value: serde_json::Value = serde_json::from_str(&result).expect("json");
assert_eq!(value["data"]["driveCount"], 0);
assert_eq!(value["data"].as_object().expect("data").len(), 4);
assert!(!result.contains("secret"));
let directory = tempfile::tempdir().expect("tempdir");
let output = directory.path().join("environment.zip");
save_signed_environment_export(&output, &export, &cancel).expect("save");
assert!(save_signed_environment_export(&output, &export, &cancel).is_err());
}
}
+3 -2
View File
@@ -125,8 +125,9 @@ pub(crate) use diagnostics::{
LocalTraceCaptureError, LocalTraceCaptureRuntime, request_local_trace_capture, spawn_local_trace_capture_runtime,
};
pub use environment::{
ENVIRONMENT_CAPABILITY, ENVIRONMENT_SCHEMA_VERSION, EnvironmentCollectionRequest, EnvironmentError,
EnvironmentFilesystemType, EnvironmentInventory, EnvironmentOsFamily, MAX_ENVIRONMENT_DURATION, collect_environment,
ENVIRONMENT_CAPABILITY, ENVIRONMENT_SCHEMA_VERSION, EnvironmentCollectionRequest, EnvironmentError, EnvironmentExportRequest,
EnvironmentFilesystemType, EnvironmentInventory, EnvironmentOsFamily, MAX_ENVIRONMENT_DURATION, SavedEnvironmentExport,
SignedEnvironmentExport, collect_environment, save_signed_environment_export, sign_environment_inventory,
};
pub use heartbeat::{CoarseNodeSummary, HeartbeatError, HeartbeatStatus};
pub use identity::{DeviceIdentity, IdentityError, RegistrationProof, RegistrationTranscript};
+7 -3
View File
@@ -165,9 +165,13 @@ impl From<EnvironmentError> for CollectorError {
EnvironmentError::TimedOut => Self::TimedOut,
EnvironmentError::TaskFailed => Self::TaskFailed,
EnvironmentError::SourceUnavailable(_) => Self::SourceUnavailable,
EnvironmentError::UnsupportedVersion | EnvironmentError::UnsupportedCapability | EnvironmentError::InvalidTimeout => {
Self::TaskFailed
}
EnvironmentError::UnsupportedVersion
| EnvironmentError::UnsupportedCapability
| EnvironmentError::InvalidTimeout
| EnvironmentError::InvalidExport
| EnvironmentError::ExportEncoding
| EnvironmentError::ExportSigning
| EnvironmentError::ExportIo => Self::TaskFailed,
}
}
}
+52 -1
View File
@@ -308,7 +308,58 @@ async fn execute_connect_environment_inventory(options: ConnectEnvironmentInvent
error => Error::other(error),
})?,
};
println!("{}", serde_json::to_string(&inventory).map_err(Error::other)?);
let Some(output) = options.output else {
println!("{}", serde_json::to_string(&inventory).map_err(Error::other)?);
return Ok(());
};
use crate::connect::{EnvironmentExportRequest, IdentityStore, save_signed_environment_export, sign_environment_inventory};
use rand::{TryRng as _, rngs::SysRng};
let required =
|value: Option<String>, name: &str| value.ok_or_else(|| Error::other(format!("--{name} is required with --output")));
let key = IdentityStore::new(options.state_dir.join("identity"))
.load()
.map_err(Error::other)?
.ok_or_else(|| Error::other("connect environment export requires an enrolled device identity"))?;
let produced_at_unix = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map_err(Error::other)
.and_then(|d| i64::try_from(d.as_secs()).map_err(Error::other))?;
let mut nonce = [0_u8; 32];
SysRng.try_fill_bytes(&mut nonce).map_err(Error::other)?;
let request = EnvironmentExportRequest {
confirmed: options.acknowledge_l1,
organization_name: required(options.organization, "organization")?,
cluster_name: required(options.cluster, "cluster")?,
device_name: required(options.device, "device")?,
run_uid: required(options.run_uid, "run-uid")?,
artifact_uid: required(options.artifact_uid, "artifact-uid")?,
consent_uid: required(options.consent_uid, "consent-uid")?,
policy_revision: options
.policy_revision
.ok_or_else(|| Error::other("--policy-revision is required with --output"))?,
produced_at_unix,
expires_at_unix: options
.expires_at_unix
.ok_or_else(|| Error::other("--expires-at is required with --output"))?,
nonce,
source_commit: crate::version::build::COMMIT_HASH.to_owned(),
executable_sha256: hash_current_executable()?,
rustfs_version: env!("CARGO_PKG_VERSION").to_owned(),
build_features: enabled_build_features(),
};
let export = sign_environment_inventory(&inventory, &request, &key, Duration::from_secs(options.timeout_seconds), &cancel)
.map_err(Error::other)?;
let writer_cancel = cancel.clone();
let receipt = tokio::task::spawn_blocking(move || save_signed_environment_export(&output, &export, &writer_cancel))
.await
.map_err(Error::other)?
.map_err(Error::other)?;
println!("tool=inventory.environment outcome=SUCCEEDED reason=COMPLETE");
println!(
"artifact={} bytes={} sha256={}",
receipt.artifact_uid, receipt.archive_size_bytes, receipt.archive_sha256
);
println!("upload=not-performed");
Ok(())
}