diff --git a/.github/workflows/_ci-relay.yml b/.github/workflows/_ci-relay.yml index b616b8e7f35..cf87d7a6cad 100644 --- a/.github/workflows/_ci-relay.yml +++ b/.github/workflows/_ci-relay.yml @@ -198,6 +198,8 @@ jobs: with: name: desktop-e2e-relay path: target/ci + - name: Restore relay executable permission + run: chmod +x ./target/ci/buzz-relay - name: PostgreSQL-backed tests env: BUZZ_POSTGRES_ADMIN_URL: postgres://buzz:${{ env.BUZZ_TEST_POSTGRES_PASSWORD }}@localhost:5432/postgres diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index ae82c131ec7..ed34ee4c8d3 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -644,7 +644,7 @@ pub enum AuthState { Pending { challenge: String }, Authenticated(AuthContext), | GET | `/.well-known/nostr.json` | NIP-05 identity | | GET | `/health` | Health check | | GET | `/_liveness` | Liveness probe | -| GET | `/_readiness` | Readiness probe | +| GET | `/_readiness` | Readiness probe — local process lifecycle only | | POST | `/events` | Submit a signed Nostr event over HTTP (same ingest path as WebSocket `EVENT`) | | POST | `/query` | Query Nostr events over HTTP with NIP-01 filters | | POST | `/count` | Count Nostr events over HTTP with NIP-45 filters | diff --git a/Justfile b/Justfile index 8bc31ab21a6..6377b290db8 100644 --- a/Justfile +++ b/Justfile @@ -472,7 +472,7 @@ test-unit: # because they live in the binary target; the nested # `tests::postgres_tests::` stays in the PostgreSQL lane. cargo nextest run -p buzz-relay --lib --bin buzz-relay \ - -E 'test(/^api::admin::/) + test(/^handlers::channel_authz::/) + test(/^handlers::moderation_authz::/) + test(/^handlers::side_effects::tests::/) + test(/^storage_sweep::tests::/) + test(/^nip_fi_http::tests::/) + test(/^nip_fi_config::tests::/) + test(/^router::tests::/) + test(/^api::parse_query_tests::/) + test(/^api::git::transport::off_mode_precedence_tests::/) + (kind(bin) & (test(/^tests::/) + test(/^composition_tests::/)) - test(/^tests::postgres_tests::/))' + -E 'test(/^api::admin::/) + test(/^handlers::channel_authz::/) + test(/^handlers::moderation_authz::/) + test(/^handlers::side_effects::tests::/) + test(/^storage_sweep::tests::/) + test(/^nip_fi_http::tests::/) + test(/^nip_fi_config::tests::/) + test(/^readiness::tests::/) + test(/^router::tests::/) + test(/^api::parse_query_tests::/) + test(/^api::git::transport::off_mode_precedence_tests::/) + test(=state::tests::neither_a_confirmed_inactive_community_nor_a_failed_lookup_admits_the_socket) + (kind(bin) & (test(/^tests::/) + test(/^composition_tests::/)) - test(/^tests::postgres_tests::/))' # ACP author-gate and queue tests protect the trust boundary between # relay events and agent prompts. They are infra-free; ignored lifecycle # tests remain excluded and run in their dedicated integration lanes. diff --git a/crates/buzz-relay/src/lib.rs b/crates/buzz-relay/src/lib.rs index 6991ebd3681..70f3ef4892a 100644 --- a/crates/buzz-relay/src/lib.rs +++ b/crates/buzz-relay/src/lib.rs @@ -45,7 +45,8 @@ pub mod operator_listener; pub mod protocol; /// Durable NIP-PL matcher and delivery worker. pub mod push_runtime; -mod readiness; +/// Readiness-probe telemetry and the per-pod dependency sampler behind `/_status`. +pub mod readiness; /// Axum router construction. pub mod router; /// Shared application state. diff --git a/crates/buzz-relay/src/main.rs b/crates/buzz-relay/src/main.rs index 6adea419da1..d2cf5fa7871 100644 --- a/crates/buzz-relay/src/main.rs +++ b/crates/buzz-relay/src/main.rs @@ -227,6 +227,10 @@ async fn run_relay_main(boot: BootTracker) -> anyhow::Result<()> { let usage_interval_secs = usage_metrics_interval_secs(); let usage_idle_timeout_secs = usage_metrics_idle_timeout_secs(usage_interval_secs); + let dependency_sample_completion_republish_interval = + buzz_relay::readiness::dependency_sample_completion_republish_interval( + usage_idle_timeout_secs, + ); let (boot, ()) = boot.run_required( StartupPhase::MetricsBind, || relay_metrics::try_install(config.metrics_port, usage_idle_timeout_secs), @@ -244,6 +248,7 @@ async fn run_relay_main(boot: BootTracker) -> anyhow::Result<()> { info!( port = config.metrics_port, idle_timeout_secs = usage_idle_timeout_secs, + completion_republish_secs = dependency_sample_completion_republish_interval.as_secs(), "Prometheus metrics exporter started" ); @@ -455,7 +460,13 @@ async fn run_relay_main(boot: BootTracker) -> anyhow::Result<()> { cfg.create_pool(Some(deadpool_redis::Runtime::Tokio1)) .map_err(|e| anyhow::anyhow!("Redis pool creation failed: {e}"))? }; - let redis_health_pool = redis_pool.clone(); // cheap Arc clone — shared with readiness handler + let redis_health_pool = redis_pool.clone(); // cheap Arc clone — shared with AppState + // One-time bootstrap gate, deliberately before AppState and therefore before + // the health listener binds. Post-start Redis failures are dependency + // failures and must never move readiness; never having connected at all is + // a broken deployment, not a blip. + buzz_relay::state::verify_redis_command_path(&redis_health_pool).await?; + info!("Redis command path connected"); let pubsub = Arc::new( PubSubManager::new(&config.redis_url, redis_pool) .await @@ -1100,6 +1111,16 @@ async fn run_relay_main(boot: BootTracker) -> anyhow::Result<()> { )); } + // Per-pod dependency diagnostics runtime: one seam starts the dependency + // sampler and its independent completion-epoch republisher together. + { + let diagnostics_state = Arc::clone(&state); + buzz_relay::readiness::start_dependency_sampler_and_completion_publisher( + diagnostics_state, + dependency_sample_completion_republish_interval, + ); + } + // Cross-pod connection-control consumer: receive disconnect commands from // Redis pub/sub (published by the pod that recorded a ban) and close any // matching local sockets. A member's live connections may land on any pod, @@ -1273,6 +1294,8 @@ async fn run_relay_main(boot: BootTracker) -> anyhow::Result<()> { serve(router, health_router, Arc::clone(&state)).await?; state.community_revalidator_cancel.cancel(); + state.dependency_sampler_cancel.cancel(); + state.dependency_completion_publisher_cancel.cancel(); // Signal the audit worker to stop accepting, flush buffered entries, and // exit. Uses a CancellationToken so it works regardless of how many diff --git a/crates/buzz-relay/src/metrics.rs b/crates/buzz-relay/src/metrics.rs index c1e72f7b75a..59bd90a0d7f 100644 --- a/crates/buzz-relay/src/metrics.rs +++ b/crates/buzz-relay/src/metrics.rs @@ -290,6 +290,7 @@ pub fn try_install(port: u16, gauge_idle_timeout_secs: u64) -> Result<(), Metric metrics::set_global_recorder(recorder) .map_err(|_error| MetricsInstallError::RecorderConflict)?; describe_readiness_metrics(); + describe_community_admission_metrics(); describe_db_pool_metrics(); describe_auth_metrics(); initialize_auth_metric_series(); @@ -306,24 +307,42 @@ pub fn install(port: u16, gauge_idle_timeout_secs: u64) { .unwrap_or_else(|error| panic!("metrics exporter must install exactly once: {error}")); } -/// Register the frozen readiness metric descriptions with the active recorder. +/// Register the frozen readiness and dependency-diagnostic metric descriptions. +/// +/// The two `buzz_readiness_*` probe families describe local process lifecycle. +/// The dependency families keep their names for dashboard continuity but are +/// published by the per-pod dependency sampler, not by the Kubernetes probe or +/// by an `/_status` request — a shared-dependency failure no longer deroutes +/// the pod, and nobody has to read the endpoint for the metrics to move. pub(crate) fn describe_readiness_metrics() { metrics::describe_counter!( "buzz_readiness_checks_total", - "Kubernetes health-listener readiness probes by terminal bounded reason" + "Kubernetes health-listener readiness probes by lifecycle reason (ready, shutting_down)" ); metrics::describe_counter!( "buzz_readiness_dependency_checks_total", - "Completed readiness dependency attempts by dependency and bounded outcome" + "Completed dependency-sampler attempts by dependency and bounded outcome" ); metrics::describe_histogram!( "buzz_readiness_check_duration_seconds", metrics::Unit::Seconds, - "Completed readiness check duration without outcome label multiplication" + "Completed dependency-sampler check duration without outcome label multiplication" ); metrics::describe_gauge!( "buzz_readiness_state", - "Latest publishable readiness state by check, where 1 is ready and 0 is not ready" + "Latest private readiness-probe observation, where 1 is ready and 0 is shutting down" + ); + metrics::describe_gauge!( + "buzz_readiness_dependency_sample_completed_timestamp_seconds", + "Unix time the cached /_status dependency report completed, absent until the first sample completes; sampler completion advances it and the publisher re-emits it" + ); +} + +/// Register the bounded community-admission contract. +pub(crate) fn describe_community_admission_metrics() { + metrics::describe_counter!( + "buzz_community_admission_checks_total", + "Durable community-active checks at socket admission by bounded outcome" ); } @@ -541,7 +560,17 @@ pub(crate) fn readiness_test_recorder() -> ( metrics_exporter_prometheus::PrometheusRecorder, metrics_exporter_prometheus::PrometheusHandle, ) { - let recorder = configured_prometheus_builder(300).build_recorder(); + readiness_test_recorder_with_idle_timeout(300) +} + +#[cfg(test)] +pub(crate) fn readiness_test_recorder_with_idle_timeout( + gauge_idle_timeout_secs: u64, +) -> ( + metrics_exporter_prometheus::PrometheusRecorder, + metrics_exporter_prometheus::PrometheusHandle, +) { + let recorder = configured_prometheus_builder(gauge_idle_timeout_secs).build_recorder(); let handle = recorder.handle(); (recorder, handle) } diff --git a/crates/buzz-relay/src/readiness.rs b/crates/buzz-relay/src/readiness.rs index 79a1a985707..f39689f20ac 100644 --- a/crates/buzz-relay/src/readiness.rs +++ b/crates/buzz-relay/src/readiness.rs @@ -1,43 +1,108 @@ -//! Readiness dependency evaluation and ordered metrics publication. +//! Readiness-probe telemetry and the dependency diagnostics behind `/_status`. //! -//! [`ReadinessCoordinator`] is process-owned. Its mutex is the linearization -//! point shared by health-probe commits and terminal shutdown, so an older -//! evaluation can never overwrite newer gauges or publish ready after shutdown. +//! Readiness is deliberately *not* a dependency question. A shared Postgres or +//! Redis failure is shared by every replica, so evaluating it in the probe took +//! the whole deployment out of the load balancer at once and left a reconnect +//! burst with nowhere to land. The Kubernetes probe therefore answers from this +//! process's own lifecycle (see [`crate::router`]), and the same dependency +//! evaluation is reported on the diagnostic `/_status` endpoint, which is never +//! wired to a probe. +//! +//! Dependency evaluation is also decoupled from requests. One per-pod loop +//! ([`run_dependency_sampler`]) evaluates on a fixed cadence, publishes the +//! dependency metrics, and caches the report; `/_status` only reads that cache. +//! Evaluating per request made the load a pressured dependency sees depend on +//! how often someone looked at the endpoint, with nothing bounding how many +//! evaluations could be in flight at once. use std::future::Future; -use std::sync::{Arc, Mutex, MutexGuard, PoisonError}; -use std::time::Duration; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; +use std::sync::{Arc, Mutex, PoisonError}; +use std::time::{Duration, SystemTime}; use buzz_db::{Db, DbError, DbReadinessOutcome}; use tokio::time::Instant; +use tokio_util::sync::CancellationToken; + +use crate::state::AppState; + +const DEPENDENCY_TIMEOUT: Duration = Duration::from_secs(2); + +/// Fixed cadence of the per-pod dependency sampling loop. +/// +/// Slow enough that a pod adds negligible load to a shared dependency, fast +/// enough that an operator opening `/_status` during an incident reads +/// something current. Documented in `deploy/charts/buzz/README.md`. +pub const DEPENDENCY_SAMPLE_INTERVAL: Duration = Duration::from_secs(30); -const READINESS_TIMEOUT: Duration = Duration::from_secs(2); +/// The completion-epoch publisher always refreshes at least this frequently. +/// +/// Main computes a cadence from the configured gauge idle timeout and caps it +/// at this value so the completion timestamp series survives idle eviction +/// without adding high-frequency noise. +const DEPENDENCY_SAMPLE_COMPLETION_REPUBLISH_MAX_INTERVAL: Duration = DEPENDENCY_SAMPLE_INTERVAL; + +/// Age past which a cached report is reported stale rather than current. +/// +/// Two cadences: one full cycle can be missed by an evaluation that consumed +/// its whole [`DEPENDENCY_TIMEOUT`] budget, so anything older than that means +/// the sampler itself is not keeping up. +const DEPENDENCY_SAMPLE_STALE_AFTER: Duration = DEPENDENCY_SAMPLE_INTERVAL.saturating_mul(2); /// Closed label set exported by `buzz_readiness_checks_total{reason}`. +/// +/// Readiness answers a local lifecycle question, so this set cannot grow with +/// the number of shared dependencies the relay talks to. #[cfg(test)] -pub(crate) const READINESS_REASON_LABELS: [&str; 12] = [ - "ready", - "shutting_down", - "postgres_pool_timeout", - "postgres_pool_error", - "postgres_query_timeout", - "postgres_query_error", - "redis_pool_timeout", - "redis_pool_error", - "deletion_catalog_timeout", - "deletion_catalog_error", - "overall_timeout", - "multiple_dependencies_failed", -]; - -/// Maximum raw Prometheus series emitted by readiness for one pod. +pub(crate) const READINESS_REASON_LABELS: [&str; 2] = ["ready", "shutting_down"]; + +/// Maximum raw Prometheus series emitted by readiness and its dependency +/// diagnostics for one pod. /// -/// - 12 overall reasons +/// - 2 probe reasons /// - 11 valid dependency/outcome pairs (Postgres 5, Redis 3, catalog 3) /// - 4 histograms x (15 configured buckets + `+Inf` + count + sum) = 72 -/// - 4 current-state gauges +/// - 2 readiness gauges (overall lifecycle + completion epoch) #[cfg(test)] -pub(crate) const READINESS_RAW_SERIES_PER_POD: usize = 12 + 11 + (4 * 18) + 4; +pub(crate) const READINESS_RAW_SERIES_PER_POD: usize = 2 + 11 + (4 * 18) + 1 + 1; + +/// Terminal outcome of one readiness probe. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum ReadinessReason { + Ready, + ShuttingDown, +} + +impl ReadinessReason { + pub(crate) fn label(self) -> &'static str { + match self { + Self::Ready => "ready", + Self::ShuttingDown => "shutting_down", + } + } + + pub(crate) fn is_ready(self) -> bool { + self == Self::Ready + } +} + +/// Records one readiness probe served by the private health listener. +/// +/// The counter and gauge describe the same immutable lifecycle observation. +/// The gauge is therefore the latest private readiness-probe observation, not +/// a transition-owned lifecycle mirror. +pub(crate) fn record_readiness_probe(reason: ReadinessReason) { + metrics::counter!( + "buzz_readiness_checks_total", + "reason" => reason.label(), + ) + .increment(1); + metrics::gauge!("buzz_readiness_state", "check" => "overall").set(if reason.is_ready() { + 1.0 + } else { + 0.0 + }); +} #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub(crate) enum PostgresOutcome { @@ -130,10 +195,14 @@ impl DeletionCatalogOutcome { } } +/// Aggregate dependency verdict reported in the `/_status` diagnostics body. +/// +/// This is a diagnostic field, never a metric label: it exists so an operator +/// reading `/_status` gets the same one-line summary the readiness body used to +/// carry. #[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub(crate) enum ReadinessReason { +pub(crate) enum DependencyReason { Ready, - ShuttingDown, PostgresPoolTimeout, PostgresPoolError, PostgresQueryTimeout, @@ -146,11 +215,10 @@ pub(crate) enum ReadinessReason { MultipleDependenciesFailed, } -impl ReadinessReason { +impl DependencyReason { pub(crate) fn label(self) -> &'static str { match self { Self::Ready => "ready", - Self::ShuttingDown => "shutting_down", Self::PostgresPoolTimeout => "postgres_pool_timeout", Self::PostgresPoolError => "postgres_pool_error", Self::PostgresQueryTimeout => "postgres_query_timeout", @@ -178,26 +246,18 @@ impl TimedOutcome { } } +/// One completed dependency evaluation. Every dependency always runs, so the +/// report carries three outcomes and never a partial shape. #[derive(Debug, Clone, Copy)] -pub(crate) struct ReadinessEvaluation { - postgres: Option>, - redis: Option>, - deletion_catalog: Option>, - pub(crate) reason: ReadinessReason, +pub(crate) struct DependencyReport { + postgres: TimedOutcome, + redis: TimedOutcome, + deletion_catalog: TimedOutcome, + pub(crate) reason: DependencyReason, total_duration: Duration, } -impl ReadinessEvaluation { - pub(crate) fn shutting_down() -> Self { - Self { - postgres: None, - redis: None, - deletion_catalog: None, - reason: ReadinessReason::ShuttingDown, - total_duration: Duration::ZERO, - } - } - +impl DependencyReport { #[cfg(test)] pub(crate) fn from_results( postgres: TimedOutcome, @@ -216,34 +276,24 @@ impl ReadinessEvaluation { ) -> Self { let reason = final_reason(postgres.outcome, redis.outcome, deletion_catalog.outcome); Self { - postgres: Some(postgres), - redis: Some(redis), - deletion_catalog: Some(deletion_catalog), + postgres, + redis, + deletion_catalog, reason, total_duration, } } - pub(crate) fn is_ready(self) -> bool { - self.reason == ReadinessReason::Ready - } - pub(crate) fn postgres_ready(self) -> bool { - self.postgres - .is_some_and(|result| result.outcome.is_success()) + self.postgres.outcome.is_success() } pub(crate) fn redis_ready(self) -> bool { - self.redis.is_some_and(|result| result.outcome.is_success()) + self.redis.outcome.is_success() } pub(crate) fn deletion_catalog_ready(self) -> bool { - self.deletion_catalog - .is_some_and(|result| result.outcome.is_success()) - } - - fn dependencies_ran(self) -> bool { - self.postgres.is_some() || self.redis.is_some() || self.deletion_catalog.is_some() + self.deletion_catalog.outcome.is_success() } } @@ -251,37 +301,39 @@ fn final_reason( postgres: PostgresOutcome, redis: RedisOutcome, deletion_catalog: DeletionCatalogOutcome, -) -> ReadinessReason { +) -> DependencyReason { let failure_count = usize::from(!postgres.is_success()) + usize::from(!redis.is_success()) + usize::from(!deletion_catalog.is_success()); if failure_count == 0 { - return ReadinessReason::Ready; + return DependencyReason::Ready; } if failure_count > 1 { let all_failures_are_timeouts = (postgres.is_success() || postgres.is_timeout()) && (redis.is_success() || redis.is_timeout()) && (deletion_catalog.is_success() || deletion_catalog.is_timeout()); return if all_failures_are_timeouts { - ReadinessReason::OverallTimeout + DependencyReason::OverallTimeout } else { - ReadinessReason::MultipleDependenciesFailed + DependencyReason::MultipleDependenciesFailed }; } match postgres { - PostgresOutcome::PoolTimeout => ReadinessReason::PostgresPoolTimeout, - PostgresOutcome::PoolError => ReadinessReason::PostgresPoolError, - PostgresOutcome::QueryTimeout => ReadinessReason::PostgresQueryTimeout, - PostgresOutcome::QueryError => ReadinessReason::PostgresQueryError, + PostgresOutcome::PoolTimeout => DependencyReason::PostgresPoolTimeout, + PostgresOutcome::PoolError => DependencyReason::PostgresPoolError, + PostgresOutcome::QueryTimeout => DependencyReason::PostgresQueryTimeout, + PostgresOutcome::QueryError => DependencyReason::PostgresQueryError, PostgresOutcome::Success => match redis { - RedisOutcome::PoolTimeout => ReadinessReason::RedisPoolTimeout, - RedisOutcome::PoolError => ReadinessReason::RedisPoolError, + RedisOutcome::PoolTimeout => DependencyReason::RedisPoolTimeout, + RedisOutcome::PoolError => DependencyReason::RedisPoolError, RedisOutcome::Success => match deletion_catalog { - DeletionCatalogOutcome::OperationTimeout => ReadinessReason::DeletionCatalogTimeout, - DeletionCatalogOutcome::OperationError => ReadinessReason::DeletionCatalogError, - DeletionCatalogOutcome::Success => ReadinessReason::Ready, + DeletionCatalogOutcome::OperationTimeout => { + DependencyReason::DeletionCatalogTimeout + } + DeletionCatalogOutcome::OperationError => DependencyReason::DeletionCatalogError, + DeletionCatalogOutcome::Success => DependencyReason::Ready, }, }, } @@ -303,7 +355,7 @@ async fn evaluate_dependencies( postgres: P, redis: R, deletion_catalog: D, -) -> ReadinessEvaluation +) -> DependencyReport where P: Future, R: Future, @@ -312,7 +364,7 @@ where let started_at = Instant::now(); let (postgres, redis, deletion_catalog) = tokio::join!(timed(postgres), timed(redis), timed(deletion_catalog),); - ReadinessEvaluation::for_dependencies(postgres, redis, deletion_catalog, started_at.elapsed()) + DependencyReport::for_dependencies(postgres, redis, deletion_catalog, started_at.elapsed()) } async fn redis_check(pool: &deadpool_redis::Pool, deadline: Instant) -> RedisOutcome { @@ -345,16 +397,16 @@ fn classify_deletion_catalog_result(result: buzz_db::Result<()>) -> DeletionCata } #[async_trait::async_trait] -pub(crate) trait ReadinessEvaluator: Send + Sync { - async fn evaluate(&self, db: &Db, redis_pool: &deadpool_redis::Pool) -> ReadinessEvaluation; +pub(crate) trait DependencyEvaluator: Send + Sync { + async fn evaluate(&self, db: &Db, redis_pool: &deadpool_redis::Pool) -> DependencyReport; } -struct ProductionReadinessEvaluator; +struct ProductionDependencyEvaluator; #[async_trait::async_trait] -impl ReadinessEvaluator for ProductionReadinessEvaluator { - async fn evaluate(&self, db: &Db, redis_pool: &deadpool_redis::Pool) -> ReadinessEvaluation { - let deadline = Instant::now() + READINESS_TIMEOUT; +impl DependencyEvaluator for ProductionDependencyEvaluator { + async fn evaluate(&self, db: &Db, redis_pool: &deadpool_redis::Pool) -> DependencyReport { + let deadline = Instant::now() + DEPENDENCY_TIMEOUT; evaluate_dependencies( async { db.readiness_check(deadline).await.into() }, redis_check(redis_pool, deadline), @@ -364,150 +416,280 @@ impl ReadinessEvaluator for ProductionReadinessEvaluator { } } +/// One completed evaluation and when it was observed. #[derive(Debug, Clone, Copy)] -pub(crate) struct ProbeTicket { - generation: u64, +struct DependencySample { + report: DependencyReport, + observed_at: Instant, } +/// What the cache can tell `/_status`. +/// +/// "No report yet" is a distinct state, not a fabricated healthy one, and a +/// report is always accompanied by its age: a cached verdict presented without +/// one would read as authoritative however long ago it was taken. #[derive(Debug, Clone, Copy)] -pub(crate) enum ProbeStart { - Evaluate(ProbeTicket), - ShuttingDown, +pub(crate) enum DependencySnapshot { + /// The sampler has not completed its first evaluation yet. + NotYetSampled, + Sampled { + report: DependencyReport, + age: Duration, + stale: bool, + }, } -#[derive(Debug, Default)] -struct PublicationState { - next_generation: u64, - latest_published_generation: u64, - shutdown_generation: Option, +/// Republish cadence for the completion-epoch gauge. +/// +/// The configured gauge idle timeout comes from `main.rs`; this returns a +/// bounded cadence that is strictly below that timeout and never slower than +/// the dependency sampler itself. +pub fn dependency_sample_completion_republish_interval(gauge_idle_timeout_secs: u64) -> Duration { + let idle_timeout_secs = gauge_idle_timeout_secs.max(1); + let refresh_secs = idle_timeout_secs.saturating_div(3).max(1); + let strict_upper_bound_secs = idle_timeout_secs.saturating_sub(1).max(1); + Duration::from_secs( + refresh_secs + .min(strict_upper_bound_secs) + .min(DEPENDENCY_SAMPLE_COMPLETION_REPUBLISH_MAX_INTERVAL.as_secs()), + ) } -/// Serializes readiness result publication with terminal process shutdown. -pub(crate) struct ReadinessCoordinator { - state: Mutex, - evaluator: Arc, +/// The per-pod owner of shared-dependency evaluation. +/// +/// The production [`run_dependency_sampler`] loop is the sole runtime owner of +/// [`Self::sample`], so at most one evaluation exists at a time and no request +/// path can start another. The companion completion publisher only reads the +/// stored epoch, while `/_status` reads [`Self::snapshot`]; neither touches a +/// dependency. +pub(crate) struct DependencyDiagnostics { + evaluator: Arc, + latest: Mutex>, + sample_completion_epoch_seconds: AtomicU64, + sample_completion_written: AtomicBool, } -impl Default for ReadinessCoordinator { +impl Default for DependencyDiagnostics { fn default() -> Self { - Self { - state: Mutex::new(PublicationState::default()), - evaluator: Arc::new(ProductionReadinessEvaluator), - } + Self::with_evaluator(Arc::new(ProductionDependencyEvaluator)) } } -impl ReadinessCoordinator { - #[cfg(test)] - pub(crate) fn with_evaluator(evaluator: Arc) -> Self { +impl DependencyDiagnostics { + pub(crate) fn with_evaluator(evaluator: Arc) -> Self { Self { - state: Mutex::new(PublicationState::default()), evaluator, + latest: Mutex::new(None), + sample_completion_epoch_seconds: AtomicU64::new(0), + sample_completion_written: AtomicBool::new(false), } } - fn lock_state(&self) -> MutexGuard<'_, PublicationState> { - self.state.lock().unwrap_or_else(PoisonError::into_inner) - } - - pub(crate) async fn evaluate( - &self, - db: &Db, - redis_pool: &deadpool_redis::Pool, - ) -> ReadinessEvaluation { - self.evaluator.evaluate(db, redis_pool).await + /// Runs one bounded evaluation, publishes its telemetry, and replaces the + /// cached report. + /// + /// The completion timestamp is written last, once the cache already serves + /// this report, so the gauge can never describe a sample `/_status` is not + /// yet answering with. + pub(crate) async fn sample(&self, db: &Db, redis_pool: &deadpool_redis::Pool) { + let report = self.evaluator.evaluate(db, redis_pool).await; + record_dependency_report(&report); + let sample = DependencySample { + report, + observed_at: Instant::now(), + }; + *self.latest.lock().unwrap_or_else(PoisonError::into_inner) = Some(sample); + self.record_dependency_sample_completion(SystemTime::now()); } - /// Allocates a health-probe generation or records a truthful shutdown fast path. - pub(crate) fn begin_probe(&self) -> ProbeStart { - let mut state = self.lock_state(); - if state.shutdown_generation.is_some() { - let evaluation = ReadinessEvaluation::shutting_down(); - record_attempt_metrics(&evaluation, ReadinessReason::ShuttingDown); - record_overall_state(false); - return ProbeStart::ShuttingDown; + /// The latest completed evaluation with its age. Starts no dependency work. + pub(crate) fn snapshot(&self) -> DependencySnapshot { + let latest = *self.latest.lock().unwrap_or_else(PoisonError::into_inner); + match latest { + None => DependencySnapshot::NotYetSampled, + Some(sample) => { + let age = sample.observed_at.elapsed(); + DependencySnapshot::Sampled { + report: sample.report, + age, + stale: age > DEPENDENCY_SAMPLE_STALE_AFTER, + } + } } - - state.next_generation = state.next_generation.saturating_add(1); - ProbeStart::Evaluate(ProbeTicket { - generation: state.next_generation, - }) } - /// Commits one completed health probe through the shared publication fence. - pub(crate) fn finish_probe( - &self, - ticket: ProbeTicket, - evaluation: ReadinessEvaluation, - ) -> ReadinessEvaluation { - let mut state = self.lock_state(); - if state.shutdown_generation.is_some() { - record_attempt_metrics(&evaluation, ReadinessReason::ShuttingDown); - return ReadinessEvaluation::shutting_down(); - } - - record_attempt_metrics(&evaluation, evaluation.reason); - if ticket.generation > state.latest_published_generation { - record_current_state(&evaluation); - state.latest_published_generation = ticket.generation; - } - evaluation + fn record_dependency_sample_completion(&self, completed_at: SystemTime) { + let epoch_seconds = dependency_sample_completion_epoch_seconds(completed_at); + self.sample_completion_epoch_seconds + .store(epoch_seconds, Ordering::Release); + self.sample_completion_written + .store(true, Ordering::Release); + publish_dependency_sample_completion_metric(epoch_seconds); } - /// Returns whether a compatibility/public readiness evaluation may start. - pub(crate) fn public_evaluation_allowed(&self) -> bool { - self.lock_state().shutdown_generation.is_none() + fn latest_sample_completion_epoch_seconds(&self) -> Option { + self.sample_completion_written + .load(Ordering::Acquire) + .then(|| self.sample_completion_epoch_seconds.load(Ordering::Acquire)) } - /// Makes shutdown dominate a public request that was already in flight. - pub(crate) fn finish_public_evaluation( - &self, - evaluation: ReadinessEvaluation, - ) -> ReadinessEvaluation { - if self.lock_state().shutdown_generation.is_some() { - ReadinessEvaluation::shutting_down() - } else { - evaluation + /// Re-emits the stored epoch, then checks it is still the stored one. + /// + /// A sample can complete and publish a newer epoch while this write is in + /// flight, and the metrics facade exposes only a bare `set`, so nothing + /// downstream rejects the older value once it lands last. Verifying after + /// the write — rather than before it — is what closes that window: a + /// republish that lost the race re-emits the newer epoch instead of + /// leaving the exported series moved backwards. Another pass costs another + /// completed sample, so this ends as soon as no completion is racing it. + fn republish_dependency_sample_completion(&self) { + let mut published = None; + while let Some(epoch_seconds) = self.latest_sample_completion_epoch_seconds() { + if published == Some(epoch_seconds) { + break; + } + publish_dependency_sample_completion_metric(epoch_seconds); + published = Some(epoch_seconds); } } +} + +/// Runs the per-pod dependency sampling loop until `cancel` fires. +/// +/// One loop, one fixed cadence, each evaluation awaited before the next tick is +/// taken, so this pod never has two evaluations in flight. `Skip` matches the +/// community revalidator: an evaluation that overruns its slot delays the next +/// cycle instead of queueing a catch-up burst into the dependency that was +/// already slow. The first tick fires immediately, so the not-yet-sampled +/// window is one evaluation long. +pub async fn run_dependency_sampler(state: Arc, cancel: CancellationToken) { + run_dependency_sampler_for_diagnostics( + Arc::clone(&state.dependency_diagnostics), + state.db.clone(), + state.redis_pool.clone(), + DEPENDENCY_SAMPLE_INTERVAL, + cancel, + ) + .await; +} - /// Commits terminal shutdown and immediately publishes overall not-ready. - pub(crate) fn begin_shutdown(&self) { - let mut state = self.lock_state(); - if state.shutdown_generation.is_none() { - let generation = state.next_generation.saturating_add(1); - state.shutdown_generation = Some(generation); - record_overall_state(false); +/// Starts the dependency sampler and completion republisher together. +/// +/// Main calls this once at startup so sampler and publisher ownership lives at +/// one seam instead of being wired independently. +pub fn start_dependency_sampler_and_completion_publisher( + state: Arc, + republish_interval: Duration, +) { + let sampler_cancel = state.dependency_sampler_cancel.clone(); + tokio::spawn(run_dependency_sampler(Arc::clone(&state), sampler_cancel)); + + let publisher_cancel = state.dependency_completion_publisher_cancel.clone(); + tokio::spawn(run_dependency_sample_completion_publisher( + state, + republish_interval, + publisher_cancel, + )); +} + +async fn run_dependency_sampler_for_diagnostics( + diagnostics: Arc, + db: Db, + redis_pool: deadpool_redis::Pool, + sample_interval: Duration, + cancel: CancellationToken, +) { + let mut interval = tokio::time::interval(sample_interval.max(Duration::from_millis(1))); + interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + loop { + tokio::select! { + biased; + _ = cancel.cancelled() => break, + _ = interval.tick() => { + diagnostics.sample(&db, &redis_pool).await; + } } } } -fn record_attempt_metrics(evaluation: &ReadinessEvaluation, reason: ReadinessReason) { - metrics::counter!( - "buzz_readiness_checks_total", - "reason" => reason.label(), +/// Re-emits the latest completion epoch so the gauge survives recorder idle +/// eviction even when no newer dependency sample completes. +pub async fn run_dependency_sample_completion_publisher( + state: Arc, + republish_interval: Duration, + cancel: CancellationToken, +) { + run_dependency_sample_completion_publisher_for_diagnostics( + Arc::clone(&state.dependency_diagnostics), + republish_interval, + cancel, ) - .increment(1); + .await; +} - if !evaluation.dependencies_ran() { - return; +async fn run_dependency_sample_completion_publisher_for_diagnostics( + diagnostics: Arc, + republish_interval: Duration, + cancel: CancellationToken, +) { + let mut interval = tokio::time::interval(republish_interval.max(Duration::from_millis(1))); + interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + loop { + tokio::select! { + biased; + _ = cancel.cancelled() => break, + _ = interval.tick() => diagnostics.republish_dependency_sample_completion(), + } } +} +/// Records one dependency evaluation. Counters and durations only — dependency +/// health has no publishable "current state" now that no probe consumes it; a +/// per-dependency gauge would read as an authoritative verdict on infrastructure +/// this pod only samples every [`DEPENDENCY_SAMPLE_INTERVAL`]. +fn record_dependency_report(report: &DependencyReport) { metrics::histogram!( "buzz_readiness_check_duration_seconds", "check" => "overall", ) - .record(evaluation.total_duration.as_secs_f64()); + .record(report.total_duration.as_secs_f64()); + + record_dependency_attempt( + "postgres", + report.postgres.outcome.label(), + report.postgres.duration, + ); + record_dependency_attempt("redis", report.redis.outcome.label(), report.redis.duration); + record_dependency_attempt( + "deletion_catalog", + report.deletion_catalog.outcome.label(), + report.deletion_catalog.duration, + ); +} - if let Some(result) = evaluation.postgres { - record_dependency_attempt("postgres", result.outcome.label(), result.duration); - } - if let Some(result) = evaluation.redis { - record_dependency_attempt("redis", result.outcome.label(), result.duration); - } - if let Some(result) = evaluation.deletion_catalog { - record_dependency_attempt("deletion_catalog", result.outcome.label(), result.duration); - } +/// Publishes the Unix time the cached report completed. +/// +/// A completion timestamp rather than an age, because age then belongs to the +/// query — `time() - buzz_readiness_dependency_sample_completed_timestamp_seconds` +/// — and grows on its own while this pod is wedged. A gauge carrying the age +/// needs a writer to advance it, so the one failure it most needs to expose, a +/// sampler that stopped running, is the one that would freeze it at its last +/// value and read as permanently fresh. A completed sample is the only thing +/// that may advance this epoch; the independent publisher only re-emits the +/// stored value so the series survives gauge idle-eviction. This keeps the +/// `buzz_storage_sweep_age_seconds` convention that absence means +/// "not yet sampled" rather than "fresh". +/// +fn dependency_sample_completion_epoch_seconds(completed_at: SystemTime) -> u64 { + completed_at + .duration_since(SystemTime::UNIX_EPOCH) + .unwrap_or_default() + .as_secs() +} + +fn publish_dependency_sample_completion_metric(epoch_seconds: u64) { + metrics::gauge!("buzz_readiness_dependency_sample_completed_timestamp_seconds") + .set(epoch_seconds as f64); } fn record_dependency_attempt(dependency: &'static str, outcome: &'static str, duration: Duration) { @@ -524,35 +706,6 @@ fn record_dependency_attempt(dependency: &'static str, outcome: &'static str, du .record(duration.as_secs_f64()); } -fn record_current_state(evaluation: &ReadinessEvaluation) { - record_overall_state(evaluation.is_ready()); - if let Some(result) = evaluation.postgres { - record_dependency_state("postgres", result.outcome.is_success()); - } - if let Some(result) = evaluation.redis { - record_dependency_state("redis", result.outcome.is_success()); - } - if let Some(result) = evaluation.deletion_catalog { - record_dependency_state("deletion_catalog", result.outcome.is_success()); - } -} - -fn record_overall_state(ready: bool) { - metrics::gauge!("buzz_readiness_state", "check" => "overall").set(if ready { - 1.0 - } else { - 0.0 - }); -} - -fn record_dependency_state(dependency: &'static str, ready: bool) { - metrics::gauge!("buzz_readiness_state", "check" => dependency).set(if ready { - 1.0 - } else { - 0.0 - }); -} - #[cfg(test)] mod tests { use metrics_util::debugging::{DebugValue, DebuggingRecorder}; @@ -560,17 +713,15 @@ mod tests { use super::*; - fn ready_evaluation() -> ReadinessEvaluation { - ReadinessEvaluation::from_results( - TimedOutcome::new(PostgresOutcome::Success, Duration::from_millis(35)), - TimedOutcome::new(RedisOutcome::Success, Duration::from_millis(10)), - TimedOutcome::new(DeletionCatalogOutcome::Success, Duration::from_millis(20)), - Duration::from_millis(35), - ) - } + type Snapshot = Vec<( + CompositeKey, + Option, + Option, + DebugValue, + )>; - fn redis_failure_evaluation() -> ReadinessEvaluation { - ReadinessEvaluation::from_results( + fn redis_failure_report() -> DependencyReport { + DependencyReport::from_results( TimedOutcome::new(PostgresOutcome::Success, Duration::from_millis(35)), TimedOutcome::new(RedisOutcome::PoolTimeout, Duration::from_secs(2)), TimedOutcome::new(DeletionCatalogOutcome::Success, Duration::from_millis(20)), @@ -579,12 +730,7 @@ mod tests { } fn exact_metric<'a>( - snapshot: &'a [( - CompositeKey, - Option, - Option, - DebugValue, - )], + snapshot: &'a Snapshot, name: &str, labels: &[(&str, &str)], ) -> Option<&'a DebugValue> { @@ -601,15 +747,7 @@ mod tests { }) } - fn gauge_value( - snapshot: &[( - CompositeKey, - Option, - Option, - DebugValue, - )], - check: &str, - ) -> f64 { + fn gauge_value(snapshot: &Snapshot, check: &str) -> f64 { let value = exact_metric(snapshot, "buzz_readiness_state", &[("check", check)]) .expect("readiness gauge"); let DebugValue::Gauge(value) = value else { @@ -620,7 +758,7 @@ mod tests { #[tokio::test(start_paused = true)] async fn evaluation_preserves_a_completed_check_when_another_times_out() { - let evaluation = evaluate_dependencies( + let report = evaluate_dependencies( async { tokio::time::sleep(Duration::from_millis(35)).await; PostgresOutcome::Success @@ -636,15 +774,9 @@ mod tests { ) .await; - assert_eq!(evaluation.reason, ReadinessReason::RedisPoolTimeout); - assert_eq!( - evaluation.postgres.map(|result| result.duration), - Some(Duration::from_millis(35)) - ); - assert_eq!( - evaluation.redis.map(|result| result.duration), - Some(Duration::from_secs(2)) - ); + assert_eq!(report.reason, DependencyReason::RedisPoolTimeout); + assert_eq!(report.postgres.duration, Duration::from_millis(35)); + assert_eq!(report.redis.duration, Duration::from_secs(2)); } #[test] @@ -655,7 +787,7 @@ mod tests { RedisOutcome::PoolTimeout, DeletionCatalogOutcome::Success, ), - ReadinessReason::OverallTimeout + DependencyReason::OverallTimeout ); } @@ -696,7 +828,7 @@ mod tests { .map(DeletionCatalogOutcome::label), ["success", "operation_timeout", "operation_error"] ); - assert_eq!(READINESS_RAW_SERIES_PER_POD, 99); + assert_eq!(READINESS_RAW_SERIES_PER_POD, 87); } #[test] @@ -712,144 +844,707 @@ mod tests { } #[test] - fn slow_older_failure_cannot_overwrite_newer_success_gauges() { - let coordinator = ReadinessCoordinator::default(); - let ProbeStart::Evaluate(slow_a) = coordinator.begin_probe() else { - panic!("serving probe A"); - }; - let ProbeStart::Evaluate(fast_b) = coordinator.begin_probe() else { - panic!("serving probe B"); - }; + fn completion_timestamp_republish_interval_stays_below_idle_timeout() { + assert_eq!( + dependency_sample_completion_republish_interval(900), + DEPENDENCY_SAMPLE_INTERVAL, + "default idle timeout should republish on the sampler cadence" + ); + assert_eq!( + dependency_sample_completion_republish_interval(15), + Duration::from_secs(5), + "the minimum configured idle timeout must still get multiple republishes" + ); + assert!(dependency_sample_completion_republish_interval(15) < Duration::from_secs(15)); + } + + /// The readiness gauge and counter use the same immutable reason sampled by + /// the private probe. A dependency evaluation — however bad — must never + /// move them, which is what let a shared outage deroute every replica at + /// once. + #[test] + fn readiness_telemetry_tracks_lifecycle_and_dependency_failure_never_moves_it() { let recorder = DebuggingRecorder::new(); let snapshotter = recorder.snapshotter(); metrics::with_local_recorder(&recorder, || { - coordinator.finish_probe(fast_b, ready_evaluation()); - coordinator.finish_probe(slow_a, redis_failure_evaluation()); + record_readiness_probe(ReadinessReason::Ready); + record_dependency_report(&redis_failure_report()); }); - let snapshot = snapshotter.snapshot().into_vec(); + let after_failure = snapshotter.snapshot().into_vec(); - assert_eq!(gauge_value(&snapshot, "overall"), 1.0); - assert_eq!(gauge_value(&snapshot, "redis"), 1.0); + assert_eq!(gauge_value(&after_failure, "overall"), 1.0); assert!(matches!( exact_metric( - &snapshot, + &after_failure, "buzz_readiness_checks_total", &[("reason", "ready")] ), Some(DebugValue::Counter(1)) )); + assert!( + matches!( + exact_metric( + &after_failure, + "buzz_readiness_dependency_checks_total", + &[("dependency", "redis"), ("outcome", "pool_timeout")] + ), + Some(DebugValue::Counter(1)) + ), + "dependency diagnostics must still be counted" + ); + for dependency in ["postgres", "redis", "deletion_catalog"] { + assert!( + exact_metric( + &after_failure, + "buzz_readiness_state", + &[("check", dependency)] + ) + .is_none(), + "{dependency} must not publish a readiness gauge" + ); + } + + metrics::with_local_recorder(&recorder, || { + record_readiness_probe(ReadinessReason::ShuttingDown); + }); + let after_shutdown = snapshotter.snapshot().into_vec(); + + assert_eq!(gauge_value(&after_shutdown, "overall"), 0.0); assert!(matches!( exact_metric( - &snapshot, + &after_shutdown, "buzz_readiness_checks_total", - &[("reason", "redis_pool_timeout")] + &[("reason", "shutting_down")] ), Some(DebugValue::Counter(1)) )); } + /// A shutdown probe records no dependency attempt or latency sample: it did + /// not evaluate anything, and fabricating a sample would misreport the + /// dependency's real health during a rollout. #[test] - fn slow_older_success_cannot_overwrite_newer_failure_gauges() { - let coordinator = ReadinessCoordinator::default(); - let ProbeStart::Evaluate(slow_a) = coordinator.begin_probe() else { - panic!("serving probe A"); - }; - let ProbeStart::Evaluate(fast_b) = coordinator.begin_probe() else { - panic!("serving probe B"); - }; + fn a_readiness_probe_never_records_dependency_attempts() { let recorder = DebuggingRecorder::new(); let snapshotter = recorder.snapshotter(); metrics::with_local_recorder(&recorder, || { - coordinator.finish_probe(fast_b, redis_failure_evaluation()); - coordinator.finish_probe(slow_a, ready_evaluation()); + record_readiness_probe(ReadinessReason::Ready); + record_readiness_probe(ReadinessReason::ShuttingDown); }); let snapshot = snapshotter.snapshot().into_vec(); - assert_eq!(gauge_value(&snapshot, "overall"), 0.0); - assert_eq!(gauge_value(&snapshot, "postgres"), 1.0); - assert_eq!(gauge_value(&snapshot, "redis"), 0.0); - assert_eq!(gauge_value(&snapshot, "deletion_catalog"), 1.0); + assert!(snapshot.iter().all(|(key, _, _, _)| { + key.key().name() != "buzz_readiness_dependency_checks_total" + && key.key().name() != "buzz_readiness_check_duration_seconds" + })); } + /// A `Db` and a Redis pool on a closed port. The scripted evaluators below + /// never touch either, so no connection is ever attempted; they exist only + /// to satisfy the production `sample` signature. + fn unreachable_dependencies() -> (Db, deadpool_redis::Pool) { + let pool = sqlx::PgPool::connect_lazy("postgres://127.0.0.1:1/buzz").expect("lazy pg pool"); + let redis_pool = deadpool_redis::Config::from_url("redis://127.0.0.1:1") + .create_pool(Some(deadpool_redis::Runtime::Tokio1)) + .expect("redis pool"); + (Db::from_pool(pool), redis_pool) + } + + struct FixedEvaluator(DependencyReport); + + #[async_trait::async_trait] + impl DependencyEvaluator for FixedEvaluator { + async fn evaluate(&self, _db: &Db, _redis_pool: &deadpool_redis::Pool) -> DependencyReport { + self.0 + } + } + + struct DelayedEvaluator { + report: DependencyReport, + started: Arc, + release: Arc, + } + + #[async_trait::async_trait] + impl DependencyEvaluator for DelayedEvaluator { + async fn evaluate(&self, _db: &Db, _redis_pool: &deadpool_redis::Pool) -> DependencyReport { + self.started.notify_one(); + self.release.notified().await; + self.report + } + } + + fn ready_report() -> DependencyReport { + DependencyReport::from_results( + TimedOutcome::new(PostgresOutcome::Success, Duration::from_millis(3)), + TimedOutcome::new(RedisOutcome::Success, Duration::from_millis(2)), + TimedOutcome::new(DeletionCatalogOutcome::Success, Duration::from_millis(1)), + Duration::from_millis(3), + ) + } + + fn ready_diagnostics() -> DependencyDiagnostics { + DependencyDiagnostics::with_evaluator(Arc::new(FixedEvaluator(ready_report()))) + } + + fn sampled(snapshot: DependencySnapshot) -> (Duration, bool) { + let DependencySnapshot::Sampled { age, stale, .. } = snapshot else { + panic!("a completed sample must be reported as sampled"); + }; + (age, stale) + } + + /// An operator reading `/_status` must be able to tell "nothing has been + /// sampled yet" from "this is current" from "this outlived the sampler". + /// A cached report presented without its age would read as authoritative + /// however old it is. + #[tokio::test(start_paused = true)] + async fn freshness_separates_not_yet_sampled_from_a_fresh_and_a_stale_report() { + let (db, redis_pool) = unreachable_dependencies(); + let diagnostics = ready_diagnostics(); + + assert!( + matches!(diagnostics.snapshot(), DependencySnapshot::NotYetSampled), + "no evaluation has completed, so there is nothing to report" + ); + + diagnostics.sample(&db, &redis_pool).await; + assert_eq!(sampled(diagnostics.snapshot()), (Duration::ZERO, false)); + + tokio::time::advance(DEPENDENCY_SAMPLE_INTERVAL).await; + assert_eq!( + sampled(diagnostics.snapshot()), + (DEPENDENCY_SAMPLE_INTERVAL, false), + "one cadence of age is the steady state, not staleness" + ); + + tokio::time::advance(DEPENDENCY_SAMPLE_INTERVAL + Duration::from_secs(1)).await; + let (age, stale) = sampled(diagnostics.snapshot()); + assert_eq!(age, DEPENDENCY_SAMPLE_INTERVAL * 2 + Duration::from_secs(1)); + assert!(stale, "a report that outlived two cadences missed a cycle"); + } + + /// The completion timestamp written since the previous snapshot, if any. + /// + /// `Snapshotter::snapshot` drains, so a window in which nothing wrote this + /// metric either omits the key entirely or, once registered, replays as + /// `0`. Zero is not a time any sample could have completed at, so folding + /// it into `None` keeps "nobody wrote this in that window" expressible — + /// the property a completion timestamp must have and an age cannot. + fn sample_completion_metric_write(snapshot: &Snapshot) -> Option { + exact_metric( + snapshot, + "buzz_readiness_dependency_sample_completed_timestamp_seconds", + &[], + ) + .map(|value| { + let DebugValue::Gauge(value) = value else { + panic!("the sample completion timestamp must be a gauge"); + }; + value.into_inner() as u64 + }) + .filter(|written| *written != 0) + } + + fn epoch_seconds_now() -> u64 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .expect("the system clock is after the Unix epoch") + .as_secs() + } + + fn sample_completion_metric_value_from_scrape(scrape: &str) -> Option { + scrape.lines().find_map(|line| { + line.strip_prefix("buzz_readiness_dependency_sample_completed_timestamp_seconds ") + .and_then(|value| value.parse::().ok()) + .map(|value| value as u64) + }) + } + + fn sample_completion_metric_type_from_scrape(scrape: &str) -> Option<&str> { + scrape.lines().find_map(|line| { + line.strip_prefix( + "# TYPE buzz_readiness_dependency_sample_completed_timestamp_seconds ", + ) + }) + } + + /// Minimal model of Datadog OpenMetrics v2 with `send_monotonic_counter:true`. + /// + /// Gauge values are forwarded as-is. Counter values become per-scrape deltas + /// (`monotonic_count` / `.count`), so the exported sample itself is no + /// longer available to monitor queries. + fn datadog_openmetrics_v2_completion_epoch( + scrape: &str, + previous_counter_raw: &mut Option, + ) -> Option { + let raw = sample_completion_metric_value_from_scrape(scrape)?; + match sample_completion_metric_type_from_scrape(scrape) { + Some("gauge") => Some(raw), + Some("counter") => { + let transformed = previous_counter_raw + .map(|previous| raw.saturating_sub(previous)) + .unwrap_or(raw); + *previous_counter_raw = Some(raw); + Some(transformed) + } + Some(other) => panic!("unexpected metric type: {other}"), + None => panic!("missing sample completion metric type"), + } + } + + /// The scrape-side half of the same question. Freshness is published as the + /// Unix time the cached report completed, written only by a completed + /// sample, so the age belongs to the query + /// (`time() - buzz_readiness_dependency_sample_completed_timestamp_seconds`) + /// and grows on its own while this pod is wedged. A gauge carrying the age + /// instead would need a writer to advance it, so the one failure it most + /// needs to expose — a sampler that stopped — is the one that would freeze + /// it at its last value and read as permanently fresh. #[test] - fn shutdown_fast_path_preserves_dependency_state_and_histograms() { - let coordinator = ReadinessCoordinator::default(); + fn the_completion_timestamp_gauge_tracks_completed_samples() { + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .start_paused(true) + .build() + .expect("paused current-thread runtime"); let recorder = DebuggingRecorder::new(); let snapshotter = recorder.snapshotter(); metrics::with_local_recorder(&recorder, || { - let ProbeStart::Evaluate(ticket) = coordinator.begin_probe() else { - panic!("initial serving probe"); - }; - coordinator.finish_probe(ticket, ready_evaluation()); - coordinator.begin_shutdown(); - assert!(matches!( - coordinator.begin_probe(), - ProbeStart::ShuttingDown - )); + runtime.block_on(async { + let (db, redis_pool) = unreachable_dependencies(); + let diagnostics = ready_diagnostics(); + + assert_eq!( + sample_completion_metric_write(&snapshotter.snapshot().into_vec()), + None, + "absence is the not-yet-sampled signal, not a fresh zero" + ); + + let before_first = epoch_seconds_now(); + diagnostics.sample(&db, &redis_pool).await; + let after_first = epoch_seconds_now(); + let first = sample_completion_metric_write(&snapshotter.snapshot().into_vec()) + .expect("a completed sample must publish when it completed"); + assert!( + (before_first..=after_first).contains(&first), + "{first} must be the wall time the sample completed, \ + not an age, and not a stale reading" + ); + + // Paused time: no task, timer, or aging loop can run here. + tokio::time::advance(DEPENDENCY_SAMPLE_STALE_AFTER + Duration::from_secs(1)).await; + assert_eq!( + sample_completion_metric_write(&snapshotter.snapshot().into_vec()), + None, + "only a completed sample may write the gauge, so the age a \ + query derives from it grows with no server-side writer" + ); + assert!( + sampled(diagnostics.snapshot()).1, + "the cached `/_status` age reports the same report as stale" + ); + + let before_second = epoch_seconds_now(); + diagnostics.sample(&db, &redis_pool).await; + let after_second = epoch_seconds_now(); + let second = sample_completion_metric_write(&snapshotter.snapshot().into_vec()) + .expect("the next completed sample republishes the timestamp"); + assert!( + (before_second..=after_second).contains(&second) && second >= first, + "{second} must re-anchor to the second completion, after {first}" + ); + assert!( + !sampled(diagnostics.snapshot()).1, + "a fresh completion clears staleness" + ); + }); }); - let after = snapshotter.snapshot().into_vec(); + } - for dependency in ["postgres", "redis", "deletion_catalog"] { - assert_eq!( - gauge_value(&after, dependency), - 1.0, - "shutdown must not fabricate {dependency} state" - ); - } - for check in ["overall", "postgres", "redis", "deletion_catalog"] { + #[test] + fn openmetrics_v2_preserves_epoch_only_when_raw_prometheus_type_is_gauge() { + let (recorder, handle) = crate::metrics::readiness_test_recorder(); + let mut previous_counter_raw = None; + + assert_eq!( + datadog_openmetrics_v2_completion_epoch(&handle.render(), &mut previous_counter_raw), + None, + "absence before the first completion remains absence after transform" + ); + + let first = 4_750_000_001_u64; + metrics::with_local_recorder(&recorder, || { + publish_dependency_sample_completion_metric(first); + }); + let after_first_completion = handle.render(); + assert_eq!( + datadog_openmetrics_v2_completion_epoch( + &after_first_completion, + &mut previous_counter_raw, + ), + Some(first), + "after first completion the transformed value must be the completion epoch" + ); + + // Same raw scrape value after idle timeout: sampler produced no new completion. + assert_eq!( + datadog_openmetrics_v2_completion_epoch( + &after_first_completion, + &mut previous_counter_raw, + ), + Some(first), + "idle-timeout survival still must preserve the full epoch" + ); + + let second = 4_750_000_005_u64; + metrics::with_local_recorder(&recorder, || { + publish_dependency_sample_completion_metric(second); + }); + let after_later_completion = handle.render(); + assert_eq!( + datadog_openmetrics_v2_completion_epoch( + &after_later_completion, + &mut previous_counter_raw, + ), + Some(second), + "a later completion must transform to the new full epoch" + ); + + // Same raw value again while only the sampler is stopped. + assert_eq!( + datadog_openmetrics_v2_completion_epoch( + &after_later_completion, + &mut previous_counter_raw, + ), + Some(second), + "sampler stoppage must not collapse the epoch to a delta" + ); + } + + async fn wait_for_sample_completion_scrape_value( + handle: &metrics_exporter_prometheus::PrometheusHandle, + predicate: impl Fn(u64) -> bool, + ) -> u64 { + let deadline = std::time::Instant::now() + Duration::from_secs(5); + loop { + if let Some(value) = sample_completion_metric_value_from_scrape(&handle.render()) { + if predicate(value) { + return value; + } + } assert!( - matches!( - exact_metric( - &after, - "buzz_readiness_check_duration_seconds", - &[("check", check)] - ), - Some(DebugValue::Histogram(values)) if values.len() == 1 - ), - "shutdown fast path must not add a {check} duration" + std::time::Instant::now() < deadline, + "timed out waiting for completion timestamp scrape value" ); + tokio::time::sleep(Duration::from_millis(25)).await; } - assert_eq!(gauge_value(&after, "overall"), 0.0); - assert!(matches!( - exact_metric( - &after, - "buzz_readiness_checks_total", - &[("reason", "shutting_down")] - ), - Some(DebugValue::Counter(1)) - )); } #[test] - fn shutdown_dominates_an_in_flight_success_without_resurrecting_gauges() { - let coordinator = ReadinessCoordinator::default(); - let ProbeStart::Evaluate(ticket) = coordinator.begin_probe() else { - panic!("serving probe"); - }; - let recorder = DebuggingRecorder::new(); - let snapshotter = recorder.snapshotter(); + fn dependency_runtime_retains_the_completion_timestamp_when_only_sampler_stops() { + let timeout = Duration::from_secs(1); + let republish_interval = Duration::from_millis(100); + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .start_paused(true) + .build() + .expect("current-thread runtime"); + let (recorder, handle) = + crate::metrics::readiness_test_recorder_with_idle_timeout(timeout.as_secs()); - let response = metrics::with_local_recorder(&recorder, || { - coordinator.begin_shutdown(); - coordinator.finish_probe(ticket, ready_evaluation()) + metrics::with_local_recorder(&recorder, || { + runtime.block_on(async { + let mut state = crate::state::tests::test_state_with_database_url( + "postgres://127.0.0.1:1/buzz", + ) + .await; + let sample_started = Arc::new(tokio::sync::Notify::new()); + let release_sample = Arc::new(tokio::sync::Notify::new()); + Arc::get_mut(&mut state) + .expect("sole readiness state") + .set_dependency_evaluator(Arc::new(DelayedEvaluator { + report: ready_report(), + started: Arc::clone(&sample_started), + release: Arc::clone(&release_sample), + })); + + start_dependency_sampler_and_completion_publisher( + Arc::clone(&state), + republish_interval, + ); + + tokio::time::timeout(Duration::from_secs(5), sample_started.notified()) + .await + .expect("dependency sampler must start its first sample within 5 seconds"); + for _ in 0..3 { + tokio::time::advance(republish_interval).await; + tokio::task::yield_now().await; + } + assert_eq!( + sample_completion_metric_value_from_scrape(&handle.render()), + None, + "republisher ticks must not fabricate an epoch before the first completion" + ); + + release_sample.notify_one(); + + let first = + wait_for_sample_completion_scrape_value(&handle, |value| value > 0).await; + + state.dependency_sampler_cancel.cancel(); + tokio::time::resume(); + + tokio::time::sleep(timeout + Duration::from_millis(200)).await; + assert_eq!( + sample_completion_metric_value_from_scrape(&handle.render()), + Some(first), + "publisher must keep the completion epoch exported after idle timeout" + ); + + tokio::time::sleep(timeout + Duration::from_millis(200)).await; + assert_eq!( + sample_completion_metric_value_from_scrape(&handle.render()), + Some(first), + "with no sampler activity, publisher must keep exporting the stored epoch" + ); + + state.dependency_completion_publisher_cancel.cancel(); + }); }); - let snapshot = snapshotter.snapshot().into_vec(); + } - assert_eq!(response.reason, ReadinessReason::ShuttingDown); - assert_eq!(gauge_value(&snapshot, "overall"), 0.0); - assert!( - exact_metric(&snapshot, "buzz_readiness_state", &[("check", "postgres")]).is_none() + /// How long a park step waits before it is a failure rather than a hang. + const PARK_SIGNAL_TIMEOUT: Duration = Duration::from_secs(10); + + /// A completion-epoch gauge write that can be stopped at the recorder + /// boundary — the last instruction of a publish, after the writer has + /// already read the epoch it is publishing. + /// + /// Arming names one epoch value and fires once, so the sampler's own write + /// passes straight through and only the republish under test parks. + struct PublishPark { + armed_for: Mutex>, + parked_tx: std::sync::mpsc::SyncSender<()>, + parked_rx: Mutex>, + release_tx: std::sync::mpsc::SyncSender<()>, + release_rx: Mutex>, + } + + impl PublishPark { + fn new() -> Self { + let (parked_tx, parked_rx) = std::sync::mpsc::sync_channel(1); + let (release_tx, release_rx) = std::sync::mpsc::sync_channel(1); + Self { + armed_for: Mutex::new(None), + parked_tx, + parked_rx: Mutex::new(parked_rx), + release_tx, + release_rx: Mutex::new(release_rx), + } + } + + fn arm(&self, epoch_seconds: u64) { + *self.armed_for.lock().expect("park state") = Some(epoch_seconds); + } + + fn park_if_armed(&self, epoch_seconds: u64) { + { + let mut armed = self.armed_for.lock().expect("park state"); + if *armed != Some(epoch_seconds) { + return; + } + *armed = None; + } + self.parked_tx.send(()).expect("announce the parked write"); + self.release_rx + .lock() + .expect("release channel") + .recv_timeout(PARK_SIGNAL_TIMEOUT) + .expect("the parked write must be released"); + } + + fn wait_until_parked(&self) { + self.parked_rx + .lock() + .expect("park channel") + .recv_timeout(PARK_SIGNAL_TIMEOUT) + .expect("the armed write must reach the recorder"); + } + + fn release(&self) { + self.release_tx.send(()).expect("release the parked write"); + } + } + + /// Delegates to the real Prometheus recorder, with the completion-epoch + /// gauge routed through [`PublishPark`]. + struct ParkingRecorder { + inner: metrics_exporter_prometheus::PrometheusRecorder, + park: Arc, + } + + struct ParkingGauge { + inner: metrics::Gauge, + park: Arc, + } + + impl metrics::GaugeFn for ParkingGauge { + fn increment(&self, value: f64) { + self.inner.increment(value); + } + + fn decrement(&self, value: f64) { + self.inner.decrement(value); + } + + fn set(&self, value: f64) { + self.park.park_if_armed(value as u64); + self.inner.set(value); + } + } + + impl metrics::Recorder for ParkingRecorder { + fn describe_counter( + &self, + key: metrics::KeyName, + unit: Option, + description: metrics::SharedString, + ) { + self.inner.describe_counter(key, unit, description); + } + + fn describe_gauge( + &self, + key: metrics::KeyName, + unit: Option, + description: metrics::SharedString, + ) { + self.inner.describe_gauge(key, unit, description); + } + + fn describe_histogram( + &self, + key: metrics::KeyName, + unit: Option, + description: metrics::SharedString, + ) { + self.inner.describe_histogram(key, unit, description); + } + + fn register_counter( + &self, + key: &metrics::Key, + metadata: &metrics::Metadata<'_>, + ) -> metrics::Counter { + self.inner.register_counter(key, metadata) + } + + fn register_gauge( + &self, + key: &metrics::Key, + metadata: &metrics::Metadata<'_>, + ) -> metrics::Gauge { + let gauge = self.inner.register_gauge(key, metadata); + if key.name() == "buzz_readiness_dependency_sample_completed_timestamp_seconds" { + metrics::Gauge::from_arc(Arc::new(ParkingGauge { + inner: gauge, + park: Arc::clone(&self.park), + })) + } else { + gauge + } + } + + fn register_histogram( + &self, + key: &metrics::Key, + metadata: &metrics::Metadata<'_>, + ) -> metrics::Histogram { + self.inner.register_histogram(key, metadata) + } + } + + /// Two tasks write this gauge — the sampler, which owns the epoch, and the + /// idle-refresh publisher, which may only re-emit it — and the metrics + /// facade offers a bare `set` with no compare-and-set to lean on. So a + /// republish that read the stored epoch before a sample completed can still + /// be inside its own write when the newer epoch lands, and plain last-write + /// -wins would leave the exported series moved backwards until the next + /// republish tick. + /// + /// This parks the republish at the recorder, the last point in its write, + /// completes a newer sample behind it — which must not be blocked by the + /// parked republish — and only then releases it. The scrape must report the + /// newer completion. Times are injected so the two epochs are exact and the + /// interleaving does not depend on the wall clock. + #[test] + fn a_republish_racing_a_completion_cannot_move_the_exported_epoch_backwards() { + const EARLIER_EPOCH_SECONDS: u64 = 1_700_000_000; + const LATER_EPOCH_SECONDS: u64 = 1_700_000_030; + + let (prometheus, handle) = crate::metrics::readiness_test_recorder(); + let park = Arc::new(PublishPark::new()); + let recorder = Arc::new(ParkingRecorder { + inner: prometheus, + park: Arc::clone(&park), + }); + let diagnostics = Arc::new(ready_diagnostics()); + + metrics::with_local_recorder(recorder.as_ref(), || { + diagnostics.record_dependency_sample_completion( + SystemTime::UNIX_EPOCH + Duration::from_secs(EARLIER_EPOCH_SECONDS), + ); + }); + assert_eq!( + sample_completion_metric_value_from_scrape(&handle.render()), + Some(EARLIER_EPOCH_SECONDS), + "the first completion owns the epoch" + ); + + park.arm(EARLIER_EPOCH_SECONDS); + let republisher = std::thread::spawn({ + let recorder = Arc::clone(&recorder); + let diagnostics = Arc::clone(&diagnostics); + move || { + metrics::with_local_recorder(recorder.as_ref(), || { + diagnostics.republish_dependency_sample_completion(); + }); + } + }); + + park.wait_until_parked(); + + metrics::with_local_recorder(recorder.as_ref(), || { + diagnostics.record_dependency_sample_completion( + SystemTime::UNIX_EPOCH + Duration::from_secs(LATER_EPOCH_SECONDS), + ); + }); + assert_eq!( + sample_completion_metric_value_from_scrape(&handle.render()), + Some(LATER_EPOCH_SECONDS), + "a completed sample must publish its epoch while a republish is still in flight" + ); + + park.release(); + republisher.join().expect("republish thread"); + + assert_eq!( + sample_completion_metric_value_from_scrape(&handle.render()), + Some(LATER_EPOCH_SECONDS), + "a republish that lost the race must not re-export the epoch it read \ + before the newer sample completed" + ); + } + + #[test] + fn readiness_reason_labels_are_the_closed_lifecycle_set() { + assert_eq!( + [ReadinessReason::Ready, ReadinessReason::ShuttingDown].map(ReadinessReason::label), + READINESS_REASON_LABELS ); - assert!(matches!( - exact_metric( - &snapshot, - "buzz_readiness_dependency_checks_total", - &[("dependency", "postgres"), ("outcome", "success")] - ), - Some(DebugValue::Counter(1)) - )); } } diff --git a/crates/buzz-relay/src/router.rs b/crates/buzz-relay/src/router.rs index 303ecff82f0..3e96c02c13d 100644 --- a/crates/buzz-relay/src/router.rs +++ b/crates/buzz-relay/src/router.rs @@ -25,7 +25,7 @@ use crate::connection::handle_connection; use crate::metrics::track_metrics; use crate::nip11::{nip11_document, relay_info_handler}; use crate::nip_fi_http::http_denial; -use crate::readiness::{self, ReadinessEvaluation, ReadinessReason}; +use crate::readiness::{self, DependencySnapshot, ReadinessReason}; use crate::state::AppState; // ── NIP-FI fail-closed assertion guard ─────────────────────────────────────── @@ -645,60 +645,48 @@ async fn liveness_handler() -> impl IntoResponse { (StatusCode::OK, "ok") } -/// Compatibility endpoint on the public listener. It evaluates dependencies -/// and preserves the existing response contract but never records rollout -/// telemetry. +/// Compatibility endpoint on the public listener. Same lifecycle answer as the +/// probe, but public traffic must never move rollout telemetry. async fn public_readiness_handler(State(state): State>) -> impl IntoResponse { - if !state.readiness.public_evaluation_allowed() { - return readiness_response(ReadinessEvaluation::shutting_down(), false); - } - - let evaluation = state.readiness.evaluate(&state.db, &state.redis_pool).await; - let evaluation = state.readiness.finish_public_evaluation(evaluation); - readiness_response(evaluation, false) + readiness_response(readiness_reason(&state)) } -/// Kubernetes health-listener endpoint. All rollout metrics flow through the -/// process-owned coordinator so shutdown and probe generations are ordered. +/// Kubernetes health-listener endpoint — the only source of rollout readiness +/// telemetry. Its single lifecycle sample determines every observable result +/// of this request: counter, gauge, HTTP status, and body. async fn kubernetes_readiness_handler(State(state): State>) -> impl IntoResponse { - let readiness::ProbeStart::Evaluate(ticket) = state.readiness.begin_probe() else { - return readiness_response(ReadinessEvaluation::shutting_down(), true); - }; + let reason = readiness_reason(&state); + readiness::record_readiness_probe(reason); + readiness_response(reason) +} - let evaluation = state.readiness.evaluate(&state.db, &state.redis_pool).await; - let evaluation = state.readiness.finish_probe(ticket, evaluation); - readiness_response(evaluation, true) +/// Readiness answers for this process only. +/// +/// Shared Postgres, Redis, and deletion-catalog health used to gate this +/// answer, which meant one shared outage removed every replica from the load +/// balancer simultaneously and left a reconnect burst with nowhere to land. +/// Those checks now report on `/_status`. The health listener does not bind +/// until the database, migrations, Redis, and pub/sub are up (see +/// `buzz-relay/src/main.rs`), so an answering process is a booted process and +/// needs no separate startup state. +fn readiness_reason(state: &AppState) -> ReadinessReason { + if state.shutting_down.load(Ordering::Acquire) { + ReadinessReason::ShuttingDown + } else { + ReadinessReason::Ready + } } -fn readiness_response( - evaluation: ReadinessEvaluation, - include_reason: bool, -) -> axum::response::Response { - if evaluation.reason == ReadinessReason::ShuttingDown { - return ( +fn readiness_response(reason: ReadinessReason) -> axum::response::Response { + match reason { + ReadinessReason::Ready => { + (StatusCode::OK, Json(json!({"status": "ready"}))).into_response() + } + ReadinessReason::ShuttingDown => ( StatusCode::SERVICE_UNAVAILABLE, Json(json!({"status": "shutting_down"})), ) - .into_response(); - } - - let pg_ok = evaluation.postgres_ready(); - let redis_ok = evaluation.redis_ready(); - let deletion_catalog_ok = evaluation.deletion_catalog_ready(); - - if evaluation.is_ready() { - (StatusCode::OK, Json(json!({"status": "ready"}))).into_response() - } else { - let mut payload = json!({ - "status": "not_ready", - "postgres": pg_ok, - "redis": redis_ok, - "deletion_catalog": deletion_catalog_ok - }); - if include_reason { - payload["reason"] = json!(evaluation.reason.label()); - } - (StatusCode::SERVICE_UNAVAILABLE, Json(payload)).into_response() + .into_response(), } } @@ -715,9 +703,44 @@ fn status_payload(uptime_secs: u64) -> serde_json::Value { }) } -/// Status endpoint — service name, version, uptime, and intrinsic build identity. +/// The dependency fields the readiness body used to carry, now a diagnostic +/// read of the sampler's cache. +/// +/// `sample` is always present so a reader can never mistake a cached verdict +/// for a current one: `not_yet_sampled` before the sampler's first evaluation +/// completes, then `fresh` or `stale` alongside the report's own age. +fn dependency_diagnostics_payload(snapshot: DependencySnapshot) -> serde_json::Value { + let interval_seconds = readiness::DEPENDENCY_SAMPLE_INTERVAL.as_secs(); + match snapshot { + DependencySnapshot::NotYetSampled => json!({ + "sample": "not_yet_sampled", + "sample_interval_seconds": interval_seconds, + }), + DependencySnapshot::Sampled { report, age, stale } => json!({ + "sample": if stale { "stale" } else { "fresh" }, + "sample_interval_seconds": interval_seconds, + "sample_age_seconds": age.as_secs(), + "postgres": report.postgres_ready(), + "redis": report.redis_ready(), + "deletion_catalog": report.deletion_catalog_ready(), + "reason": report.reason.label(), + }), + } +} + +/// Status endpoint — service name, version, uptime, intrinsic build identity, +/// and the cached shared-dependency diagnostics. +/// +/// Health-listener only, and never wired to a Kubernetes probe: this is where +/// an operator looks to tell "the pod is fine, Postgres is not" apart from "the +/// pod is broken". It reads only what +/// [`readiness::run_dependency_sampler`] has already cached, so however often +/// it is polled it adds no load to the shared pools. async fn status_handler(State(state): State>) -> impl IntoResponse { - Json(status_payload(state.started_at.elapsed().as_secs())) + let mut payload = status_payload(state.started_at.elapsed().as_secs()); + payload["dependencies"] = + dependency_diagnostics_payload(state.dependency_diagnostics.snapshot()); + Json(payload) } /// `/_mesh` — live mesh status: peer table, connection/phi state, per-peer @@ -761,7 +784,6 @@ fn build_cors_layer(cors_origins: &[String]) -> CorsLayer { #[cfg(test)] mod tests { use std::collections::VecDeque; - use std::sync::atomic::{AtomicUsize, Ordering as AtomicOrdering}; use std::sync::{Mutex, PoisonError}; use std::time::Duration; @@ -770,91 +792,63 @@ mod tests { use opentelemetry::trace::TracerProvider as _; use opentelemetry_sdk::trace::{InMemorySpanExporter, SdkTracerProvider}; use tokio::net::TcpListener; - use tokio::sync::{mpsc, Notify}; + use tokio::sync::mpsc; use tokio_tungstenite::{connect_async, tungstenite::Message}; use tower::ServiceBuilder; use tracing::Instrument as _; use tracing_subscriber::prelude::*; use super::*; + use crate::readiness::DependencyReport; - struct ScriptedReadinessEvaluator { - evaluations: Mutex>, + struct ScriptedDependencyEvaluator { + evaluations: Mutex>, + evaluations_started: std::sync::atomic::AtomicUsize, } - impl ScriptedReadinessEvaluator { - fn new(evaluations: impl IntoIterator) -> Self { + impl ScriptedDependencyEvaluator { + fn new(evaluations: impl IntoIterator) -> Self { Self { evaluations: Mutex::new(evaluations.into_iter().collect()), + evaluations_started: std::sync::atomic::AtomicUsize::new(0), } } - fn push(&self, evaluation: ReadinessEvaluation) { + fn push(&self, evaluation: DependencyReport) { self.evaluations .lock() .unwrap_or_else(PoisonError::into_inner) .push_back(evaluation); } + + /// How many times a caller actually reached the shared dependencies. + fn evaluations_started(&self) -> usize { + self.evaluations_started.load(Ordering::SeqCst) + } } #[async_trait::async_trait] - impl readiness::ReadinessEvaluator for ScriptedReadinessEvaluator { + impl readiness::DependencyEvaluator for ScriptedDependencyEvaluator { async fn evaluate( &self, _db: &buzz_db::Db, _redis_pool: &deadpool_redis::Pool, - ) -> ReadinessEvaluation { + ) -> DependencyReport { + self.evaluations_started.fetch_add(1, Ordering::SeqCst); self.evaluations .lock() .unwrap_or_else(PoisonError::into_inner) .pop_front() - .expect("scripted readiness evaluation") - } - } - - struct BarrierReadinessEvaluator { - calls: AtomicUsize, - first_started: Notify, - release_first: Notify, - first: ReadinessEvaluation, - second: ReadinessEvaluation, - } - - impl BarrierReadinessEvaluator { - fn new(first: ReadinessEvaluation, second: ReadinessEvaluation) -> Self { - Self { - calls: AtomicUsize::new(0), - first_started: Notify::new(), - release_first: Notify::new(), - first, - second, - } + .expect("scripted dependency report") } } - #[async_trait::async_trait] - impl readiness::ReadinessEvaluator for BarrierReadinessEvaluator { - async fn evaluate( - &self, - _db: &buzz_db::Db, - _redis_pool: &deadpool_redis::Pool, - ) -> ReadinessEvaluation { - if self.calls.fetch_add(1, AtomicOrdering::SeqCst) == 0 { - self.first_started.notify_waiters(); - self.release_first.notified().await; - self.first - } else { - self.second - } - } - } - - fn readiness_evaluation( + fn dependency_report( postgres: readiness::PostgresOutcome, redis: readiness::RedisOutcome, deletion_catalog: readiness::DeletionCatalogOutcome, - ) -> ReadinessEvaluation { - ReadinessEvaluation::from_results( + ) -> DependencyReport { + DependencyReport::from_results( readiness::TimedOutcome::new(postgres, Duration::from_millis(35)), readiness::TimedOutcome::new(redis, Duration::from_millis(20)), readiness::TimedOutcome::new(deletion_catalog, Duration::from_millis(15)), @@ -862,8 +856,8 @@ mod tests { ) } - fn ready_evaluation() -> ReadinessEvaluation { - readiness_evaluation( + fn ready_report() -> DependencyReport { + dependency_report( readiness::PostgresOutcome::Success, readiness::RedisOutcome::Success, readiness::DeletionCatalogOutcome::Success, @@ -945,10 +939,19 @@ mod tests { Arc::new(state) } - async fn readiness_state(evaluator: Arc) -> Arc { + async fn readiness_state(evaluator: Arc) -> Arc { + let mut state = unreachable_dependency_state().await; + Arc::get_mut(&mut state) + .expect("sole reference") + .set_dependency_evaluator(evaluator); + state + } + + /// A relay process whose shared Postgres and Redis are both unroutable. + async fn unreachable_dependency_state() -> Arc { let mut config = crate::config::Config::from_env().expect("default config loads"); config.require_relay_membership = false; - config.database_url = "postgres://buzz:buzz_dev@127.0.0.1:1/buzz".to_string(); + config.database_url = "postgres://buzz:buzz_dev@127.0.0.1:1/buzz".to_string(); // sadscan:disable np.postgres.1 -- local test-only credentials on a closed port config.redis_url = "redis://127.0.0.1:1".to_string(); let pool = sqlx::PgPool::connect_lazy(&config.database_url).expect("lazy pg pool"); let db = buzz_db::Db::from_pool(pool.clone()); @@ -968,7 +971,7 @@ mod tests { buzz_workflow::WorkflowConfig::default(), )); let media_storage = buzz_media::MediaStorage::new(&config.media).expect("media storage"); - let (mut state, _audit_shutdown) = AppState::new( + let (state, _audit_shutdown) = AppState::new( config, db, redis_pool, @@ -980,7 +983,6 @@ mod tests { nostr::Keys::generate(), media_storage, ); - state.set_readiness_evaluator(evaluator); Arc::new(state) } @@ -1001,6 +1003,300 @@ mod tests { (status, payload) } + async fn status_request(router: Router) -> (StatusCode, serde_json::Value) { + let response = router + .oneshot( + Request::get("/_status") + .body(Body::empty()) + .expect("status request"), + ) + .await + .expect("status response"); + let status = response.status(); + let body = axum::body::to_bytes(response.into_body(), 64 * 1024) + .await + .expect("status response body"); + let payload = serde_json::from_slice(&body).expect("status JSON"); + (status, payload) + } + + /// The incident regression. Shared Postgres and Redis pressure took every + /// replica out of the load balancer at once, so a reconnect burst had + /// nowhere to land. Readiness answers for this process only: a pod whose + /// shared dependencies are unreachable is still a healthy pod, and only a + /// local shutdown may withdraw it. + #[tokio::test] + async fn readiness_answers_from_local_lifecycle_not_shared_dependencies() { + let state = unreachable_dependency_state().await; + + for router in [ + build_health_router(state.clone()), + build_router(state.clone()), + ] { + assert_eq!( + readiness_request(router).await, + (StatusCode::OK, json!({"status": "ready"})), + "unreachable shared dependencies must not deroute a healthy pod" + ); + } + + state.begin_shutdown(); + + for router in [ + build_health_router(state.clone()), + build_router(state.clone()), + ] { + assert_eq!( + readiness_request(router).await, + ( + StatusCode::SERVICE_UNAVAILABLE, + json!({"status": "shutting_down"}) + ), + "a draining pod must still withdraw itself" + ); + } + } + + /// Dependency health did not disappear with the probe — it moved to the + /// diagnostic endpoint, which is never wired to a Kubernetes probe. The + /// fields the readiness body used to carry are still there, now qualified + /// by how old the sample behind them is. + #[tokio::test] + async fn status_retains_dependency_diagnostics_off_the_probe_path() { + let evaluator = Arc::new(ScriptedDependencyEvaluator::new([dependency_report( + readiness::PostgresOutcome::Success, + readiness::RedisOutcome::PoolTimeout, + readiness::DeletionCatalogOutcome::Success, + )])); + let state = readiness_state(evaluator).await; + state + .dependency_diagnostics + .sample(&state.db, &state.redis_pool) + .await; + + let (status, payload) = status_request(build_health_router(state)).await; + + assert_eq!(status, StatusCode::OK); + assert_eq!(payload["service"], "buzz-relay"); + assert_eq!( + payload["dependencies"], + json!({ + "sample": "fresh", + "sample_interval_seconds": 30, + "sample_age_seconds": 0, + "postgres": true, + "redis": false, + "deletion_catalog": true, + "reason": "redis_pool_timeout" + }) + ); + } + + /// `/_status` is an operator diagnostic, not a dependency driver. Evaluating + /// per request let operator curiosity — and anything that polls the + /// endpoint — add Postgres, Redis, and deletion-catalog work to a shared + /// dependency that is already under pressure, with no bound on how many + /// evaluations could be in flight at once. The endpoint reads the cached + /// report the per-pod sampler owns and starts nothing. + #[tokio::test] + async fn status_reads_the_cached_report_and_never_starts_a_dependency_check() { + let evaluator = Arc::new(ScriptedDependencyEvaluator::new([ready_report()])); + let state = readiness_state(evaluator.clone()).await; + let health = build_health_router(state.clone()); + + for _ in 0..3 { + let (status, payload) = status_request(health.clone()).await; + assert_eq!(status, StatusCode::OK); + assert_eq!( + payload["dependencies"], + json!({ + "sample": "not_yet_sampled", + "sample_interval_seconds": 30, + }), + "before the first sample completes there is no report to serve" + ); + } + + assert_eq!( + evaluator.evaluations_started(), + 0, + "a status request must never reach the shared dependencies" + ); + + // Once the sampler has a report, and only then, the endpoint serves it. + state + .dependency_diagnostics + .sample(&state.db, &state.redis_pool) + .await; + let (status, payload) = status_request(health).await; + + assert_eq!(status, StatusCode::OK); + assert_eq!(payload["dependencies"]["sample"], json!("fresh")); + assert_eq!(payload["dependencies"]["reason"], json!("ready")); + assert_eq!( + evaluator.evaluations_started(), + 1, + "the sampler is the only caller that evaluates" + ); + } + + /// Always answers, recording how many evaluations started and the peak + /// number in flight, so a loop test can assert cadence and single-flight + /// without a scripted queue to exhaust. + struct ObservedDependencyEvaluator { + report: DependencyReport, + duration: Duration, + started: std::sync::atomic::AtomicUsize, + in_flight: std::sync::atomic::AtomicUsize, + peak_in_flight: std::sync::atomic::AtomicUsize, + } + + impl ObservedDependencyEvaluator { + fn new(report: DependencyReport, duration: Duration) -> Self { + Self { + report, + duration, + started: std::sync::atomic::AtomicUsize::new(0), + in_flight: std::sync::atomic::AtomicUsize::new(0), + peak_in_flight: std::sync::atomic::AtomicUsize::new(0), + } + } + + fn started(&self) -> usize { + self.started.load(Ordering::SeqCst) + } + + fn peak_in_flight(&self) -> usize { + self.peak_in_flight.load(Ordering::SeqCst) + } + } + + #[async_trait::async_trait] + impl readiness::DependencyEvaluator for ObservedDependencyEvaluator { + async fn evaluate( + &self, + _db: &buzz_db::Db, + _redis_pool: &deadpool_redis::Pool, + ) -> DependencyReport { + self.started.fetch_add(1, Ordering::SeqCst); + let in_flight = self.in_flight.fetch_add(1, Ordering::SeqCst) + 1; + self.peak_in_flight.fetch_max(in_flight, Ordering::SeqCst); + if !self.duration.is_zero() { + tokio::time::sleep(self.duration).await; + } + self.in_flight.fetch_sub(1, Ordering::SeqCst); + self.report + } + } + + /// Runs the production sampler on a paused clock for `window`, then cancels + /// it and returns the evaluator's observations. + async fn run_sampler_for( + evaluator: Arc, + window: Duration, + ) -> Arc { + let state = readiness_state(evaluator.clone()).await; + let cancel = state.dependency_sampler_cancel.clone(); + let sampler = tokio::spawn(readiness::run_dependency_sampler( + state.clone(), + cancel.clone(), + )); + tokio::time::sleep(window).await; + cancel.cancel(); + sampler.await.expect("sampler task"); + evaluator + } + + /// Dependency telemetry must keep describing the shared dependencies whether + /// or not anyone reads `/_status`. Request-driven evaluation meant a quiet + /// endpoint produced a flat dashboard during the exact outage it existed to + /// explain. + #[test] + fn the_dependency_sampler_emits_telemetry_without_any_request() { + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .start_paused(true) + .build() + .expect("paused current-thread runtime"); + let (recorder, handle) = crate::metrics::readiness_test_recorder(); + + metrics::with_local_recorder(&recorder, || { + crate::metrics::describe_readiness_metrics(); + runtime.block_on(async { + // The first tick fires immediately, then one per cadence. + let evaluator = run_sampler_for( + Arc::new(ObservedDependencyEvaluator::new( + ready_report(), + Duration::ZERO, + )), + readiness::DEPENDENCY_SAMPLE_INTERVAL * 3 + Duration::from_secs(1), + ) + .await; + + assert_eq!(evaluator.started(), 4); + let rendered = handle.render(); + assert_eq!( + metric_value( + &rendered, + "buzz_readiness_dependency_checks_total{dependency=\"postgres\",outcome=\"success\"}" + ), + 4.0, + "every sampling cycle must publish its dependency outcomes" + ); + assert_eq!( + metric_value( + &rendered, + "buzz_readiness_check_duration_seconds_count{check=\"overall\"}" + ), + 4.0 + ); + assert!( + !rendered.contains("buzz_readiness_checks_total{"), + "sampling is not a readiness probe and must not move probe telemetry" + ); + }); + }); + } + + /// The bound that replaces the request-driven design's lack of one. The + /// sampler awaits each evaluation before taking the next tick, so a + /// dependency slower than the cadence lowers the sampling rate instead of + /// stacking probes on top of the slowness that caused it. + #[test] + fn the_dependency_sampler_never_runs_two_evaluations_at_once() { + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .start_paused(true) + .build() + .expect("paused current-thread runtime"); + + runtime.block_on(async { + let slow = readiness::DEPENDENCY_SAMPLE_INTERVAL * 2 + Duration::from_secs(1); + let window = readiness::DEPENDENCY_SAMPLE_INTERVAL * 10; + let evaluator = run_sampler_for( + Arc::new(ObservedDependencyEvaluator::new(ready_report(), slow)), + window, + ) + .await; + + assert_eq!( + evaluator.peak_in_flight(), + 1, + "the sampler must own the only in-flight evaluation" + ); + // 61-second evaluations run back to back from t=0 in a 300-second + // window: five, not the ten ticks the cadence offered. An + // evaluation started per tick regardless of the last one would + // have started ten and held several open at once. + assert_eq!(evaluator.started(), 5); + assert!( + evaluator.started() + < (window.as_secs() / readiness::DEPENDENCY_SAMPLE_INTERVAL.as_secs()) as usize, + "a slow dependency must throttle sampling, not be sampled on every tick" + ); + }); + } + fn readiness_metric_lines(rendered: &str) -> Vec<&str> { rendered .lines() @@ -1008,15 +1304,6 @@ mod tests { .collect() } - fn sorted_readiness_metric_lines(rendered: &str) -> Vec { - let mut lines = readiness_metric_lines(rendered) - .into_iter() - .map(str::to_owned) - .collect::>(); - lines.sort(); - lines - } - fn metric_value(rendered: &str, exact_prefix: &str) -> f64 { rendered .lines() @@ -1028,16 +1315,23 @@ mod tests { .unwrap_or_else(|| panic!("missing metric line: {exact_prefix}")) } + /// The frozen telemetry contract for the health listener. + /// + /// Readiness is lifecycle-only: its counter carries exactly two reasons and + /// its gauge is the latest private readiness-probe observation, never a + /// dependency or a transition-owned lifecycle mirror. Dependency families + /// are still exported, but only by the per-pod sampler, and neither + /// public-listener traffic nor an `/_status` request moves anything. #[test] - fn production_readiness_routes_export_the_frozen_health_only_contract() { + fn production_health_routes_export_the_frozen_telemetry_contract() { let runtime = tokio::runtime::Builder::new_current_thread() .enable_all() .build() .expect("current-thread runtime"); - let evaluator = Arc::new(ScriptedReadinessEvaluator::new(std::iter::repeat_n( - ready_evaluation(), - 4, - ))); + // Seeded with the first sampling cycle only; the coverage loop below + // pushes the rest, one per cycle, so the evaluator never serves a + // report the assertions did not choose. + let evaluator = Arc::new(ScriptedDependencyEvaluator::new([ready_report()])); let (recorder, handle) = crate::metrics::readiness_test_recorder(); metrics::with_local_recorder(&recorder, || { @@ -1062,138 +1356,130 @@ mod tests { readiness_request(health.clone()).await, (StatusCode::OK, json!({"status": "ready"})) ); - let first_scrape = handle.render(); - - assert!(first_scrape.contains("# TYPE buzz_readiness_checks_total counter")); - assert!(first_scrape - .contains("# TYPE buzz_readiness_dependency_checks_total counter")); - assert!(first_scrape - .contains("# TYPE buzz_readiness_check_duration_seconds histogram")); - assert!(first_scrape.contains("# TYPE buzz_readiness_state gauge")); + let after_probe = handle.render(); + + assert!(after_probe.contains("# TYPE buzz_readiness_checks_total counter")); + assert!(after_probe.contains("# TYPE buzz_readiness_state gauge")); assert_eq!( - metric_value( - &first_scrape, - "buzz_readiness_checks_total{reason=\"ready\"}" - ), + metric_value(&after_probe, "buzz_readiness_checks_total{reason=\"ready\"}"), 1.0 ); assert_eq!( - metric_value( - &first_scrape, - "buzz_readiness_dependency_checks_total{dependency=\"postgres\",outcome=\"success\"}" - ), + metric_value(&after_probe, "buzz_readiness_state{check=\"overall\"}"), 1.0 ); + assert!( + !after_probe.contains("buzz_readiness_dependency_checks_total{"), + "the probe must not touch a shared dependency" + ); + assert!( + !after_probe.contains("buzz_readiness_check_duration_seconds_count"), + "the probe must not record a dependency latency sample" + ); + + // Dependency telemetry now belongs to the sampler; the + // endpoint only reads what the sampler cached. + state + .dependency_diagnostics + .sample(&state.db, &state.redis_pool) + .await; + let (status, payload) = status_request(health.clone()).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(payload["dependencies"]["sample"], json!("fresh")); + assert_eq!(payload["dependencies"]["reason"], json!("ready")); + let after_status = handle.render(); + assert!( + after_status.contains("# TYPE buzz_readiness_dependency_checks_total counter") + ); + assert!( + after_status.contains("# TYPE buzz_readiness_check_duration_seconds histogram") + ); assert_eq!( metric_value( - &first_scrape, - "buzz_readiness_state{check=\"overall\"}" + &after_status, + "buzz_readiness_dependency_checks_total{dependency=\"postgres\",outcome=\"success\"}" ), 1.0 ); for bucket in ["2", "2.5", "+Inf"] { - assert!(first_scrape.contains(&format!( + assert!(after_status.contains(&format!( "buzz_readiness_check_duration_seconds_bucket{{check=\"overall\",le=\"{bucket}\"}}" ))); } - assert!(!first_scrape.contains("result=")); - assert!(!first_scrape + assert!(!after_status.contains("result=")); + assert!(!after_status .lines() .filter(|line| line.starts_with("buzz_readiness_check_duration_seconds")) .any(|line| line.contains("outcome="))); + for dependency in ["postgres", "redis", "deletion_catalog"] { + assert!( + !after_status + .contains(&format!("buzz_readiness_state{{check=\"{dependency}\"}}")), + "dependency health has no publishable readiness gauge" + ); + } - let before_public_failure = sorted_readiness_metric_lines(&first_scrape); - evaluator.push(readiness_evaluation( - readiness::PostgresOutcome::Success, - readiness::RedisOutcome::PoolTimeout, - readiness::DeletionCatalogOutcome::Success, - )); - assert_eq!( - readiness_request(public.clone()).await, + // A failing dependency is reported and changes nothing about + // whether this pod stays in the load balancer. This set also + // covers every valid dependency/outcome pair, so the series + // total below is exact rather than merely bounded. + let coverage = [ ( - StatusCode::SERVICE_UNAVAILABLE, - json!({ - "status": "not_ready", - "postgres": true, - "redis": false, - "deletion_catalog": true - }) - ) - ); - assert_eq!( - sorted_readiness_metric_lines(&handle.render()), - before_public_failure - ); - - let contract_evaluations = [ - readiness_evaluation( readiness::PostgresOutcome::PoolTimeout, - readiness::RedisOutcome::Success, - readiness::DeletionCatalogOutcome::Success, + readiness::RedisOutcome::PoolTimeout, + readiness::DeletionCatalogOutcome::OperationTimeout, ), - readiness_evaluation( + ( readiness::PostgresOutcome::PoolError, - readiness::RedisOutcome::Success, - readiness::DeletionCatalogOutcome::Success, + readiness::RedisOutcome::PoolError, + readiness::DeletionCatalogOutcome::OperationError, ), - readiness_evaluation( + ( readiness::PostgresOutcome::QueryTimeout, readiness::RedisOutcome::Success, readiness::DeletionCatalogOutcome::Success, ), - readiness_evaluation( + ( readiness::PostgresOutcome::QueryError, readiness::RedisOutcome::Success, readiness::DeletionCatalogOutcome::Success, ), - readiness_evaluation( - readiness::PostgresOutcome::Success, - readiness::RedisOutcome::PoolTimeout, - readiness::DeletionCatalogOutcome::Success, - ), - readiness_evaluation( - readiness::PostgresOutcome::Success, - readiness::RedisOutcome::PoolError, - readiness::DeletionCatalogOutcome::Success, - ), - readiness_evaluation( - readiness::PostgresOutcome::Success, - readiness::RedisOutcome::Success, - readiness::DeletionCatalogOutcome::OperationTimeout, - ), - readiness_evaluation( - readiness::PostgresOutcome::Success, - readiness::RedisOutcome::Success, - readiness::DeletionCatalogOutcome::OperationError, - ), - readiness_evaluation( - readiness::PostgresOutcome::PoolTimeout, - readiness::RedisOutcome::PoolTimeout, - readiness::DeletionCatalogOutcome::OperationTimeout, - ), - readiness_evaluation( - readiness::PostgresOutcome::PoolError, - readiness::RedisOutcome::PoolError, - readiness::DeletionCatalogOutcome::Success, - ), ]; - for evaluation in contract_evaluations { - evaluator.push(evaluation); - let (status, payload) = readiness_request(health.clone()).await; - assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE); - assert_eq!(payload["reason"], json!(evaluation.reason.label())); + for (index, (postgres, redis, deletion_catalog)) in + coverage.into_iter().enumerate() + { + evaluator.push(dependency_report(postgres, redis, deletion_catalog)); + state + .dependency_diagnostics + .sample(&state.db, &state.redis_pool) + .await; + let (status, degraded) = status_request(health.clone()).await; + assert_eq!(status, StatusCode::OK); + if index == 0 { + assert_eq!( + degraded["dependencies"], + json!({ + "sample": "fresh", + "sample_interval_seconds": 30, + "sample_age_seconds": 0, + "postgres": false, + "redis": false, + "deletion_catalog": false, + "reason": "overall_timeout" + }) + ); + } + assert_eq!( + readiness_request(health.clone()).await, + (StatusCode::OK, json!({"status": "ready"})), + "a failing dependency must never deroute this pod" + ); } - let before_shutdown = handle.render(); - let histogram_counts_before = ["overall", "postgres", "redis", "deletion_catalog"] - .map(|check| { - metric_value( - &before_shutdown, - &format!( - "buzz_readiness_check_duration_seconds_count{{check=\"{check}\"}}" - ), - ) - }); + let histogram_count_before = metric_value( + &handle.render(), + "buzz_readiness_check_duration_seconds_count{check=\"overall\"}", + ); state.begin_shutdown(); assert_eq!( readiness_request(public).await, @@ -1202,10 +1488,18 @@ mod tests { json!({"status": "shutting_down"}) ) ); - let after_public_shutdown = handle.render(); - assert!(after_public_shutdown + assert!(handle + .render() .lines() .all(|line| !line.contains("reason=\"shutting_down\""))); + assert_eq!( + metric_value( + &handle.render(), + "buzz_readiness_state{check=\"overall\"}" + ), + 1.0, + "shutdown and public traffic must not update the private probe gauge" + ); assert_eq!( readiness_request(health).await, @@ -1215,31 +1509,43 @@ mod tests { ) ); let final_scrape = handle.render(); - let histogram_counts_after = ["overall", "postgres", "redis", "deletion_catalog"] - .map(|check| { - metric_value( - &final_scrape, - &format!( - "buzz_readiness_check_duration_seconds_count{{check=\"{check}\"}}" - ), - ) - }); - assert_eq!(histogram_counts_after, histogram_counts_before); assert_eq!( metric_value( &final_scrape, - "buzz_readiness_checks_total{reason=\"shutting_down\"}" + "buzz_readiness_check_duration_seconds_count{check=\"overall\"}" ), - 1.0 + histogram_count_before, + "shutdown must not fabricate a dependency latency sample" ); assert_eq!( metric_value( &final_scrape, - "buzz_readiness_state{check=\"overall\"}" + "buzz_readiness_checks_total{reason=\"shutting_down\"}" ), + 1.0 + ); + assert_eq!( + metric_value(&final_scrape, "buzz_readiness_state{check=\"overall\"}"), 0.0 ); assert!(!final_scrape.contains("sensitive-sql-or-url")); + // Freshness is part of the frozen contract: one unlabelled + // gauge carrying when the cached report completed. The sampler + // advances it and the publisher re-emits it, so + // `time() - ` ages a stalled sampler out from a scrape + // alone. + assert!(final_scrape.contains( + "# TYPE buzz_readiness_dependency_sample_completed_timestamp_seconds gauge" + )); + assert_eq!( + final_scrape + .lines() + .filter(|line| line.starts_with( + "buzz_readiness_dependency_sample_completed_timestamp_seconds" + )) + .count(), + 1 + ); let exported_reasons = final_scrape .lines() @@ -1249,138 +1555,7 @@ mod tests { assert_eq!( readiness_metric_lines(&final_scrape).len(), readiness::READINESS_RAW_SERIES_PER_POD, - "readiness series contract must stay at or below its 99-series cap" - ); - }); - }); - } - - fn run_out_of_order_route_case( - first: ReadinessEvaluation, - second: ReadinessEvaluation, - ) -> (serde_json::Value, serde_json::Value, String) { - let runtime = tokio::runtime::Builder::new_current_thread() - .enable_all() - .build() - .expect("current-thread runtime"); - let evaluator = Arc::new(BarrierReadinessEvaluator::new(first, second)); - let (recorder, handle) = crate::metrics::readiness_test_recorder(); - - metrics::with_local_recorder(&recorder, || { - runtime.block_on(async { - let state = readiness_state(evaluator.clone()).await; - let health = build_health_router(state); - let first_started = evaluator.first_started.notified(); - let slow_first = tokio::spawn(readiness_request(health.clone())); - first_started.await; - - let (_, second_payload) = readiness_request(health).await; - evaluator.release_first.notify_one(); - let (_, first_payload) = slow_first.await.expect("slow first probe task"); - (first_payload, second_payload, handle.render()) - }) - }) - } - - #[test] - fn real_health_route_generation_fence_covers_both_completion_orders() { - let failure = readiness_evaluation( - readiness::PostgresOutcome::Success, - readiness::RedisOutcome::PoolTimeout, - readiness::DeletionCatalogOutcome::Success, - ); - - let (older_failure, newer_success, success_scrape) = - run_out_of_order_route_case(failure, ready_evaluation()); - assert_eq!(older_failure["reason"], json!("redis_pool_timeout")); - assert_eq!(newer_success, json!({"status": "ready"})); - assert_eq!( - metric_value(&success_scrape, "buzz_readiness_state{check=\"overall\"}"), - 1.0 - ); - assert_eq!( - metric_value(&success_scrape, "buzz_readiness_state{check=\"redis\"}"), - 1.0 - ); - - let (older_success, newer_failure, failure_scrape) = - run_out_of_order_route_case(ready_evaluation(), failure); - assert_eq!(older_success, json!({"status": "ready"})); - assert_eq!(newer_failure["reason"], json!("redis_pool_timeout")); - assert_eq!( - metric_value(&failure_scrape, "buzz_readiness_state{check=\"overall\"}"), - 0.0 - ); - assert_eq!( - metric_value(&failure_scrape, "buzz_readiness_state{check=\"redis\"}"), - 0.0 - ); - for scrape in [&success_scrape, &failure_scrape] { - assert_eq!( - metric_value(scrape, "buzz_readiness_checks_total{reason=\"ready\"}"), - 1.0 - ); - assert_eq!( - metric_value( - scrape, - "buzz_readiness_checks_total{reason=\"redis_pool_timeout\"}" - ), - 1.0 - ); - } - } - - #[test] - fn real_health_route_shutdown_fence_dominates_an_in_flight_success() { - let runtime = tokio::runtime::Builder::new_current_thread() - .enable_all() - .build() - .expect("current-thread runtime"); - let evaluator = Arc::new(BarrierReadinessEvaluator::new( - ready_evaluation(), - ready_evaluation(), - )); - let (recorder, handle) = crate::metrics::readiness_test_recorder(); - - metrics::with_local_recorder(&recorder, || { - runtime.block_on(async { - let state = readiness_state(evaluator.clone()).await; - let health = build_health_router(state.clone()); - let first_started = evaluator.first_started.notified(); - let in_flight = tokio::spawn(readiness_request(health)); - first_started.await; - - state.begin_shutdown(); - evaluator.release_first.notify_one(); - assert_eq!( - in_flight.await.expect("in-flight readiness task"), - ( - StatusCode::SERVICE_UNAVAILABLE, - json!({"status": "shutting_down"}) - ) - ); - - let scrape = handle.render(); - assert_eq!( - metric_value(&scrape, "buzz_readiness_state{check=\"overall\"}"), - 0.0 - ); - assert!(scrape - .lines() - .all(|line| !line.starts_with("buzz_readiness_state{check=\"postgres\"}"))); - assert_eq!( - metric_value( - &scrape, - "buzz_readiness_checks_total{reason=\"shutting_down\"}" - ), - 1.0 - ); - assert_eq!( - metric_value( - &scrape, - "buzz_readiness_dependency_checks_total{dependency=\"postgres\",outcome=\"success\"}" - ), - 1.0 + "readiness series contract must stay at or below its 87-series cap" ); }); }); diff --git a/crates/buzz-relay/src/state.rs b/crates/buzz-relay/src/state.rs index 72376402001..0bcab5d63a4 100644 --- a/crates/buzz-relay/src/state.rs +++ b/crates/buzz-relay/src/state.rs @@ -182,10 +182,99 @@ impl Drop for CommunityConnectionGuard { } } +/// Message reported when the one-time Redis bootstrap gate rejects startup. +/// +/// Bounded and stable so operators and the boot regression test can match on +/// it without parsing the underlying driver error. +pub const REDIS_BOOTSTRAP_FAILURE: &str = "Redis command path unavailable at startup"; + +/// Budget for the one-time bootstrap PING. A refused port answers immediately; +/// this only bounds a blackholed address, where hanging forever would be worse +/// than exiting. +const REDIS_BOOTSTRAP_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5); + +/// Keep the Redis driver's per-command timeout outside the relay-owned startup +/// budget so [`REDIS_BOOTSTRAP_TIMEOUT`] remains the one authoritative bound. +const REDIS_BOOTSTRAP_DRIVER_TIMEOUT: std::time::Duration = + REDIS_BOOTSTRAP_TIMEOUT.saturating_mul(2); + +/// Proves once, during startup, that the Redis command path this pod will serve +/// from can actually be reached. +/// +/// `deadpool_redis` pools dial lazily and `PubSubManager::new` only allocates +/// channels, so without this nothing in boot ever opened a command connection: +/// a relay came up against a dead Redis, bound its health listener, and — since +/// readiness reports local lifecycle only — advertised ready forever. Binding +/// that listener is a one-way latch, so the check has to happen before it, and +/// it is deliberately a *startup* gate: once serving, a Redis blip is a +/// dependency failure and must never change readiness. +pub async fn verify_redis_command_path(pool: &deadpool_redis::Pool) -> anyhow::Result<()> { + let ping = async { + let connection = pool + .get() + .await + .map_err(|error| anyhow::anyhow!("{REDIS_BOOTSTRAP_FAILURE}: {error}"))?; + // This one-shot connection is removed from the pool so extending its + // driver timeout cannot leak into normal serving traffic. The relay's + // outer timeout below must bound both lazy checkout and PING. + let mut connection = deadpool_redis::Connection::take(connection); + connection.set_response_timeout(REDIS_BOOTSTRAP_DRIVER_TIMEOUT); + redis::cmd("PING") + .query_async::(&mut connection) + .await + .map_err(|error| anyhow::anyhow!("{REDIS_BOOTSTRAP_FAILURE}: {error}")) + }; + + match tokio::time::timeout(REDIS_BOOTSTRAP_TIMEOUT, ping).await { + Err(_) => Err(anyhow::anyhow!( + "{REDIS_BOOTSTRAP_FAILURE}: no response within {REDIS_BOOTSTRAP_TIMEOUT:?}" + )), + Ok(result) => result.map(|_| ()), + } +} + +/// Bounded outcome of the durable community-active check run when a socket is +/// admitted. +/// +/// `outcome` is the only dimension. Community, tenant, connection, and error +/// text are request-controlled and deliberately absent from the label set. +#[derive(Debug, Clone, Copy)] +enum AdmissionOutcome { + Active, + Inactive, + CheckError, +} + +impl AdmissionOutcome { + fn label(self) -> &'static str { + match self { + Self::Active => "active", + Self::Inactive => "inactive", + Self::CheckError => "check_error", + } + } +} + +fn record_admission_check(outcome: AdmissionOutcome) { + metrics::counter!( + "buzz_community_admission_checks_total", + "outcome" => outcome.label(), + ) + .increment(1); +} + /// Registers a socket, durably revalidates its community, then runs it. /// /// The ordering is the archival admission invariant: archive-before-query is /// observed by the query, while archive-after-registration sees the token. +/// +/// Admission is fail-closed: only an affirmative `Ok(true)` may serve. Both +/// `Ok(false)` and a lookup `Err` cancel, because neither proves this tenant is +/// currently admitted, and `docs/multi-tenant-relay.md` I5 +/// (`Inv_AdmissionFence`) grants capability only to an actor *currently* +/// admitted to that community. The two are still told apart in telemetry +/// (`buzz_community_admission_checks_total{outcome}`) so an operator can +/// separate archival from database pressure. pub(crate) async fn run_registered_community_connection( registry: &CommunityConnectionRegistry, connection_id: Uuid, @@ -201,9 +290,29 @@ pub(crate) async fn run_registered_community_connection record_admission_check(AdmissionOutcome::Active), + Ok(false) => { + record_admission_check(AdmissionOutcome::Inactive); + cancel.cancel(); + return; + } + Err(error) => { + // A lookup failure is not an answer, so it cannot authorize one. + // Admitting here would begin serving AUTH and REQ for a tenant + // whose lifecycle is unknown, and the adjacent host-binding seam + // already refuses on exactly this evidence (see + // `router::nip11_or_ws_handler`). The client sees an ordinary dial + // failure and retries. + record_admission_check(AdmissionOutcome::CheckError); + tracing::warn!( + %community_id, + %error, + "community active check failed; refusing the socket" + ); + cancel.cancel(); + return; + } } if cancel.is_cancelled() { return; @@ -719,8 +828,14 @@ pub struct AppState { pub audio_rooms: Arc, /// Set to `true` on SIGTERM — readiness probe returns 503. pub shutting_down: Arc, - /// Orders readiness gauge publication against terminal shutdown. - pub(crate) readiness: Arc, + /// Cached shared-dependency evaluation behind the diagnostic `/_status` + /// endpoint, owned by [`crate::readiness::run_dependency_sampler`]. Never + /// consulted by a Kubernetes probe, and never evaluated by a request. + pub(crate) dependency_diagnostics: Arc, + /// Stops only the periodic dependency sampler during graceful shutdown. + pub dependency_sampler_cancel: CancellationToken, + /// Stops only the completion-epoch publisher during graceful shutdown. + pub dependency_completion_publisher_cancel: CancellationToken, /// Process start time — used by `/_status` endpoint. pub started_at: Instant, /// Shared, community-scoped NIP-98 replay prevention. @@ -946,7 +1061,9 @@ impl AppState { git_pack_cache, audio_rooms: Arc::new(AudioRoomManager::new()), shutting_down: Arc::new(AtomicBool::new(false)), - readiness: Arc::new(crate::readiness::ReadinessCoordinator::default()), + dependency_diagnostics: Arc::new(crate::readiness::DependencyDiagnostics::default()), + dependency_sampler_cancel: CancellationToken::new(), + dependency_completion_publisher_cancel: CancellationToken::new(), started_at: Instant::now(), nip98_replay, gif_http_client, @@ -990,21 +1107,21 @@ impl AppState { ) } - /// Atomically closes readiness publication before exposing shutdown to - /// the relay's other fast-path lifecycle checks. + /// Withdraws this pod from routing. The lifecycle flag is authoritative for + /// `/_readiness`; the private probe publishes its sampled observation to + /// the readiness gauge on its next request. pub fn begin_shutdown(&self) { - self.readiness.begin_shutdown(); self.shutting_down.store(true, Ordering::Release); } #[cfg(test)] - pub(crate) fn set_readiness_evaluator( + pub(crate) fn set_dependency_evaluator( &mut self, - evaluator: Arc, + evaluator: Arc, ) { - self.readiness = Arc::new(crate::readiness::ReadinessCoordinator::with_evaluator( - evaluator, - )); + self.dependency_diagnostics = Arc::new( + crate::readiness::DependencyDiagnostics::with_evaluator(evaluator), + ); } /// Inter-relay mesh handle. `None` ⇒ mesh-off / single-instance: callers @@ -2096,6 +2213,144 @@ pub(crate) mod tests { assert!(!started_during.load(Ordering::SeqCst)); } + /// Reads one `buzz_community_admission_checks_total` series by exact label set. + fn admission_counter( + snapshot: &[( + metrics_util::CompositeKey, + Option, + Option, + metrics_util::debugging::DebugValue, + )], + outcome: &str, + ) -> Option { + snapshot.iter().find_map(|(key, _, _, value)| { + let labels = key + .key() + .labels() + .map(|label| (label.key(), label.value())) + .collect::>(); + if key.key().name() != "buzz_community_admission_checks_total" + || labels != [("outcome", outcome)] + { + return None; + } + match value { + metrics_util::debugging::DebugValue::Counter(count) => Some(*count), + _ => panic!("community admission checks must be a counter"), + } + }) + } + + /// Admission is fail-closed on both non-affirmative outcomes. A confirmed + /// `Ok(false)` and a lookup `Err` are different diagnoses — the counter + /// keeps them apart — but neither is proof of current admission, and + /// `docs/multi-tenant-relay.md` I5 (`Inv_AdmissionFence`) grants read or + /// membership capability only to an actor *currently* admitted to that + /// community. Serving AUTH/REQ on an unproven tenant lifecycle is the + /// failure this guards. + #[test] + fn neither_a_confirmed_inactive_community_nor_a_failed_lookup_admits_the_socket() { + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("current-thread runtime"); + let recorder = metrics_util::debugging::DebuggingRecorder::new(); + let snapshotter = recorder.snapshotter(); + + let (inactive_cancel, inactive_started, error_cancel, error_started, active_started) = + metrics::with_local_recorder(&recorder, || { + runtime.block_on(async { + let registry = CommunityConnectionRegistry::new(); + let community = CommunityId::from_uuid(Uuid::from_u128(0xa)); + + let inactive_cancel = CancellationToken::new(); + let inactive_started = Arc::new(AtomicBool::new(false)); + let started = Arc::clone(&inactive_started); + run_registered_community_connection( + ®istry, + Uuid::new_v4(), + community, + CommunityConnectionControl::new(inactive_cancel.clone()), + || async { Ok(false) }, + move |_| async move { started.store(true, Ordering::SeqCst) }, + ) + .await; + + let error_cancel = CancellationToken::new(); + let error_started = Arc::new(AtomicBool::new(false)); + let started = Arc::clone(&error_started); + run_registered_community_connection( + ®istry, + Uuid::new_v4(), + community, + CommunityConnectionControl::new(error_cancel.clone()), + || async { Err(buzz_db::DbError::Sqlx(sqlx::Error::PoolTimedOut)) }, + move |_| async move { started.store(true, Ordering::SeqCst) }, + ) + .await; + + let active_started = Arc::new(AtomicBool::new(false)); + let started = Arc::clone(&active_started); + run_registered_community_connection( + ®istry, + Uuid::new_v4(), + community, + CommunityConnectionControl::new(CancellationToken::new()), + || async { Ok(true) }, + move |_| async move { started.store(true, Ordering::SeqCst) }, + ) + .await; + + ( + inactive_cancel, + inactive_started, + error_cancel, + error_started, + active_started, + ) + }) + }); + + assert!( + inactive_cancel.is_cancelled(), + "a confirmed-inactive community must still cancel its socket" + ); + assert!( + !inactive_started.load(Ordering::SeqCst), + "a confirmed-inactive community must never start the socket body" + ); + assert!( + error_cancel.is_cancelled(), + "a failed active check must cancel its socket, not admit it" + ); + assert!( + !error_started.load(Ordering::SeqCst), + "a failed active check must never start serving AUTH/REQ on an unproven tenant" + ); + assert!(active_started.load(Ordering::SeqCst)); + + let snapshot = snapshotter.snapshot().into_vec(); + assert_eq!(admission_counter(&snapshot, "inactive"), Some(1)); + assert_eq!(admission_counter(&snapshot, "check_error"), Some(1)); + assert_eq!(admission_counter(&snapshot, "active"), Some(1)); + + let label_sets = snapshot + .iter() + .filter(|(key, _, _, _)| key.key().name() == "buzz_community_admission_checks_total") + .map(|(key, _, _, _)| { + key.key() + .labels() + .map(|label| label.key().to_owned()) + .collect::>() + }) + .collect::>(); + assert_eq!(label_sets.len(), 3, "outcome is the only dimension"); + assert!( + label_sets.iter().all(|labels| labels == &["outcome"]), + "admission telemetry must never carry community, tenant, or error labels: {label_sets:?}" + ); + } + #[tokio::test] async fn revalidation_continues_after_one_community_lookup_failure() { let registry = CommunityConnectionRegistry::new(); diff --git a/crates/buzz-relay/tests/boot_lifecycle.rs b/crates/buzz-relay/tests/boot_lifecycle.rs index 17762bbdeab..1d43fc3236e 100644 --- a/crates/buzz-relay/tests/boot_lifecycle.rs +++ b/crates/buzz-relay/tests/boot_lifecycle.rs @@ -3,6 +3,7 @@ use std::{ io::{Read as _, Write as _}, net::{TcpListener, TcpStream}, process::{Child, Command, ExitStatus, Output, Stdio}, + sync::mpsc, thread::{self, JoinHandle}, time::{Duration, Instant}, }; @@ -10,11 +11,14 @@ use std::{ use serde_json::Value; use buzz_relay::lifecycle::StartupPhase; +use buzz_relay::state::REDIS_BOOTSTRAP_FAILURE; const VALID_RELAY_PRIVATE_KEY: &str = "0000000000000000000000000000000000000000000000000000000000000001"; const CHILD_TIMEOUT: Duration = Duration::from_secs(10); const MAX_CAPTURE_BYTES: u64 = 1024 * 1024; +/// Bound on every wait for the relay to export a metric family. +const METRICS_SCRAPE_DEADLINE: Duration = Duration::from_secs(8); struct RelayProcess { child: Option, @@ -156,20 +160,27 @@ fn scrape_metrics(port: u16) -> std::io::Result { } fn wait_for_relay_metrics(process: &mut RelayProcess, port: u16) -> String { - let deadline = Instant::now() + Duration::from_secs(8); + wait_for_scraped_metric(process, port, "buzz_audit_enabled") +} + +/// Polls the relay's own `/metrics` until `needle` appears, bounded by +/// [`METRICS_SCRAPE_DEADLINE`]. A relay that exits first is a failure, not a +/// timeout, so the panic names the real cause. +fn wait_for_scraped_metric(process: &mut RelayProcess, port: u16, needle: &str) -> String { + let deadline = Instant::now() + METRICS_SCRAPE_DEADLINE; loop { assert!( process.try_wait().is_none(), - "relay exited before its metrics endpoint became usable" + "relay exited before exporting {needle}" ); if let Ok(response) = scrape_metrics(port) { - if response.contains("buzz_audit_enabled") { + if response.contains(needle) { return response; } } assert!( Instant::now() < deadline, - "relay metrics did not become scrapeable within 8s" + "relay did not export {needle} within {METRICS_SCRAPE_DEADLINE:?}" ); thread::sleep(Duration::from_millis(20)); } @@ -500,3 +511,334 @@ fn successful_main_emits_complete_lifecycle_without_startup_metrics() { assert_terminal(&events, "metrics_bind", "succeeded", None); assert_terminal(&events, "process_telemetry", "succeeded", None); } + +/// Boot gates that need a live Postgres to reach the code under test. Named +/// `postgres_tests` so `.config/nextest.toml`'s `postgres-ci` default filter +/// discovers them structurally; the wrapper hands each test its own database +/// through `DATABASE_URL`. +mod postgres_tests { + use super::*; + + const REDIS_BOOTSTRAP_BUDGET: Duration = Duration::from_secs(5); + const REDIS_BOOTSTRAP_SCHEDULING_SLACK: Duration = Duration::from_secs(3); + + /// A TCP peer that completes Redis's metadata handshake, then reads and + /// holds PING without replying. This distinguishes a bounded checkout + + /// PING from a refused connection, which returns before the bootstrap + /// timeout is exercised. + struct HangingRedisPeer { + redis_url: String, + request_received: mpsc::Receiver>, + stop: mpsc::Sender<()>, + worker: Option>, + } + + impl HangingRedisPeer { + fn spawn() -> Self { + let listener = TcpListener::bind(("127.0.0.1", 0)).expect("bind fake Redis peer"); + listener + .set_nonblocking(true) + .expect("set fake Redis listener nonblocking"); + let port = listener.local_addr().expect("fake Redis address").port(); + let (request_tx, request_received) = mpsc::channel(); + let (stop, stop_rx) = mpsc::channel(); + let worker = thread::spawn(move || loop { + if stop_rx.try_recv().is_ok() { + return; + } + let (mut stream, _) = match listener.accept() { + Ok(accepted) => accepted, + Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => { + thread::sleep(Duration::from_millis(10)); + continue; + } + Err(error) => panic!("accept fake Redis connection: {error}"), + }; + stream + .set_read_timeout(Some(Duration::from_millis(100))) + .expect("bound fake Redis read"); + let mut request = Vec::new(); + loop { + let mut chunk = [0_u8; 4096]; + match stream.read(&mut chunk) { + Ok(0) => return, + Ok(read) => { + request.extend_from_slice(&chunk[..read]); + assert!( + request.len() <= 4096, + "fake Redis peer received an oversized request" + ); + let setinfo_commands = request + .windows(b"SETINFO".len()) + .filter(|window| *window == b"SETINFO") + .count(); + if setinfo_commands >= 2 { + // redis-rs pipelines CLIENT SETINFO lib-name + // and lib-ver while establishing a connection. + // Complete that handshake so the relay reaches + // its explicit bootstrap PING, then hold it. + stream + .write_all(b"+OK\r\n+OK\r\n") + .expect("reply to Redis client handshake"); + request.clear(); + continue; + } + if request + .windows(b"PING".len()) + .any(|window| window == b"PING") + { + let _ = request_tx.send(request); + let _ = stop_rx.recv_timeout(CHILD_TIMEOUT); + return; + } + } + Err(error) + if matches!( + error.kind(), + std::io::ErrorKind::WouldBlock | std::io::ErrorKind::TimedOut + ) => + { + if stop_rx.try_recv().is_ok() { + return; + } + } + Err(error) => panic!("read fake Redis request: {error}"), + } + } + }); + Self { + redis_url: format!("redis://127.0.0.1:{port}"), + request_received, + stop, + worker: Some(worker), + } + } + + fn redis_url(&self) -> &str { + &self.redis_url + } + + fn assert_request_received(&self, logs: &str) { + let request = self + .request_received + .recv_timeout(Duration::from_secs(1)) + .unwrap_or_else(|_| panic!("relay must reach the fake Redis peer: {logs}")); + assert!(!request.is_empty(), "fake Redis peer read an empty request"); + } + } + + impl Drop for HangingRedisPeer { + fn drop(&mut self) { + let _ = self.stop.send(()); + if let Some(worker) = self.worker.take() { + worker.join().expect("fake Redis peer must not panic"); + } + } + } + + fn reserve_closed_port() -> u16 { + let reserved = TcpListener::bind(("127.0.0.1", 0)).expect("reserve port"); + let port = reserved.local_addr().expect("reserved address").port(); + drop(reserved); + port + } + + /// Runs the relay until it exits on its own, or kills it once `timeout` + /// passes. Unlike `run_relay`, a relay that keeps serving is a result to + /// assert on rather than a panic, which is the whole point here. + fn run_until_exit(environment: &[(&str, &str)], timeout: Duration) -> (bool, Output) { + let mut process = RelayProcess::spawn(environment); + let deadline = Instant::now() + timeout; + while Instant::now() < deadline { + if process.try_wait().is_some() { + return (true, process.wait(Duration::from_secs(2))); + } + thread::sleep(Duration::from_millis(20)); + } + (false, process.terminate()) + } + + /// Redis is required for pub/sub fan-out, presence, and typing, but nothing + /// in boot ever opened a command connection: `deadpool_redis` pools dial + /// lazily and `PubSubManager::new` only allocates channels, so "Redis + /// pub/sub connected" was logged against a dead port. A relay could + /// therefore boot with Redis unreachable, bind its health listener, and — + /// now that readiness answers from local lifecycle alone — advertise ready + /// for the rest of its life. The bootstrap gate is the one-time proof that + /// the command path has connected at least once, and it has to land before + /// the listener binds, because binding is the one-way latch that makes this + /// pod routable. + /// + /// The git conformance probe is disabled so the only remaining startup-fatal + /// gate is the one under test. + #[test] + #[ignore = "requires PostgreSQL"] + fn unreachable_redis_fails_boot_before_the_health_listener_binds() { + let database_url = std::env::var("DATABASE_URL") + .expect("postgres lane provides DATABASE_URL for each test process"); + let redis_url = format!("redis://127.0.0.1:{}", reserve_closed_port()); + let metrics_port = reserve_closed_port().to_string(); + let health_port = reserve_closed_port(); + let health_port_value = health_port.to_string(); + + let (exited, output) = run_until_exit( + &[ + ("BUZZ_RELAY_PRIVATE_KEY", VALID_RELAY_PRIVATE_KEY), + ("BUZZ_METRICS_PORT", &metrics_port), + ("BUZZ_HEALTH_PORT", &health_port_value), + ("DATABASE_URL", &database_url), + ("REDIS_URL", &redis_url), + ("BUZZ_GIT_CONFORMANCE_PROBE", "false"), + ], + CHILD_TIMEOUT, + ); + let logs = format!( + "{}{}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); + + assert!( + !logs.contains("Health probe listener started"), + "the Redis bootstrap gate must run before the health listener binds: {logs}" + ); + assert!( + exited && !output.status.success(), + "an unreachable Redis command path must be startup-fatal: {logs}" + ); + assert!( + logs.contains(REDIS_BOOTSTRAP_FAILURE), + "the failure must name the gate that rejected boot: {logs}" + ); + assert!( + TcpListener::bind(("0.0.0.0", health_port)).is_ok(), + "the health port must never have been bound" + ); + } + + /// A peer that accepts the socket but withholds its Redis response exercises + /// the outer timeout around both lazy pool checkout and PING. Removing or + /// narrowing that timeout makes this test kill a still-running relay at the + /// deadline instead of observing a startup failure. + #[test] + #[ignore = "requires PostgreSQL"] + fn hanging_redis_peer_times_out_before_the_health_listener_binds() { + let database_url = std::env::var("DATABASE_URL") + .expect("postgres lane provides DATABASE_URL for each test process"); + let redis_peer = HangingRedisPeer::spawn(); + let metrics_port = reserve_closed_port().to_string(); + let health_port = reserve_closed_port(); + let health_port_value = health_port.to_string(); + let timeout = REDIS_BOOTSTRAP_BUDGET + REDIS_BOOTSTRAP_SCHEDULING_SLACK; + let started_at = Instant::now(); + + let (exited, output) = run_until_exit( + &[ + ("BUZZ_RELAY_PRIVATE_KEY", VALID_RELAY_PRIVATE_KEY), + ("BUZZ_METRICS_PORT", &metrics_port), + ("BUZZ_HEALTH_PORT", &health_port_value), + ("DATABASE_URL", &database_url), + ("REDIS_URL", redis_peer.redis_url()), + ("BUZZ_GIT_CONFORMANCE_PROBE", "false"), + ], + timeout, + ); + let elapsed = started_at.elapsed(); + let logs = format!( + "{}{}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); + redis_peer.assert_request_received(&logs); + + assert!( + elapsed >= REDIS_BOOTSTRAP_BUDGET, + "the fake peer must hold the request through the bootstrap budget: {elapsed:?}: {logs}" + ); + assert!( + exited && elapsed < timeout && !output.status.success(), + "the outer bootstrap timeout must terminate the relay within scheduling slack: {elapsed:?}: {logs}" + ); + assert!( + logs.contains(REDIS_BOOTSTRAP_FAILURE), + "the timeout must report the bounded Redis bootstrap failure: {logs}" + ); + assert!( + !logs.contains("Health probe listener started"), + "the health listener must not bind before Redis bootstrap succeeds: {logs}" + ); + assert!( + TcpListener::bind(("0.0.0.0", health_port)).is_ok(), + "the health port must never have been bound" + ); + } + + /// Startup is the only owner of dependency evaluation. `/_status` just reads + /// the cache, a readiness probe records no dependency attempt at all, and no + /// request path may start a check — so if `main` stops spawning the sampler, + /// this pod evaluates Postgres, Redis, and the deletion catalog exactly + /// never: the dependency families stay absent from its scrape and `/_status` + /// answers `not_yet_sampled` for the pod's whole life. + /// + /// No in-process test can fail on that, because each one drives `sample` + /// itself. This one boots the real binary against real dependencies and + /// reads only the relay's own `/metrics` — no probe, no `/_status`, nothing + /// that could evaluate a dependency on the test's behalf. The first tick + /// fires immediately, so the wait is bounded by + /// [`METRICS_SCRAPE_DEADLINE`] and never a fixed sleep. + /// + /// The completion-timestamp gauge is the needle because `sample` writes it + /// last, after the cache already serves the report: observing it proves the + /// whole cycle ran, with no window where the counters have landed but the + /// timestamp has not. + #[test] + #[ignore = "requires PostgreSQL"] + fn startup_owns_the_dependency_sampler() { + let database_url = std::env::var("DATABASE_URL") + .expect("postgres lane provides DATABASE_URL for each test process"); + let redis_url = std::env::var("REDIS_URL") + .expect("the postgres lane runs alongside Redis and exports REDIS_URL"); + let metrics_port = reserve_closed_port(); + let metrics_port_value = metrics_port.to_string(); + let health_port_value = reserve_closed_port().to_string(); + let bind_addr = format!("127.0.0.1:{}", reserve_closed_port()); + + let mut process = RelayProcess::spawn(&[ + ("BUZZ_RELAY_PRIVATE_KEY", VALID_RELAY_PRIVATE_KEY), + ("BUZZ_METRICS_PORT", &metrics_port_value), + ("BUZZ_HEALTH_PORT", &health_port_value), + ("BUZZ_BIND_ADDR", &bind_addr), + ("DATABASE_URL", &database_url), + ("REDIS_URL", &redis_url), + ("BUZZ_GIT_CONFORMANCE_PROBE", "false"), + ]); + let scrape = wait_for_scraped_metric( + &mut process, + metrics_port, + "buzz_readiness_dependency_sample_completed_timestamp_seconds", + ); + let output = process.terminate(); + let logs = format!( + "{}{}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); + + assert!( + scrape.contains("buzz_readiness_dependency_checks_total{"), + "the same cycle must publish its per-dependency outcomes: {scrape}" + ); + assert!( + scrape.contains("buzz_readiness_check_duration_seconds"), + "a completed evaluation must publish its latency too: {scrape}" + ); + assert!( + !scrape.contains("buzz_readiness_checks_total{"), + "no readiness probe was sent, so the sampler alone produced this: {scrape}" + ); + assert!( + logs.contains("Health probe listener started"), + "the sampler must be owned by a relay that finished booting: {logs}" + ); + } +} diff --git a/deploy/charts/buzz/README.md b/deploy/charts/buzz/README.md index a6b74e87572..f5075fd19d5 100644 --- a/deploy/charts/buzz/README.md +++ b/deploy/charts/buzz/README.md @@ -110,10 +110,11 @@ SigV4 signing, so do not put the bucket into `s3.endpoint`; pass Railway's base Object storage is contacted during relay startup only when `BUZZ_GIT_CONFORMANCE_PROBE` is enabled (the relay default). A probe failure is -startup-fatal, so Kubernetes readiness never opens. If an operator explicitly -disables that probe through `relay.extraEnv`, `/_readiness` does not test object -storage; configuration is still parsed strictly, but reachability and addressing -errors surface on the first storage operation. +startup-fatal, so the process exits and Kubernetes readiness never opens. If an +operator explicitly disables that probe through `relay.extraEnv`, configuration +is still parsed strictly, but reachability and addressing errors surface on the +first storage operation. `/_readiness` tests no external dependency in either +case — see the readiness contract below. ### Early-startup telemetry contract @@ -125,26 +126,114 @@ These phases intentionally do not emit metrics. Most run before the Prometheus exporter exists, and one uniform log-only contract preserves every phase's real event time and failure without assigning an eventual scrape time to earlier work. +### Readiness contract + +**`/_readiness` reports local process lifecycle only.** It performs no +Postgres, Redis, or deletion-catalog I/O: `shutting_down` returns 503, and any +other state returns 200. Shared dependencies are shared by every replica, so +gating the probe on them removed the whole deployment from the load balancer at +once and left a reconnect burst with nowhere to land. There is no separate +"starting" state — the health listener does not bind until the database, +migrations, Redis, and pub/sub are up, so a process that can answer has booted. + +Shared-dependency health moved to **`/_status`** on the same private health +listener, under a `dependencies` object carrying the `postgres`, `redis`, +`deletion_catalog`, and aggregate `reason` fields the readiness body used to +return. Do not wire `/_status` to a Kubernetes probe. + +**`/_status` is a cached read.** It performs no Postgres, Redis, or +deletion-catalog I/O of its own. One background loop per pod evaluates the +three dependencies **every 30 seconds**, awaiting each evaluation before taking +the next tick, so a pod never has more than one evaluation in flight no matter +how often — or how rarely — the endpoint is read. The loop publishes the +dependency metrics below and caches the report `/_status` serves. Its first +cycle runs at startup, and each evaluation is bounded by a two-second budget. + +Every `dependencies` object therefore states how old its report is: + +| Field | Meaning | +|-------|---------| +| `sample: "not_yet_sampled"` | the first cycle has not completed; no `postgres`/`redis`/`deletion_catalog`/`reason` fields are present, because there is no observation to report | +| `sample: "fresh"` | the report is at most two cadences (60s) old | +| `sample: "stale"` | the report outlived two cadences, so the sampler missed at least one cycle — read the verdict as history, not as current state | +| `sample_age_seconds` | age of the report at request time (absent when `not_yet_sampled`) | +| `sample_interval_seconds` | the sampling cadence, `30` | + +Polling `/_status` more often than the cadence returns the same cached report; +it does not make the data fresher and adds no dependency load. + ### Readiness telemetry contract Only requests served by the private health listener (`BUZZ_HEALTH_PORT`) emit -rollout readiness telemetry. The compatibility `/_readiness` route on the public -app listener returns health but does not change these metrics. +rollout telemetry. The compatibility `/_readiness` route on the public app +listener returns the same lifecycle answer but does not change these metrics. + +| Metric | Type | Labels | Source | +|--------|------|--------|--------| +| `buzz_readiness_checks_total` | counter | `reason` ∈ {`ready`, `shutting_down`} | `/_readiness` | +| `buzz_readiness_state` | gauge | `check="overall"`; latest private probe observation, 1 ready or 0 shutting down | `/_readiness` | +| `buzz_readiness_dependency_checks_total` | counter | `dependency`, typed bounded `outcome` | dependency sampler | +| `buzz_readiness_check_duration_seconds` | histogram | `check` only | dependency sampler | +| `buzz_readiness_dependency_sample_completed_timestamp_seconds` | gauge | none; Unix time the cached report completed | completion publisher | + +The three dependency families keep their `buzz_readiness_*` names for dashboard +continuity, but nothing about them is request-driven any more: the 30-second +sampler publishes them whether or not anyone reads `/_status`, so a quiet +endpoint no longer produces a flat dashboard during the outage it exists to +explain. + +`buzz_readiness_dependency_sample_completed_timestamp_seconds` carries **when +the cached report completed**, in Unix seconds. The relay sampler is the only +owner allowed to advance that epoch, and it writes it immediately after the +cache is replaced. A separate bounded publisher re-emits the stored epoch often +enough to survive local gauge idle-timeout; republishing never advances the +timestamp. If a republish races a newer completion, one scrape can briefly see +the older epoch, but the republish path verifies after writing and repairs to the +newer stored epoch before that republish call returns. So the value stands still +when sampling stops, and the time elapsed since it was written is whatever the +reader computes at read time. + +Following the `buzz_storage_sweep_age_seconds` convention, the series is not +emitted until the first sample completes: its absence means "not yet sampled", +not "fresh". + +That is the whole server-side contract. Freshness alerting is built from this +gauge in the monitoring provider, and the monitor query, thresholds, and +per-pod tag grouping belong with the deployment's monitor configuration rather +than in this chart — they depend on the provider's query grammar and on the +tags its agent attaches, neither of which this repo owns. + +Two properties of the neighboring families are worth knowing when building +that alerting. `buzz_readiness_dependency_checks_total` and +`buzz_readiness_check_duration_seconds` are cumulative, so a sampler that stops +leaves their last values exported and scraped indefinitely: those series stay +present and flat rather than disappearing. And per-report freshness for a human +reading a single pod is already on the `sample`, `sample_age_seconds`, and +`sample_interval_seconds` fields of `/_status` above. + +The schema has a ceiling of 87 raw Prometheus series per pod: 2 probe reasons, +11 valid dependency/outcome pairs, 72 histogram series, and 2 gauges (overall lifecycle + completion timestamp). Do not +add pod, ReplicaSet, version, rollout, error text, SQL, URL, tenant, user, +community, pubkey, header, query, or other request-controlled labels. A +readiness probe records no dependency attempt or latency sample at all. +The readiness gauge is not a monotonic lifecycle mirror: shutdown changes the +authoritative lifecycle flag, and the next private readiness probe observes and +publishes that state. + +### Community admission telemetry | Metric | Type | Labels | |--------|------|--------| -| `buzz_readiness_checks_total` | counter | `reason` from the closed readiness-reason set | -| `buzz_readiness_dependency_checks_total` | counter | `dependency`, typed bounded `outcome` | -| `buzz_readiness_check_duration_seconds` | histogram | `check` only | -| `buzz_readiness_state` | gauge | `check` only; latest publishable generation | - -The schema has a ceiling of 99 raw Prometheus series per pod: 12 overall -reasons, 11 valid dependency/outcome pairs, 72 histogram series, and 4 gauges. -Do not add pod, ReplicaSet, version, rollout, error text, SQL, URL, tenant, -user, community, pubkey, header, query, or other request-controlled labels. -Shutdown without dependency evaluation increments only -`buzz_readiness_checks_total{reason="shutting_down"}` and sets the overall -state to zero; it does not fabricate dependency failures or latency samples. +| `buzz_community_admission_checks_total` | counter | `outcome` ∈ {`active`, `inactive`, `check_error`} | + +Counts the durable community-active check run before a socket is admitted. +Admission is fail-closed: only `active` serves. `inactive` (a confirmed +archival answer) and `check_error` (the lookup itself failed, so the tenant +lifecycle is unknown) both refuse the socket before any AUTH or REQ frame is +read; the client sees an ordinary dial failure and retries. The two outcomes +stay distinct so a rise in `check_error` reads as database pressure rather than +archival. `outcome` is the only dimension — community and error text are +request-controlled and must never become labels. ### Operation-aware database pool acquisition contract diff --git a/docs/deployment-identity.md b/docs/deployment-identity.md index 32d8c426e76..f5bc8710d10 100644 --- a/docs/deployment-identity.md +++ b/docs/deployment-identity.md @@ -47,6 +47,15 @@ The relay health listener exposes intrinsic build identity at `/_status`: "source_sha": "<40-character-source-sha>", "id": "github-actions::", "url": "https://github.com/block/buzz/actions/runs//attempts/" + }, + "dependencies": { + "sample": "fresh", + "sample_interval_seconds": 30, + "sample_age_seconds": 12, + "postgres": true, + "redis": true, + "deletion_catalog": true, + "reason": "ready" } } ``` @@ -54,6 +63,16 @@ The relay health listener exposes intrinsic build identity at `/_status`: Non-CI builds report stable `unknown` or `local` fallback values instead of claiming provenance they do not have. +`dependencies` is a cached diagnostic snapshot of shared-dependency health. A +per-pod background loop evaluates the dependencies every 30 seconds; the +endpoint only reads the latest report and never contacts a dependency itself, +so polling it costs nothing. `sample` is always present and reports whether +that cached verdict is `fresh`, `stale`, or `not_yet_sampled` — before the +first cycle completes the health fields are absent rather than defaulted. +`/_readiness` does not consult any of this — see +[the readiness contract](../deploy/charts/buzz/README.md#readiness-contract) — +so this endpoint must never be wired to a Kubernetes probe. + ## Helm digest pinning Buzz chart `0.1.8` and newer accept an immutable image digest: diff --git a/scripts/run-tests.sh b/scripts/run-tests.sh index 3732cdc476b..32bfa1c94bd 100755 --- a/scripts/run-tests.sh +++ b/scripts/run-tests.sh @@ -155,9 +155,9 @@ run_unit_tests() { run_test_step "buzz-acp unit tests" \ cargo test -p buzz-acp --lib -- --nocapture - # Mirror the three infra-free relay handler modules in `just test-unit`'s - # nextest expression. Keep the side-effects filter pinned to `::tests::` so - # it does not select the sibling Postgres-backed test module. + # Mirror the relay filters from `just test-unit`: the three handler modules, + # storage-snapshot helpers, readiness and router unit suites, and the single + # scoped admission regression in state::tests. run_test_step "buzz-relay channel authorization tests" \ cargo test -p buzz-relay --lib handlers::channel_authz:: -- --nocapture @@ -169,6 +169,15 @@ run_unit_tests() { run_test_step "buzz-relay storage snapshot tests" \ cargo test -p buzz-relay --lib storage_sweep::tests:: -- --nocapture + + run_test_step "buzz-relay readiness tests" \ + cargo test -p buzz-relay --lib readiness::tests:: -- --nocapture + + run_test_step "buzz-relay router tests" \ + cargo test -p buzz-relay --lib router::tests:: -- --nocapture + + run_test_step "buzz-relay admission regression test" \ + cargo test -p buzz-relay --lib state::tests::neither_a_confirmed_inactive_community_nor_a_failed_lookup_admits_the_socket -- --nocapture } # ---- DB / integration tests (infra required) --------------------------------