Skip to main content

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}