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}