mj_controller/server/api/wait_policy.rs
1use super::*;
2
3/// Classify a harness stop reason.
4///
5/// Stop reasons are free text the harness chooses, so the comparison is
6/// case-insensitive and tolerates both `end_turn` and `endTurn`. Anything
7/// unrecognized is an error carrying the raw reason, because silently calling
8/// an unknown ending "finished" would tell the caller its work succeeded when
9/// nobody knows that it did.
10pub fn map_stop_reason(stop_reason: &str) -> (WaitOutcome, Option<String>) {
11 use mj_core::state::{PromptCompletion, classify_prompt_completion};
12
13 match classify_prompt_completion(stop_reason) {
14 PromptCompletion::InputRequired => (WaitOutcome::InputRequired, None),
15 PromptCompletion::Finished => (WaitOutcome::Finished, None),
16 PromptCompletion::Cancelled => (WaitOutcome::Cancelled, None),
17 PromptCompletion::QuotaLimit => (WaitOutcome::QuotaLimit, None),
18 PromptCompletion::Error => (WaitOutcome::Error, Some(stop_reason.to_owned())),
19 }
20}
21
22/// Everything one pass of the wait loop knows about a session.
23#[derive(Debug, Clone, Default, PartialEq)]
24pub struct WaitObservation {
25 pub checking_continuation: bool,
26 pub background_work: Option<ApiBackgroundWork>,
27 pub pending_elicitations: Vec<mj_core::elicitation::ElicitationRequest>,
28 pub lifecycle: Option<ViewerLifecycleCategory>,
29 /// A resume operation owns this session now. Its durable record still says
30 /// stopped — it stays stopped until the archive has been verified — so
31 /// without this a wait would answer `stopped` for a session that is on its
32 /// way up.
33 pub resuming: bool,
34 /// A close owns this session now. Like `resuming`, this ends nothing: the
35 /// wait follows the close until it finishes.
36 pub closing: bool,
37 /// The reason a close recorded on a session that is alive again, which is
38 /// what a close that failed leaves behind. It is published only until the
39 /// next action or transition for the session succeeds, so it always refers
40 /// to a close nobody has recovered from.
41 pub close_failure: Option<String>,
42 /// The session cannot take a prompt yet: it is still provisioning, a
43 /// lifecycle operation owns it, or its worker is not attached. A wait with
44 /// no target turn keeps waiting, because "finished" would invite a prompt
45 /// that is then refused.
46 pub cannot_take_prompt: bool,
47 /// A recorded launch failure names this session.
48 pub launch_failed: bool,
49 /// Why the launch failed, when a reason was recorded.
50 pub launch_error: Option<String>,
51 pub execution: MaterializedExecutionState,
52 pub active_turn: Option<MaterializedTurn>,
53 pub last_turn_outcome: Option<MaterializedTurnOutcome>,
54 pub queued: usize,
55 pub capacity_retry: Option<CapacityRetry>,
56 pub retry_assessment_pending: bool,
57 pub quota_recovery: Option<mj_core::continuation::QuotaRecovery>,
58 pub start_status: Option<StartStatus>,
59 /// This session is a Mjolnir sub-agent child, whose answer says where its
60 /// report came from.
61 pub subagent: bool,
62 /// The finished turn, by command id, whose child still owes its report.
63 /// Mjolnir reminds the child to hand it back, so like an armed retry the
64 /// turn is not an ending yet.
65 pub report_pending_for: Option<String>,
66 /// The report a child handed back, and the turn it answers.
67 pub handback: Option<(String, String)>,
68}
69
70impl WaitObservation {
71 /// Fold a child's recorded report into this observation, judged against
72 /// the turn this observation saw finish.
73 pub fn apply_subagent_report(
74 &mut self,
75 handback_tool: bool,
76 report: &mj_core::subagent::SubagentReport,
77 now_ms: i64,
78 ) {
79 self.subagent = true;
80 let Some(turn) = self.last_turn_outcome.as_ref() else {
81 return;
82 };
83 let in_flight = self
84 .active_turn
85 .iter()
86 .map(|turn| turn.command_id.as_str())
87 .collect::<Vec<_>>();
88 match mj_core::subagent::report_state(handback_tool, report, Some(turn), &in_flight, now_ms)
89 {
90 mj_core::subagent::ReportState::Pending { .. } => {
91 self.report_pending_for = Some(turn.command_id.clone());
92 }
93 mj_core::subagent::ReportState::Delivered(message) => {
94 self.handback = Some((turn.command_id.clone(), message));
95 }
96 mj_core::subagent::ReportState::Fallback => {}
97 }
98 }
99}
100
101/// The transcript positions one finished turn covers.
102///
103/// A turn is a span, not a starting point. The session keeps recording after a
104/// turn ends — a harness resume notice arrives as an agent message of its own —
105/// and only what falls inside the span is that turn's work.
106#[derive(Debug, Clone, Copy, PartialEq, Eq)]
107pub struct TurnSpan {
108 pub start_position: u64,
109 pub completed_position: u64,
110}
111
112/// What one pass of the wait loop concluded, before the turn summary is read.
113#[derive(Debug, Clone, PartialEq, Eq)]
114pub struct WaitDecision {
115 pub outcome: WaitOutcome,
116 pub stop_reason: Option<String>,
117 pub message: Option<String>,
118 pub turn_id: Option<u64>,
119 /// Which transcript positions the finished turn covers, so its summary can
120 /// be read.
121 pub turn: Option<TurnSpan>,
122}
123
124impl WaitDecision {
125 pub(super) fn simple(outcome: WaitOutcome, message: Option<String>) -> Self {
126 Self {
127 outcome,
128 stop_reason: None,
129 message,
130 turn_id: None,
131 turn: None,
132 }
133 }
134
135 pub(super) fn from_outcome(outcome: &MaterializedTurnOutcome) -> Self {
136 let (kind, stop_reason, message) = match &outcome.outcome {
137 TurnOutcomeKind::Completed { stop_reason } => {
138 let (kind, message) = map_stop_reason(stop_reason);
139 (
140 kind,
141 Some(stop_reason.clone()),
142 outcome
143 .diagnostic
144 .as_ref()
145 .map(|d| d.message.clone())
146 .or(message),
147 )
148 }
149 TurnOutcomeKind::Rejected { message } => {
150 (WaitOutcome::Error, None, Some(message.clone()))
151 }
152 TurnOutcomeKind::Interrupted { message } => {
153 (WaitOutcome::Error, None, Some(message.clone()))
154 }
155 };
156 Self {
157 outcome: kind,
158 stop_reason,
159 message,
160 turn_id: outcome.accepted_ordinal,
161 turn: outcome.turn_start_position.map(|start_position| TurnSpan {
162 start_position,
163 completed_position: outcome.completed_ordinal,
164 }),
165 }
166 }
167}
168
169/// Decide whether this observation ends the wait.
170///
171/// A wait answers for one turn, so only the turn's own fate ends it. In
172/// particular a session that is carrying an error from some earlier, unrelated
173/// action is not a reason to fail the turn the caller asked about: the session
174/// error badge has no expiry, and reporting it here made every later wait on
175/// that session return `error` while the turn ran on perfectly well.
176///
177/// The rules run in order, and the order is the point:
178///
179/// 0. A resume running for this session ends nothing: it is a session coming
180/// up, and its durable record says stopped until the archive is verified.
181/// A close running for it ends nothing either, for the same reason in
182/// reverse: the wait follows it and reports how it ended. A close that
183/// left the session alive failed, and this is the only place left to say
184/// so, because the request that asked for it was answered when it was
185/// admitted. That reason outlives the wait that started it, on purpose: a
186/// close can fail before the next command has even connected, and it is
187/// cleared as soon as anything for the session succeeds.
188/// 1. A stopped or stopping session ends the wait as `stopped`, superseding
189/// any initialization result that raced with the close request.
190/// 2. A launch failure or failed initialization is reported before a turn; a durable
191/// failed lifecycle ends it as `error` even after a daemon restart.
192/// 3. Otherwise the wait has a target turn: the caller's explicit `turn_id`,
193/// else the turn a create-with-prompt call submitted, else "the newest
194/// one", which additionally requires the session to be idle with an empty
195/// queue — with queued prompts, "idle" alone would return an earlier
196/// prompt's outcome — and able to take a prompt, so a session that is
197/// still provisioning or reattaching is not reported as finished. A
198/// prompt handed over at creation that has not become a turn yet is
199/// queued work too: until it is submitted there is no turn to target,
200/// and an idle session would otherwise read as finished before it.
201/// 4. A completed turn under server assessment or with a retry armed is not an
202/// ending: the worker may submit the retry itself, so the wait keeps waiting.
203///
204/// A turn that really did fail still reports `error`: a rejected or interrupted
205/// turn, and an unrecognized stop reason, all come back through the turn record
206/// in rule 3.
207pub fn resolve_wait(observation: &WaitObservation, request: &WaitRequest) -> Option<WaitDecision> {
208 let stopping = matches!(
209 observation.lifecycle,
210 Some(ViewerLifecycleCategory::Suspended | ViewerLifecycleCategory::Suspending)
211 ) || matches!(
212 observation.execution,
213 MaterializedExecutionState::Closing | MaterializedExecutionState::Closed
214 );
215 // A resume owns the session: nothing about it has settled yet, and its
216 // durable record still says stopped. The wait keeps waiting; its own
217 // deadline still bounds it.
218 if observation.resuming {
219 return None;
220 }
221 if let Some(reason) = &observation.close_failure {
222 return Some(WaitDecision::simple(
223 WaitOutcome::Error,
224 Some(reason.clone()),
225 ));
226 }
227 // The close owns the session; its own deadline still bounds this wait.
228 if observation.closing {
229 return None;
230 }
231 if stopping {
232 return Some(WaitDecision::simple(
233 WaitOutcome::Stopped,
234 // A resume that failed rolled the record back to stopped and left
235 // its reason there. Reporting it is the difference between "the
236 // session is stopped" and knowing why it did not come up.
237 Some(
238 observation
239 .launch_error
240 .clone()
241 .unwrap_or_else(|| "the session is stopped or stopping".to_owned()),
242 ),
243 ));
244 }
245 if observation.launch_failed {
246 return Some(WaitDecision::simple(
247 WaitOutcome::Error,
248 Some(
249 observation
250 .launch_error
251 .clone()
252 .unwrap_or_else(|| "the session failed to launch".to_owned()),
253 ),
254 ));
255 }
256 if let Some(StartStatus::Failed { message }) = &observation.start_status {
257 return Some(WaitDecision::simple(
258 WaitOutcome::Error,
259 Some(message.clone()),
260 ));
261 }
262 if observation.lifecycle == Some(ViewerLifecycleCategory::Failed) {
263 return Some(WaitDecision::simple(
264 WaitOutcome::Error,
265 // A close that left the session dead recorded why; saying only
266 // that it failed would throw that away.
267 Some(
268 observation
269 .launch_error
270 .clone()
271 .unwrap_or_else(|| "the session is in a failed state".to_owned()),
272 ),
273 ));
274 }
275 let retry_pending = |outcome: &MaterializedTurnOutcome| {
276 observation.report_pending_for.as_deref() == Some(outcome.command_id.as_str())
277 || observation.retry_assessment_pending
278 || observation.capacity_retry.is_some()
279 || observation.quota_recovery.as_ref().is_some_and(|r| {
280 r.retry_at_ms.is_some() && r.completed_command_id == outcome.command_id
281 })
282 };
283 if let Some(recovery) = &observation.quota_recovery
284 && recovery.retry_at_ms.is_none()
285 && observation
286 .last_turn_outcome
287 .as_ref()
288 .is_some_and(|t| t.command_id == recovery.completed_command_id)
289 && request.turn_id.is_none_or(|target| {
290 observation
291 .last_turn_outcome
292 .as_ref()
293 .and_then(|t| t.accepted_ordinal)
294 .is_some_and(|a| a >= target)
295 })
296 {
297 return Some(WaitDecision::simple(
298 WaitOutcome::QuotaLimit,
299 Some(recovery.notice.clone()),
300 ));
301 }
302 let target = request.turn_id.or(match &observation.start_status {
303 Some(StartStatus::Submitted { turn_id }) => Some(*turn_id),
304 _ => None,
305 });
306 let target_finished = target.is_some_and(|target| {
307 observation
308 .last_turn_outcome
309 .as_ref()
310 .is_some_and(|outcome| {
311 outcome
312 .accepted_ordinal
313 .is_some_and(|ordinal| ordinal >= target)
314 && !retry_pending(outcome)
315 })
316 });
317 if request.return_on_input && !target_finished && !observation.pending_elicitations.is_empty() {
318 return Some(WaitDecision {
319 outcome: WaitOutcome::InputRequired,
320 stop_reason: None,
321 message: Some("the harness needs a response to a structured input request".into()),
322 turn_id: observation
323 .active_turn
324 .as_ref()
325 .and_then(|turn| turn.accepted_ordinal),
326 turn: None,
327 });
328 }
329 match target {
330 Some(target) => {
331 let outcome = observation.last_turn_outcome.as_ref()?;
332 if outcome
333 .accepted_ordinal
334 .is_none_or(|ordinal| ordinal < target)
335 {
336 return None;
337 }
338 if retry_pending(outcome) {
339 return None;
340 }
341 Some(WaitDecision::from_outcome(outcome))
342 }
343 None => {
344 if observation.checking_continuation
345 || observation.cannot_take_prompt
346 || matches!(observation.start_status, Some(StartStatus::Pending))
347 || observation.execution != MaterializedExecutionState::Idle
348 || observation.active_turn.is_some()
349 || observation.queued > 0
350 {
351 return None;
352 }
353 match observation.last_turn_outcome.as_ref() {
354 Some(outcome) if retry_pending(outcome) => None,
355 Some(outcome) => Some(WaitDecision::from_outcome(outcome)),
356 // Idle with nothing queued and nothing ever finished: there is
357 // no turn to wait for, so say so immediately rather than block
358 // for the full timeout.
359 None => Some(WaitDecision::simple(WaitOutcome::Finished, None)),
360 }
361 }
362 }
363}
364
365// ---------------------------------------------------------------------------
366// Router
367// ---------------------------------------------------------------------------
368
369#[cfg(test)]
370mod tests {
371 use super::*;
372
373 #[test]
374 fn awaiting_input_is_a_successful_wait_outcome() {
375 assert_eq!(
376 map_stop_reason(mj_core::acp::AWAITING_INPUT_STOP_REASON),
377 (WaitOutcome::InputRequired, None)
378 );
379 }
380}