use std::path::PathBuf;
use std::time::Duration;
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
use super::super::{CheckResult, CheckStatus};
use super::{check_daemon_health_at, interpret_health_body, summarize};
#[derive(Clone, Copy)]
enum Reply {
After(Duration),
Frame(&'static str),
Never,
}
async fn stand_in(reply: Reply) -> (PathBuf, tempfile::TempDir, tokio::task::JoinHandle<()>) {
let tmp = tempfile::tempdir().expect("tempdir");
let socket = tmp.path().join("sockets").join("stand-in.sock");
let listener = trusty_common::uds::bind_hardened(&socket).expect("bind hardened socket");
let task = tokio::spawn(async move {
let mut held = Vec::new();
while let Ok((stream, _)) = listener.accept().await {
match reply {
Reply::Never => held.push(stream),
Reply::After(delay) => {
let body = serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"result": {
"status": "ok",
"daemon_state": "ready",
"worker": {
"in_flight": 0,
"wedged": false,
"stall_tracking_ok": true,
},
},
});
tokio::spawn(answer_once(stream, delay, body.to_string()));
}
Reply::Frame(raw) => {
tokio::spawn(answer_once(stream, Duration::ZERO, raw.to_string()));
}
}
}
});
(socket, tmp, task)
}
async fn answer_once(stream: tokio::net::UnixStream, delay: Duration, frame: String) {
let mut reader = BufReader::new(stream);
let mut line = String::new();
if reader.read_line(&mut line).await.is_err() {
return;
}
tokio::time::sleep(delay).await;
let _ = reader
.get_mut()
.write_all(format!("{frame}\n").as_bytes())
.await;
}
async fn probe(socket: &std::path::Path, budget: Duration) -> CheckResult {
tokio::time::timeout(
budget + Duration::from_secs(5),
check_daemon_health_at("daemon socket".to_string(), socket, budget),
)
.await
.expect("the probe must finish within its own budget")
}
#[tokio::test]
async fn a_socket_that_accepts_and_never_answers_is_unknown_within_the_budget() {
let (socket, _tmp, task) = stand_in(Reply::Never).await;
let result = probe(&socket, Duration::from_millis(200)).await;
task.abort();
assert_eq!(result.status, CheckStatus::Unknown, "{result:?}");
let detail = result.detail.as_deref().unwrap_or_default();
assert!(
detail.contains("did not answer") && detail.contains("could not be determined"),
"the timeout must be named, not dressed as a verdict: {detail}"
);
assert!(
!summarize(&[result]).healthy,
"a timed-out probe must not end the run healthy"
);
}
#[tokio::test]
async fn an_answer_without_a_health_body_is_never_a_healthy_run() {
let cases = [
(
r#"{"jsonrpc":"2.0","id":1,"error":{"code":-32603,"message":"boom"}}"#,
CheckStatus::Fail,
),
(r#"{"jsonrpc":"2.0","id":1}"#, CheckStatus::Unknown),
("not json", CheckStatus::Fail),
];
for (frame, expected) in cases {
let (socket, _tmp, task) = stand_in(Reply::Frame(frame)).await;
let result = probe(&socket, Duration::from_secs(2)).await;
task.abort();
assert_eq!(result.status, expected, "frame {frame}: {result:?}");
assert!(
!summarize(&[result]).healthy,
"frame {frame} must not end the run healthy"
);
}
}
#[tokio::test]
async fn a_responsive_daemon_passes() {
for delay in [Duration::ZERO, Duration::from_millis(300)] {
let (socket, _tmp, task) = stand_in(Reply::After(delay)).await;
let result = probe(&socket, Duration::from_secs(2)).await;
task.abort();
assert_eq!(
result.status,
CheckStatus::Pass,
"a daemon answering after {delay:?} must pass: {result:?}"
);
assert!(summarize(&[result]).healthy, "delay {delay:?}");
}
}
#[tokio::test]
async fn a_wedge_from_a_held_handle_lock_fails_and_names_the_lock() {
let frame = r#"{"jsonrpc":"2.0","id":1,"result":{"status":"wedged","daemon_state":"ready","worker":{"in_flight":0,"wedged":true,"stalled_lock":{"palace":"trusty-tools","lock":"write","age_secs":131}}}}"#;
let (socket, _tmp, task) = stand_in(Reply::Frame(frame)).await;
let result = probe(&socket, Duration::from_secs(2)).await;
task.abort();
assert_eq!(result.status, CheckStatus::Fail, "{result:?}");
let detail = result.detail.as_deref().unwrap_or_default();
assert!(
detail.contains("write lock of palace trusty-tools") && detail.contains("at least 131s"),
"the held lock must be named in prose, not as JSON: {detail}"
);
assert!(!summarize(&[result]).healthy);
}
#[test]
fn stopped_stall_tracking_is_undetermined_not_a_warning() {
let result = interpret_health_body(
"HTTP daemon".to_string(),
"http://x/health",
200,
Some(&serde_json::json!({
"status": "degraded",
"detail": "palace lock stall ticker has not run for 412s (interval 30000ms); a \
held lock may go unnoticed between health polls (#4001)",
"daemon_state": "ready",
"worker": {"in_flight": 0, "wedged": false, "stall_tracking_ok": false},
})),
);
assert_eq!(result.status, CheckStatus::Unknown, "{result:?}");
let detail = result.detail.as_deref().unwrap_or_default();
assert!(
detail.contains("UNKNOWN") && detail.contains("stall tracking is not running"),
"the dead detector must be named: {detail}"
);
assert!(
!summarize(&[CheckResult::pass("a", "fine"), result]).healthy,
"a daemon that stopped watching its palace locks must not end the run green"
);
}
#[test]
fn a_daemon_with_no_stall_detector_is_undetermined_not_a_pass() {
let result = interpret_health_body(
"HTTP daemon".to_string(),
"http://x/health",
200,
Some(&serde_json::json!({
"status": "ok",
"daemon_state": "ready",
"worker": {"in_flight": 0, "oldest_age_secs": 1, "wedged": false},
})),
);
assert_eq!(result.status, CheckStatus::Unknown, "{result:?}");
let detail = result.detail.as_deref().unwrap_or_default();
assert!(
detail.contains("no palace-lock stall detector") && detail.contains("UNKNOWN"),
"the missing detector must be named: {detail}"
);
assert!(
!summarize(&[CheckResult::pass("a", "fine"), result]).healthy,
"the incident build must not end the run green"
);
}
#[test]
fn a_pool_wedge_beside_a_benign_stamp_names_the_pool() {
let result = interpret_health_body(
"HTTP daemon".to_string(),
"http://x/health",
200,
Some(&serde_json::json!({
"status": "wedged",
"daemon_state": "ready",
"worker": {
"in_flight": 4,
"oldest_age_secs": 900,
"wedged": true,
"wedged_reason": "pool",
"stalled_lock": {"palace": "scratch", "lock": "commit", "age_secs": 3},
},
})),
);
assert_eq!(result.status, CheckStatus::Fail, "{result:?}");
let detail = result.detail.as_deref().unwrap_or_default();
assert!(
detail.contains("WEDGED worker pool")
&& detail.contains("900s")
&& detail.contains("4 in flight"),
"the pool must be the subject: {detail}"
);
assert!(
!detail.contains("scratch"),
"a sub-threshold stamp is not the wedge: {detail}"
);
}
#[test]
fn an_undetermined_check_is_not_a_healthy_run() {
let summary = summarize(&[
CheckResult::pass("a", "fine"),
CheckResult::unknown("daemon socket", "did not answer"),
]);
assert!(!summary.healthy);
assert_eq!(
summary.line,
"1 passed, 0 warnings, 1 undetermined, 0 failed."
);
}
#[test]
fn a_non_health_undetermined_check_also_ends_the_run_unhealthy() {
let summary = summarize(&[
CheckResult::pass("daemon socket", "workers progressing"),
CheckResult::unknown(
"Tier S facts",
"a daemon predating #4890 does not report `affirmed_at`",
),
CheckResult::unknown(
"MCP registrations",
"could not resolve the home directory, so no client config was read",
),
]);
assert!(!summary.healthy);
assert_eq!(
summary.line,
"1 passed, 0 warnings, 2 undetermined, 0 failed."
);
}
#[test]
fn warnings_alone_are_a_healthy_run() {
let summary = summarize(&[
CheckResult::pass("a", "fine"),
CheckResult::warn("b", "minor"),
]);
assert!(summary.healthy);
assert!(!summarize(&[CheckResult::fail("c", "broken")]).healthy);
}