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    /// A recorded launch failure names this session.
43    pub launch_failed: bool,
44    /// Why the launch failed, when a reason was recorded.
45    pub launch_error: Option<String>,
46    pub execution: MaterializedExecutionState,
47    pub active_turn: Option<MaterializedTurn>,
48    pub last_turn_outcome: Option<MaterializedTurnOutcome>,
49    pub queued: usize,
50    pub capacity_retry: Option<CapacityRetry>,
51    pub quota_recovery: Option<mj_core::continuation::QuotaRecovery>,
52    pub start_status: Option<StartStatus>,
53}
54
55/// The transcript positions one finished turn covers.
56///
57/// A turn is a span, not a starting point. The session keeps recording after a
58/// turn ends — a harness resume notice arrives as an agent message of its own —
59/// and only what falls inside the span is that turn's work.
60#[derive(Debug, Clone, Copy, PartialEq, Eq)]
61pub struct TurnSpan {
62    pub start_position: u64,
63    pub completed_position: u64,
64}
65
66/// What one pass of the wait loop concluded, before the turn summary is read.
67#[derive(Debug, Clone, PartialEq, Eq)]
68pub struct WaitDecision {
69    pub outcome: WaitOutcome,
70    pub stop_reason: Option<String>,
71    pub message: Option<String>,
72    pub turn_id: Option<u64>,
73    /// Which transcript positions the finished turn covers, so its summary can
74    /// be read.
75    pub turn: Option<TurnSpan>,
76}
77
78impl WaitDecision {
79    pub(super) fn simple(outcome: WaitOutcome, message: Option<String>) -> Self {
80        Self {
81            outcome,
82            stop_reason: None,
83            message,
84            turn_id: None,
85            turn: None,
86        }
87    }
88
89    pub(super) fn from_outcome(outcome: &MaterializedTurnOutcome) -> Self {
90        let (kind, stop_reason, message) = match &outcome.outcome {
91            TurnOutcomeKind::Completed { stop_reason } => {
92                let (kind, message) = map_stop_reason(stop_reason);
93                (
94                    kind,
95                    Some(stop_reason.clone()),
96                    outcome
97                        .diagnostic
98                        .as_ref()
99                        .map(|d| d.message.clone())
100                        .or(message),
101                )
102            }
103            TurnOutcomeKind::Rejected { message } => {
104                (WaitOutcome::Error, None, Some(message.clone()))
105            }
106            TurnOutcomeKind::Interrupted { message } => {
107                (WaitOutcome::Error, None, Some(message.clone()))
108            }
109        };
110        Self {
111            outcome: kind,
112            stop_reason,
113            message,
114            turn_id: outcome.accepted_ordinal,
115            turn: outcome.turn_start_position.map(|start_position| TurnSpan {
116                start_position,
117                completed_position: outcome.completed_ordinal,
118            }),
119        }
120    }
121}
122
123/// Decide whether this observation ends the wait.
124///
125/// A wait answers for one turn, so only the turn's own fate ends it. In
126/// particular a session that is carrying an error from some earlier, unrelated
127/// action is not a reason to fail the turn the caller asked about: the session
128/// error badge has no expiry, and reporting it here made every later wait on
129/// that session return `error` while the turn ran on perfectly well.
130///
131/// The rules run in order, and the order is the point:
132///
133/// 0. A resume running for this session ends nothing: it is a session coming
134///    up, and its durable record says stopped until the archive is verified.
135///    A close running for it ends nothing either, for the same reason in
136///    reverse: the wait follows it and reports how it ended. A close that
137///    left the session alive failed, and this is the only place left to say
138///    so, because the request that asked for it was answered when it was
139///    admitted. That reason outlives the wait that started it, on purpose: a
140///    close can fail before the next command has even connected, and it is
141///    cleared as soon as anything for the session succeeds.
142/// 1. A stopped or stopping session ends the wait as `stopped`, superseding
143///    any initialization result that raced with the close request.
144/// 2. A launch failure or failed initialization is reported before a turn; a durable
145///    failed lifecycle ends it as `error` even after a daemon restart.
146/// 3. Otherwise the wait has a target turn: the caller's explicit `turn_id`,
147///    else the turn a create-with-prompt call submitted, else "the newest
148///    one", which additionally requires the session to be idle with an empty
149///    queue — with queued prompts, "idle" alone would return an earlier
150///    prompt's outcome.
151/// 4. A capacity outcome with a retry armed is not an ending: the worker will
152///    submit the retry itself, so the wait keeps waiting.
153///
154/// A turn that really did fail still reports `error`: a rejected or interrupted
155/// turn, and an unrecognized stop reason, all come back through the turn record
156/// in rule 3.
157pub fn resolve_wait(observation: &WaitObservation, request: &WaitRequest) -> Option<WaitDecision> {
158    let stopping = matches!(
159        observation.lifecycle,
160        Some(ViewerLifecycleCategory::Suspended | ViewerLifecycleCategory::Suspending)
161    ) || matches!(
162        observation.execution,
163        MaterializedExecutionState::Closing | MaterializedExecutionState::Closed
164    );
165    // A resume owns the session: nothing about it has settled yet, and its
166    // durable record still says stopped. The wait keeps waiting; its own
167    // deadline still bounds it.
168    if observation.resuming {
169        return None;
170    }
171    if let Some(reason) = &observation.close_failure {
172        return Some(WaitDecision::simple(
173            WaitOutcome::Error,
174            Some(reason.clone()),
175        ));
176    }
177    // The close owns the session; its own deadline still bounds this wait.
178    if observation.closing {
179        return None;
180    }
181    if stopping {
182        return Some(WaitDecision::simple(
183            WaitOutcome::Stopped,
184            // A resume that failed rolled the record back to stopped and left
185            // its reason there. Reporting it is the difference between "the
186            // session is stopped" and knowing why it did not come up.
187            Some(
188                observation
189                    .launch_error
190                    .clone()
191                    .unwrap_or_else(|| "the session is stopped or stopping".to_owned()),
192            ),
193        ));
194    }
195    if observation.launch_failed {
196        return Some(WaitDecision::simple(
197            WaitOutcome::Error,
198            Some(
199                observation
200                    .launch_error
201                    .clone()
202                    .unwrap_or_else(|| "the session failed to launch".to_owned()),
203            ),
204        ));
205    }
206    if let Some(StartStatus::Failed { message }) = &observation.start_status {
207        return Some(WaitDecision::simple(
208            WaitOutcome::Error,
209            Some(message.clone()),
210        ));
211    }
212    if observation.lifecycle == Some(ViewerLifecycleCategory::Failed) {
213        return Some(WaitDecision::simple(
214            WaitOutcome::Error,
215            // A close that left the session dead recorded why; saying only
216            // that it failed would throw that away.
217            Some(
218                observation
219                    .launch_error
220                    .clone()
221                    .unwrap_or_else(|| "the session is in a failed state".to_owned()),
222            ),
223        ));
224    }
225    let retry_pending = |outcome: &MaterializedTurnOutcome| {
226        observation.quota_recovery.as_ref().is_some_and(|r| {
227            r.retry_at_ms.is_some() && r.completed_command_id == outcome.command_id
228        }) || observation.capacity_retry.is_some()
229            && matches!(
230                &outcome.outcome,
231                TurnOutcomeKind::Completed { stop_reason } if is_capacity_stop_reason(stop_reason)
232            )
233    };
234    if let Some(recovery) = &observation.quota_recovery
235        && recovery.retry_at_ms.is_none()
236        && observation
237            .last_turn_outcome
238            .as_ref()
239            .is_some_and(|t| t.command_id == recovery.completed_command_id)
240        && request.turn_id.is_none_or(|target| {
241            observation
242                .last_turn_outcome
243                .as_ref()
244                .and_then(|t| t.accepted_ordinal)
245                .is_some_and(|a| a >= target)
246        })
247    {
248        return Some(WaitDecision::simple(
249            WaitOutcome::QuotaLimit,
250            Some(recovery.notice.clone()),
251        ));
252    }
253    let target = request.turn_id.or(match &observation.start_status {
254        Some(StartStatus::Submitted { turn_id }) => Some(*turn_id),
255        _ => None,
256    });
257    let target_finished = target.is_some_and(|target| {
258        observation
259            .last_turn_outcome
260            .as_ref()
261            .is_some_and(|outcome| {
262                outcome
263                    .accepted_ordinal
264                    .is_some_and(|ordinal| ordinal >= target)
265                    && !retry_pending(outcome)
266            })
267    });
268    if request.return_on_input && !target_finished && !observation.pending_elicitations.is_empty() {
269        return Some(WaitDecision {
270            outcome: WaitOutcome::InputRequired,
271            stop_reason: None,
272            message: Some("the harness needs a response to a structured input request".into()),
273            turn_id: observation
274                .active_turn
275                .as_ref()
276                .and_then(|turn| turn.accepted_ordinal),
277            turn: None,
278        });
279    }
280    match target {
281        Some(target) => {
282            let outcome = observation.last_turn_outcome.as_ref()?;
283            if outcome
284                .accepted_ordinal
285                .is_none_or(|ordinal| ordinal < target)
286            {
287                return None;
288            }
289            if retry_pending(outcome) {
290                return None;
291            }
292            Some(WaitDecision::from_outcome(outcome))
293        }
294        None => {
295            if observation.checking_continuation
296                || observation.execution != MaterializedExecutionState::Idle
297                || observation.active_turn.is_some()
298                || observation.queued > 0
299            {
300                return None;
301            }
302            match observation.last_turn_outcome.as_ref() {
303                Some(outcome) if retry_pending(outcome) => None,
304                Some(outcome) => Some(WaitDecision::from_outcome(outcome)),
305                // Idle with nothing queued and nothing ever finished: there is
306                // no turn to wait for, so say so immediately rather than block
307                // for the full timeout.
308                None => Some(WaitDecision::simple(WaitOutcome::Finished, None)),
309            }
310        }
311    }
312}
313
314// ---------------------------------------------------------------------------
315// Router
316// ---------------------------------------------------------------------------
317
318#[cfg(test)]
319mod tests {
320    use super::*;
321
322    #[test]
323    fn awaiting_input_is_a_successful_wait_outcome() {
324        assert_eq!(
325            map_stop_reason(mj_core::acp::AWAITING_INPUT_STOP_REASON),
326            (WaitOutcome::InputRequired, None)
327        );
328    }
329}