feat(connect): capture signed offline service health (#8219)

This commit is contained in:
Chris
2026-09-29 10:50:23 +08:00
committed by GitHub
parent c62979a45b
commit a1724f3dbe
8 changed files with 536 additions and 16 deletions
+48
View File
@@ -137,6 +137,8 @@ pub enum ConnectCommands {
Performance(ConnectPerformanceOpts),
/// Capture a consent-bound local profile and write a signed export
Profile(ConnectProfileOpts),
/// Capture bounded service health in the running server and write a signed export
Health(ConnectHealthOpts),
/// Capture allow-listed local log events and write a signed export
Logs(ConnectLogsOpts),
/// Record, forward, or replay consent-bound telemetry
@@ -1027,6 +1029,50 @@ pub enum ConnectLogsMode {
Live,
}
/// `connect health` options.
#[derive(Args, Clone)]
pub struct ConnectHealthOpts {
/// Owner-only server state directory containing an enrolled offline identity
#[arg(long = "state-dir")]
pub state_dir: PathBuf,
/// SHA-256 of the enrolled offline public key
#[arg(long = "offline-key-id", value_parser = NonEmptyStringValueParser::new())]
pub offline_key_id: String,
/// New local archive path; an existing file is never replaced
#[arg(long)]
pub output: PathBuf,
/// Organization resource name bound to the export
#[arg(long, value_parser = NonEmptyStringValueParser::new())]
pub organization: String,
/// Cluster resource name bound to the export
#[arg(long, value_parser = NonEmptyStringValueParser::new())]
pub cluster: String,
/// Cluster-device resource name bound to the export
#[arg(long, value_parser = NonEmptyStringValueParser::new())]
pub device: String,
/// UUIDv7 diagnostic run identifier issued by Connect
#[arg(long = "run-uid", value_parser = NonEmptyStringValueParser::new())]
pub run_uid: String,
/// UUIDv7 artifact identifier issued by Connect
#[arg(long = "artifact-uid", value_parser = NonEmptyStringValueParser::new())]
pub artifact_uid: String,
/// UUIDv7 consent identifier issued by Connect
#[arg(long = "consent-uid", value_parser = NonEmptyStringValueParser::new())]
pub consent_uid: String,
/// Consent policy revision bound to this capture
#[arg(long = "policy-revision")]
pub policy_revision: u64,
/// Consent expiry as UTC Unix seconds
#[arg(long = "consent-expires-at")]
pub consent_expires_at_unix: i64,
/// Artifact expiry as UTC Unix seconds
#[arg(long = "expires-at")]
pub expires_at_unix: i64,
/// Confirm this explicit local L0 health capture
#[arg(long = "acknowledge-l0", required = true, action = clap::ArgAction::SetTrue)]
pub acknowledge_l0: bool,
}
/// `connect profile` options.
#[derive(Args, Clone)]
pub struct ConnectProfileOpts {
@@ -1589,6 +1635,8 @@ pub enum CommandResult {
ConnectSiteReplicationPerformance(Box<ConnectSiteReplicationPerformanceOpts>),
/// Consent-bound local Connect profile export
ConnectProfile(ConnectProfileOpts),
/// Consent-bound serving-process health export
ConnectHealth(ConnectHealthOpts),
/// Consent-bound local Connect log export
ConnectLogs(ConnectLogsOpts),
/// Consent-bound local Connect telemetry operation
+1
View File
@@ -52,6 +52,7 @@ mod config_test;
// Re-export public types
#[cfg(test)]
pub(crate) use cli::Cli;
pub use cli::ConnectHealthOpts;
pub use cli::ConnectSiteReplicationPerformanceOpts;
pub use cli::{CommandResult, InfoOpts, InfoType};
pub use cli::{
+1
View File
@@ -159,6 +159,7 @@ impl Opt {
}
},
ConnectCommands::Profile(opts) => Ok(CommandResult::ConnectProfile(opts)),
ConnectCommands::Health(opts) => Ok(CommandResult::ConnectHealth(opts)),
ConnectCommands::Logs(opts) => Ok(CommandResult::ConnectLogs(opts)),
ConnectCommands::Telemetry(opts) => Ok(CommandResult::ConnectTelemetry(opts.command)),
ConnectCommands::Top(opts) => Ok(CommandResult::ConnectTop(opts.command)),
+94
View File
@@ -19,7 +19,9 @@
//! raw admin responses and object data never enter the result.
use std::collections::BTreeSet;
use std::fs::{self, File, OpenOptions};
use std::io::{Cursor, Write as _};
use std::path::Path;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Instant;
@@ -302,6 +304,20 @@ pub enum HealthError {
Encoding,
}
#[derive(Debug, Error)]
pub enum HealthSaveError {
#[error("health_export_invalid_request")]
InvalidRequest,
#[error("health_export_limit_exceeded")]
LimitExceeded,
#[error("health_export_cancelled")]
Cancelled,
#[error("health_export_io_failed")]
Io(#[source] std::io::Error),
#[error("health_export_durability_failed_after_commit")]
DurabilityAfterCommit(#[source] std::io::Error),
}
struct CollectorLease;
impl CollectorLease {
@@ -598,6 +614,55 @@ pub fn sign_health_export(
})
}
pub fn save_signed_health_export(
output: &Path,
archive: &[u8],
artifact_uid: &str,
cancel: &CancellationToken,
) -> Result<(), HealthSaveError> {
if cancel.is_cancelled() {
return Err(HealthSaveError::Cancelled);
}
if archive.is_empty() || archive.len() as u64 > MAX_HEALTH_OUTPUT_BYTES {
return Err(HealthSaveError::LimitExceeded);
}
let parent = output
.parent()
.filter(|path| !path.as_os_str().is_empty())
.unwrap_or_else(|| Path::new("."));
let filename = output.file_name().ok_or(HealthSaveError::InvalidRequest)?.to_string_lossy();
let temporary = parent.join(format!(".{filename}.{artifact_uid}.partial"));
let mut options = OpenOptions::new();
options.write(true).create_new(true);
#[cfg(unix)]
{
use std::os::unix::fs::OpenOptionsExt as _;
options.mode(0o600);
}
let mut file = options.open(&temporary).map_err(HealthSaveError::Io)?;
let result = (|| {
file.write_all(archive).map_err(HealthSaveError::Io)?;
if cancel.is_cancelled() {
return Err(HealthSaveError::Cancelled);
}
file.sync_all().map_err(HealthSaveError::Io)?;
if cancel.is_cancelled() {
return Err(HealthSaveError::Cancelled);
}
fs::hard_link(&temporary, output).map_err(HealthSaveError::Io)?;
fs::remove_file(&temporary).map_err(HealthSaveError::DurabilityAfterCommit)?;
#[cfg(unix)]
File::open(parent)
.and_then(|directory| directory.sync_all())
.map_err(HealthSaveError::DurabilityAfterCommit)?;
Ok(())
})();
if result.is_err() {
let _ = fs::remove_file(&temporary);
}
result
}
fn result_reason(freshness: HealthFreshness, observation: &HealthSourceObservation) -> HealthResultReason {
match freshness {
HealthFreshness::Stale => HealthResultReason::EvidenceStale,
@@ -824,4 +889,33 @@ mod tests {
drop(active);
CollectorLease::acquire().expect("collector lease should be released");
}
#[test]
fn offline_health_archive_is_private_and_never_replaces_existing_output() {
let directory = tempfile::tempdir().unwrap();
let output = directory.path().join("health.zip");
let cancel = CancellationToken::new();
save_signed_health_export(&output, b"signed archive", "019e3ae0-0000-7000-8000-000000000025", &cancel).unwrap();
assert_eq!(std::fs::read(&output).unwrap(), b"signed archive");
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt as _;
assert_eq!(std::fs::metadata(&output).unwrap().permissions().mode() & 0o777, 0o600);
}
assert!(matches!(
save_signed_health_export(&output, b"different", "019e3ae0-0000-7000-8000-000000000026", &cancel),
Err(HealthSaveError::Io(error)) if error.kind() == std::io::ErrorKind::AlreadyExists
));
assert_eq!(std::fs::read(&output).unwrap(), b"signed archive");
cancel.cancel();
assert!(matches!(
save_signed_health_export(
&directory.path().join("cancelled.zip"),
b"data",
"019e3ae0-0000-7000-8000-000000000027",
&cancel
),
Err(HealthSaveError::Cancelled)
));
}
}
+5 -4
View File
@@ -68,10 +68,10 @@ mod trace_runtime;
pub use health::{
HEALTH_CATALOG_CHECKS, HEALTH_SCHEMA_VERSION, HEALTH_SERVICE_CAPABILITY, HEALTH_TIMEOUT_SECONDS, HealthCheckResult,
HealthDiagnosticResult, HealthError, HealthFreshness, HealthOutcome, HealthResultReason, HealthRuleOutcome,
HealthDiagnosticResult, HealthError, HealthFreshness, HealthOutcome, HealthResultReason, HealthRuleOutcome, HealthSaveError,
HealthServiceRequest, HealthSourceObservation, LocalHealthConsent, MAX_EVIDENCE_AGE_SECONDS, MAX_HEALTH_CPU_MILLIS,
MAX_HEALTH_MEMORY_BYTES, MAX_HEALTH_OUTPUT_BYTES, SignedHealthExport, collect_runtime_health, evaluate_health_observation,
sign_health_export,
save_signed_health_export, sign_health_export,
};
pub use inspect::{
INSPECT_CAPABILITY, INSPECT_SCHEMA_VERSION, InspectArtifactConsent, InspectDiagnosticResult, InspectError, InspectFinding,
@@ -172,8 +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_native_threads_profile, request_local_runtime_profile,
request_local_top_disk, request_local_trace_capture, spawn_local_trace_capture_runtime,
LocalHealthRequest, LocalTraceCaptureError, LocalTraceCaptureRuntime, request_local_health,
request_local_native_threads_profile, 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;
+319 -1
View File
@@ -14,7 +14,7 @@
//! Owner-only local transport between diagnostic CLI commands and the running server.
//!
//! Version 1 accepts TRACE_RECORD, RUNTIME_PROFILE, NATIVE_THREADS_PROFILE and TOP_DISK. Signed requests select
//! Version 1 accepts TRACE_RECORD, RUNTIME_PROFILE, NATIVE_THREADS_PROFILE, TOP_DISK and HEALTH. 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,6 +66,8 @@ pub(crate) enum LocalTraceCaptureError {
Producer(#[from] TelemetryProducerError),
#[error("runtime profile request rejected: {0}")]
RuntimeProfile(String),
#[error("health request rejected: {0}")]
Health(String),
#[error("top disk request rejected: {0}")]
TopDisk(String),
#[error("local diagnostic cancellation was not acknowledged")]
@@ -118,6 +120,34 @@ enum CaptureRequest {
protocol_version: u16,
request: super::top_disk::LocalTopDiskRequest,
},
Health {
protocol_version: u16,
request: LocalHealthRequest,
},
}
#[derive(Clone, Debug, Deserialize, Serialize)]
#[serde(deny_unknown_fields, rename_all = "camelCase")]
pub(crate) struct LocalHealthRequest {
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 schema_version: u16,
pub capability: String,
pub consent_uid: String,
pub policy_revision: u64,
pub consent_expires_at_unix: i64,
pub acknowledge_l0: bool,
pub expires_at_unix: i64,
}
pub(crate) struct LocalHealthArchive {
pub artifact_uid: String,
pub archive_bytes: Vec<u8>,
pub archive_sha256: String,
}
#[derive(Clone, Copy, Debug, Deserialize, Serialize)]
@@ -193,6 +223,14 @@ enum CaptureResponse {
TopDiskError {
code: RuntimeErrorCode,
},
HealthOk {
archive_base64: String,
archive_sha256: String,
artifact_uid: String,
},
HealthError {
code: RuntimeErrorCode,
},
}
pub(crate) fn spawn_local_trace_capture_runtime(
@@ -313,6 +351,89 @@ pub(crate) async fn request_local_runtime_profile(
request_local_profile_capture(state_root, request, cancel, false).await
}
pub(crate) async fn request_local_health(
state_root: &Path,
request: LocalHealthRequest,
cancel: &CancellationToken,
) -> Result<LocalHealthArchive, LocalTraceCaptureError> {
if request.capability != super::health::HEALTH_SERVICE_CAPABILITY {
return Err(LocalTraceCaptureError::Protocol);
}
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 mut bytes = serde_json::to_vec(&CaptureRequest::Health {
protocol_version: PROTOCOL_VERSION,
request,
})
.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)?;
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(Duration::from_secs(super::health::HEALTH_TIMEOUT_SECONDS + 5), response);
tokio::pin!(response);
let response = tokio::select! {
biased;
_ = cancel.cancelled() => {
writer.shutdown().await.map_err(|_| LocalTraceCaptureError::CancellationUnconfirmed)?;
let acknowledged = matches!(tokio::time::timeout(Duration::from_secs(2), &mut response).await,
Ok(Ok(Ok(CaptureResponse::HealthError { .. } | CaptureResponse::HealthOk { .. }))));
if !acknowledged { return Err(LocalTraceCaptureError::CancellationUnconfirmed); }
return Err(LocalTraceCaptureError::Health("CANCELLED".to_owned()));
}
response = &mut response => response.map_err(|_| LocalTraceCaptureError::Protocol)??,
};
match response {
CaptureResponse::HealthOk {
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() as u64 > super::health::MAX_HEALTH_OUTPUT_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(LocalHealthArchive {
artifact_uid,
archive_bytes,
archive_sha256,
})
}
CaptureResponse::HealthError { code } => Err(LocalTraceCaptureError::Health(format!("{code:?}"))),
_ => Err(LocalTraceCaptureError::Protocol),
}
}
pub(crate) async fn request_local_native_threads_profile(
state_root: &Path,
request: LocalRuntimeProfileRequest,
@@ -545,6 +666,56 @@ async fn handle_runtime_profile(
}
}
async fn handle_health(
mut reader: BufReader<tokio::net::unix::OwnedReadHalf>,
mut writer: tokio::net::unix::OwnedWriteHalf,
state_root: &Path,
protocol_version: u16,
request: LocalHealthRequest,
shutdown: CancellationToken,
) {
let cancel = shutdown.child_token();
let capture = async {
match tokio::time::timeout(
Duration::from_secs(super::health::HEALTH_TIMEOUT_SECONDS),
capture_local_health(state_root, protocol_version, request, &cancel),
)
.await
{
Ok(result) => result,
Err(_) => {
cancel.cancel();
Err(RuntimeErrorCode::TimedOut)
}
}
};
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::HealthOk {
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::HealthError { 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 handle_top_disk(
mut reader: BufReader<tokio::net::unix::OwnedReadHalf>,
mut writer: tokio::net::unix::OwnedWriteHalf,
@@ -582,6 +753,87 @@ async fn handle_top_disk(
}
}
async fn capture_local_health(
state_root: &Path,
protocol_version: u16,
input: LocalHealthRequest,
cancel: &CancellationToken,
) -> Result<super::health::SignedHealthExport, RuntimeErrorCode> {
use super::health::{HEALTH_SERVICE_CAPABILITY, HealthServiceRequest, LocalHealthConsent};
use rand::{TryRng as _, rngs::SysRng};
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))
|| input.schema_version != 1
|| input.capability != HEALTH_SERVICE_CAPABILITY
{
return Err(RuntimeErrorCode::InvalidRequest);
}
if cancel.is_cancelled() {
return Err(RuntimeErrorCode::Cancelled);
}
if !input.acknowledge_l0 || input.policy_revision == 0 {
return Err(RuntimeErrorCode::ConsentRequired);
}
let now = unix_now().map_err(|_| RuntimeErrorCode::CollectionFailed)?;
if input.consent_expires_at_unix <= now || input.expires_at_unix > input.consent_expires_at_unix {
return Err(RuntimeErrorCode::ConsentExpired);
}
if input.expires_at_unix <= now {
return Err(RuntimeErrorCode::Expired);
}
let key = load_offline_key(state_root, &input.offline_key_id)?;
let provenance = super::job_delivery::executable_provenance()
.await
.map_err(|_| RuntimeErrorCode::CollectionFailed)?;
let mut nonce = [0_u8; 32];
SysRng
.try_fill_bytes(&mut nonce)
.map_err(|_| RuntimeErrorCode::CollectionFailed)?;
let request = HealthServiceRequest {
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,
schema_version: input.schema_version,
capability: input.capability,
consent: LocalHealthConsent {
consent_uid: input.consent_uid,
policy_revision: input.policy_revision,
expires_at_unix: input.consent_expires_at_unix,
active: input.acknowledge_l0,
},
produced_at_unix: unix_now().map_err(|_| RuntimeErrorCode::CollectionFailed)?,
expires_at_unix: input.expires_at_unix,
nonce,
max_evidence_age_seconds: 300,
provenance,
};
super::health::collect_runtime_health(&request, &key, cancel)
.await
.map_err(health_error)
}
fn health_error(error: super::health::HealthError) -> RuntimeErrorCode {
use super::health::HealthError;
match error {
HealthError::ConsentRequired => RuntimeErrorCode::ConsentRequired,
HealthError::ConsentExpired => RuntimeErrorCode::ConsentExpired,
HealthError::Expired => RuntimeErrorCode::Expired,
HealthError::LimitExceeded => RuntimeErrorCode::LimitExceeded,
HealthError::Busy => RuntimeErrorCode::Busy,
HealthError::Cancelled => RuntimeErrorCode::Cancelled,
HealthError::Unsupported | HealthError::InvalidRequest => RuntimeErrorCode::InvalidRequest,
HealthError::SourceUnavailable | HealthError::CollectionFailed | HealthError::Signing | HealthError::Encoding => {
RuntimeErrorCode::CollectionFailed
}
}
}
async fn capture_local_top_disk(
state_root: &Path,
protocol_version: u16,
@@ -957,6 +1209,13 @@ async fn handle_connection(stream: UnixStream, state_root: PathBuf, shutdown: Ca
handle_top_disk(reader, writer, &state_root, protocol_version, request, shutdown).await;
return;
}
CaptureRequest::Health {
protocol_version,
request,
} => {
handle_health(reader, writer, &state_root, protocol_version, request, shutdown).await;
return;
}
CaptureRequest::TraceRecord {
protocol_version,
consent_expires_at_unix,
@@ -1307,6 +1566,65 @@ mod tests {
}
}
fn health_request(state: &std::path::Path) -> super::LocalHealthRequest {
let (request, _) = runtime_request(state);
super::LocalHealthRequest {
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,
schema_version: 1,
capability: "health.check.service@1".to_owned(),
consent_uid: request.consent_uid,
policy_revision: request.policy_revision,
consent_expires_at_unix: request.consent_expires_at_unix,
acknowledge_l0: true,
expires_at_unix: request.expires_at_unix,
}
}
#[tokio::test]
async fn local_health_requires_consent_bounds_and_separate_offline_identity() {
let state = tempfile::tempdir().unwrap();
std::fs::set_permissions(state.path(), std::fs::Permissions::from_mode(0o700)).unwrap();
let mut request = health_request(state.path());
std::fs::remove_file(crate::connect::OfflineKeyStore::new(state.path()).key_path()).unwrap();
request.acknowledge_l0 = false;
assert!(matches!(
super::capture_local_health(state.path(), 1, request.clone(), &CancellationToken::new()).await,
Err(super::RuntimeErrorCode::ConsentRequired)
));
request.acknowledge_l0 = true;
let valid_consent_expiry = request.consent_expires_at_unix;
request.consent_expires_at_unix = request.expires_at_unix - 1;
assert!(matches!(
super::capture_local_health(state.path(), 1, request.clone(), &CancellationToken::new()).await,
Err(super::RuntimeErrorCode::ConsentExpired)
));
request.consent_expires_at_unix = valid_consent_expiry;
let valid_expiry = request.expires_at_unix;
request.expires_at_unix = request.consent_expires_at_unix - 121;
assert!(matches!(
super::capture_local_health(state.path(), 1, request.clone(), &CancellationToken::new()).await,
Err(super::RuntimeErrorCode::Expired)
));
request.expires_at_unix = valid_expiry;
assert!(matches!(
super::capture_local_health(state.path(), 1, request.clone(), &CancellationToken::new()).await,
Err(super::RuntimeErrorCode::IdentityUnavailable)
));
assert!(
super::request_local_health(state.path(), request.clone(), &CancellationToken::new())
.await
.is_err()
);
let mut value = serde_json::to_value(&request).unwrap();
value["command"] = serde_json::json!("shell.exec");
assert!(serde_json::from_value::<super::LocalHealthRequest>(value).is_err());
}
#[tokio::test]
async fn local_top_disk_rejects_invalid_requests_without_identity_or_socket_fallback() {
let state = tempfile::tempdir().unwrap();
+7 -6
View File
@@ -88,10 +88,15 @@ pub use diagnostics::{
};
pub use diagnostics::{
HEALTH_CATALOG_CHECKS, HEALTH_SCHEMA_VERSION, HEALTH_SERVICE_CAPABILITY, HEALTH_TIMEOUT_SECONDS, HealthCheckResult,
HealthDiagnosticResult, HealthError, HealthFreshness, HealthOutcome, HealthResultReason, HealthRuleOutcome,
HealthDiagnosticResult, HealthError, HealthFreshness, HealthOutcome, HealthResultReason, HealthRuleOutcome, HealthSaveError,
HealthServiceRequest, HealthSourceObservation, LocalHealthConsent, MAX_EVIDENCE_AGE_SECONDS, MAX_HEALTH_CPU_MILLIS,
MAX_HEALTH_MEMORY_BYTES, MAX_HEALTH_OUTPUT_BYTES, SignedHealthExport, collect_runtime_health, evaluate_health_observation,
sign_health_export,
save_signed_health_export, sign_health_export,
};
pub(crate) use diagnostics::{
LocalHealthRequest, LocalTraceCaptureError, LocalTraceCaptureRuntime, request_local_health,
request_local_native_threads_profile, request_local_runtime_profile, request_local_top_disk, request_local_trace_capture,
spawn_local_trace_capture_runtime,
};
pub use diagnostics::{
LocalNetworkConsent, MAX_NETWORK_ARCHIVE_BYTES, MAX_NETWORK_BANDWIDTH_BYTES_PER_SECOND, MAX_NETWORK_DECOMPRESSED_BYTES,
@@ -128,10 +133,6 @@ pub use diagnostics::{
TopRpcData, capture_top_api, capture_top_disk, capture_top_locks, capture_top_net, capture_top_rpc, evaluate_disk_window,
evaluate_network_window, save_signed_top_export, sign_top_export,
};
pub(crate) use diagnostics::{
LocalTraceCaptureError, LocalTraceCaptureRuntime, request_local_native_threads_profile, 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,
EnvironmentFilesystemType, EnvironmentInventory, EnvironmentOsFamily, MAX_ENVIRONMENT_DURATION, SavedEnvironmentExport,
+61 -5
View File
@@ -15,11 +15,11 @@
use crate::{
config::{
CommandResult, Config, ConnectClientPerformanceOperation, ConnectClientPerformanceOpts, ConnectDrivePerformanceOpts,
ConnectEnvironmentInventoryOpts, ConnectInspectObjectOpts, ConnectLicenseCommands, ConnectLicenseScopeOpts,
ConnectLogsMode, ConnectLogsOpts, ConnectObjectPerformanceOperation, ConnectObjectPerformanceOpts, ConnectProfileOpts,
ConnectProfileTool, ConnectRelayMaterialKind, ConnectRelayOpts, ConnectReportUploadOpts,
ConnectSiteReplicationPerformanceOpts, ConnectTelemetryArtifactOpts, ConnectTelemetryCommands, ConnectThreadProfileScope,
ConnectTopCommands, Opt,
ConnectEnvironmentInventoryOpts, ConnectHealthOpts, ConnectInspectObjectOpts, ConnectLicenseCommands,
ConnectLicenseScopeOpts, ConnectLogsMode, ConnectLogsOpts, ConnectObjectPerformanceOperation,
ConnectObjectPerformanceOpts, ConnectProfileOpts, ConnectProfileTool, ConnectRelayMaterialKind, ConnectRelayOpts,
ConnectReportUploadOpts, ConnectSiteReplicationPerformanceOpts, ConnectTelemetryArtifactOpts, ConnectTelemetryCommands,
ConnectThreadProfileScope, ConnectTopCommands, Opt,
},
startup_lifecycle::{StartupRuntimeLifecycle, run_startup_runtime_lifecycle},
startup_preflight::{StartupServerPreflightError, bootstrap_external_prefix_compat, init_startup_server_preflight},
@@ -149,6 +149,7 @@ async fn async_main() -> Result<()> {
return execute_connect_site_replication_performance(*options).await;
}
CommandResult::ConnectProfile(options) => return execute_connect_profile(options).await,
CommandResult::ConnectHealth(options) => return execute_connect_health(options).await,
CommandResult::ConnectLogs(options) => return execute_connect_logs(options).await,
CommandResult::ConnectTelemetry(command) => return execute_connect_telemetry(command).await,
CommandResult::ConnectTop(command) => return execute_connect_top(command).await,
@@ -1394,6 +1395,61 @@ async fn execute_connect_drive_performance(options: ConnectDrivePerformanceOpts)
Ok(())
}
async fn execute_connect_health(options: ConnectHealthOpts) -> Result<()> {
let cancel = CancellationToken::new();
let request = crate::connect::LocalHealthRequest {
offline_key_id: options.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,
schema_version: 1,
capability: crate::connect::HEALTH_SERVICE_CAPABILITY.to_owned(),
consent_uid: options.consent_uid,
policy_revision: options.policy_revision,
consent_expires_at_unix: options.consent_expires_at_unix,
acknowledge_l0: options.acknowledge_l0,
expires_at_unix: options.expires_at_unix,
};
let capture = crate::connect::request_local_health(&options.state_dir, request, &cancel);
tokio::pin!(capture);
let archive = 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("health collection cancelled")),
};
}
result = &mut capture => result.map_err(Error::other)?,
};
let size = archive.archive_bytes.len();
let sha256 = archive.archive_sha256.clone();
let artifact_uid = archive.artifact_uid.clone();
let output = options.output;
let writer_cancel = cancel.clone();
let mut writer = tokio::task::spawn_blocking(move || {
crate::connect::save_signed_health_export(&output, &archive.archive_bytes, &archive.artifact_uid, &writer_cancel)
});
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)?;
return Err(Error::other("health export cancelled"));
}
result = &mut writer => result.map_err(Error::other)?.map_err(Error::other)?,
}
println!("tool=health.check outcome=PARTIAL");
println!("artifact={artifact_uid} bytes={size} sha256={sha256}");
println!("upload=not-performed");
Ok(())
}
async fn execute_connect_profile(options: ConnectProfileOpts) -> Result<()> {
use crate::connect::{
IdentityStore, LocalProfileConsent, ProfileCaptureRequest, ProfileProvenance, ThreadProfileScope, export_cpu_profile,