返回 CodeWhale
test_cases_03.rs
根目录 / crates / tui / src / core / engine / tests / test_cases_03.rs
1
2
3 #[tokio::test]
4 async fn queued_failed_turn_cancels_older_goal_continuation_without_third_call() {
5 let objective = "stop after the intervening failure";
6 let model = std::sync::Arc::new(FailingGoalModelClient {
7 calls: std::sync::atomic::AtomicUsize::new(0),
8 message: "deterministic queued turn failure".to_string(),
9 });
10 let config = goal_custom_route_config();
11 let client: crate::core::model_client::SharedModelClient = model.clone();
12 let (mut engine, handle) = Engine::new_with_model_client(
13 EngineConfig {
14 model: "local-model".to_string(),
15 snapshots_enabled: false,
16 terminal_chrome_enabled: false,
17 goal_objective: Some(objective.to_string()),
18 ..EngineConfig::default()
19 },
20 &config,
21 client,
22 );
23 let goal_state = engine.config.goal_state.clone();
24
25 handle
26 .send(active_goal_message_op(
27 &config,
28 "queued turn that will fail",
29 objective,
30 None,
31 ))
32 .await
33 .expect("queue failing ordinary turn");
34 engine.schedule_goal_continuation(Vec::new()).await;
35 assert!(engine.has_scheduled_goal_continuation());
36 let run_task = tokio::spawn(engine.run());
37
38 let session = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot())
39 .await
40 .expect("queued failure did not settle")
41 .expect("post-failure session snapshot");
42 assert_eq!(
43 model.calls.load(std::sync::atomic::Ordering::SeqCst),
44 1,
45 "the stale synthetic token must not make a third provider call"
46 );
47 let goal = goal_state.lock().expect("goal lock").snapshot();
48 assert_eq!(goal.status, "blocked");
49 assert!(
50 goal.blocker
51 .as_deref()
52 .is_some_and(|blocker| blocker.contains("deterministic queued turn failure")),
53 "{goal:?}"
54 );
55 let prompt = system_prompt_text(session.system_prompt.expect("blocked system prompt"));
56 assert!(!prompt.contains("<session_goal>"), "{prompt}");
57
58 handle.send(Op::Shutdown).await.expect("shutdown engine");
59 run_task.await.expect("engine task");
60 }
61
62 #[tokio::test]
63 async fn configured_goal_delay_is_cancellable_without_starting_another_turn() {
64 let model = std::sync::Arc::new(FailingGoalModelClient {
65 calls: std::sync::atomic::AtomicUsize::new(0),
66 message: "the delayed provider turn must not start".to_string(),
67 });
68 let config = goal_custom_route_config();
69 let client: crate::core::model_client::SharedModelClient = model.clone();
70 let (mut engine, handle) = Engine::new_with_model_client(
71 EngineConfig {
72 model: "local-model".to_string(),
73 snapshots_enabled: false,
74 terminal_chrome_enabled: false,
75 goal_continuation_delay_seconds: 300,
76 ..EngineConfig::default()
77 },
78 &config,
79 client,
80 );
81
82 engine.schedule_goal_continuation(Vec::new()).await;
83 let run_task = tokio::spawn(engine.run());
84
85 {
86 let mut events = handle.rx_event.write().await;
87 let waiting = tokio::time::timeout(model_turn_event_timeout(), events.recv())
88 .await
89 .expect("missing continuation wait event")
90 .expect("engine event channel closed");
91 assert!(matches!(
92 waiting,
93 Event::GoalContinuationWaiting { delay_seconds: 300 }
94 ));
95 }
96
97 handle.cancel();
98 {
99 let mut events = handle.rx_event.write().await;
100 let ended = tokio::time::timeout(model_turn_event_timeout(), async {
101 loop {
102 if let Some(Event::GoalContinuationWaitEnded { interrupted }) = events.recv().await
103 {
104 break interrupted;
105 }
106 }
107 })
108 .await
109 .expect("cancel did not end the continuation delay");
110 assert!(
111 ended,
112 "the wait receipt must identify an explicit interrupt"
113 );
114 }
115
116 let _ = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot())
117 .await
118 .expect("engine did not accept controls after cancelling the delay")
119 .expect("session snapshot after cancelled delay");
120 assert_eq!(
121 model.calls.load(std::sync::atomic::Ordering::SeqCst),
122 0,
123 "cancelling the quiet period must happen before provider dispatch"
124 );
125
126 handle.send(Op::Shutdown).await.expect("shutdown engine");
127 tokio::time::timeout(model_turn_event_timeout(), run_task)
128 .await
129 .expect("engine did not shut down after delay cancellation")
130 .expect("engine task");
131 }
132
133 #[tokio::test]
134 async fn sync_session_boundary_discards_delayed_goal_and_runtime_mcp_capabilities() {
135 let (mut engine, _handle) = Engine::new(
136 EngineConfig {
137 goal_continuation_delay_seconds: 300,
138 ..EngineConfig::default()
139 },
140 &Config::default(),
141 );
142 engine.session.id = "session-a".to_string();
143 engine.schedule_goal_continuation(Vec::new()).await;
144 assert!(engine.has_scheduled_goal_continuation());
145 engine.ensure_mcp_pool().await.expect("initialize MCP pool");
146 assert!(engine.mcp_pool.is_some());
147
148 assert_eq!(
149 engine.install_synced_session_id("session-a".to_string()),
150 None
151 );
152 assert!(engine.has_scheduled_goal_continuation());
153 assert!(
154 engine.mcp_pool.is_some(),
155 "same-id reload keeps runtime state"
156 );
157
158 assert_eq!(
159 engine.install_synced_session_id("session-b".to_string()),
160 Some("session-a".to_string())
161 );
162 assert!(!engine.has_scheduled_goal_continuation());
163 assert!(
164 engine.mcp_pool.is_none(),
165 "same-workspace B must not inherit A's runtime-added MCP servers"
166 );
167 assert_eq!(
168 engine.install_synced_session_id("session-a".to_string()),
169 Some("session-b".to_string())
170 );
171 assert!(
172 !engine.has_scheduled_goal_continuation() && engine.mcp_pool.is_none(),
173 "A -> B -> A must not resurrect process-local state from the first A"
174 );
175 }
176
177 #[tokio::test]
178 async fn sync_session_boundary_rejects_an_already_enqueued_goal_token() {
179 let (mut engine, _handle) = Engine::new(
180 EngineConfig {
181 goal_continuation_delay_seconds: 0,
182 ..EngineConfig::default()
183 },
184 &Config::default(),
185 );
186 engine.session.id = "session-a".to_string();
187 engine.schedule_goal_continuation(Vec::new()).await;
188 assert!(
189 engine
190 .scheduled_goal_continuation
191 .as_ref()
192 .is_some_and(|scheduled| scheduled.enqueued),
193 "fixture must place A's token in the engine mailbox"
194 );
195
196 engine.install_synced_session_id("session-b".to_string());
197 let input = engine
198 .next_run_input(false)
199 .await
200 .expect("queued continuation token");
201 let EngineRunInput::Operation(op) = input else {
202 panic!("expected queued continuation operation");
203 };
204 let Op::ContinueGoal {
205 dynamic_tools,
206 engine_schedule_id,
207 } = *op
208 else {
209 panic!("expected queued continuation token");
210 };
211 assert!(
212 engine
213 .take_scheduled_goal_continuation(engine_schedule_id, dynamic_tools)
214 .is_none(),
215 "B must reject A's already-enqueued synthetic turn token"
216 );
217 }
218
219 #[tokio::test]
220 async fn cancellation_after_delay_expiry_beats_queued_continuation_dispatch() {
221 let model = std::sync::Arc::new(FailingGoalModelClient {
222 calls: std::sync::atomic::AtomicUsize::new(0),
223 message: "the raced provider turn must not start".to_string(),
224 });
225 let config = goal_custom_route_config();
226 let client: crate::core::model_client::SharedModelClient = model.clone();
227 let (mut engine, handle) = Engine::new_with_model_client(
228 EngineConfig {
229 model: "local-model".to_string(),
230 snapshots_enabled: false,
231 terminal_chrome_enabled: false,
232 goal_continuation_delay_seconds: 300,
233 ..EngineConfig::default()
234 },
235 &config,
236 client,
237 );
238
239 engine.schedule_goal_continuation(Vec::new()).await;
240 // Deterministically place the fixture at the timer/mailbox boundary: the
241 // quiet period expired and its one coalesced token is already queued, but
242 // the engine has not consumed it yet.
243 engine
244 .scheduled_goal_continuation
245 .as_mut()
246 .expect("scheduled continuation")
247 .ready_at = None;
248 engine.try_flush_pending_goal_continuation();
249 assert!(
250 engine
251 .scheduled_goal_continuation
252 .as_ref()
253 .is_some_and(|scheduled| scheduled.enqueued)
254 );
255 handle.cancel();
256 let run_task = tokio::spawn(engine.run());
257
258 let _ = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot())
259 .await
260 .expect("engine did not settle the delay-expiry cancellation race")
261 .expect("session snapshot after expiry race");
262 assert_eq!(
263 model.calls.load(std::sync::atomic::Ordering::SeqCst),
264 0,
265 "cancelled delayed token must be discarded before provider dispatch"
266 );
267
268 handle.send(Op::Shutdown).await.expect("shutdown engine");
269 tokio::time::timeout(model_turn_event_timeout(), run_task)
270 .await
271 .expect("engine did not shut down after expiry race")
272 .expect("engine task");
273 }
274
275 #[tokio::test]
276 async fn configured_goal_delay_expires_into_exactly_one_continuation() {
277 let model = std::sync::Arc::new(FailingGoalModelClient {
278 calls: std::sync::atomic::AtomicUsize::new(0),
279 message: "stop after proving delayed dispatch".to_string(),
280 });
281 let config = goal_custom_route_config();
282 let client: crate::core::model_client::SharedModelClient = model.clone();
283 let (mut engine, handle) = Engine::new_with_model_client(
284 EngineConfig {
285 model: "local-model".to_string(),
286 snapshots_enabled: false,
287 terminal_chrome_enabled: false,
288 goal_objective: Some("dispatch once after the cadence".to_string()),
289 goal_continuation_delay_seconds: 1,
290 ..EngineConfig::default()
291 },
292 &config,
293 client,
294 );
295 engine
296 .config
297 .goal_state
298 .lock()
299 .expect("goal lock")
300 .sync_from_host_status(
301 Some("dispatch once after the cadence"),
302 None,
303 crate::tools::goal::GoalStatus::Active,
304 );
305 engine.schedule_goal_continuation(Vec::new()).await;
306 let run_task = tokio::spawn(engine.run());
307
308 let (mut saw_waiting, mut saw_ready, mut saw_started) = (false, false, false);
309 {
310 let mut events = handle.rx_event.write().await;
311 tokio::time::timeout(Duration::from_secs(3), async {
312 while let Some(event) = events.recv().await {
313 match event {
314 Event::GoalContinuationWaiting { delay_seconds: 1 } => saw_waiting = true,
315 Event::GoalContinuationWaitEnded { interrupted: false } => saw_ready = true,
316 Event::TurnStarted { .. } => {
317 saw_started = true;
318 break;
319 }
320 _ => {}
321 }
322 }
323 })
324 .await
325 .expect("configured delay did not dispatch its continuation");
326 }
327 assert!(saw_waiting && saw_ready && saw_started);
328
329 let _ = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot())
330 .await
331 .expect("delayed failing turn did not settle")
332 .expect("session snapshot after delayed dispatch");
333 assert_eq!(
334 model.calls.load(std::sync::atomic::Ordering::SeqCst),
335 1,
336 "one delayed schedule must create exactly one provider request"
337 );
338
339 handle.send(Op::Shutdown).await.expect("shutdown engine");
340 tokio::time::timeout(model_turn_event_timeout(), run_task)
341 .await
342 .expect("engine did not shut down after delayed dispatch")
343 .expect("engine task");
344 }
345
346 #[tokio::test]
347 async fn goal_pause_during_configured_delay_cancels_pending_continuation() {
348 let model = std::sync::Arc::new(FailingGoalModelClient {
349 calls: std::sync::atomic::AtomicUsize::new(0),
350 message: "the paused provider turn must not start".to_string(),
351 });
352 let config = goal_custom_route_config();
353 let client: crate::core::model_client::SharedModelClient = model.clone();
354 let (mut engine, handle) = Engine::new_with_model_client(
355 EngineConfig {
356 model: "local-model".to_string(),
357 snapshots_enabled: false,
358 terminal_chrome_enabled: false,
359 goal_objective: Some("coordinate until paused".to_string()),
360 goal_continuation_delay_seconds: 300,
361 ..EngineConfig::default()
362 },
363 &config,
364 client,
365 );
366 let goal_state = engine.config.goal_state.clone();
367 goal_state.lock().expect("goal lock").sync_from_host_status(
368 Some("coordinate until paused"),
369 None,
370 crate::tools::goal::GoalStatus::Active,
371 );
372 engine.schedule_goal_continuation(Vec::new()).await;
373 let run_task = tokio::spawn(engine.run());
374
375 handle
376 .send(Op::SetGoalStatus {
377 goal_id: None,
378 status: crate::tools::goal::GoalStatus::Paused,
379 clear: false,
380 })
381 .await
382 .expect("pause delayed goal");
383 {
384 let mut events = handle.rx_event.write().await;
385 let interrupted = tokio::time::timeout(model_turn_event_timeout(), async {
386 loop {
387 if let Some(Event::GoalContinuationWaitEnded { interrupted }) = events.recv().await
388 {
389 break interrupted;
390 }
391 }
392 })
393 .await
394 .expect("goal pause did not end the continuation delay");
395 assert!(
396 interrupted,
397 "a pause is an explicit interruption, not a ready-to-run receipt"
398 );
399 }
400 let _ = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot())
401 .await
402 .expect("goal pause did not settle during delay")
403 .expect("session snapshot after goal pause");
404
405 assert_eq!(
406 goal_state.lock().expect("goal lock").snapshot().status,
407 "paused"
408 );
409 assert_eq!(
410 model.calls.load(std::sync::atomic::Ordering::SeqCst),
411 0,
412 "a goal status control must beat the delayed continuation"
413 );
414
415 handle.send(Op::Shutdown).await.expect("shutdown engine");
416 tokio::time::timeout(model_turn_event_timeout(), run_task)
417 .await
418 .expect("engine did not shut down after pausing delayed goal")
419 .expect("engine task");
420 }
421
422 #[tokio::test]
423 async fn host_injected_goal_continuation_waits_out_the_quiet_period() {
424 // Host-managed sessions never run the engine-owned scheduler
425 // (schedule_goal_continuation is gated on !host_managed_turns), so the
426 // host injects ContinueGoal with engine_schedule_id None and this arm is
427 // the only dispatch site. A positive configured delay must be awaited
428 // before the provider request starts.
429 let model = std::sync::Arc::new(FailingGoalModelClient {
430 calls: std::sync::atomic::AtomicUsize::new(0),
431 message: "the host-injected continuation must wait first".to_string(),
432 });
433 let config = goal_custom_route_config();
434 let client: crate::core::model_client::SharedModelClient = model.clone();
435 let (engine, handle) = Engine::new_with_model_client(
436 EngineConfig {
437 model: "local-model".to_string(),
438 snapshots_enabled: false,
439 terminal_chrome_enabled: false,
440 goal_objective: Some("wait out the cadence".to_string()),
441 goal_continuation_delay_seconds: 1,
442 ..EngineConfig::default()
443 },
444 &config,
445 client,
446 );
447 engine
448 .config
449 .goal_state
450 .lock()
451 .expect("goal lock")
452 .sync_from_host_status(
453 Some("wait out the cadence"),
454 None,
455 crate::tools::goal::GoalStatus::Active,
456 );
457 let run_task = tokio::spawn(engine.run());
458
459 let queued_at = Instant::now();
460 handle
461 .send(Op::ContinueGoal {
462 dynamic_tools: Vec::new(),
463 engine_schedule_id: None,
464 })
465 .await
466 .expect("queue host-injected continuation");
467 {
468 let mut events = handle.rx_event.write().await;
469 tokio::time::timeout(Duration::from_secs(3), async {
470 while let Some(event) = events.recv().await {
471 if matches!(event, Event::TurnStarted { .. }) {
472 break;
473 }
474 }
475 })
476 .await
477 .expect("host-injected continuation did not dispatch after the quiet period");
478 }
479 let elapsed = queued_at.elapsed();
480 assert!(
481 elapsed >= Duration::from_millis(900),
482 "dispatch must wait out the configured quiet period, dispatched after {elapsed:?}"
483 );
484
485 let _ = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot())
486 .await
487 .expect("host-injected delayed turn did not settle")
488 .expect("session snapshot after delayed host-injected dispatch");
489 assert_eq!(
490 model.calls.load(std::sync::atomic::Ordering::SeqCst),
491 1,
492 "the host-injected token must dispatch exactly one provider request"
493 );
494
495 handle.send(Op::Shutdown).await.expect("shutdown engine");
496 tokio::time::timeout(model_turn_event_timeout(), run_task)
497 .await
498 .expect("engine did not shut down after host-injected delay")
499 .expect("engine task");
500 }
501
502 #[tokio::test]
503 async fn host_injected_goal_continuation_with_zero_delay_dispatches_immediately() {
504 let model = std::sync::Arc::new(FailingGoalModelClient {
505 calls: std::sync::atomic::AtomicUsize::new(0),
506 message: "zero-delay host-injected dispatch proof".to_string(),
507 });
508 let config = goal_custom_route_config();
509 let client: crate::core::model_client::SharedModelClient = model.clone();
510 let (engine, handle) = Engine::new_with_model_client(
511 EngineConfig {
512 model: "local-model".to_string(),
513 snapshots_enabled: false,
514 terminal_chrome_enabled: false,
515 goal_objective: Some("dispatch without a cadence".to_string()),
516 goal_continuation_delay_seconds: 0,
517 ..EngineConfig::default()
518 },
519 &config,
520 client,
521 );
522 engine
523 .config
524 .goal_state
525 .lock()
526 .expect("goal lock")
527 .sync_from_host_status(
528 Some("dispatch without a cadence"),
529 None,
530 crate::tools::goal::GoalStatus::Active,
531 );
532 let run_task = tokio::spawn(engine.run());
533
534 handle
535 .send(Op::ContinueGoal {
536 dynamic_tools: Vec::new(),
537 engine_schedule_id: None,
538 })
539 .await
540 .expect("queue zero-delay host-injected continuation");
541 {
542 let mut events = handle.rx_event.write().await;
543 tokio::time::timeout(Duration::from_secs(1), async {
544 while let Some(event) = events.recv().await {
545 if matches!(event, Event::TurnStarted { .. }) {
546 break;
547 }
548 }
549 })
550 .await
551 .expect("a zero-delay host-injected continuation must dispatch immediately");
552 }
553
554 let _ = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot())
555 .await
556 .expect("zero-delay host-injected turn did not settle")
557 .expect("session snapshot after zero-delay dispatch");
558 assert_eq!(
559 model.calls.load(std::sync::atomic::Ordering::SeqCst),
560 1,
561 "zero delay must not suppress the host-injected dispatch"
562 );
563
564 handle.send(Op::Shutdown).await.expect("shutdown engine");
565 tokio::time::timeout(model_turn_event_timeout(), run_task)
566 .await
567 .expect("engine did not shut down after zero-delay dispatch")
568 .expect("engine task");
569 }
570
571 #[tokio::test]
572 async fn cancellation_during_host_injected_continuation_wait_never_dispatches() {
573 let model = std::sync::Arc::new(FailingGoalModelClient {
574 calls: std::sync::atomic::AtomicUsize::new(0),
575 message: "a cancelled host-injected wait must never start".to_string(),
576 });
577 let config = goal_custom_route_config();
578 let client: crate::core::model_client::SharedModelClient = model.clone();
579 let (engine, handle) = Engine::new_with_model_client(
580 EngineConfig {
581 model: "local-model".to_string(),
582 snapshots_enabled: false,
583 terminal_chrome_enabled: false,
584 goal_objective: Some("cancel mid-cadence".to_string()),
585 goal_continuation_delay_seconds: 1,
586 ..EngineConfig::default()
587 },
588 &config,
589 client,
590 );
591 engine
592 .config
593 .goal_state
594 .lock()
595 .expect("goal lock")
596 .sync_from_host_status(
597 Some("cancel mid-cadence"),
598 None,
599 crate::tools::goal::GoalStatus::Active,
600 );
601 let run_task = tokio::spawn(engine.run());
602
603 handle
604 .send(Op::ContinueGoal {
605 dynamic_tools: Vec::new(),
606 engine_schedule_id: None,
607 })
608 .await
609 .expect("queue host-injected continuation");
610 // Enter the quiet period, then cancel: the biased wait must drop the
611 // pending pass instead of dispatching when the period would have elapsed.
612 tokio::time::sleep(Duration::from_millis(100)).await;
613 handle.cancel();
614
615 // A dispatch bug would surface a TurnStarted once the 1s period expires;
616 // a correct cancellation settles back into the mailbox loop silently.
617 let dispatched = {
618 let mut events = handle.rx_event.write().await;
619 tokio::time::timeout(Duration::from_secs(2), async {
620 loop {
621 let Some(event) = events.recv().await else {
622 break false;
623 };
624 if matches!(event, Event::TurnStarted { .. }) {
625 break true;
626 }
627 }
628 })
629 .await
630 .unwrap_or(false)
631 };
632 assert!(
633 !dispatched,
634 "cancellation during the host-injected quiet period must never dispatch"
635 );
636
637 // The engine must keep accepting controls after the cancelled wait.
638 let _ = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot())
639 .await
640 .expect("engine did not accept controls after the cancelled wait")
641 .expect("session snapshot after cancelled wait");
642 assert_eq!(
643 model.calls.load(std::sync::atomic::Ordering::SeqCst),
644 0,
645 "the cancelled host-injected token must never reach the provider"
646 );
647
648 handle.send(Op::Shutdown).await.expect("shutdown engine");
649 tokio::time::timeout(model_turn_event_timeout(), run_task)
650 .await
651 .expect("engine did not shut down after cancelled wait")
652 .expect("engine task");
653 }
654
655 #[tokio::test]
656 async fn queued_not_started_turn_cancels_older_goal_continuation() {
657 let objective = "stop when the queued turn cannot start";
658 let model = std::sync::Arc::new(FailingGoalModelClient {
659 calls: std::sync::atomic::AtomicUsize::new(0),
660 message: "model must never be called".to_string(),
661 });
662 let config = goal_custom_route_config();
663 let client: crate::core::model_client::SharedModelClient = model.clone();
664 let (mut engine, handle) = Engine::new_with_model_client(
665 EngineConfig {
666 model: "local-model".to_string(),
667 snapshots_enabled: false,
668 terminal_chrome_enabled: false,
669 goal_objective: Some(objective.to_string()),
670 ..EngineConfig::default()
671 },
672 &config,
673 client,
674 );
675 let goal_state = engine.config.goal_state.clone();
676 engine.model_client = None;
677 engine.codewhale_client_error = Some("deterministic missing model client".to_string());
678
679 handle
680 .send(active_goal_message_op(
681 &config,
682 "queued turn that cannot start",
683 objective,
684 None,
685 ))
686 .await
687 .expect("queue not-started ordinary turn");
688 engine.schedule_goal_continuation(Vec::new()).await;
689 let run_task = tokio::spawn(engine.run());
690
691 let session = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot())
692 .await
693 .expect("not-started turn did not settle")
694 .expect("post-rejection session snapshot");
695 assert_eq!(
696 model.calls.load(std::sync::atomic::Ordering::SeqCst),
697 0,
698 "neither the rejected turn nor its stale token may call the provider"
699 );
700 let goal = goal_state.lock().expect("goal lock").snapshot();
701 assert_eq!(goal.status, "blocked");
702 assert!(
703 goal.blocker
704 .as_deref()
705 .is_some_and(|blocker| blocker.contains("could not be started")),
706 "{goal:?}"
707 );
708 let prompt = system_prompt_text(session.system_prompt.expect("blocked system prompt"));
709 assert!(!prompt.contains("<session_goal>"), "{prompt}");
710
711 handle.send(Op::Shutdown).await.expect("shutdown engine");
712 run_task.await.expect("engine task");
713 }
714
715 #[tokio::test]
716 async fn queued_interrupted_turn_cancels_older_goal_continuation_without_third_call() {
717 let objective = "pause after the intervening cancellation";
718 let request_entered = std::sync::Arc::new(tokio::sync::Notify::new());
719 let release_request = std::sync::Arc::new(tokio::sync::Notify::new());
720 let model = std::sync::Arc::new(FirstRequestGatedGoalModelClient {
721 calls: std::sync::atomic::AtomicUsize::new(0),
722 request_entered: std::sync::Arc::clone(&request_entered),
723 release_request,
724 });
725 let config = goal_custom_route_config();
726 let client: crate::core::model_client::SharedModelClient = model.clone();
727 let (mut engine, handle) = Engine::new_with_model_client(
728 EngineConfig {
729 model: "local-model".to_string(),
730 snapshots_enabled: false,
731 terminal_chrome_enabled: false,
732 goal_objective: Some(objective.to_string()),
733 ..EngineConfig::default()
734 },
735 &config,
736 client,
737 );
738 let goal_state = engine.config.goal_state.clone();
739
740 handle
741 .send(active_goal_message_op(
742 &config,
743 "queued turn that will be cancelled",
744 objective,
745 None,
746 ))
747 .await
748 .expect("queue interruptible ordinary turn");
749 engine.schedule_goal_continuation(Vec::new()).await;
750 let run_task = tokio::spawn(engine.run());
751 tokio::time::timeout(model_turn_event_timeout(), request_entered.notified())
752 .await
753 .expect("queued request was never entered");
754 handle.cancel();
755
756 let session = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot())
757 .await
758 .expect("interrupted turn did not settle")
759 .expect("post-interruption session snapshot");
760 assert_eq!(
761 model.calls.load(std::sync::atomic::Ordering::SeqCst),
762 1,
763 "the stale synthetic token must not make a third provider call"
764 );
765 let goal = goal_state.lock().expect("goal lock").snapshot();
766 // Interrupted ordinary turns cancel stale auto-continuation only; the goal
767 // stays Active so the next user message continues without /goal resume.
768 assert_eq!(goal.status, "active");
769 assert_eq!(goal.blocker, None);
770 assert_eq!(goal.pause_reason, None);
771 let prompt = system_prompt_text(session.system_prompt.expect("active system prompt"));
772 assert!(prompt.contains("<session_goal>"), "{prompt}");
773
774 handle.send(Op::Shutdown).await.expect("shutdown engine");
775 run_task.await.expect("engine task");
776 }
777
778 #[tokio::test]
779 async fn initial_goal_failure_projects_blocked_state() {
780 let objective = "block the initial failed goal turn";
781 let leaked_secret = "sk-initial-goal-secret-123456";
782 let model = std::sync::Arc::new(FailingGoalModelClient {
783 calls: std::sync::atomic::AtomicUsize::new(0),
784 message: format!("initial provider failure: {leaked_secret}"),
785 });
786 let config = goal_custom_route_config();
787 let client: crate::core::model_client::SharedModelClient = model.clone();
788 let (engine, handle) = Engine::new_with_model_client(
789 EngineConfig {
790 model: "local-model".to_string(),
791 snapshots_enabled: false,
792 terminal_chrome_enabled: false,
793 ..EngineConfig::default()
794 },
795 &config,
796 client,
797 );
798 let goal_state = engine.config.goal_state.clone();
799 let run_task = tokio::spawn(engine.run());
800
801 handle
802 .send(active_goal_message_op(
803 &config,
804 "start a goal whose first turn fails",
805 objective,
806 None,
807 ))
808 .await
809 .expect("send initial goal turn");
810 let session = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot())
811 .await
812 .expect("initial goal failure did not settle")
813 .expect("post-failure session snapshot");
814
815 assert_eq!(model.calls.load(std::sync::atomic::Ordering::SeqCst), 1);
816 let goal = goal_state.lock().expect("goal lock").snapshot();
817 assert_eq!(goal.objective.as_deref(), Some(objective));
818 assert_eq!(goal.status, "blocked");
819 let blocker = goal.blocker.as_deref().expect("failure blocker");
820 assert!(blocker.contains("initial provider failure"), "{blocker}");
821 assert!(!blocker.contains(leaked_secret), "{blocker}");
822 let prompt = system_prompt_text(session.system_prompt.expect("blocked system prompt"));
823 assert!(!prompt.contains("<session_goal>"), "{prompt}");
824
825 handle.send(Op::Shutdown).await.expect("shutdown engine");
826 run_task.await.expect("engine task");
827 }
828
829 /// A goal the runtime stopped (its turn failed) resumes when the person
830 /// writes again; the host still reports Blocked because it only learns of
831 /// the resume from this turn's GoalUpdated.
832 #[tokio::test]
833 async fn user_message_resumes_a_goal_only_the_runtime_blocked() {
834 let objective = "resume after a runtime stop";
835 let model = std::sync::Arc::new(FailingGoalModelClient {
836 calls: std::sync::atomic::AtomicUsize::new(0),
837 message: "turn deadline elapsed".to_string(),
838 });
839 let config = goal_custom_route_config();
840 let client: crate::core::model_client::SharedModelClient = model.clone();
841 let (engine, handle) = Engine::new_with_model_client(
842 EngineConfig {
843 model: "local-model".to_string(),
844 snapshots_enabled: false,
845 terminal_chrome_enabled: false,
846 ..EngineConfig::default()
847 },
848 &config,
849 client,
850 );
851 let goal_state = engine.config.goal_state.clone();
852 let run_task = tokio::spawn(engine.run());
853 let settle = || async {
854 tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot())
855 .await
856 .expect("turn did not settle")
857 .expect("session snapshot")
858 };
859
860 handle
861 .send(active_goal_message_op(&config, "start", objective, None))
862 .await
863 .expect("send goal turn");
864 settle().await;
865 let blocked = goal_state.lock().expect("goal lock").snapshot();
866 assert_eq!(blocked.status, "blocked");
867
868 let Op::SendMessage(mut spec) = active_goal_message_op(&config, "continue", objective, None)
869 else {
870 unreachable!()
871 };
872 spec.goal_status = crate::tools::goal::GoalStatus::Blocked;
873 handle
874 .send(Op::SendMessage(spec))
875 .await
876 .expect("send continue");
877 settle().await;
878 assert_eq!(model.calls.load(std::sync::atomic::Ordering::SeqCst), 2);
879 let resumed = goal_state.lock().expect("goal lock").snapshot();
880 assert_ne!(
881 resumed.goal_id, blocked.goal_id,
882 "the continue turn ran as a resumed goal revision"
883 );
884
885 // A blocker the model reported is a judgement: the next message is an
886 // ordinary turn and the goal stays blocked on that report.
887 goal_state
888 .lock()
889 .expect("goal lock")
890 .mark_blocked("needs the staging credentials".to_string())
891 .unwrap();
892 let reported = goal_state.lock().expect("goal lock").snapshot();
893 let Op::SendMessage(mut spec) = active_goal_message_op(&config, "continue", objective, None)
894 else {
895 unreachable!()
896 };
897 spec.goal_status = crate::tools::goal::GoalStatus::Blocked;
898 handle
899 .send(Op::SendMessage(spec))
900 .await
901 .expect("send ordinary turn");
902 settle().await;
903 let after = goal_state.lock().expect("goal lock").snapshot();
904 assert_eq!(after.status, "blocked");
905 assert_eq!(after.goal_id, reported.goal_id);
906 assert_eq!(
907 after.blocker.as_deref(),
908 Some("needs the staging credentials")
909 );
910
911 handle.send(Op::Shutdown).await.expect("shutdown engine");
912 run_task.await.expect("engine task");
913 }
914
915 #[tokio::test]
916 async fn initial_goal_interruption_keeps_goal_active() {
917 let objective = "keep goal active after interrupted turn";
918 let request_entered = std::sync::Arc::new(tokio::sync::Notify::new());
919 let release_request = std::sync::Arc::new(tokio::sync::Notify::new());
920 let model = std::sync::Arc::new(FirstRequestGatedGoalModelClient {
921 calls: std::sync::atomic::AtomicUsize::new(0),
922 request_entered: std::sync::Arc::clone(&request_entered),
923 release_request,
924 });
925 let config = goal_custom_route_config();
926 let client: crate::core::model_client::SharedModelClient = model.clone();
927 let (engine, handle) = Engine::new_with_model_client(
928 EngineConfig {
929 model: "local-model".to_string(),
930 snapshots_enabled: false,
931 terminal_chrome_enabled: false,
932 ..EngineConfig::default()
933 },
934 &config,
935 client,
936 );
937 let goal_state = engine.config.goal_state.clone();
938 let run_task = tokio::spawn(engine.run());
939
940 handle
941 .send(active_goal_message_op(
942 &config,
943 "start a goal whose first turn is cancelled",
944 objective,
945 None,
946 ))
947 .await
948 .expect("send initial goal turn");
949 tokio::time::timeout(model_turn_event_timeout(), request_entered.notified())
950 .await
951 .expect("initial request was never entered");
952 handle.cancel();
953
954 let session = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot())
955 .await
956 .expect("initial interruption did not settle")
957 .expect("post-interruption session snapshot");
958 assert_eq!(model.calls.load(std::sync::atomic::Ordering::SeqCst), 1);
959 let goal = goal_state.lock().expect("goal lock").snapshot();
960 assert_eq!(goal.objective.as_deref(), Some(objective));
961 assert_eq!(goal.status, "active");
962 assert_eq!(goal.blocker, None);
963 assert_eq!(goal.pause_reason, None);
964 let prompt = system_prompt_text(session.system_prompt.expect("active system prompt"));
965 // Durable goals stay in the prompt after interrupt so the next turn continues.
966 assert!(prompt.contains("<session_goal>"), "{prompt}");
967 assert!(prompt.contains(objective), "{prompt}");
968
969 handle.send(Op::Shutdown).await.expect("shutdown engine");
970 run_task.await.expect("engine task");
971 }
972
973 #[tokio::test]
974 async fn saturated_goal_controls_run_before_ready_idle_child_completion() {
975 use crate::tools::subagent::SubAgentCompletion;
976
977 let stale_tool = DynamicToolSpec {
978 namespace: Some("goal-regression".to_string()),
979 name: "stale".to_string(),
980 description: "stale tool catalog".to_string(),
981 input_schema: json!({"type": "object"}),
982 defer_loading: false,
983 };
984 let fresh_tool = DynamicToolSpec {
985 name: "fresh".to_string(),
986 description: "fresh tool catalog".to_string(),
987 ..stale_tool.clone()
988 };
989 let config = goal_custom_route_config();
990 let (mut engine, handle) = Engine::new(
991 EngineConfig {
992 snapshots_enabled: false,
993 terminal_chrome_enabled: false,
994 ..EngineConfig::default()
995 },
996 &config,
997 );
998
999 for index in 0..ENGINE_OP_CHANNEL_CAPACITY {
1000 let status = if index + 1 == ENGINE_OP_CHANNEL_CAPACITY {
1001 crate::tools::goal::GoalStatus::Paused
1002 } else {
1003 crate::tools::goal::GoalStatus::Active
1004 };
1005 handle
1006 .tx_op
1007 .try_send(Op::SetGoalStatus {
1008 goal_id: None,
1009 status,
1010 clear: false,
1011 })
1012 .unwrap_or_else(|error| panic!("fill ordering mailbox slot {index}: {error}"));
1013 }
1014 assert_eq!(handle.tx_op.capacity(), 0, "fixture must saturate mailbox");
1015 engine
1016 .tx_subagent_completion
1017 .try_send(SubAgentCompletion {
1018 owner_session_id: engine.session.id.clone(),
1019 agent_id: "agent_ready_during_backpressure".to_string(),
1020 payload: "ready child completion".to_string(),
1021 })
1022 .expect("queue ready idle child completion");
1023 engine.schedule_goal_continuation(vec![stale_tool]).await;
1024 assert!(
1025 engine.has_scheduled_goal_continuation(),
1026 "a live schedule must activate temporary op priority"
1027 );
1028
1029 // Inspect the exact production receive helper without running handlers:
1030 // all controls that filled the mailbox, including the final pause, must be
1031 // selected before the already-ready idle child completion.
1032 for index in 0..ENGINE_OP_CHANNEL_CAPACITY {
1033 let input = tokio::time::timeout(model_turn_event_timeout(), engine.next_run_input(false))
1034 .await
1035 .expect("backpressured operation receive timed out")
1036 .expect("engine input");
1037 let EngineRunInput::Operation(op) = input else {
1038 panic!("idle child completion beat queued control {index}");
1039 };
1040 let Op::SetGoalStatus { status, clear, .. } = *op else {
1041 panic!("unexpected operation before queued control {index}");
1042 };
1043 assert!(!clear);
1044 let expected = if index + 1 == ENGINE_OP_CHANNEL_CAPACITY {
1045 crate::tools::goal::GoalStatus::Paused
1046 } else {
1047 crate::tools::goal::GoalStatus::Active
1048 };
1049 assert_eq!(status, expected);
1050 if index == 0 {
1051 // Refresh after capacity opens. The existing token must retain its
1052 // FIFO position behind the remaining controls while carrying the
1053 // newest runtime tool catalog when it is eventually consumed.
1054 engine
1055 .schedule_goal_continuation(vec![fresh_tool.clone()])
1056 .await;
1057 }
1058 }
1059
1060 let token = engine
1061 .next_run_input(false)
1062 .await
1063 .expect("backpressured continuation token");
1064 let EngineRunInput::Operation(token) = token else {
1065 panic!("idle child completion beat the backpressured continuation token");
1066 };
1067 let Op::ContinueGoal {
1068 dynamic_tools,
1069 engine_schedule_id,
1070 } = *token
1071 else {
1072 panic!("expected engine-owned continuation token");
1073 };
1074 assert!(dynamic_tools.is_empty());
1075 let continued_tools = engine
1076 .take_scheduled_goal_continuation(engine_schedule_id, dynamic_tools)
1077 .expect("engine-owned continuation token must consume its schedule marker");
1078 assert_eq!(continued_tools, vec![fresh_tool]);
1079 assert!(!engine.has_scheduled_goal_continuation());
1080
1081 let child = engine
1082 .next_run_input(false)
1083 .await
1084 .expect("ready idle child completion");
1085 let EngineRunInput::SubAgentCompletion(child) = child else {
1086 panic!("unexpected operation after backpressure drain");
1087 };
1088 assert_eq!(child.agent_id, "agent_ready_during_backpressure");
1089 }
1090
1091 #[tokio::test]
1092 async fn unsaturated_goal_control_runs_before_ready_idle_child_completion() {
1093 use crate::tools::subagent::SubAgentCompletion;
1094
1095 let config = goal_custom_route_config();
1096 let (mut engine, handle) = Engine::new(
1097 EngineConfig {
1098 snapshots_enabled: false,
1099 terminal_chrome_enabled: false,
1100 ..EngineConfig::default()
1101 },
1102 &config,
1103 );
1104
1105 handle
1106 .tx_op
1107 .try_send(Op::SetGoalStatus {
1108 goal_id: None,
1109 status: crate::tools::goal::GoalStatus::Paused,
1110 clear: false,
1111 })
1112 .expect("queue unsaturated pause");
1113 engine
1114 .tx_subagent_completion
1115 .try_send(SubAgentCompletion {
1116 owner_session_id: engine.session.id.clone(),
1117 agent_id: "agent_ready_without_backpressure".to_string(),
1118 payload: "ready child completion".to_string(),
1119 })
1120 .expect("queue ready idle child completion");
1121 engine.schedule_goal_continuation(Vec::new()).await;
1122 assert!(engine.has_scheduled_goal_continuation());
1123 assert!(
1124 handle.tx_op.capacity() > 0,
1125 "fixture must leave the mailbox unsaturated"
1126 );
1127
1128 let first = engine
1129 .next_run_input(false)
1130 .await
1131 .expect("queued pause must be selected");
1132 let EngineRunInput::Operation(first) = first else {
1133 panic!("ready child completion beat an unsaturated queued pause");
1134 };
1135 assert!(matches!(
1136 *first,
1137 Op::SetGoalStatus {
1138 goal_id: None,
1139 status: crate::tools::goal::GoalStatus::Paused,
1140 clear: false
1141 }
1142 ));
1143
1144 let token = engine
1145 .next_run_input(false)
1146 .await
1147 .expect("scheduled continuation token");
1148 let EngineRunInput::Operation(token) = token else {
1149 panic!("ready child completion beat the live continuation token");
1150 };
1151 let Op::ContinueGoal {
1152 dynamic_tools,
1153 engine_schedule_id,
1154 } = *token
1155 else {
1156 panic!("expected continuation token behind pause");
1157 };
1158 engine
1159 .take_scheduled_goal_continuation(engine_schedule_id, dynamic_tools)
1160 .expect("consume live continuation schedule");
1161 assert!(!engine.has_scheduled_goal_continuation());
1162
1163 let child = engine
1164 .next_run_input(false)
1165 .await
1166 .expect("idle child completion after schedule consumption");
1167 let EngineRunInput::SubAgentCompletion(child) = child else {
1168 panic!("normal child fairness did not resume after schedule consumption");
1169 };
1170 assert_eq!(child.agent_id, "agent_ready_without_backpressure");
1171 }
1172
1173 #[tokio::test]
1174 async fn cross_turn_token_budget_exhaustion_does_not_pause_goal() {
1175 use crate::llm_client::mock::{MockLlmClient, canned};
1176
1177 let budget_turn = vec![
1178 canned::message_start("mock_goal_budget"),
1179 canned::text_block_start(0),
1180 canned::text_delta(0, "budget spent"),
1181 canned::block_stop(0),
1182 canned::message_delta(
1183 "end_turn",
1184 Some(Usage {
1185 input_tokens: 8,
1186 output_tokens: 3,
1187 ..Usage::default()
1188 }),
1189 ),
1190 canned::message_stop(),
1191 ];
1192 let model = std::sync::Arc::new(MockLlmClient::new(vec![
1193 budget_turn,
1194 canned::simple_text_turn("the cross-turn continuation runs past the exhausted budget"),
1195 ]));
1196 let client: crate::core::model_client::SharedModelClient = model.clone();
1197 let config = Config::default();
1198 let engine_config = EngineConfig {
1199 max_steps: 1,
1200 snapshots_enabled: false,
1201 subagents_enabled: false,
1202 terminal_chrome_enabled: false,
1203 goal_objective: Some("finish within budget".to_string()),
1204 goal_token_budget: Some(10),
1205 ..EngineConfig::default()
1206 };
1207 let (engine, handle) = Engine::new_with_model_client(engine_config, &config, client);
1208 let goal_state = engine.config.goal_state.clone();
1209 let run_task = tokio::spawn(engine.run());
1210
1211 handle
1212 .send(Op::SendMessage(TurnSpec {
1213 max_output_tokens: None,
1214 content: "start budgeted goal".to_string(),
1215 images: Vec::new(),
1216 mode: AppMode::Agent,
1217 route: resolved_route_for_test(&config, crate::config::DEFAULT_TEXT_MODEL),
1218 compaction: Box::new(CompactionConfig::default()),
1219 initial_routed_usage: Box::default(),
1220 goal_objective: Some("finish within budget".to_string()),
1221 goal_token_budget: Some(10),
1222 goal_status: crate::tools::goal::GoalStatus::Active,
1223 reasoning_effort: None,
1224 reasoning_effort_auto: false,
1225 auto_model: false,
1226 allow_shell: false,
1227 trust_mode: false,
1228 auto_approve: false,
1229 approval_mode: ApprovalMode::Suggest,
1230 translation_enabled: false,
1231 allowed_tools: None,
1232 dynamic_tools: Vec::new(),
1233 hook_executor: None,
1234 verbosity: None,
1235 provenance: UserInputProvenance::ExternalUser,
1236 submission_id: None,
1237 }))
1238 .await
1239 .expect("send budgeted goal turn");
1240
1241 let mut starts = 0;
1242 let mut completed_turns = 0;
1243 let mut saw_blocked_goal = false;
1244 while completed_turns < 2 || !saw_blocked_goal {
1245 let event = tokio::time::timeout(model_turn_event_timeout(), async {
1246 handle.rx_event.write().await.recv().await
1247 })
1248 .await
1249 .expect("budget goal event timeout")
1250 .expect("budget goal event");
1251 match event {
1252 Event::TurnStarted { .. } => starts += 1,
1253 Event::TurnComplete { status, error, .. } => {
1254 if status == TurnOutcomeStatus::Completed {
1255 completed_turns += 1;
1256 } else {
1257 // The fixture mock is exhausted after the continuation
1258 // turns; that provider failure blocks the goal (a
1259 // legitimate terminal) — it is not a budget pause.
1260 assert!(
1261 error
1262 .as_deref()
1263 .is_some_and(|e| e.contains("no canned turn queued")),
1264 "unexpected non-completed turn: {status:?} {error:?}"
1265 );
1266 }
1267 }
1268 Event::GoalUpdated { snapshot } if snapshot.status == "paused" => {
1269 panic!(
1270 "budgets are telemetry-only in unbounded goal mode; the goal must not \
1271 pause on budget (pause_reason={:?})",
1272 snapshot.pause_reason
1273 );
1274 }
1275 Event::GoalUpdated { snapshot } if snapshot.status == "blocked" => {
1276 saw_blocked_goal = true;
1277 }
1278 _ => {}
1279 }
1280 }
1281
1282 let snapshot = goal_state.lock().expect("goal lock").snapshot();
1283 assert_eq!(snapshot.status, "blocked");
1284 assert_eq!(
1285 snapshot.pause_reason, None,
1286 "budget must never be the pause reason in unbounded goal mode"
1287 );
1288 assert_eq!(snapshot.tokens_used, 11);
1289 assert_eq!(snapshot.token_budget, Some(10));
1290 assert!(
1291 starts >= 2,
1292 "budget exhaustion must not stop the cross-turn continuation (starts={starts})"
1293 );
1294 assert!(
1295 model.call_count() >= 2,
1296 "the continuation must issue a second provider call (calls={})",
1297 model.call_count()
1298 );
1299
1300 handle.send(Op::Shutdown).await.expect("shutdown engine");
1301 run_task.await.expect("engine task");
1302 }
1303
1304 #[tokio::test]
1305 async fn current_turn_usage_does_not_stop_budgeted_goal_after_one_provider_call() {
1306 use crate::llm_client::mock::{MockLlmClient, canned};
1307
1308 let objective = "stop before an intra-turn budget overspend";
1309 let budget_turn = vec![
1310 canned::message_start("mock_goal_current_turn_budget"),
1311 canned::text_block_start(0),
1312 canned::text_delta(0, "budget spent"),
1313 canned::block_stop(0),
1314 canned::message_delta(
1315 "end_turn",
1316 Some(Usage {
1317 input_tokens: 8,
1318 output_tokens: 3,
1319 ..Usage::default()
1320 }),
1321 ),
1322 canned::message_stop(),
1323 ];
1324 let model = std::sync::Arc::new(MockLlmClient::new(vec![
1325 budget_turn,
1326 canned::simple_text_turn("the continuation runs past the exhausted budget"),
1327 ]));
1328 let client: crate::core::model_client::SharedModelClient = model.clone();
1329 let config = goal_custom_route_config();
1330 let (engine, handle) = Engine::new_with_model_client(
1331 EngineConfig {
1332 model: "local-model".to_string(),
1333 snapshots_enabled: false,
1334 subagents_enabled: false,
1335 terminal_chrome_enabled: false,
1336 goal_objective: Some(objective.to_string()),
1337 goal_token_budget: Some(10),
1338 ..EngineConfig::default()
1339 },
1340 &config,
1341 client,
1342 );
1343 let goal_state = engine.config.goal_state.clone();
1344 let run_task = tokio::spawn(engine.run());
1345
1346 handle
1347 .send(active_goal_message_op(
1348 &config,
1349 "start the budgeted goal",
1350 objective,
1351 Some(10),
1352 ))
1353 .await
1354 .expect("send budgeted goal turn");
1355 let mut starts = 0;
1356 let mut saw_blocked_goal = false;
1357 while !saw_blocked_goal {
1358 let event = tokio::time::timeout(model_turn_event_timeout(), async {
1359 handle.rx_event.write().await.recv().await
1360 })
1361 .await
1362 .expect("current-turn budget continuation did not settle")
1363 .expect("current-turn budget event");
1364 match event {
1365 Event::TurnStarted { .. } => starts += 1,
1366 Event::TurnComplete { status, error, .. } => {
1367 if status != TurnOutcomeStatus::Completed {
1368 assert!(
1369 error
1370 .as_deref()
1371 .is_some_and(|e| e.contains("no canned turn queued")),
1372 "unexpected non-completed turn: {status:?} {error:?}"
1373 );
1374 }
1375 }
1376 Event::GoalUpdated { snapshot } if snapshot.status == "paused" => {
1377 panic!(
1378 "budgets are telemetry-only in unbounded goal mode; the goal must not \
1379 pause on budget (pause_reason={:?})",
1380 snapshot.pause_reason
1381 );
1382 }
1383 Event::GoalUpdated { snapshot } if snapshot.status == "blocked" => {
1384 saw_blocked_goal = true;
1385 }
1386 _ => {}
1387 }
1388 }
1389 assert!(
1390 model.call_count() >= 2,
1391 "current-turn usage must not stop additional provider calls (calls={})",
1392 model.call_count()
1393 );
1394 assert_eq!(model.remaining_turns(), 0);
1395 let goal = goal_state.lock().expect("goal lock").snapshot();
1396 assert_eq!(goal.status, "blocked");
1397 assert_eq!(
1398 goal.pause_reason, None,
1399 "budget must never be the pause reason in unbounded goal mode"
1400 );
1401 assert_eq!(goal.tokens_used, 11);
1402 assert_eq!(goal.token_budget, Some(10));
1403 assert!(
1404 starts >= 1,
1405 "the initial goal turn must start before its intra-turn continuation (starts={starts})"
1406 );
1407
1408 handle.send(Op::Shutdown).await.expect("shutdown engine");
1409 run_task.await.expect("engine task");
1410 }
1411
1412 #[tokio::test]
1413 async fn tool_response_crossing_goal_budget_issues_second_provider_request() {
1414 use crate::llm_client::mock::{MockLlmClient, canned};
1415
1416 let objective = "stop after a budget-crossing goal tool response";
1417 let budget_tool_turn = vec![
1418 canned::message_start("mock_goal_tool_budget"),
1419 canned::tool_use_block_start(0, "call-get-goal", "get_goal"),
1420 canned::tool_input_delta(0, "{}"),
1421 canned::block_stop(0),
1422 canned::message_delta(
1423 "tool_use",
1424 Some(Usage {
1425 input_tokens: 8,
1426 output_tokens: 3,
1427 ..Usage::default()
1428 }),
1429 ),
1430 canned::message_stop(),
1431 ];
1432 let model = std::sync::Arc::new(MockLlmClient::new(vec![
1433 budget_tool_turn,
1434 canned::simple_text_turn(
1435 "this second provider response is issued past the exhausted budget",
1436 ),
1437 ]));
1438 let client: crate::core::model_client::SharedModelClient = model.clone();
1439 let config = goal_custom_route_config();
1440 let (engine, handle) = Engine::new_with_model_client(
1441 EngineConfig {
1442 model: "local-model".to_string(),
1443 snapshots_enabled: false,
1444 subagents_enabled: false,
1445 terminal_chrome_enabled: false,
1446 goal_objective: Some(objective.to_string()),
1447 goal_token_budget: Some(10),
1448 ..EngineConfig::default()
1449 },
1450 &config,
1451 client,
1452 );
1453 let goal_state = engine.config.goal_state.clone();
1454 let run_task = tokio::spawn(engine.run());
1455
1456 handle
1457 .send(active_goal_message_op(
1458 &config,
1459 "inspect the goal without overspending",
1460 objective,
1461 Some(10),
1462 ))
1463 .await
1464 .expect("send budgeted goal tool turn");
1465
1466 let mut saw_get_goal = false;
1467 let mut saw_blocked_goal = false;
1468 while !saw_blocked_goal {
1469 let event = tokio::time::timeout(model_turn_event_timeout(), async {
1470 handle.rx_event.write().await.recv().await
1471 })
1472 .await
1473 .expect("budget-crossing goal tool turn did not settle")
1474 .expect("budget-crossing goal tool event");
1475 match event {
1476 Event::ToolCallComplete { name, result, .. } if name == "get_goal" => {
1477 assert!(result.expect("get_goal result").success);
1478 saw_get_goal = true;
1479 }
1480 Event::TurnComplete { status, error, .. } => {
1481 if status != TurnOutcomeStatus::Completed {
1482 assert!(
1483 error
1484 .as_deref()
1485 .is_some_and(|e| e.contains("no canned turn queued")),
1486 "unexpected non-completed turn: {status:?} {error:?}"
1487 );
1488 }
1489 }
1490 Event::GoalUpdated { snapshot } if snapshot.status == "paused" => {
1491 panic!(
1492 "budgets are telemetry-only in unbounded goal mode; the goal must not \
1493 pause on budget (pause_reason={:?})",
1494 snapshot.pause_reason
1495 );
1496 }
1497 Event::GoalUpdated { snapshot } if snapshot.status == "blocked" => {
1498 saw_blocked_goal = true;
1499 }
1500 _ => {}
1501 }
1502 }
1503
1504 assert!(saw_get_goal, "the first response's goal tool must execute");
1505 assert!(
1506 model.call_count() >= 2,
1507 "the provider-request boundary must not stop the second model call (calls={})",
1508 model.call_count()
1509 );
1510 assert_eq!(model.remaining_turns(), 0);
1511 let goal = goal_state.lock().expect("goal lock").snapshot();
1512 assert_eq!(goal.status, "blocked");
1513 assert_eq!(goal.tokens_used, 11, "usage must be durably recorded once");
1514 assert_eq!(goal.token_budget, Some(10));
1515 assert_eq!(
1516 goal.pause_reason, None,
1517 "budget must never be the pause reason in unbounded goal mode"
1518 );
1519
1520 handle.send(Op::Shutdown).await.expect("shutdown engine");
1521 run_task.await.expect("engine task");
1522 }
1522 lines RUST