use std::sync::atomic::Ordering;
use std::sync::Arc;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use crate::session::ServerState;
use super::rpc::{CoderSessionEntry, CoderSessionMap};
use super::session::{needs_you_from, CoderEventKind, CoderState};
pub const STALL_GRACE_SECS: u64 = 60;
pub const STALL_POLL_SECS: u64 = 30;
const STALLED_FAILURE_KIND: &str = "stalled";
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct LivenessSnapshot {
pub id: String,
pub state: CoderState,
pub last_activity_at: u64,
pub deadline_secs: Option<u64>,
pub waiting_on_human: bool,
}
pub(crate) fn stalled_sessions(sessions: &[LivenessSnapshot], now: u64) -> Vec<String> {
sessions
.iter()
.filter(|session| !session.state.is_terminal())
.filter(|session| {
!matches!(
session.state,
CoderState::NeedsApproval | CoderState::ContractProposed
)
})
.filter(|session| !session.waiting_on_human)
.filter_map(|session| {
let deadline = session.deadline_secs?;
let idle = now.saturating_sub(session.last_activity_at);
(idle > deadline.saturating_add(STALL_GRACE_SECS)).then(|| session.id.clone())
})
.collect()
}
async fn snapshot(entry: &Arc<CoderSessionEntry>) -> LivenessSnapshot {
let event_at = entry.sink.last_event_at();
let session = entry.session.lock().await;
let wall_secs = entry.session_wall_secs.load(Ordering::SeqCst);
LivenessSnapshot {
id: session.id.clone(),
state: session.state,
last_activity_at: event_at.max(session.updated_at),
deadline_secs: (wall_secs > 0).then_some(wall_secs),
waiting_on_human: needs_you_from(
session.state,
entry.user_input.is_pending(),
entry.attention.auth_outstanding(),
entry.attention.approval_kind(),
)
.is_some(),
}
}
async fn fail_if_still_stalled(entry: &Arc<CoderSessionEntry>, now: u64) -> bool {
fail_if_still_stalled_with(entry, now, || std::future::ready(())).await
}
async fn fail_if_still_stalled_with<F, Fut>(
entry: &Arc<CoderSessionEntry>,
now: u64,
before_abort: F,
) -> bool
where
F: FnOnce() -> Fut,
Fut: std::future::Future<Output = ()>,
{
if stalled_sessions(&[snapshot(entry).await], now).is_empty() {
return false;
}
before_abort().await;
let mut session = entry.session.lock().await;
let wall_secs = entry.session_wall_secs.load(Ordering::SeqCst);
let waiting_on_human = needs_you_from(
session.state,
entry.user_input.is_pending(),
entry.attention.auth_outstanding(),
entry.attention.approval_kind(),
)
.is_some();
let last_activity_at = entry.sink.last_event_at().max(session.updated_at);
let still_stalled = wall_secs > 0
&& !session.state.is_terminal()
&& !waiting_on_human
&& now.saturating_sub(last_activity_at) > wall_secs.saturating_add(STALL_GRACE_SECS);
if !still_stalled {
tracing::debug!(
session_id = %session.id,
state = session.state.as_str(),
waiting_on_human,
"coder session changed after watchdog selection; leaving it fully intact"
);
return false;
}
entry.cancel.store(true, Ordering::SeqCst);
if let Some(handle) = entry.task.lock().expect("task slot poisoned").take() {
handle.abort();
}
let reason = format!(
"coder session stalled: no progress for more than its {wall_secs}s deadline plus \
{STALL_GRACE_SECS}s grace"
);
session.error = Some(reason.clone());
session.failure_kind = Some(STALLED_FAILURE_KIND.to_string());
entry.sink.emit(CoderEventKind::Error {
message: reason.clone(),
});
if let Err(error) = session.transition(CoderState::Failed, &entry.sink) {
tracing::warn!(session_id = %session.id, %error, "stalled coder session transition failed");
return false;
}
true
}
pub(crate) async fn sweep_stalled_sessions(
sessions: &tokio::sync::Mutex<CoderSessionMap>,
now: u64,
) -> usize {
let entries: Vec<Arc<CoderSessionEntry>> = sessions.lock().await.values().cloned().collect();
let mut snapshots = Vec::with_capacity(entries.len());
for entry in &entries {
snapshots.push(snapshot(entry).await);
}
let selected = stalled_sessions(&snapshots, now);
let mut failed = 0;
for (entry, snapshot) in entries.iter().zip(&snapshots) {
if selected.iter().any(|id| id == &snapshot.id) {
failed += usize::from(fail_if_still_stalled(entry, now).await);
}
}
failed
}
fn now_secs() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs()
}
pub fn spawn_coder_session_watchdog(state: &Arc<ServerState>) {
let state = Arc::downgrade(state);
tokio::spawn(async move {
let cadence = Duration::from_secs(STALL_POLL_SECS);
let start = tokio::time::Instant::now() + cadence;
let mut ticker = tokio::time::interval_at(start, cadence);
ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
loop {
ticker.tick().await;
let Some(state) = state.upgrade() else {
return;
};
let failed = sweep_stalled_sessions(&state.coder_sessions, now_secs()).await;
if failed > 0 {
tracing::warn!(failed, "failed stalled coder sessions");
}
}
});
}
#[cfg(test)]
mod tests {
use std::collections::VecDeque;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use async_trait::async_trait;
use car_inference::{GenerateRequest, InferenceResult};
use super::*;
use crate::coder::contract::CheckResult;
use crate::coder::native_loop::{LoopFailure, LoopOutcome, TurnGenerator};
use crate::coder::router::EngineChoice;
use crate::coder::rpc::{AttentionState, CoderSessionEntry};
use crate::coder::session::{CoderEventKind, CoderSession, EventSink, UserInputGate};
use crate::coder::skill_memory::RepairMemory;
fn snapshot(id: &str, state: CoderState, last_activity_at: u64) -> LivenessSnapshot {
LivenessSnapshot {
id: id.into(),
state,
last_activity_at,
deadline_secs: Some(100),
waiting_on_human: false,
}
}
#[test]
fn stalled_selection_targets_only_old_unattended_sessions() {
assert_eq!(
stalled_sessions(&[snapshot("stalled", CoderState::Running, 1)], 1_000),
vec!["stalled"]
);
}
#[test]
fn stalled_selection_leaves_fresh_terminal_unbounded_and_human_waits_alone() {
let fresh = snapshot("fresh", CoderState::Running, 950);
let approval = snapshot("approval", CoderState::NeedsApproval, 1);
let mut question = snapshot("question", CoderState::Running, 1);
question.waiting_on_human = true;
let mut unlimited = snapshot("unlimited", CoderState::Running, 1);
unlimited.deadline_secs = None;
let sessions = vec![
fresh,
approval,
question,
unlimited,
snapshot("failed", CoderState::Failed, 1),
];
assert!(stalled_sessions(&sessions, 1_000).is_empty());
}
struct FailingGenerator;
#[async_trait]
impl TurnGenerator for FailingGenerator {
async fn generate(&self, _req: GenerateRequest) -> Result<InferenceResult, String> {
Err("unused".into())
}
}
fn test_entry(
now: u64,
) -> (
Arc<CoderSessionEntry>,
Arc<Mutex<Vec<crate::coder::CoderEvent>>>,
) {
let (sink, emitted) = EventSink::collecting("watchdog-action");
let mut session = CoderSession::new(".", "watch me", EngineChoice::Native, 1, None);
session.id = "watchdog-action".into();
session.state = CoderState::Running;
session.updated_at = now - 1_000;
(
Arc::new(CoderSessionEntry {
session: Arc::new(tokio::sync::Mutex::new(session)),
events: Arc::new(tokio::sync::Mutex::new(VecDeque::new())),
cancel: Arc::new(AtomicBool::new(false)),
preparation: tokio::sync::RwLock::new(()),
session_wall_secs: AtomicU64::new(100),
sink: Arc::new(sink),
infra: car_multi::SharedInfra::new(),
generator: Arc::new(FailingGenerator),
routing_exclusions: Vec::new(),
memory: RepairMemory::disabled(),
mcp_endpoint: None,
mcp_config_dir: None,
user_input: Arc::new(UserInputGate::new()),
attention: Arc::new(AttentionState::default()),
next_seq: Arc::new(AtomicU64::new(0)),
task: std::sync::Mutex::new(None),
fleet: std::sync::Mutex::new(None),
}),
emitted,
)
}
#[tokio::test]
async fn a_sweep_tick_fails_and_cancels_a_stalled_session() {
let sessions = tokio::sync::Mutex::new(CoderSessionMap::new());
let now = 10_000;
let (entry, emitted) = test_entry(now);
*entry.task.lock().expect("task slot poisoned") = Some(tokio::spawn(async {
std::future::pending::<()>().await;
}));
sessions
.lock()
.await
.insert("watchdog-action".into(), entry.clone());
assert_eq!(sweep_stalled_sessions(&sessions, now).await, 1);
assert!(entry.cancel.load(Ordering::SeqCst));
assert!(
entry.task.lock().expect("task slot poisoned").is_none(),
"the real task handle must be taken and aborted"
);
let session = entry.session.lock().await;
assert_eq!(session.state, CoderState::Failed);
assert_eq!(session.failure_kind.as_deref(), Some("stalled"));
drop(session);
assert!(emitted
.lock()
.expect("events poisoned")
.iter()
.any(|event| {
matches!(
&event.kind,
CoderEventKind::Error { message } if message.contains("stalled")
)
}));
}
#[tokio::test]
async fn a_fresh_event_behind_a_contended_buffer_keeps_the_session_live() {
let sessions = tokio::sync::Mutex::new(CoderSessionMap::new());
let now = now_secs();
let (entry, emitted) = test_entry(now);
*entry.task.lock().expect("task slot poisoned") = Some(tokio::spawn(async {
std::future::pending::<()>().await;
}));
let event = entry
.sink
.emit(CoderEventKind::IterationStarted { n: 1, max: 2 });
entry.events.lock().await.push_back(event);
sessions
.lock()
.await
.insert("watchdog-action".into(), entry.clone());
let buffer_guard = entry.events.lock().await;
assert_eq!(sweep_stalled_sessions(&sessions, now).await, 0);
assert!(!entry.cancel.load(Ordering::SeqCst));
assert!(entry.task.lock().expect("task slot poisoned").is_some());
assert_eq!(entry.session.lock().await.state, CoderState::Running);
assert!(!emitted
.lock()
.expect("events poisoned")
.iter()
.any(|event| {
matches!(
&event.kind,
CoderEventKind::Error { message } if message.contains("stalled")
)
}));
drop(buffer_guard);
entry
.task
.lock()
.expect("task slot poisoned")
.take()
.expect("healthy task remains installed")
.abort();
}
#[tokio::test]
async fn a_completion_racing_the_tick_keeps_its_real_terminal() {
let now = 10_000;
let (entry, emitted) = test_entry(now);
let worktree = tempfile::tempdir().expect("worktree");
let release = Arc::new(tokio::sync::Notify::new());
let (done_tx, done_rx) = tokio::sync::oneshot::channel();
let task_entry = entry.clone();
let task_worktree = worktree.path().to_path_buf();
let task_release = release.clone();
let real_result = CheckResult {
credentials_allowed: false,
name: "real-finalizer".into(),
passed: false,
exit_code: Some(1),
output_tail: "the loop finished first".into(),
duration_ms: 1,
timed_out: false,
deadline_clamped: false,
};
let task_result = real_result.clone();
*entry.task.lock().expect("task slot poisoned") = Some(tokio::spawn(async move {
task_release.notified().await;
crate::coder::rpc::finalize_outcome_for_watchdog_test(
&task_entry,
&task_worktree,
LoopOutcome::lost(
LoopFailure::Infrastructure,
Some("real loop terminal".into()),
7,
vec![task_result],
),
)
.await;
let _ = done_tx.send(());
}));
let action_release = release.clone();
let acted = fail_if_still_stalled_with(&entry, now, move || async move {
action_release.notify_one();
done_rx.await.expect("real finalizer must finish");
})
.await;
assert!(!acted, "a real terminal must beat the stalled rewrite");
assert!(
entry.task.lock().expect("task slot poisoned").is_some(),
"a completed session rejected by the recheck keeps its task slot intact"
);
assert!(!entry.cancel.load(Ordering::SeqCst));
let session = entry.session.lock().await;
assert_eq!(session.state, CoderState::Failed);
assert_eq!(session.iterations, 7);
assert_eq!(session.failure_kind.as_deref(), Some("infrastructure"));
assert_eq!(session.last_check_results, vec![real_result]);
assert_eq!(session.error.as_deref(), Some("real loop terminal"));
drop(session);
{
let emitted = emitted.lock().expect("events poisoned");
assert!(!emitted.iter().any(|event| matches!(
&event.kind,
CoderEventKind::Error { message } if message.contains("stalled")
)));
assert!(!emitted
.iter()
.any(|event| matches!(event.kind, CoderEventKind::DiffReady { .. })));
}
let handle = entry
.task
.lock()
.expect("task slot poisoned")
.take()
.expect("completed task remains installed");
handle.await.expect("real finalizer task");
}
#[tokio::test]
async fn a_session_terminal_at_action_time_is_untouched() {
let now = 10_000;
let (entry, emitted) = test_entry(now);
*entry.task.lock().expect("task slot poisoned") = Some(tokio::spawn(async {
std::future::pending::<()>().await;
}));
let action_entry = entry.clone();
let acted = fail_if_still_stalled_with(&entry, now, move || async move {
let mut session = action_entry.session.lock().await;
session.state = CoderState::Failed;
session.error = Some("completed independently".into());
session.failure_kind = Some("infrastructure".into());
session.iterations = 9;
})
.await;
assert!(!acted);
assert!(
entry.task.lock().expect("task slot poisoned").is_some(),
"a session rejected by the under-lock recheck must keep its task"
);
assert!(
!entry.cancel.load(Ordering::SeqCst),
"a session rejected by the under-lock recheck must stay uncancelled"
);
let session = entry.session.lock().await;
assert_eq!(session.error.as_deref(), Some("completed independently"));
assert_eq!(session.failure_kind.as_deref(), Some("infrastructure"));
assert_eq!(session.iterations, 9);
drop(session);
assert!(emitted.lock().expect("events poisoned").is_empty());
entry
.task
.lock()
.expect("task slot poisoned")
.take()
.expect("untouched task remains installed")
.abort();
}
#[tokio::test]
async fn a_session_that_becomes_human_waiting_at_action_time_is_untouched() {
let now = 10_000;
let (entry, emitted) = test_entry(now);
*entry.task.lock().expect("task slot poisoned") = Some(tokio::spawn(async {
std::future::pending::<()>().await;
}));
let action_entry = entry.clone();
let acted = fail_if_still_stalled_with(&entry, now, move || async move {
let _answer = action_entry.user_input.park("Which option?");
})
.await;
assert!(!acted);
assert!(entry.task.lock().expect("task slot poisoned").is_some());
assert!(!entry.cancel.load(Ordering::SeqCst));
assert!(entry.user_input.is_pending());
assert_eq!(entry.session.lock().await.state, CoderState::Running);
assert!(emitted.lock().expect("events poisoned").is_empty());
entry.user_input.clear();
entry
.task
.lock()
.expect("task slot poisoned")
.take()
.expect("untouched task remains installed")
.abort();
}
#[test]
fn daemon_start_spawns_the_coder_session_watchdog() {
let main = include_str!("../../../car-server/src/main.rs");
assert!(
main.contains("spawn_coder_session_watchdog(&server_state)"),
"car-server startup must spawn the coder session watchdog"
);
}
}