use crate::obs::log::Logger;
use serde_json::json;
use std::borrow::Cow;
use std::io::{Read, Write};
use std::net::{TcpListener, TcpStream};
use std::thread;
use std::time::Duration;
const STALE_AFTER_MS: u64 = crate::obs::health::LIVENESS_STALE_AFTER_MS;
pub fn spawn(addr: &str, log: Logger) -> std::io::Result<()> {
let listener = TcpListener::bind(normalize_bind_addr(addr).as_ref())?;
let local = listener.local_addr().ok();
let bound = local
.map(|a| a.to_string())
.unwrap_or_else(|| addr.to_string());
log.info(
"metrics.serving",
json!({"addr": bound, "endpoints": ["/metrics", "/healthz", "/readyz"]}),
);
if local.is_some_and(|a| a.ip().is_unspecified()) {
log.warn(
"metrics.bound_all_interfaces",
json!({"addr": bound, "note": "read-only probe surface reachable on all interfaces; restrict via firewall/NetworkPolicy or bind 127.0.0.1:PORT"}),
);
}
thread::Builder::new()
.name("metrics-http".into())
.spawn(move || {
for s in listener.incoming().flatten() {
let _ = handle(s);
}
})?;
Ok(())
}
fn normalize_bind_addr(addr: &str) -> Cow<'_, str> {
if addr.starts_with(':') {
Cow::Owned(format!("0.0.0.0{addr}"))
} else {
Cow::Borrowed(addr)
}
}
fn handle(mut stream: TcpStream) -> std::io::Result<()> {
let _ = stream.set_read_timeout(Some(Duration::from_secs(5)));
let path = read_request_target(&mut stream)?;
let (status, ctype, body) = route(&path);
let head = format!(
"HTTP/1.1 {status}\r\nContent-Type: {ctype}\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
body.len()
);
stream.write_all(head.as_bytes())?;
stream.write_all(body.as_bytes())?;
stream.flush()
}
fn route(path: &str) -> (&'static str, &'static str, String) {
let path = path.split('?').next().unwrap_or(path);
match path {
"/metrics" => (
"200 OK",
"text/plain; version=0.0.4",
super::metrics::render_prometheus(),
),
"/healthz" => health_response(),
"/readyz" => readiness_response(),
_ => ("404 Not Found", "text/plain", "not found\n".into()),
}
}
fn readiness_response() -> (&'static str, &'static str, String) {
let lame_duck = crate::signals::lame_duck();
let draining = crate::signals::draining();
let intel_all_down = crate::signals::intel_all_down();
if lame_duck || draining || intel_all_down {
(
"503 Service Unavailable",
"text/plain",
format!(
"not ready lame_duck={lame_duck} draining={draining} intel_all_down={intel_all_down}\n"
),
)
} else {
("200 OK", "text/plain", "ready\n".into())
}
}
fn health_response() -> (&'static str, &'static str, String) {
let draining = crate::signals::draining();
let age = crate::obs::health::tick_age_ms();
if !draining && age < STALE_AFTER_MS {
("200 OK", "text/plain", format!("ok tick_age_ms={age}\n"))
} else {
(
"503 Service Unavailable",
"text/plain",
format!("unhealthy draining={draining} tick_age_ms={age}\n"),
)
}
}
fn read_request_target(stream: &mut TcpStream) -> std::io::Result<String> {
let mut buf = [0u8; 1024];
let n = stream.read(&mut buf)?;
let head = String::from_utf8_lossy(&buf[..n]);
let line = head.lines().next().unwrap_or("");
Ok(line.split_whitespace().nth(1).unwrap_or("/").to_string())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn routes_known_and_unknown_paths() {
let (s, ct, body) = route("/metrics");
assert!(s.starts_with("200"));
assert!(ct.contains("version=0.0.4"));
assert!(body.contains("agent_runs_started_total"));
let (s, _, _) = route("/nope");
assert!(s.starts_with("404"));
}
#[test]
fn readyz_flips_503_under_lame_duck_then_clears() {
let _g = crate::signals::test_guard();
crate::signals::set_lame_duck(false);
if !crate::signals::draining() {
let (s, _, body) = route("/readyz");
assert_eq!(s, "200 OK", "baseline /readyz is ready");
assert_eq!(body, "ready\n");
}
crate::signals::set_lame_duck(true);
let (s, _, body) = route("/readyz");
assert_eq!(s, "503 Service Unavailable", "lame-duck → NotReady");
assert!(body.contains("lame_duck=true"), "body: {body}");
crate::signals::set_lame_duck(false);
if !crate::signals::draining() {
let (s, _, _) = route("/readyz");
assert_eq!(s, "200 OK", "clearing lame-duck restores Ready");
}
}
#[test]
fn readyz_flips_503_under_intel_all_down_then_clears() {
let _g = crate::signals::test_guard();
if crate::signals::draining() {
return; }
let (s, _, _) = route("/readyz");
assert_eq!(s, "200 OK", "baseline /readyz is ready");
crate::signals::set_intel_all_down(true);
let (s, _, body) = route("/readyz");
assert_eq!(s, "503 Service Unavailable", "intel all-down → NotReady");
assert!(body.contains("intel_all_down=true"), "body: {body}");
crate::signals::set_intel_all_down(false);
let (s, _, _) = route("/readyz");
assert_eq!(s, "200 OK", "intel recovery restores Ready");
}
#[test]
fn normalize_bind_addr_expands_bare_port_only() {
assert_eq!(normalize_bind_addr(":9090"), "0.0.0.0:9090");
assert_eq!(normalize_bind_addr("0.0.0.0:9090"), "0.0.0.0:9090");
assert_eq!(normalize_bind_addr("127.0.0.1:9090"), "127.0.0.1:9090");
assert_eq!(normalize_bind_addr("localhost:9090"), "localhost:9090");
assert_eq!(normalize_bind_addr("[::]:9090"), "[::]:9090");
}
#[test]
fn bare_port_actually_binds() {
let l = TcpListener::bind(normalize_bind_addr(":0").as_ref());
assert!(l.is_ok(), "bare :port must bind after normalisation: {l:?}");
}
#[test]
fn metrics_path_ignores_query_string() {
let (s, _, _) = route("/metrics?foo=bar");
assert!(s.starts_with("200"));
}
#[test]
fn healthz_reflects_only_the_reactor_heartbeat() {
let _sig = crate::signals::test_guard();
let _g = crate::obs::health::HEARTBEAT_TEST_LOCK.lock().unwrap();
crate::signals::set_lame_duck(false); if crate::signals::draining() {
return; }
crate::obs::health::tick();
let (s, _, body) = route("/healthz");
assert_eq!(s, "200 OK", "fresh reactor tick → live");
assert!(body.starts_with("ok tick_age_ms="), "body: {body}");
crate::obs::health::set_tick_age_for_test(STALE_AFTER_MS + 1_000);
let (s, _, body) = route("/healthz");
assert_eq!(s, "503 Service Unavailable", "wedged reactor → unhealthy");
assert!(body.starts_with("unhealthy"), "body: {body}");
crate::obs::health::tick();
let (s, _, _) = route("/healthz");
assert_eq!(s, "200 OK", "a renewed reactor tick restores liveness");
}
}