leviath_runtime/pipeline/wedge.rs
1//! The wedge watchdog: fail a run that no system can ever look at again.
2//!
3//! Every non-terminal agent rests between ticks holding exactly one phase
4//! marker, and each marker is a claim some system has on it: `ReadyToInfer` is
5//! dispatch's, `AwaitingInference` is the inference collector's, `FanOutWaiting`
6//! is the fan-out collector's, and so on. The marker is what makes the agent
7//! reachable. An agent that is non-terminal and holds *none* of them is in a
8//! state no query matches, so nothing will ever touch it again, and it stays
9//! `running` in `meta.json` for the life of the daemon.
10//!
11//! That is not hypothetical. `PipelineWorld` already logs "a pipeline system
12//! panicked outside any agent's scope; the daemon survived (an agent may be
13//! wedged - cancel it via `lev cancel <run-id>`)": the runtime knows it can
14//! strand a run and asks a person to clean up. Issue #202 is what happens when
15//! the "person" is an unattended harness. It counted the run as occupying a slot
16//! for ever, and its factory ran down to zero free slots over a few hours.
17//!
18//! ## Why this cannot produce a false positive
19//!
20//! The trigger is structural, not temporal. It is not "this run looks old" -
21//! that is the shape of the misdiagnosis in issues #184, #189, #190 and #197,
22//! where a fresh `updated_at`, a bare `waiting`, and a `pid` field each got read
23//! as evidence they were never able to give. It is "the pipeline's own
24//! invariants say this state cannot exist", and the timeout only absorbs the
25//! transient windows inside a tick.
26//!
27//! Those invariants hold on both sides. Every site that removes a phase marker
28//! either inserts a successor or sets a terminal status, and `spawn_agent_seeded`
29//! always lands `Active + ReadyToInfer`, so no ordinary path arrives here.
30//!
31//! What is *not* touched, and why:
32//!
33//! - A long inference holds `AwaitingInference` and `InFlightWork`, and is
34//! bounded twice over besides (the provider's job timeout, and the lane
35//! supervisor that turns a dead task into an ordinary error outcome). A
36//! fifteen-minute call is never a candidate.
37//! - A full inference pool leaves the agent `ReadyToInfer` with a
38//! [`DispatchStall`](super::DispatchStall). That is backpressure working as
39//! designed, and issue #190's watchdog already declines to fail it.
40//! - A tool batch holds `AwaitingTools` and is deliberately unbounded: it may
41//! park off-lane on a tool approval, an `ask_user`, or a `wait_for_agent` that
42//! ends only when some other run does. Any clock-based rule would have to
43//! guess a bound here. This one does not have to.
44//! - A run blocked on a person holds an interaction marker and is left entirely
45//! alone. That is issue #204's territory, and deliberately so: killing a run
46//! somebody is about to answer is a worse failure than leaking a slot, and two
47//! timeouts racing over one status is how #184 happened.
48//! - `Paused` is skipped before the clock is even read, and loses any record it
49//! was carrying, so resuming never finds a run the watchdog had started
50//! counting.
51//!
52//! One residual class this does not cover, stated plainly: an agent that *holds*
53//! a marker but is missing some component the matching system's query also
54//! requires, so that query never sees it. No production path builds an agent
55//! that way today. Covering it would need exactly the "looks old" reasoning the
56//! rest of this module exists to avoid, so it is left uncovered rather than
57//! guessed at.
58//!
59//! ## It composes upward
60//!
61//! Failing a wedged child is enough to free its parent. The fan-out collector
62//! already counts an `Error` worker as finished, and a `requires_children` gate
63//! releases on any terminal child, so there is no parent case to special-case.
64
65use super::*;
66
67/// An agent found in a state no system can reach, and when it was first seen
68/// that way.
69///
70/// One field, unlike [`DispatchStall`](super::DispatchStall), which also carries
71/// a freshness stamp. That record is written by the dispatch systems and read by
72/// a different one, so it has to cope with its writer going away. This one has a
73/// single owner: the watchdog inserts it when the condition holds, keeps the
74/// original `since` while it keeps holding, and removes it the moment the agent
75/// becomes reachable again. There is nothing to go stale.
76#[derive(Component, Debug, Clone, Copy, PartialEq, Eq)]
77pub struct Wedged {
78 /// Unix seconds when the agent was first seen unreachable.
79 pub since: i64,
80}
81
82/// How long an agent may sit unreachable before the run is failed.
83///
84/// A world resource rather than a constant because the daemon serves it from
85/// `[limits] wedge_timeout_secs`. Zero disables the watchdog.
86#[derive(Resource, Debug, Clone, Copy, PartialEq, Eq)]
87pub struct WedgeTimeout(pub u64);
88
89impl Default for WedgeTimeout {
90 fn default() -> Self {
91 Self(DEFAULT_WEDGE_TIMEOUT_SECS)
92 }
93}
94
95/// Default grace period before an unreachable agent fails its run: `0`, meaning
96/// the watchdog is off unless an operator turns it on.
97///
98/// Off by default because this fails runs, and a daemon that starts killing work
99/// after an upgrade nobody asked for is a worse outcome than the leak it
100/// prevents. `300` is the value to set once you want it: five minutes is ten of
101/// the daemon's thirty-second re-drives, the same span `dead_cycles_before_relief`
102/// already treats as "long enough to act on".
103pub const DEFAULT_WEDGE_TIMEOUT_SECS: u64 = 0;
104
105/// Query filter: the agent holds nothing that will cause any system to look at
106/// it again.
107///
108/// Every entry is a component whose *presence* means some system has this agent
109/// queued. This is the load-bearing list in the module, and the one thing here
110/// that needs maintaining: **a new phase marker must be added here**, or an agent
111/// resting on it will be mistaken for an unreachable one. The table-driven test
112/// at the bottom of this file is what catches that.
113///
114/// Note that [`PipelineWorld::fingerprint`](crate::world::PipelineWorld) counts a
115/// subset of these markers for a different purpose (deciding whether a tick
116/// changed anything). The two lists answer different questions and are not
117/// interchangeable.
118///
119/// Components that are per-agent *data* rather than a claim on the agent -
120/// `StageCursor`, `StageProgress`, `DynamicTools`, the auto-approve markers, and
121/// `DispatchStall` itself - are deliberately absent. Their presence says nothing
122/// about whether anything is going to run.
123pub type Unreachable = (
124 (
125 Without<ReadyToInfer>,
126 Without<AwaitingInference>,
127 Without<ProcessResponse>,
128 Without<ReadyForTools>,
129 Without<AwaitingTools>,
130 Without<ReadyForTransition>,
131 Without<ResolveTransition>,
132 Without<ToolsNeedRefresh>,
133 Without<StageJustEntered>,
134 Without<AwaitingTransitionChoice>,
135 Without<AwaitingTransitionResponse>,
136 Without<WaitingForChildren>,
137 ),
138 (
139 Without<AwaitingCompaction>,
140 Without<PendingEdgeCompact>,
141 Without<crate::context_transform::AwaitingContentSummary>,
142 Without<crate::context_transform::PendingContentSummary>,
143 Without<crate::title::PendingTitle>,
144 Without<crate::title::AwaitingTitle>,
145 Without<crate::components::AwaitingInteraction>,
146 Without<crate::gate_prompt::AwaitingGatePrompt>,
147 Without<crate::gate_prompt::GateResolved>,
148 Without<crate::interaction_points::ReadyForInteractionPoint>,
149 Without<crate::interaction_points::AwaitingInteractionPoint>,
150 Without<crate::fanout::FanOutWaiting>,
151 ),
152 (
153 Without<InFlightWork>,
154 Without<crate::tick_scope::PanickedInParallel>,
155 ),
156);
157
158/// What `fail_wedged_runs` selects.
159///
160/// `&'static` is bevy's `WorldQuery` convention, not a claim about
161/// lifetimes: the borrow is bound when the query is fetched.
162type WedgedRunQuery = (
163 Entity,
164 Option<&'static Wedged>,
165 &'static mut AgentState,
166 Option<&'static mut StageIoBuffer>,
167);
168
169/// Wedge watchdog: fail any non-terminal agent that has been unreachable for
170/// longer than [`WedgeTimeout`].
171///
172/// See the module documentation for why this is safe. In short: an agent matches
173/// [`Unreachable`] only in a state the rest of the pipeline guarantees it never
174/// leaves an agent in, so anything that matches is already lost.
175pub fn fail_wedged_runs(
176 mut agents: Query<WedgedRunQuery, Unreachable>,
177 timeout: Option<Res<WedgeTimeout>>,
178 mut commands: Commands,
179) {
180 crate::tick_scope::clear();
181 let limit = timeout.map(|t| t.0).unwrap_or(DEFAULT_WEDGE_TIMEOUT_SECS);
182 if limit == 0 {
183 return; // watchdog disabled
184 }
185 let now = chrono::Utc::now().timestamp();
186 for (entity, wedged, mut state, buffer) in agents.iter_mut() {
187 crate::tick_scope::enter(entity);
188 // A terminal agent is the goal, not a problem, and the host reaps it. A
189 // paused one is stopped because somebody stopped it, and must come back
190 // with no clock already running against it.
191 if is_terminal(&state.status) || state.status == AgentStatus::Paused {
192 if wedged.is_some() {
193 commands.entity(entity).remove::<Wedged>();
194 }
195 continue;
196 }
197 let since = wedged.map(|w| w.since).unwrap_or(now);
198 if now.saturating_sub(since) < limit as i64 {
199 // Inside the grace period: record it (keeping the original start) so
200 // the next tick measures the whole wait rather than restarting.
201 commands.entity(entity).insert(Wedged { since });
202 continue;
203 }
204 let waited = now.saturating_sub(since);
205 let message = format!(
206 "run stopped being driven in stage '{}': {waited}s with nothing in flight and \
207 no marker any system acts on, so it can never move again. Failing it releases \
208 the capacity it was holding. This is a bug in Leviath - please report it with \
209 the daemon log",
210 state.current_stage
211 );
212 tracing::error!(
213 stage = %state.current_stage,
214 wedged_secs = waited,
215 "failing a run that nothing can drive"
216 );
217 if let Some(mut buffer) = buffer {
218 buffer.logs.push((0, format!("[wedged] {message}")));
219 }
220 state.status = AgentStatus::Error { message };
221 commands.entity(entity).remove::<Wedged>();
222 }
223}
224
225/// Whether an agent has stopped for good. Terminal agents are left to the host,
226/// which reaps them once their state has been persisted and reported.
227fn is_terminal(status: &AgentStatus) -> bool {
228 matches!(
229 status,
230 AgentStatus::Complete | AgentStatus::Error { .. } | AgentStatus::Cancelled
231 )
232}
233
234#[cfg(test)]
235mod tests;