Skip to main content

toolkit/runtime/
readiness.rs

1//! Framework-managed readiness state for the `OoP` bootstrap (`cpt-cf-fr-eventual-readiness`).
2//!
3//! An `OoP` gear becomes *live* the moment its HTTP server binds (`/healthz`),
4//! but only becomes *ready* (`/readyz`) once startup is complete, every critical
5//! dependency has been resolved, **and** the gear's registered healthchecks
6//! report it can serve traffic (Spring Boot-style health groups per
7//! `cpt-cf-fr-eventual-readiness`).
8//!
9//! Readiness reuses the framework's standard healthcheck mechanism
10//! ([`crate::healthcheck`]): a gear expresses readiness once, via
11//! [`RestApiCapability::healthcheck`](crate::contracts::RestApiCapability::healthcheck),
12//! and it is honored identically whether the gear is hosted in-process by the
13//! `api-gateway` or run `OoP`. This module layers the three `OoP`-only concerns
14//! the gateway path lacks — startup completion, critical-dependency resolution
15//! gating, and the graceful-drain readiness flip — on top of that shared
16//! healthcheck report.
17//!
18//! The healthcheck report itself is fanned out concurrently, timeout-bounded,
19//! panic-isolated, and cached inside the [`RestHealthcheckRegistry`], so a burst
20//! of probe traffic cannot storm the registered checks. The dependency and
21//! draining state layered on top are cheap live reads.
22
23use std::collections::BTreeMap;
24use std::sync::Arc;
25use std::sync::atomic::{AtomicBool, Ordering};
26use std::time::Duration;
27
28use parking_lot::RwLock;
29
30use crate::healthcheck::{HealthcheckReport, HealthcheckStatus, RestHealthcheckRegistry};
31
32/// Sticky, dynamic dependency-resolution tracker shared across serving modes
33/// (in-process host and `OoP`).
34///
35/// **Passive:** the resolution loop / consumer-wiring feeds it
36/// ([`register_dep`](Self::register_dep), [`mark_resolved`](Self::mark_resolved));
37/// it never runs a check itself — unlike the active
38/// [`Healthcheck`](crate::healthcheck::Healthcheck) registry. It answers exactly
39/// one question — *are all declared dependencies resolved?* — and carries the
40/// process-wide draining flag both `/readyz` front-ends consult (the in-process
41/// [`ReadinessHealthcheck`] leaf check and the `OoP` [`ReadinessState`]
42/// aggregator).
43///
44/// Startup-gating and **sticky**: once a dependency is resolved it stays
45/// resolved. A provider vanishing later does NOT revert `/readyz` — that is
46/// runtime churn, handled lazily by the directory-resolving client at call time
47/// (liveness is `/healthz`, not `/readyz`).
48#[derive(Debug, Default)]
49pub struct DependencyChecker {
50    /// Graceful-shutdown flag; when set, `/readyz` reports `503` regardless of deps.
51    draining: AtomicBool,
52    /// `dep_gear` → resolved. Populated by the proxy-wiring / consumer-wiring
53    /// phase from each `#[toolkit::consumes]` registration; flipped `true` once
54    /// the provider endpoint resolves.
55    deps: RwLock<BTreeMap<String, bool>>,
56}
57
58impl DependencyChecker {
59    /// Create an empty checker (no deps → trivially resolved / ready).
60    #[must_use]
61    pub fn new() -> Self {
62        Self::default()
63    }
64
65    /// Declare a consumed dependency that gates readiness (idempotent; keeps an
66    /// already-resolved entry resolved).
67    pub fn register_dep(&self, dep_gear: impl Into<String>) {
68        self.deps.write().entry(dep_gear.into()).or_insert(false);
69    }
70
71    /// Mark a previously-registered dependency as resolved. Unknown names are
72    /// ignored. Returns `true` only on the `false → true` transition (so callers
73    /// can log the resolution exactly once).
74    pub fn mark_resolved(&self, dep_gear: &str) -> bool {
75        if let Some(resolved) = self.deps.write().get_mut(dep_gear) {
76            let was = *resolved;
77            *resolved = true;
78            !was
79        } else {
80            false
81        }
82    }
83
84    /// Names of consumed dependency gears not yet resolved.
85    #[must_use]
86    pub fn unresolved_deps(&self) -> Vec<String> {
87        self.deps
88            .read()
89            .iter()
90            .filter(|(_, resolved)| !**resolved)
91            .map(|(name, _)| name.clone())
92            .collect()
93    }
94
95    /// Whether every declared dependency has resolved (or none were declared).
96    #[must_use]
97    pub fn all_resolved(&self) -> bool {
98        self.deps.read().values().all(|resolved| *resolved)
99    }
100
101    /// Begin (or clear) draining. Setting it makes `/readyz` report `503`
102    /// regardless of dependency state so upstreams drain the instance.
103    pub fn set_draining(&self, draining: bool) {
104        self.draining.store(draining, Ordering::SeqCst);
105    }
106
107    /// Whether shutdown has begun.
108    #[must_use]
109    pub fn is_draining(&self) -> bool {
110        self.draining.load(Ordering::SeqCst)
111    }
112
113    /// Whether the process is serving traffic: all deps resolved and not draining.
114    #[must_use]
115    pub fn is_ready(&self) -> bool {
116        !self.is_draining() && self.all_resolved()
117    }
118}
119
120/// Bridges a [`DependencyChecker`] into the
121/// [`RestHealthcheckRegistry`](crate::healthcheck::RestHealthcheckRegistry) as a
122/// single leaf check, for the **in-process** (`api-gateway`-hosted) `/readyz`
123/// path where readiness is one check among many in a shared registry.
124///
125/// The gateway `/readyz` reports `503` whenever any registered check is
126/// `Unhealthy`, so mapping draining / unresolved-deps → `Unhealthy` gates the
127/// probe exactly as ADR-0007 intends. `/healthz` (liveness) is a static handler
128/// and is unaffected. The registry caches reports (~2s), so a readiness/drain
129/// transition surfaces on `/readyz` with up to that lag.
130///
131/// The `OoP` path instead uses [`ReadinessState::evaluate`] directly (it owns
132/// its registry), so this leaf adapter is not used there — avoiding the
133/// re-aggregation recursion that owning the same registry would cause.
134pub struct ReadinessHealthcheck {
135    deps: Arc<DependencyChecker>,
136}
137
138impl ReadinessHealthcheck {
139    /// Wrap the shared dependency checker.
140    #[must_use]
141    pub fn new(deps: Arc<DependencyChecker>) -> Self {
142        Self { deps }
143    }
144}
145
146#[async_trait::async_trait]
147impl crate::healthcheck::Healthcheck for ReadinessHealthcheck {
148    fn name(&self) -> &'static str {
149        "readiness"
150    }
151
152    async fn check(&self) -> crate::healthcheck::HealthcheckResult {
153        use crate::healthcheck::HealthcheckResult;
154        if self.deps.is_draining() {
155            return HealthcheckResult::unhealthy("draining").with_code("draining");
156        }
157        let unresolved = self.deps.unresolved_deps();
158        if unresolved.is_empty() {
159            HealthcheckResult::healthy()
160        } else {
161            HealthcheckResult::unhealthy(format!(
162                "unresolved dependencies: {}",
163                unresolved.join(", ")
164            ))
165            .with_code("deps_unresolved")
166        }
167    }
168}
169
170/// Default per-check timeout for readiness healthchecks.
171///
172/// Matches the `api-gateway` `healthcheck_timeout_ms` default so a gear's
173/// healthcheck behaves identically in-process and `OoP`.
174pub const DEFAULT_HEALTHCHECK_TIMEOUT: Duration = Duration::from_millis(500);
175
176/// Lifecycle state reported on `/readyz` (`cpt-cf-adr-eventual-readiness`).
177///
178/// Serialized lowercase; the four variants are a stable wire contract.
179#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)]
180#[serde(rename_all = "lowercase")]
181pub enum ReadinessLifecycle {
182    /// Not yet able to serve — startup is not complete, critical deps are
183    /// unresolved, or a healthcheck is `Unhealthy`. Maps to `503`.
184    Starting,
185    /// Fully serving traffic. Maps to `200`.
186    Ready,
187    /// Serving with reduced functionality — a healthcheck reported `Degraded`
188    /// (e.g. an optional backend is down but a fallback is acceptable). Kept in
189    /// rotation: maps to `200`.
190    Degraded,
191    /// Graceful shutdown in progress; upstreams should stop routing. Maps to
192    /// `503`.
193    Draining,
194}
195
196/// The aggregate readiness report rendered as the `/readyz` response body.
197///
198/// Readiness still depends on the aggregated [`HealthcheckReport`] internally,
199/// but the detailed per-component report is intentionally not echoed here; it
200/// belongs on the separate `/health` endpoint (see `oop_serve.rs`).
201#[derive(Debug, Clone, serde::Serialize)]
202pub struct ReadinessReport {
203    /// Lifecycle state — the primary readiness signal (`cpt-cf-adr-eventual-readiness`).
204    pub state: ReadinessLifecycle,
205    /// Whether the gear is ready to receive traffic (`true` → `200`, else `503`).
206    /// Convenience mirror of `state ∈ {ready, degraded}` for probes/clients that
207    /// do not want to know the `state → status` mapping.
208    pub ready: bool,
209    /// Critical dependencies not yet resolved via `DirectoryService` / DNS.
210    /// Non-empty only while `starting`. Omitted from the body when empty.
211    #[serde(skip_serializing_if = "Vec::is_empty")]
212    pub unresolved_deps: Vec<String>,
213}
214
215/// Shared, framework-owned readiness state for an `OoP` gear instance.
216///
217/// Created by the `OoP` bootstrap with the gear's critical dependency names and
218/// a shared [`RestHealthcheckRegistry`] (populated from each gear's
219/// [`RestApiCapability::healthcheck`](crate::contracts::RestApiCapability::healthcheck)).
220/// Cloned as an `Arc` into the probe router and the dependency-resolution task.
221pub struct ReadinessState {
222    /// Shared dependency-resolution core (deps + draining). Owned here for the
223    /// `OoP` aggregator; the in-process path shares the same [`DependencyChecker`]
224    /// via [`Self::dependency_checker`].
225    deps: Arc<DependencyChecker>,
226    /// Whether the gear has finished startup and is actually serving traffic.
227    /// Remains `false` until the bootstrap publishes the composed routes.
228    startup_complete: AtomicBool,
229    /// Shared gear healthcheck registry; supplies the "custom checks" dimension
230    /// of readiness (fan-out, timeout, panic isolation, and caching live here).
231    healthchecks: Arc<RestHealthcheckRegistry>,
232    /// Per-check timeout passed to [`RestHealthcheckRegistry::report`].
233    check_timeout: Duration,
234}
235
236impl std::fmt::Debug for ReadinessState {
237    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
238        f.debug_struct("ReadinessState")
239            .field("unresolved_deps", &self.deps.unresolved_deps())
240            .field("draining", &self.deps.is_draining())
241            .field(
242                "startup_complete",
243                &self.startup_complete.load(Ordering::Relaxed),
244            )
245            .field("check_timeout", &self.check_timeout)
246            .finish_non_exhaustive()
247    }
248}
249
250impl ReadinessState {
251    /// Create a new readiness state seeded with the gear's critical dependency
252    /// names and the shared healthcheck registry. All listed deps start
253    /// unresolved; the gear is not ready until startup is complete, each dep is
254    /// marked resolved via [`mark_dep_resolved`](Self::mark_dep_resolved), and
255    /// the healthchecks pass. Uses [`DEFAULT_HEALTHCHECK_TIMEOUT`] as the
256    /// per-check timeout.
257    #[must_use]
258    pub fn new<I, S>(critical_deps: I, healthchecks: Arc<RestHealthcheckRegistry>) -> Arc<Self>
259    where
260        I: IntoIterator<Item = S>,
261        S: Into<String>,
262    {
263        Self::with_check_timeout(critical_deps, healthchecks, DEFAULT_HEALTHCHECK_TIMEOUT)
264    }
265
266    /// Like [`new`](Self::new) but with an explicit per-check timeout.
267    #[must_use]
268    pub fn with_check_timeout<I, S>(
269        critical_deps: I,
270        healthchecks: Arc<RestHealthcheckRegistry>,
271        check_timeout: Duration,
272    ) -> Arc<Self>
273    where
274        I: IntoIterator<Item = S>,
275        S: Into<String>,
276    {
277        let deps = Arc::new(DependencyChecker::new());
278        for dep in critical_deps {
279            deps.register_dep(dep);
280        }
281        Self::from_checker(deps, healthchecks, check_timeout)
282    }
283
284    /// Build an aggregator around an **existing** [`DependencyChecker`], so the
285    /// in-process consumer-wiring path and the `OoP` `/readyz` aggregator share
286    /// one dependency-resolution core (the checker is fed once, read by both).
287    #[must_use]
288    pub fn from_checker(
289        deps: Arc<DependencyChecker>,
290        healthchecks: Arc<RestHealthcheckRegistry>,
291        check_timeout: Duration,
292    ) -> Arc<Self> {
293        Arc::new(Self {
294            deps,
295            startup_complete: AtomicBool::new(false),
296            healthchecks,
297            check_timeout,
298        })
299    }
300
301    /// The shared dependency-resolution core, e.g. to feed it from consumer
302    /// wiring or to back an in-process [`ReadinessHealthcheck`] leaf.
303    #[must_use]
304    pub fn dependency_checker(&self) -> Arc<DependencyChecker> {
305        Arc::clone(&self.deps)
306    }
307
308    /// Declare a critical dependency dynamically (idempotent; preserves an
309    /// already-resolved entry). Complements the up-front `critical_deps` seed for
310    /// deps discovered during consumer wiring.
311    pub fn register_dep(&self, dep_gear: impl Into<String>) {
312        self.deps.register_dep(dep_gear);
313    }
314
315    /// Mark startup as complete. Idempotent; subsequent calls are ignored.
316    /// `/readyz` will not report `Ready` or `Degraded` until this is called.
317    pub fn mark_startup_complete(&self) {
318        self.startup_complete.store(true, Ordering::SeqCst);
319    }
320
321    /// Whether startup is complete.
322    #[must_use]
323    pub fn is_startup_complete(&self) -> bool {
324        self.startup_complete.load(Ordering::SeqCst)
325    }
326
327    /// Mark a critical dependency as resolved. Idempotent; unknown names are
328    /// ignored.
329    pub fn mark_dep_resolved(&self, name: &str) {
330        if self.deps.mark_resolved(name) {
331            tracing::info!(dep = %name, "critical dependency resolved");
332        }
333    }
334
335    /// Whether all critical dependencies have been resolved.
336    #[must_use]
337    pub fn all_deps_resolved(&self) -> bool {
338        self.deps.all_resolved()
339    }
340
341    /// Set (or clear) the draining flag. Setting it flips `/readyz` to `503`
342    /// immediately so upstreams pull the instance out of rotation while
343    /// in-flight requests drain.
344    pub fn set_draining(&self, draining: bool) {
345        self.deps.set_draining(draining);
346    }
347
348    /// Whether the gear is currently draining.
349    #[must_use]
350    pub fn is_draining(&self) -> bool {
351        self.deps.is_draining()
352    }
353
354    /// Run the registered healthchecks and return the aggregated report.
355    ///
356    /// Used by `/health` to expose full per-component detail and by
357    /// [`Self::evaluate`] to decide readiness state. The registry caches the
358    /// report, so repeated calls within the cache window do not re-run checks.
359    pub async fn health_report(&self) -> HealthcheckReport {
360        self.healthchecks.report(self.check_timeout).await
361    }
362
363    /// Evaluate the aggregate readiness.
364    ///
365    /// The gear is ready when it is not draining, startup is complete, all
366    /// critical deps are resolved, and the aggregated healthcheck report is not
367    /// `Unhealthy`. `Degraded` healthchecks keep the gear ready (`state =
368    /// degraded`, `ready = true`) but the detailed per-component messages belong
369    /// on `/health`, not `/readyz`. The healthcheck fan-out is cached inside the
370    /// registry, so repeated probe traffic does not re-run checks; dependency and
371    /// draining state are read live so transitions take effect immediately.
372    pub async fn evaluate(&self) -> ReadinessReport {
373        let health = self.health_report().await;
374        let draining = self.deps.is_draining();
375        let startup_complete = self.is_startup_complete();
376        let unresolved_deps: Vec<String> = self.deps.unresolved_deps();
377
378        // Draining wins; then any not-ready condition (startup not complete,
379        // unresolved deps, or an Unhealthy check) is `Starting`; then `Degraded`;
380        // else `Ready`.
381        let state = if draining {
382            ReadinessLifecycle::Draining
383        } else if !startup_complete
384            || !unresolved_deps.is_empty()
385            || health.status == HealthcheckStatus::Unhealthy
386        {
387            ReadinessLifecycle::Starting
388        } else if health.status == HealthcheckStatus::Degraded {
389            ReadinessLifecycle::Degraded
390        } else {
391            ReadinessLifecycle::Ready
392        };
393
394        let ready = matches!(
395            state,
396            ReadinessLifecycle::Ready | ReadinessLifecycle::Degraded
397        );
398
399        ReadinessReport {
400            state,
401            ready,
402            unresolved_deps,
403        }
404    }
405}
406
407#[cfg(test)]
408#[cfg_attr(coverage_nightly, coverage(off))]
409#[path = "readiness_tests.rs"]
410mod tests;
411
412#[cfg(test)]
413#[cfg_attr(coverage_nightly, coverage(off))]
414mod dep_checker_tests {
415    use super::DependencyChecker;
416
417    #[test]
418    fn no_deps_is_ready() {
419        let c = DependencyChecker::new();
420        assert!(c.all_resolved());
421        assert!(c.is_ready());
422        assert!(c.unresolved_deps().is_empty());
423    }
424
425    #[test]
426    fn unresolved_dep_lists_and_blocks() {
427        let c = DependencyChecker::new();
428        c.register_dep("billing");
429        c.register_dep("inventory");
430        assert!(!c.is_ready());
431        assert_eq!(
432            c.unresolved_deps(),
433            vec!["billing".to_owned(), "inventory".to_owned()]
434        );
435    }
436
437    #[test]
438    fn ready_once_all_deps_resolved() {
439        let c = DependencyChecker::new();
440        c.register_dep("billing");
441        c.register_dep("inventory");
442        assert!(c.mark_resolved("billing")); // false -> true
443        assert!(!c.is_ready());
444        assert!(c.mark_resolved("inventory"));
445        assert!(c.is_ready());
446    }
447
448    #[test]
449    fn register_dep_is_idempotent_and_preserves_resolved() {
450        let c = DependencyChecker::new();
451        c.register_dep("billing");
452        assert!(c.mark_resolved("billing"));
453        c.register_dep("billing"); // must NOT reset to unresolved
454        assert!(c.is_ready());
455        assert!(!c.mark_resolved("billing")); // already resolved -> no transition
456    }
457
458    #[test]
459    fn mark_resolved_unknown_dep_is_noop() {
460        let c = DependencyChecker::new();
461        assert!(!c.mark_resolved("nope")); // no panic, no transition, no entry
462        assert!(c.is_ready());
463    }
464
465    #[test]
466    fn draining_overrides_ready() {
467        let c = DependencyChecker::new();
468        assert!(c.is_ready());
469        c.set_draining(true);
470        assert!(c.is_draining());
471        assert!(!c.is_ready());
472        c.set_draining(false);
473        assert!(c.is_ready());
474    }
475
476    #[tokio::test]
477    async fn readiness_healthcheck_maps_state_to_status() {
478        use super::ReadinessHealthcheck;
479        use crate::healthcheck::{Healthcheck, HealthcheckStatus};
480        use std::sync::Arc;
481
482        let checker = Arc::new(DependencyChecker::new());
483        checker.register_dep("billing");
484        let hc = ReadinessHealthcheck::new(checker.clone());
485        assert_eq!(hc.name(), "readiness");
486
487        // Unresolved dep → Unhealthy (503), stable code, dep listed.
488        let starting = hc.check().await;
489        assert_eq!(starting.status, HealthcheckStatus::Unhealthy);
490        assert_eq!(starting.code.as_deref(), Some("deps_unresolved"));
491        assert!(starting.message.unwrap().contains("billing"));
492
493        // Resolved → Healthy (200).
494        checker.mark_resolved("billing");
495        assert_eq!(hc.check().await.status, HealthcheckStatus::Healthy);
496
497        // Draining wins regardless of deps.
498        checker.set_draining(true);
499        let draining = hc.check().await;
500        assert_eq!(draining.status, HealthcheckStatus::Unhealthy);
501        assert_eq!(draining.code.as_deref(), Some("draining"));
502    }
503}