mj_controller/server/api/
wait_policy.rs1use super::*;
2
3pub 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#[derive(Debug, Clone, Default, PartialEq)]
24pub struct WaitObservation {
25 pub assessment: Option<mj_core::assessment::Summary>,
26 pub activity: mj_core::activity::ActivityState,
28 pub turn_completion: Option<mj_core::activity::verdict::TurnCompletion>,
29 pub background_work: Option<ApiBackgroundWork>,
30 pub pending_elicitations: Vec<mj_core::elicitation::ElicitationRequest>,
31 pub lifecycle: Option<ViewerLifecycleCategory>,
32 pub resuming: bool,
37 pub closing: bool,
40 pub close_failure: Option<String>,
45 pub cannot_take_prompt: bool,
50 pub launch_failed: bool,
52 pub launch_error: Option<String>,
54 pub execution: MaterializedExecutionState,
55 pub active_turn: Option<MaterializedTurn>,
56 pub last_turn_outcome: Option<MaterializedTurnOutcome>,
57 pub queued: usize,
58 pub capacity_retry: Option<CapacityRetry>,
59 pub retry_assessment_pending: bool,
60 pub quota_recovery: Option<mj_core::continuation::QuotaRecovery>,
61 pub start_status: Option<StartStatus>,
62 pub subagent: bool,
65 pub report_pending_for: Option<String>,
69 pub handback: Option<(String, String)>,
71}
72
73impl WaitObservation {
74 fn completed_decision(&self, outcome: &MaterializedTurnOutcome) -> Option<WaitDecision> {
75 use mj_core::activity::verdict::Decision;
76 let mut result = WaitDecision::from_outcome(outcome);
77 if result.outcome != WaitOutcome::Finished {
80 return Some(result);
81 }
82 let completion = self
83 .turn_completion
84 .as_ref()
85 .filter(|completion| completion.command_id == outcome.command_id);
86 if let Some(completion) = completion
87 && matches!(
88 completion.decision,
89 Decision::InferIdle | Decision::AwaitingInput
90 )
91 && let Some(span) = result.turn.as_mut()
92 {
93 span.completed_position = span.completed_position.max(completion.completed_ordinal);
94 }
95 match completion.map(|completion| completion.decision) {
96 Some(Decision::ExpectContinuation) => None,
97 Some(Decision::InferIdle) => Some(result),
98 Some(Decision::AwaitingInput) => {
99 result.outcome = WaitOutcome::InputRequired;
100 result.stop_reason = Some(mj_core::acp::AWAITING_INPUT_STOP_REASON.into());
101 Some(result)
102 }
103 Some(Decision::KeepCurrent) | None => {
104 let later_turn = self.active_turn.as_ref().is_some_and(|active| {
105 active
106 .accepted_ordinal
107 .is_some_and(|ordinal| ordinal > outcome.completed_ordinal)
108 });
109 (self.activity.is_idle() || later_turn).then_some(result)
110 }
111 }
112 }
113
114 pub fn apply_subagent_report(
117 &mut self,
118 handback_tool: bool,
119 report: &mj_core::subagent::SubagentReport,
120 now_ms: i64,
121 ) {
122 self.subagent = true;
123 let Some(turn) = self.last_turn_outcome.as_ref() else {
124 return;
125 };
126 let in_flight = self
127 .active_turn
128 .iter()
129 .map(|turn| turn.command_id.as_str())
130 .collect::<Vec<_>>();
131 match mj_core::subagent::report_state(handback_tool, report, Some(turn), &in_flight, now_ms)
132 {
133 mj_core::subagent::ReportState::Pending { .. } => {
134 self.report_pending_for = Some(turn.command_id.clone());
135 }
136 mj_core::subagent::ReportState::Delivered(message) => {
137 self.handback = Some((turn.command_id.clone(), message));
138 }
139 mj_core::subagent::ReportState::Fallback => {}
140 }
141 }
142}
143
144#[derive(Debug, Clone, Copy, PartialEq, Eq)]
150pub struct TurnSpan {
151 pub start_position: u64,
152 pub completed_position: u64,
153}
154
155#[derive(Debug, Clone, PartialEq, Eq)]
157pub struct WaitDecision {
158 pub outcome: WaitOutcome,
159 pub stop_reason: Option<String>,
160 pub message: Option<String>,
161 pub turn_id: Option<u64>,
162 pub turn: Option<TurnSpan>,
165}
166
167impl WaitDecision {
168 pub(super) fn simple(outcome: WaitOutcome, message: Option<String>) -> Self {
169 Self {
170 outcome,
171 stop_reason: None,
172 message,
173 turn_id: None,
174 turn: None,
175 }
176 }
177
178 pub(super) fn from_outcome(outcome: &MaterializedTurnOutcome) -> Self {
179 use mj_core::event_outcome::{OutcomeReason, TurnResultKind};
180 let result = outcome.result();
181 let kind = match result.kind {
182 TurnResultKind::Completed => WaitOutcome::Finished,
183 TurnResultKind::InputRequired => WaitOutcome::InputRequired,
184 TurnResultKind::Cancelled => WaitOutcome::Cancelled,
185 TurnResultKind::Failed if result.reason == Some(OutcomeReason::QuotaLimit) => {
186 WaitOutcome::QuotaLimit
187 }
188 TurnResultKind::Rejected | TurnResultKind::Interrupted | TurnResultKind::Failed => {
189 WaitOutcome::Error
190 }
191 };
192 let stop_reason = result.stop_reason;
193 let message = result.message;
194 Self {
195 outcome: kind,
196 stop_reason,
197 message,
198 turn_id: outcome.accepted_ordinal,
199 turn: outcome.turn_start_position.map(|start_position| TurnSpan {
200 start_position,
201 completed_position: outcome.completed_ordinal,
202 }),
203 }
204 }
205}
206
207pub fn resolve_wait(observation: &WaitObservation, request: &WaitRequest) -> Option<WaitDecision> {
250 let stopping = matches!(
251 observation.lifecycle,
252 Some(ViewerLifecycleCategory::Suspended | ViewerLifecycleCategory::Suspending)
253 ) || matches!(
254 observation.execution,
255 MaterializedExecutionState::Closing | MaterializedExecutionState::Closed
256 );
257 if observation.resuming {
261 return None;
262 }
263 if let Some(reason) = &observation.close_failure {
264 return Some(WaitDecision::simple(
265 WaitOutcome::Error,
266 Some(reason.clone()),
267 ));
268 }
269 if observation.closing {
271 return None;
272 }
273 if stopping {
274 return Some(WaitDecision::simple(
275 WaitOutcome::Stopped,
276 Some(
280 observation
281 .launch_error
282 .clone()
283 .unwrap_or_else(|| "the session is stopped or stopping".to_owned()),
284 ),
285 ));
286 }
287 if observation.launch_failed {
288 return Some(WaitDecision::simple(
289 WaitOutcome::Error,
290 Some(
291 observation
292 .launch_error
293 .clone()
294 .unwrap_or_else(|| "the session failed to launch".to_owned()),
295 ),
296 ));
297 }
298 if let Some(StartStatus::Failed { message }) = &observation.start_status {
299 return Some(WaitDecision::simple(
300 WaitOutcome::Error,
301 Some(message.clone()),
302 ));
303 }
304 if observation.lifecycle == Some(ViewerLifecycleCategory::Failed) {
305 return Some(WaitDecision::simple(
306 WaitOutcome::Error,
307 Some(
310 observation
311 .launch_error
312 .clone()
313 .unwrap_or_else(|| "the session is in a failed state".to_owned()),
314 ),
315 ));
316 }
317 let retry_pending = |outcome: &MaterializedTurnOutcome| {
318 observation.report_pending_for.as_deref() == Some(outcome.command_id.as_str())
319 || observation.retry_assessment_pending
320 || observation
321 .assessment
322 .as_ref()
323 .is_some_and(|a| a.status == mj_core::assessment::Status::Deferred)
324 || observation.capacity_retry.is_some()
325 || observation.quota_recovery.as_ref().is_some_and(|r| {
326 r.retry_at_ms.is_some() && r.completed_command_id == outcome.command_id
327 })
328 };
329 if let Some(recovery) = &observation.quota_recovery
330 && recovery.retry_at_ms.is_none()
331 && observation
332 .last_turn_outcome
333 .as_ref()
334 .is_some_and(|t| t.command_id == recovery.completed_command_id)
335 && request.turn_id.is_none_or(|target| {
336 observation
337 .last_turn_outcome
338 .as_ref()
339 .and_then(|t| t.accepted_ordinal)
340 .is_some_and(|a| a >= target)
341 })
342 {
343 return Some(WaitDecision::simple(
344 WaitOutcome::QuotaLimit,
345 Some(recovery.notice.clone()),
346 ));
347 }
348 let target = request.turn_id.or(match &observation.start_status {
349 Some(StartStatus::Submitted { turn_id }) => Some(*turn_id),
350 _ => None,
351 });
352 let target_finished = target.is_some_and(|target| {
353 observation
354 .last_turn_outcome
355 .as_ref()
356 .is_some_and(|outcome| {
357 outcome
358 .accepted_ordinal
359 .is_some_and(|ordinal| ordinal >= target)
360 && !retry_pending(outcome)
361 && observation.completed_decision(outcome).is_some()
362 })
363 });
364 if request.return_on_input && !target_finished && !observation.pending_elicitations.is_empty() {
365 return Some(WaitDecision {
366 outcome: WaitOutcome::InputRequired,
367 stop_reason: None,
368 message: Some("the harness needs a response to a structured input request".into()),
369 turn_id: observation
370 .active_turn
371 .as_ref()
372 .and_then(|turn| turn.accepted_ordinal),
373 turn: None,
374 });
375 }
376 if request.turn_id.is_none() && !observation.activity.is_idle() {
380 return None;
381 }
382 match target {
383 Some(target) => {
384 let outcome = observation.last_turn_outcome.as_ref()?;
385 if outcome
386 .accepted_ordinal
387 .is_none_or(|ordinal| ordinal < target)
388 {
389 return None;
390 }
391 if retry_pending(outcome) {
392 return None;
393 }
394 observation.completed_decision(outcome)
395 }
396 None => {
397 if observation.cannot_take_prompt
398 || matches!(observation.start_status, Some(StartStatus::Pending))
399 || observation.execution != MaterializedExecutionState::Idle
400 || observation.active_turn.is_some()
401 || observation.queued > 0
402 {
403 return None;
404 }
405 match observation.last_turn_outcome.as_ref() {
406 Some(outcome) if retry_pending(outcome) => None,
407 Some(outcome) => observation.completed_decision(outcome),
408 None => Some(WaitDecision::simple(WaitOutcome::Finished, None)),
412 }
413 }
414 }
415}
416
417#[cfg(test)]
422mod tests {
423 use super::*;
424
425 #[test]
426 fn awaiting_input_is_a_successful_wait_outcome() {
427 assert_eq!(
428 map_stop_reason(mj_core::acp::AWAITING_INPUT_STOP_REASON),
429 (WaitOutcome::InputRequired, None)
430 );
431 }
432}