From bfec547d36344e485d71c272a8175c7ffc6c6ce7 Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Sat, 1 Aug 2026 19:15:12 +0800 Subject: [PATCH] feat(health): drive KMS readiness from the probe status (#5552) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * feat(health): drive KMS readiness from the probe status Readiness reported the KMS ready on the service status bit alone, which says the manager started, not that the backend can still serve a request. Consume the background probe snapshot instead: a running service is withdrawn only when a fresh snapshot shows failures at the threshold. The readiness path stays free of backend calls — it reads the lock-free snapshot — so a struggling KMS cannot be amplified into probe traffic of its own. Everything the probe cannot speak to (no worker, no completed round, a snapshot older than three probe intervals) leaves the previous status-bit verdict standing, and the check remains off by default. Refs rustfs/backlog#1584 (part of rustfs/backlog#1562) * test(health): pin the readiness default at compile time --- rustfs/src/server/health.rs | 257 ++++++++++++++++++++++++++++++++++++ rustfs/src/server/layer.rs | 5 +- rustfs/src/server/mod.rs | 5 +- 3 files changed, 265 insertions(+), 2 deletions(-) diff --git a/rustfs/src/server/health.rs b/rustfs/src/server/health.rs index 8fe315dc9..5cac23e3f 100644 --- a/rustfs/src/server/health.rs +++ b/rustfs/src/server/health.rs @@ -18,7 +18,10 @@ use super::{ collect_cluster_read_health_report, collect_cluster_write_health_report, collect_node_readiness_report, }; use http::{Method, StatusCode}; +use rustfs_kms::ProbeStatus; +use rustfs_kms::probe::{DEFAULT_PROBE_INTERVAL, ENV_KMS_PROBE_INTERVAL_SECS, MIN_PROBE_INTERVAL}; use serde_json::{Value, json}; +use std::time::Duration; #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub(crate) struct HealthCheckState { @@ -122,6 +125,86 @@ pub(crate) fn health_minimal_response_enabled() -> bool { ) } +/// Consecutive failed probe rounds tolerated before readiness withdraws the KMS. +/// +/// One failed round is routinely an external hiccup — a Vault leader election, a +/// dropped connection, a rotated token — and reporting not ready on the first one +/// would turn that hiccup into a rolling restart of every node that shares the +/// backend. Requiring the failure to survive three rounds keeps the signal while +/// costing at most three intervals of detection latency. +const KMS_PROBE_FAILURE_THRESHOLD: u32 = 3; + +/// Probe rounds a snapshot may age before readiness stops acting on it. +/// +/// Expressed in rounds rather than seconds so the bound follows the configured +/// interval: a deployment probing every ten minutes must not read every snapshot +/// as stale. Three rounds absorbs a round that overruns its slot or a worker that +/// misses a tick without discarding a signal that is merely late. +const KMS_PROBE_STALENESS_ROUNDS: u32 = 3; + +/// Upper bound on the age of a probe snapshot that readiness is willing to act on. +pub(crate) fn kms_probe_staleness_limit() -> Duration { + configured_kms_probe_interval().saturating_mul(KMS_PROBE_STALENESS_ROUNDS) +} + +/// Probe interval currently in force. +/// +/// Readiness reads the same variable as the probe worker because it only needs +/// the value to scale [`kms_probe_staleness_limit`]; the worker owns the interval +/// itself and keeps its parsing private. +fn configured_kms_probe_interval() -> Duration { + kms_probe_interval_from_env(std::env::var(ENV_KMS_PROBE_INTERVAL_SECS).ok().as_deref()) +} + +/// Interval implied by a raw environment value, mirroring the worker's own rules: +/// unset or unparsable means the default, and anything below the floor is raised +/// to it. `0` disables the probe, which publishes no snapshot at all — the bound +/// is then never consulted, so the default stands in. +fn kms_probe_interval_from_env(value: Option<&str>) -> Duration { + let Some(seconds) = value.and_then(|value| value.trim().parse::().ok()) else { + return DEFAULT_PROBE_INTERVAL; + }; + if seconds == 0 { + return DEFAULT_PROBE_INTERVAL; + } + Duration::from_secs(seconds).max(MIN_PROBE_INTERVAL) +} + +/// KMS readiness verdict from the service status bit and the background probe. +/// +/// Reads only the published snapshot and never calls the backend: a readiness +/// request must not become load on a KMS that is already struggling, and a +/// Kubernetes probe timing out on an external dependency is how one slow Vault +/// becomes a fleet-wide restart. +/// +/// The probe can only ever withdraw readiness from a service that is otherwise +/// running, and only on evidence that is both fresh and repeated. Everything the +/// probe cannot speak to — no worker (a backend without a data key round trip, or +/// the probe disabled), no completed round yet, or a snapshot older than +/// `staleness_limit` — leaves the status-bit verdict standing rather than +/// reporting a failure nobody observed. +pub(crate) fn kms_ready_from_probe(service_running: bool, probe: Option<&ProbeStatus>, staleness_limit: Duration) -> bool { + if !service_running { + return false; + } + + let Some(probe) = probe else { + return true; + }; + + let Some(age) = probe.last_round_age() else { + return true; + }; + + if age > staleness_limit { + return true; + } + + // Every non-failure result resets the counter, so an unsupported backend or a + // recovered round can never reach the threshold. + probe.consecutive_failures < KMS_PROBE_FAILURE_THRESHOLD +} + pub(crate) fn build_component_details( storage_ready: bool, iam_ready: bool, @@ -277,3 +360,177 @@ pub(crate) fn build_health_payload(ctx: HealthPayloadContext<'_>) -> Value { payload } + +#[cfg(test)] +mod tests { + use super::super::readiness::DependencyReadiness; + use super::*; + use rustfs_kms::{ProbeFailureKind, ProbeResult}; + use serial_test::serial; + use temp_env::with_var; + use tokio::time::Instant; + + const STALENESS_LIMIT: Duration = Duration::from_secs(180); + + fn snapshot(result: ProbeResult, last_round_at: Option, consecutive_failures: u32) -> ProbeStatus { + ProbeStatus { + result, + last_round_at, + last_success_at: last_round_at.filter(|_| result == ProbeResult::Success), + last_success_unix_secs: None, + consecutive_failures, + } + } + + fn ready_report() -> DependencyReadinessReport { + DependencyReadinessReport { + readiness: DependencyReadiness { + storage_ready: true, + iam_ready: true, + lock_quorum_ready: true, + peer_health_ready: true, + }, + degraded_reasons: Vec::new(), + } + } + + #[tokio::test(start_paused = true)] + async fn a_fresh_successful_round_keeps_the_service_ready() { + let round_at = Instant::now(); + tokio::time::advance(Duration::from_secs(1)).await; + + let status = snapshot(ProbeResult::Success, Some(round_at), 0); + assert!(kms_ready_from_probe(true, Some(&status), STALENESS_LIMIT)); + } + + #[tokio::test(start_paused = true)] + async fn a_single_failed_round_does_not_withdraw_readiness() { + let round_at = Instant::now(); + tokio::time::advance(Duration::from_secs(1)).await; + + let status = snapshot( + ProbeResult::Failure(ProbeFailureKind::Decrypt), + Some(round_at), + KMS_PROBE_FAILURE_THRESHOLD - 1, + ); + assert!(kms_ready_from_probe(true, Some(&status), STALENESS_LIMIT)); + } + + #[tokio::test(start_paused = true)] + async fn failures_reaching_the_threshold_withdraw_readiness() { + let round_at = Instant::now(); + tokio::time::advance(Duration::from_secs(1)).await; + + let status = snapshot( + ProbeResult::Failure(ProbeFailureKind::Generate), + Some(round_at), + KMS_PROBE_FAILURE_THRESHOLD, + ); + assert!(!kms_ready_from_probe(true, Some(&status), STALENESS_LIMIT)); + } + + #[tokio::test(start_paused = true)] + async fn a_stale_snapshot_falls_back_to_the_status_bit() { + let round_at = Instant::now(); + tokio::time::advance(STALENESS_LIMIT + Duration::from_secs(1)).await; + + let status = snapshot( + ProbeResult::Failure(ProbeFailureKind::Mismatch), + Some(round_at), + KMS_PROBE_FAILURE_THRESHOLD * 10, + ); + assert!( + kms_ready_from_probe(true, Some(&status), STALENESS_LIMIT), + "a snapshot nobody refreshed is unknown, not a failure" + ); + } + + #[tokio::test(start_paused = true)] + async fn an_unsupported_backend_is_not_a_failure() { + let round_at = Instant::now(); + tokio::time::advance(Duration::from_secs(1)).await; + + let status = snapshot(ProbeResult::Unsupported, Some(round_at), 0); + assert!(kms_ready_from_probe(true, Some(&status), STALENESS_LIMIT)); + } + + #[tokio::test(start_paused = true)] + async fn a_stopped_service_stays_not_ready_despite_a_successful_round() { + let round_at = Instant::now(); + tokio::time::advance(Duration::from_secs(1)).await; + + let status = snapshot(ProbeResult::Success, Some(round_at), 0); + assert!(!kms_ready_from_probe(false, Some(&status), STALENESS_LIMIT)); + } + + #[test] + fn a_service_without_a_probe_keeps_the_status_bit_verdict() { + assert!(kms_ready_from_probe(true, None, STALENESS_LIMIT)); + assert!(!kms_ready_from_probe(false, None, STALENESS_LIMIT)); + } + + #[test] + fn a_probe_with_no_completed_round_keeps_the_service_ready() { + let status = snapshot(ProbeResult::Pending, None, 0); + assert!(kms_ready_from_probe(true, Some(&status), STALENESS_LIMIT)); + } + + #[test] + fn the_probe_interval_follows_the_worker_parsing_rules() { + assert_eq!(kms_probe_interval_from_env(None), DEFAULT_PROBE_INTERVAL); + assert_eq!(kms_probe_interval_from_env(Some("not-a-number")), DEFAULT_PROBE_INTERVAL); + assert_eq!(kms_probe_interval_from_env(Some("0")), DEFAULT_PROBE_INTERVAL); + assert_eq!(kms_probe_interval_from_env(Some(" 600 ")), Duration::from_secs(600)); + assert_eq!(kms_probe_interval_from_env(Some("1")), MIN_PROBE_INTERVAL); + } + + #[test] + #[serial] + fn the_staleness_bound_spans_three_probe_rounds() { + with_var(ENV_KMS_PROBE_INTERVAL_SECS, Some("30"), || { + assert_eq!(kms_probe_staleness_limit(), Duration::from_secs(90)); + }); + } + + /// Flipping this default is a deployment behaviour change, not a code + /// change: it would start failing readiness on clusters that never asked + /// for the KMS to gate it. Pinned at compile time so the constant cannot + /// drift without this line being edited too. + const _: () = assert!(!rustfs_config::DEFAULT_HEALTH_COMPAT_KMS_READY_CHECK_ENABLE); + + #[test] + #[serial] + fn a_withdrawn_kms_degrades_the_readiness_response() { + with_var(rustfs_config::ENV_HEALTH_MINIMAL_RESPONSE_ENABLE, Some("false"), || { + let report = ready_report(); + let parts = build_health_response_parts( + Method::GET, + HealthProbe::Readiness, + Some(&report), + "rustfs-endpoint", + None, + Some(false), + ); + + assert_eq!(parts.status_code, StatusCode::SERVICE_UNAVAILABLE); + let payload = parts.payload.expect("GET should include payload"); + assert_eq!(payload["details"]["kms"]["ready"], false); + assert_eq!(payload["degradedReasons"], json!(["kms_not_ready"])); + }); + } + + #[test] + #[serial] + fn a_disabled_check_leaves_the_readiness_response_untouched() { + with_var(rustfs_config::ENV_HEALTH_MINIMAL_RESPONSE_ENABLE, Some("false"), || { + let report = ready_report(); + let parts = + build_health_response_parts(Method::GET, HealthProbe::Readiness, Some(&report), "rustfs-endpoint", None, None); + + assert_eq!(parts.status_code, StatusCode::OK); + let payload = parts.payload.expect("GET should include payload"); + assert!(payload["details"].get("kms").is_none()); + assert_eq!(payload["degradedReasons"], json!([])); + }); + } +} diff --git a/rustfs/src/server/layer.rs b/rustfs/src/server/layer.rs index 538dd1de5..9d02b0e39 100644 --- a/rustfs/src/server/layer.rs +++ b/rustfs/src/server/layer.rs @@ -23,6 +23,7 @@ use crate::server::{ MINIO_ADMIN_V3_PREFIX, MINIO_HEALTH_CLUSTER_PATH, MINIO_HEALTH_CLUSTER_READ_PATH, MINIO_HEALTH_LIVE_PATH, MINIO_HEALTH_READY_PATH, PROFILE_CPU_PATH, PROFILE_MEMORY_PATH, RPC_PREFIX, RUSTFS_ADMIN_PREFIX, active_http_requests, build_health_response_parts, collect_probe_readiness, has_path_prefix, is_admin_path, is_table_catalog_path, + kms_probe_staleness_limit, kms_ready_from_probe, }; use crate::storage_api::server::layer::apply_cors_headers; use crate::storage_api::server::layer::request_context::{RequestContext, extract_request_id_from_headers, spawn_traced}; @@ -1281,7 +1282,9 @@ async fn health_kms_ready() -> bool { return true; }; - matches!(service_manager.get_status().await, rustfs_kms::KmsServiceStatus::Running) + let service_running = matches!(service_manager.get_status().await, rustfs_kms::KmsServiceStatus::Running); + let probe_status = service_manager.probe_status(); + kms_ready_from_probe(service_running, probe_status.as_deref(), kms_probe_staleness_limit()) } async fn build_public_health_http_response( diff --git a/rustfs/src/server/mod.rs b/rustfs/src/server/mod.rs index 8c2af988a..bcc00cdc4 100644 --- a/rustfs/src/server/mod.rs +++ b/rustfs/src/server/mod.rs @@ -55,7 +55,10 @@ pub(crate) use health::{ HealthPayloadContext, HealthReadinessSource, build_component_details, build_health_payload, health_check_state, readiness_source_for_probe, }; -pub(crate) use health::{HealthProbe, build_health_response_parts, collect_probe_readiness, probe_from_path}; +pub(crate) use health::{ + HealthProbe, build_health_response_parts, collect_probe_readiness, kms_probe_staleness_limit, kms_ready_from_probe, + probe_from_path, +}; pub(crate) use http::HeaderMapCarrier; pub(crate) use http::active_http_requests; pub(crate) use layer::{RequestContextLayer, is_sts_query_request};