use std::future::Future;
use std::pin::Pin;
use std::sync::{
atomic::{AtomicBool, Ordering},
Arc, RwLock,
};
use axum::{
extract::{RawQuery, State},
http::{header, StatusCode},
response::{IntoResponse, Response},
routing::get,
Router,
};
pub const LIVE_PATH: &str = "/livez";
pub const READY_PATH: &str = "/readyz";
pub const LEGACY_HEALTH_PATH: &str = "/healthz";
pub trait ReadyCheck: Send + Sync {
fn name(&self) -> &'static str;
fn ready(&self) -> Pin<Box<dyn Future<Output = bool> + Send + '_>>;
}
pub struct Gate {
name: &'static str,
check: Box<dyn Fn() -> bool + Send + Sync>,
}
impl Gate {
pub fn new<F>(name: &'static str, check: F) -> Self
where
F: Fn() -> bool + Send + Sync + 'static,
{
Self {
name,
check: Box::new(check),
}
}
}
impl ReadyCheck for Gate {
fn name(&self) -> &'static str {
self.name
}
fn ready(&self) -> Pin<Box<dyn Future<Output = bool> + Send + '_>> {
let verdict = (self.check)();
Box::pin(async move { verdict })
}
}
#[derive(Default)]
pub struct Health {
started: AtomicBool,
draining: AtomicBool,
checks: RwLock<Vec<Arc<dyn ReadyCheck>>>,
}
impl Health {
pub fn new() -> Arc<Self> {
Arc::new(Self::default())
}
pub fn ready() -> Arc<Self> {
let h = Self::default();
h.started.store(true, Ordering::Release);
Arc::new(h)
}
pub fn set_checks(&self, checks: Vec<Arc<dyn ReadyCheck>>) {
*self.checks.write().expect("health checks poisoned") = checks;
}
pub fn set_gates(&self, gates: Vec<Gate>) {
self.set_checks(
gates
.into_iter()
.map(|g| Arc::new(g) as Arc<dyn ReadyCheck>)
.collect(),
);
}
pub fn mark_started(&self) {
self.started.store(true, Ordering::Release);
}
pub fn begin_drain(&self) {
self.draining.store(true, Ordering::Release);
}
pub fn is_draining(&self) -> bool {
self.draining.load(Ordering::Acquire)
}
async fn readiness(&self) -> (bool, String) {
let mut ok = true;
let mut body = String::new();
let line = |name: &str, pass: bool, body: &mut String| {
body.push_str(if pass { "[+]" } else { "[-]" });
body.push_str(name);
body.push_str(if pass { " ok\n" } else { " failed\n" });
};
let started = self.started.load(Ordering::Acquire);
line("started", started, &mut body);
ok &= started;
let live = !self.is_draining();
line("shutdown", live, &mut body);
ok &= live;
let checks: Vec<Arc<dyn ReadyCheck>> = self
.checks
.read()
.expect("health checks poisoned")
.iter()
.cloned()
.collect();
for check in checks {
let pass = check.ready().await;
line(check.name(), pass, &mut body);
ok &= pass;
}
(ok, body)
}
}
pub fn probe_routes(health: Arc<Health>) -> Router {
Router::new()
.route(LIVE_PATH, get(livez))
.route(READY_PATH, get(readyz))
.route(LEGACY_HEALTH_PATH, get(readyz))
.route(crate::HEALTH_PATH, get(livez))
.with_state(health)
}
async fn livez(State(health): State<Arc<Health>>, RawQuery(q): RawQuery) -> Response {
if !is_verbose(q.as_deref()) {
return plaintext(StatusCode::OK, "ok".into());
}
let mut body = String::from("[+]ping ok\n");
if health.is_draining() {
body.push_str("[+]draining ok\n");
}
body.push_str("livez check passed\n");
plaintext(StatusCode::OK, body)
}
async fn readyz(State(health): State<Arc<Health>>, RawQuery(q): RawQuery) -> Response {
let (ok, detail) = health.readiness().await;
let status = if ok {
StatusCode::OK
} else {
StatusCode::SERVICE_UNAVAILABLE
};
if !is_verbose(q.as_deref()) {
let body = if ok { "ok\n" } else { "readyz check failed\n" };
return plaintext(status, body.into());
}
let mut body = detail;
body.push_str(if ok {
"readyz check passed\n"
} else {
"readyz check failed\n"
});
plaintext(status, body)
}
fn is_verbose(query: Option<&str>) -> bool {
query.is_some_and(|q| q.split('&').any(|p| p == "verbose" || p.starts_with("verbose=")))
}
fn plaintext(status: StatusCode, body: String) -> Response {
(
status,
[
(header::CONTENT_TYPE, "text/plain; charset=utf-8"),
(header::CACHE_CONTROL, "no-cache, no-store, must-revalidate"),
],
body,
)
.into_response()
}
#[cfg(test)]
mod tests {
use super::*;
use axum::{body::Body, http::Request};
use tower::ServiceExt;
async fn probe(app: Router, uri: &str) -> (StatusCode, String) {
let resp = app
.oneshot(Request::builder().uri(uri).body(Body::empty()).unwrap())
.await
.unwrap();
let status = resp.status();
let bytes = axum::body::to_bytes(resp.into_body(), 64 * 1024).await.unwrap();
(status, String::from_utf8(bytes.to_vec()).unwrap())
}
#[tokio::test]
async fn readyz_fails_before_started_but_livez_passes() {
let health = Health::new();
let app = probe_routes(health.clone());
let (status, _) = probe(app.clone(), READY_PATH).await;
assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE);
let (status, _) = probe(app.clone(), LIVE_PATH).await;
assert_eq!(status, StatusCode::OK);
health.mark_started();
let (status, body) = probe(app, READY_PATH).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(body, "ok\n");
}
#[tokio::test]
async fn a_failing_gate_fails_readiness_only() {
let health = Health::ready();
health.set_gates(vec![Gate::new("dist", || false)]);
let app = probe_routes(health);
let (status, _) = probe(app.clone(), READY_PATH).await;
assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE);
let (status, _) = probe(app, LIVE_PATH).await;
assert_eq!(status, StatusCode::OK);
}
#[tokio::test]
async fn draining_fails_readiness_and_keeps_liveness() {
let health = Health::ready();
health.begin_drain();
let app = probe_routes(health);
let (status, _) = probe(app.clone(), READY_PATH).await;
assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE);
let (status, _) = probe(app, LIVE_PATH).await;
assert_eq!(
status,
StatusCode::OK,
"a draining pod is healthy; restarting it mid-drain drops the connections \
the drain exists to protect"
);
}
#[tokio::test]
async fn verbose_lists_every_check_not_just_the_first_failure() {
let health = Health::ready();
health.set_gates(vec![Gate::new("dist", || false), Gate::new("ssr", || false)]);
let (status, body) = probe(probe_routes(health), "/readyz?verbose").await;
assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE);
assert_eq!(
body,
"[+]started ok\n[+]shutdown ok\n[-]dist failed\n[-]ssr failed\nreadyz check failed\n",
);
}
#[tokio::test]
async fn set_gates_replaces_rather_than_accumulates() {
let health = Health::ready();
health.set_gates(vec![Gate::new("dist", || false)]);
health.set_gates(vec![Gate::new("dist", || true)]);
let (status, body) = probe(probe_routes(health), "/readyz?verbose").await;
assert_eq!(status, StatusCode::OK);
assert_eq!(body.matches("dist").count(), 1, "rebuilt router duplicated a gate");
}
#[tokio::test]
async fn legacy_paths_keep_their_liveness_meaning() {
let health = Health::new();
health.set_gates(vec![Gate::new("dist", || false)]);
let app = probe_routes(health);
let (status, _) = probe(app.clone(), crate::HEALTH_PATH).await;
assert_eq!(status, StatusCode::OK);
let (status, _) = probe(app, LEGACY_HEALTH_PATH).await;
assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE);
}
#[test]
fn verbose_matches_only_the_flag() {
assert!(is_verbose(Some("verbose")));
assert!(is_verbose(Some("verbose=1")));
assert!(is_verbose(Some("exclude=ssr&verbose")));
assert!(!is_verbose(Some("verbosely=1")));
assert!(!is_verbose(None));
}
}