use std::collections::BTreeMap;
use std::time::{Duration, Instant};
pub(super) const HARNESS_READINESS_TIMEOUT: Duration =
crate::controller::NATIVE_SESSION_STARTUP_TIMEOUT;
pub(super) struct ReadinessObservation {
pub(super) session_id: String,
pub(super) harness_ready: bool,
pub(super) record_age: Option<Duration>,
pub(super) updated_at: String,
pub(super) detail: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(super) struct UnreadySession {
pub(super) session_id: String,
pub(super) observed_updated_at: String,
pub(super) waited: Duration,
pub(super) detail: Option<String>,
}
impl UnreadySession {
pub(super) fn cause(&self) -> String {
let mut cause = format!(
"the harness never advertised its configuration within {}s, so this session cannot be used; destroy it and create a replacement",
self.waited.as_secs()
);
if let Some(detail) = &self.detail {
cause.push_str("; the daemon last saw: ");
cause.push_str(detail);
}
cause
}
}
struct Waiting {
since: Instant,
reported: bool,
}
#[derive(Default)]
pub(super) struct HarnessReadinessWatch {
waiting: BTreeMap<String, Waiting>,
}
impl HarnessReadinessWatch {
pub(super) fn observe(
&mut self,
now: Instant,
observations: Vec<ReadinessObservation>,
) -> Vec<UnreadySession> {
self.waiting.retain(|session_id, _| {
observations
.iter()
.any(|observation| &observation.session_id == session_id)
});
let mut unready = Vec::new();
for observation in observations {
if observation.harness_ready {
self.waiting.remove(&observation.session_id);
continue;
}
if !self.waiting.contains_key(&observation.session_id) {
let Some(age) = observation
.record_age
.filter(|age| *age < HARNESS_READINESS_TIMEOUT)
else {
continue;
};
self.waiting.insert(
observation.session_id.clone(),
Waiting {
since: now.checked_sub(age).unwrap_or(now),
reported: false,
},
);
}
let waiting = self
.waiting
.get_mut(&observation.session_id)
.expect("the session was armed above");
let waited = now.saturating_duration_since(waiting.since);
if waiting.reported || waited < HARNESS_READINESS_TIMEOUT {
continue;
}
waiting.reported = true;
unready.push(UnreadySession {
session_id: observation.session_id,
observed_updated_at: observation.updated_at,
waited,
detail: observation.detail,
});
}
unready
}
}
#[cfg(test)]
mod tests {
use super::*;
fn observation(session_id: &str, harness_ready: bool) -> ReadinessObservation {
ReadinessObservation {
session_id: session_id.to_owned(),
harness_ready,
record_age: Some(Duration::ZERO),
updated_at: "2026-09-18T15:28:35Z".to_owned(),
detail: None,
}
}
#[test]
fn a_session_that_never_advertises_its_harness_fails_once_the_wait_runs_out() {
let mut watch = HarnessReadinessWatch::default();
let start = Instant::now();
assert!(
watch
.observe(start, vec![observation("session-1", false)])
.is_empty(),
"a session that has just reported ready is still waiting"
);
assert!(
watch
.observe(
start + HARNESS_READINESS_TIMEOUT - Duration::from_secs(1),
vec![observation("session-1", false)],
)
.is_empty(),
"the wait is not cut short"
);
let unready = watch.observe(
start + HARNESS_READINESS_TIMEOUT,
vec![ReadinessObservation {
detail: Some("relay worker is unreachable".to_owned()),
..observation("session-1", false)
}],
);
assert_eq!(
unready,
vec![UnreadySession {
session_id: "session-1".to_owned(),
observed_updated_at: "2026-09-18T15:28:35Z".to_owned(),
waited: HARNESS_READINESS_TIMEOUT,
detail: Some("relay worker is unreachable".to_owned()),
}]
);
assert!(
unready[0].cause().starts_with(&format!(
"the harness never advertised its configuration within {}s",
HARNESS_READINESS_TIMEOUT.as_secs()
)),
"{}",
unready[0].cause()
);
assert!(
unready[0].cause().contains("relay worker is unreachable"),
"the reason carries what the daemon last saw"
);
assert!(
watch
.observe(
start + HARNESS_READINESS_TIMEOUT + Duration::from_secs(30),
vec![observation("session-1", false)],
)
.is_empty(),
"one failure per wait, however long the record takes to change"
);
}
#[test]
fn a_session_whose_harness_arrives_in_time_is_never_failed() {
let mut watch = HarnessReadinessWatch::default();
let start = Instant::now();
assert!(
watch
.observe(start, vec![observation("session-1", false)])
.is_empty()
);
assert!(
watch
.observe(
start + Duration::from_secs(20),
vec![observation("session-1", true)],
)
.is_empty()
);
assert!(
watch
.observe(
start + HARNESS_READINESS_TIMEOUT + Duration::from_secs(60),
vec![observation("session-1", true)],
)
.is_empty(),
"a harness that answered is never failed for the wait it ended"
);
}
#[test]
fn the_wait_is_measured_from_the_moment_the_session_reported_ready() {
let mut watch = HarnessReadinessWatch::default();
let start = Instant::now();
let already_waited = HARNESS_READINESS_TIMEOUT - Duration::from_secs(60);
assert!(
watch
.observe(
start,
vec![ReadinessObservation {
record_age: Some(already_waited),
..observation("session-1", false)
}],
)
.is_empty()
);
let unready = watch.observe(
start + Duration::from_secs(60),
vec![ReadinessObservation {
record_age: Some(already_waited + Duration::from_secs(60)),
..observation("session-1", false)
}],
);
assert_eq!(unready.len(), 1, "the record's own clock starts the wait");
assert_eq!(unready[0].waited, HARNESS_READINESS_TIMEOUT);
}
#[test]
fn a_session_the_daemon_did_not_see_start_is_left_alone() {
let mut watch = HarnessReadinessWatch::default();
let start = Instant::now();
let stale = ReadinessObservation {
record_age: Some(HARNESS_READINESS_TIMEOUT + Duration::from_secs(1)),
..observation("session-1", false)
};
assert!(watch.observe(start, vec![stale]).is_empty());
assert!(
watch
.observe(
start + HARNESS_READINESS_TIMEOUT * 2,
vec![ReadinessObservation {
record_age: Some(HARNESS_READINESS_TIMEOUT * 3),
..observation("session-1", false)
}],
)
.is_empty(),
"a session that was already running before the daemon saw it keeps reconnecting"
);
}
#[test]
fn a_session_that_leaves_the_sweep_starts_its_wait_again() {
let mut watch = HarnessReadinessWatch::default();
let start = Instant::now();
assert!(
watch
.observe(start, vec![observation("session-1", false)])
.is_empty()
);
assert!(
watch
.observe(start + Duration::from_secs(30), Vec::new())
.is_empty()
);
assert!(
watch
.observe(
start + HARNESS_READINESS_TIMEOUT,
vec![observation("session-1", false)],
)
.is_empty(),
"the wait restarts once the session is the daemon's to watch again"
);
}
}