feat(connect): select offline key for client performance (#8242)

This commit is contained in:
Chris
2026-09-29 20:49:59 +08:00
committed by GitHub
parent b0a54f0e5c
commit 1cca0dd25c
8 changed files with 218 additions and 10 deletions
+4
View File
@@ -613,6 +613,10 @@ pub struct ConnectClientPerformanceOpts {
#[arg(long = "state-dir")]
pub state_dir: PathBuf,
/// Select an existing separately enrolled offline key by its SPKI SHA-256 identifier
#[arg(long = "offline-key-id", value_parser = NonEmptyStringValueParser::new())]
pub offline_key_id: Option<String>,
/// RustFS deployment endpoint
#[arg(long, value_parser = NonEmptyStringValueParser::new())]
pub endpoint: String,
+49
View File
@@ -128,6 +128,55 @@ mod tests {
assert!(Opt::parse_command(approved.into_iter().chain(["--peer-address", "example.com"])).is_err());
}
#[test]
fn connect_client_cli_selects_offline_key_only_when_requested() {
let base = [
"rustfs",
"connect",
"performance",
"client",
"--state-dir",
"/state",
"--endpoint",
"https://storage.example",
"--access-key-file",
"/access",
"--secret-key-file",
"/secret",
"--output",
"/output.zip",
"--organization",
"org",
"--cluster",
"cluster",
"--device",
"device",
"--run-uid",
"run",
"--artifact-uid",
"artifact",
"--consent-uid",
"consent",
"--policy-revision",
"1",
"--consent-expires-at",
"10",
"--expires-at",
"10",
"--operation",
"get",
"--acknowledge-l1",
];
assert!(matches!(
Opt::parse_command(base),
Ok(CommandResult::ConnectClientPerformance(options)) if options.offline_key_id.is_none()
));
assert!(matches!(
Opt::parse_command(base.into_iter().chain(["--offline-key-id", "selected"])),
Ok(CommandResult::ConnectClientPerformance(options)) if options.offline_key_id.as_deref() == Some("selected")
));
}
#[test]
#[serial]
fn test_tls_inspect_subcommand_parses_tls_path_alias() {
+3 -3
View File
@@ -172,9 +172,9 @@ pub use trace_record::{
};
pub use trace_replay::{LocallyReviewedTraceArtifact, ReplayedTrace, TraceReplayError, replay_trace, replay_trace_result};
pub(crate) use trace_runtime::{
LocalHealthRequest, LocalNetworkRequest, LocalTraceCaptureError, LocalTraceCaptureRuntime, request_local_health,
request_local_native_threads_profile, request_local_network, request_local_runtime_profile, request_local_top_api,
request_local_top_disk, request_local_top_locks, request_local_top_rpc, request_local_trace_capture,
LocalHealthRequest, LocalNetworkRequest, LocalTraceCaptureError, LocalTraceCaptureRuntime, load_selected_offline_key,
request_local_health, request_local_native_threads_profile, request_local_network, request_local_runtime_profile,
request_local_top_api, request_local_top_disk, request_local_top_locks, request_local_top_rpc, request_local_trace_capture,
spawn_local_trace_capture_runtime,
};
@@ -93,6 +93,8 @@ pub(crate) enum LocalTraceCaptureError {
Network(String),
#[error("local diagnostic cancellation was not acknowledged")]
CancellationUnconfirmed,
#[error("selected offline identity is unavailable or unsafe")]
OfflineIdentity,
}
pub(crate) struct LocalTraceCaptureRuntime {
@@ -1441,6 +1443,20 @@ fn load_offline_key(state_root: &Path, offline_key_id: &str) -> Result<crate::co
Ok(key)
}
pub(crate) fn load_selected_offline_key(
state_root: &Path,
offline_key_id: &str,
) -> Result<crate::connect::DeviceIdentity, LocalTraceCaptureError> {
if offline_key_id.len() != 64
|| !offline_key_id
.bytes()
.all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
{
return Err(LocalTraceCaptureError::OfflineIdentity);
}
load_offline_key(state_root, offline_key_id).map_err(|_| LocalTraceCaptureError::OfflineIdentity)
}
async fn capture_local_runtime_profile(
state_root: &Path,
protocol_version: u16,
@@ -2100,6 +2116,48 @@ mod tests {
}
}
#[test]
fn selected_offline_key_never_falls_back_or_follows_unsafe_state() {
use std::os::unix::fs::symlink;
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 store = crate::connect::OfflineKeyStore::new(state.path());
assert_eq!(
super::load_selected_offline_key(state.path(), &request.offline_key_id)
.unwrap()
.public_key_der(),
key.public_key_der()
);
assert!(super::load_selected_offline_key(state.path(), &"0".repeat(64)).is_err());
let path = store.key_path();
std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o644)).unwrap();
assert!(super::load_selected_offline_key(state.path(), &request.offline_key_id).is_err());
std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o600)).unwrap();
std::fs::set_permissions(state.path(), std::fs::Permissions::from_mode(0o755)).unwrap();
assert!(super::load_selected_offline_key(state.path(), &request.offline_key_id).is_err());
std::fs::set_permissions(state.path(), std::fs::Permissions::from_mode(0o700)).unwrap();
let target = state.path().join("key-target");
std::fs::rename(&path, &target).unwrap();
crate::connect::IdentityStore::new(state.path().join("identity"))
.load_or_create()
.unwrap();
assert!(super::load_selected_offline_key(state.path(), &request.offline_key_id).is_err());
symlink(&target, &path).unwrap();
assert!(super::load_selected_offline_key(state.path(), &request.offline_key_id).is_err());
std::fs::remove_file(&path).unwrap();
std::fs::rename(&target, &path).unwrap();
std::fs::set_permissions(state.path().join("offline"), std::fs::Permissions::from_mode(0o777)).unwrap();
assert!(super::load_selected_offline_key(state.path(), &request.offline_key_id).is_err());
std::fs::set_permissions(state.path().join("offline"), std::fs::Permissions::from_mode(0o755)).unwrap();
let linked_root = state.path().join("linked-root");
symlink(state.path(), &linked_root).unwrap();
assert!(super::load_selected_offline_key(&linked_root, &request.offline_key_id).is_err());
}
#[tokio::test]
#[serial]
async fn local_network_rejects_unapproved_expired_and_unbounded_work_before_using_key() {
@@ -50,6 +50,13 @@ pub(crate) async fn request_local_network(
Err(LocalTraceCaptureError::RuntimeUnavailable)
}
pub(crate) fn load_selected_offline_key(
_state_root: &Path,
_offline_key_id: &str,
) -> Result<crate::connect::DeviceIdentity, LocalTraceCaptureError> {
Err(LocalTraceCaptureError::RuntimeUnavailable)
}
#[derive(Debug, Error)]
pub(crate) enum LocalTraceCaptureError {
#[error("telemetry server runtime is unavailable")]
+3 -3
View File
@@ -94,9 +94,9 @@ pub use diagnostics::{
save_signed_health_export, sign_health_export,
};
pub(crate) use diagnostics::{
LocalHealthRequest, LocalNetworkRequest, LocalTraceCaptureError, LocalTraceCaptureRuntime, request_local_health,
request_local_native_threads_profile, request_local_network, request_local_runtime_profile, request_local_top_api,
request_local_top_disk, request_local_top_locks, request_local_top_rpc, request_local_trace_capture,
LocalHealthRequest, LocalNetworkRequest, LocalTraceCaptureError, LocalTraceCaptureRuntime, load_selected_offline_key,
request_local_health, request_local_native_threads_profile, request_local_network, request_local_runtime_profile,
request_local_top_api, request_local_top_disk, request_local_top_locks, request_local_top_rpc, request_local_trace_capture,
spawn_local_trace_capture_runtime,
};
pub use diagnostics::{
+8 -4
View File
@@ -867,10 +867,14 @@ async fn execute_connect_client_performance(options: ConnectClientPerformanceOpt
let duration = Duration::from_millis(options.duration_millis);
validate_client_limits(duration, options.traffic_bytes).map_err(Error::other)?;
let key = IdentityStore::new(options.state_dir.join("identity"))
.load()
.map_err(Error::other)?
.ok_or_else(|| Error::other("connect client performance requires an enrolled device identity"))?;
let key = if let Some(offline_key_id) = options.offline_key_id.as_deref() {
crate::connect::load_selected_offline_key(&options.state_dir, offline_key_id).map_err(Error::other)?
} else {
IdentityStore::new(options.state_dir.join("identity"))
.load()
.map_err(Error::other)?
.ok_or_else(|| Error::other("connect client performance requires an enrolled device identity"))?
};
let access_key = read_protected_client_credential(&options.access_key_file).map_err(Error::other)?;
let secret_key = read_protected_client_credential(&options.secret_key_file).map_err(Error::other)?;
let session_token = options
+86
View File
@@ -547,6 +547,92 @@ async fn production_cli_writes_a_verifiable_export_with_exact_binary_provenance(
.expect("public key")
.verify(&signed, &signature)
.expect("valid ES256 signature");
#[cfg(unix)]
{
// The offline key is separate from the online identity and is selected explicitly.
fs::set_permissions(&state, fs::Permissions::from_mode(0o700)).expect("private offline state root");
let offline = rustfs::connect::OfflineKeyStore::new(&state)
.load_or_create()
.expect("existing offline identity");
let offline_key_id = hex_lower(&Sha256::digest(offline.public_key_der()));
let offline_output = temp.path().join("client-offline.zip");
let mut command = Command::new(env!("CARGO_BIN_EXE_rustfs"));
command
.args(["connect", "performance", "client", "--state-dir"])
.arg(&state)
.args([
"--offline-key-id",
&offline_key_id,
"--endpoint",
&server.endpoint(),
"--access-key-file",
])
.arg(&access_key_file)
.arg("--secret-key-file")
.arg(&secret_key_file)
.arg("--output")
.arg(&offline_output)
.args(["--organization", organization, "--cluster", &cluster, "--device"])
.arg(format!("{cluster}/clusterDevices/019e3ae0-0000-7000-8000-000000000012"))
.args([
"--run-uid",
"019e3ae0-0000-7000-8000-000000000016",
"--artifact-uid",
"019e3ae0-0000-7000-8000-000000000017",
"--consent-uid",
"019e3ae0-0000-7000-8000-000000000018",
"--policy-revision",
"7",
"--consent-expires-at",
&(now() + 120).to_string(),
"--expires-at",
&(now() + 60).to_string(),
"--operation",
"get",
"--traffic-bytes",
"65536",
"--duration-millis",
"1000",
"--acknowledge-l1",
]);
let result = tokio::task::spawn_blocking(move || command.output())
.await
.expect("offline CLI task")
.expect("run production rustfs binary");
assert!(result.status.success(), "stderr: {}", String::from_utf8_lossy(&result.stderr));
assert!(String::from_utf8_lossy(&result.stdout).contains("upload=not-performed\n"));
let bytes = fs::read(&offline_output).expect("offline signed archive");
let mut archive = zip::ZipArchive::new(Cursor::new(bytes.as_slice())).expect("offline archive");
let envelope_bytes = read_archive_member(&mut archive, "envelope.json");
let signature_bytes = read_archive_member(&mut archive, "envelope.sig");
let result_bytes = read_archive_member(&mut archive, "result.json");
let envelope: serde_json::Value = serde_json::from_slice(&envelope_bytes).expect("offline envelope");
let signed_result: serde_json::Value = serde_json::from_slice(&result_bytes).expect("offline result");
assert_eq!(envelope["classification"], "L1");
assert_eq!(envelope["deviceKeyId"], offline_key_id);
assert_eq!(
signed_result["provenance"]["executableSha256"],
sha256_file(Path::new(env!("CARGO_BIN_EXE_rustfs")))
);
let signature_document: serde_json::Value = serde_json::from_slice(&signature_bytes).expect("offline signature");
let raw = URL_SAFE_NO_PAD
.decode_to_vec(signature_document["value"].as_str().expect("signature value"))
.expect("base64url signature");
let signature = Signature::from_slice(&raw).expect("P-256 signature");
let mut signed = b"rustfs-diagnostic-envelope-v1\0".to_vec();
signed.extend_from_slice(&envelope_bytes);
VerifyingKey::from_public_key_der(&offline.public_key_der())
.expect("offline public key")
.verify(&signed, &signature)
.expect("offline ES256 signature");
assert!(
VerifyingKey::from_public_key_der(&identity.public_key_der())
.expect("online public key")
.verify(&signed, &signature)
.is_err()
);
}
}
fn write_credential(path: &Path, value: &str) {