返回 CodeWhale
test_cases_04.rs
根目录 / crates / tui / src / core / engine / tests / test_cases_04.rs
1
2
3 #[tokio::test]
4 async fn queued_goal_clear_refreshes_prompt_and_cancels_stale_continuation() {
5 let config = Config::default();
6 let engine_config = EngineConfig {
7 snapshots_enabled: false,
8 terminal_chrome_enabled: false,
9 goal_objective: Some("clear this goal".to_string()),
10 goal_token_budget: Some(42_000),
11 ..EngineConfig::default()
12 };
13 let (engine, handle) = Engine::new(engine_config, &config);
14 let goal_state = engine.config.goal_state.clone();
15 let run_task = tokio::spawn(engine.run());
16
17 // Model the mailbox order produced when the user clears a goal while its
18 // prior turn is still finishing: the control reaches the queue before the
19 // synthetic continuation that TurnComplete schedules.
20 handle
21 .send(Op::SetGoalStatus {
22 goal_id: None,
23 status: crate::tools::goal::GoalStatus::Active,
24 clear: true,
25 })
26 .await
27 .expect("queue goal clear");
28 handle
29 .send(Op::ContinueGoal {
30 dynamic_tools: Vec::new(),
31 engine_schedule_id: None,
32 })
33 .await
34 .expect("queue stale continuation");
35
36 // This receipt sits behind both operations. Once it arrives, a stale
37 // continuation has either incorrectly started a turn or been consumed.
38 let session = handle
39 .get_session_snapshot()
40 .await
41 .expect("post-clear session snapshot");
42 let prompt = match session.system_prompt.expect("post-clear system prompt") {
43 SystemPrompt::Text(text) => text,
44 SystemPrompt::Blocks(blocks) => blocks
45 .into_iter()
46 .map(|block| block.text)
47 .collect::<Vec<_>>()
48 .join("\n"),
49 };
50 assert!(
51 !prompt.contains("<session_goal>"),
52 "cleared config fallback must not restore the goal prompt: {prompt}"
53 );
54 let snapshot = goal_state.lock().expect("goal lock").snapshot();
55 assert_eq!(snapshot.objective, None);
56 assert_eq!(snapshot.status, "none");
57 assert_eq!(snapshot.token_budget, None);
58
59 let mut saw_clear_session = false;
60 let mut saw_clear_goal = false;
61 let mut saw_clear_status = false;
62 {
63 let mut events = handle.rx_event.write().await;
64 while let Ok(event) = events.try_recv() {
65 match event {
66 Event::TurnStarted { .. } => {
67 panic!("queued clear must prevent a stale goal continuation")
68 }
69 Event::SessionUpdated { system_prompt, .. } => {
70 let prompt = match system_prompt.expect("clear SessionUpdated prompt") {
71 SystemPrompt::Text(text) => text,
72 SystemPrompt::Blocks(blocks) => blocks
73 .into_iter()
74 .map(|block| block.text)
75 .collect::<Vec<_>>()
76 .join("\n"),
77 };
78 assert!(!prompt.contains("<session_goal>"), "{prompt}");
79 saw_clear_session = true;
80 }
81 Event::GoalUpdated { snapshot } => {
82 assert_eq!(snapshot.objective, None);
83 assert_eq!(snapshot.status, "none");
84 saw_clear_goal = true;
85 }
86 Event::Status { message } if message == "Goal cleared." => {
87 saw_clear_status = true;
88 }
89 _ => {}
90 }
91 }
92 }
93 assert!(
94 saw_clear_session,
95 "clear must refresh persisted prompt state"
96 );
97 assert!(saw_clear_goal, "clear must emit a canonical empty snapshot");
98 assert!(saw_clear_status, "clear must remain user-visible");
99
100 handle.send(Op::Shutdown).await.expect("shutdown engine");
101 run_task.await.expect("engine task");
102 }
103
104 fn goal_custom_route_config() -> Config {
105 let mut custom = HashMap::new();
106 custom.insert(
107 "custom-a".to_string(),
108 crate::config::ProviderConfig {
109 kind: Some("openai-compatible".to_string()),
110 base_url: Some("http://127.0.0.1:18181/v1".to_string()),
111 model: Some("local-model".to_string()),
112 api_key: Some("local-test-key".to_string()),
113 ..crate::config::ProviderConfig::default()
114 },
115 );
116 Config {
117 provider: Some("custom-a".to_string()),
118 providers: Some(crate::config::ProvidersConfig {
119 custom,
120 ..crate::config::ProvidersConfig::default()
121 }),
122 ..Config::default()
123 }
124 }
125
126 #[tokio::test]
127 async fn ordinary_prose_never_activates_a_goal() {
128 let request_entered = std::sync::Arc::new(tokio::sync::Notify::new());
129 let release_request = std::sync::Arc::new(tokio::sync::Notify::new());
130 let model = std::sync::Arc::new(FirstRequestGatedGoalModelClient {
131 calls: std::sync::atomic::AtomicUsize::new(0),
132 request_entered: std::sync::Arc::clone(&request_entered),
133 release_request: std::sync::Arc::clone(&release_request),
134 });
135 let config = goal_custom_route_config();
136 let client: crate::core::model_client::SharedModelClient = model.clone();
137 let (engine, handle) = Engine::new_with_model_client(
138 EngineConfig {
139 model: "local-model".to_string(),
140 max_steps: 1,
141 snapshots_enabled: false,
142 terminal_chrome_enabled: false,
143 ..EngineConfig::default()
144 },
145 &config,
146 client,
147 );
148 let goal_state = engine.config.goal_state.clone();
149 let run_task = tokio::spawn(engine.run());
150
151 handle
152 .send(Op::SendMessage(TurnSpec {
153 max_output_tokens: None,
154 content: "hello - take over and make it your /goal to solve navier stokes".to_string(),
155 images: Vec::new(),
156 mode: AppMode::Agent,
157 route: resolved_route_for_test(&config, "local-model"),
158 compaction: Box::new(CompactionConfig::default()),
159 initial_routed_usage: Box::default(),
160 goal_objective: None,
161 goal_token_budget: None,
162 goal_status: crate::tools::goal::GoalStatus::Active,
163 reasoning_effort: None,
164 reasoning_effort_auto: false,
165 auto_model: false,
166 allow_shell: false,
167 trust_mode: false,
168 auto_approve: false,
169 approval_mode: ApprovalMode::Suggest,
170 translation_enabled: false,
171 allowed_tools: None,
172 dynamic_tools: Vec::new(),
173 hook_executor: None,
174 verbosity: None,
175 provenance: UserInputProvenance::ExternalUser,
176 submission_id: None,
177 }))
178 .await
179 .expect("send explicit natural goal turn");
180
181 // #6290 rework: the natural-language `/goal` prose parser is gone. The
182 // same wording that used to activate a goal is now an ordinary turn;
183 // only the model (`create_goal`) or the `/goal` command creates one.
184 let mut saw_goal = false;
185 loop {
186 let event = tokio::time::timeout(model_turn_event_timeout(), async {
187 handle.rx_event.write().await.recv().await
188 })
189 .await
190 .expect("prose goal event timeout")
191 .expect("prose goal event");
192 match event {
193 Event::GoalUpdated { .. } => {
194 saw_goal = true;
195 }
196 Event::TurnStarted { .. } => {
197 assert!(
198 !saw_goal,
199 "ordinary prose must not publish a goal before provider work starts"
200 );
201 break;
202 }
203 _ => {}
204 }
205 }
206
207 tokio::time::timeout(model_turn_event_timeout(), request_entered.notified())
208 .await
209 .expect("provider request was never entered");
210 release_request.notify_one();
211 let _ = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot())
212 .await
213 .expect("turn did not settle")
214 .expect("post-turn session snapshot");
215 let snapshot = goal_state.lock().expect("goal lock").snapshot();
216 assert_eq!(snapshot.objective.as_deref(), None);
217 assert!(!snapshot.is_active());
218 assert_eq!(model.calls.load(std::sync::atomic::Ordering::SeqCst), 1);
219
220 handle.send(Op::Shutdown).await.expect("shutdown engine");
221 run_task.await.expect("engine task");
222 }
223
224 /// Drive one turn in `mode` and report what the goal path did before the
225 /// provider was called: the objective published by `GoalUpdated` (if any),
226 /// whether the engine goal state is active, and how many Operate contract
227 /// messages the session log holds afterwards.
228 async fn operate_goal_probe(mode: AppMode, prompt: &str) -> (Option<String>, bool, usize) {
229 let request_entered = std::sync::Arc::new(tokio::sync::Notify::new());
230 let release_request = std::sync::Arc::new(tokio::sync::Notify::new());
231 let model = std::sync::Arc::new(FirstRequestGatedGoalModelClient {
232 calls: std::sync::atomic::AtomicUsize::new(0),
233 request_entered: std::sync::Arc::clone(&request_entered),
234 release_request: std::sync::Arc::clone(&release_request),
235 });
236 let config = goal_custom_route_config();
237 let client: crate::core::model_client::SharedModelClient = model.clone();
238 let (engine, handle) = Engine::new_with_model_client(
239 EngineConfig {
240 model: "local-model".to_string(),
241 max_steps: 1,
242 snapshots_enabled: false,
243 terminal_chrome_enabled: false,
244 ..EngineConfig::default()
245 },
246 &config,
247 client,
248 );
249 let goal_state = engine.config.goal_state.clone();
250 let run_task = tokio::spawn(engine.run());
251
252 handle
253 .send(Op::SendMessage(TurnSpec {
254 max_output_tokens: None,
255 content: prompt.to_string(),
256 images: Vec::new(),
257 mode,
258 route: resolved_route_for_test(&config, "local-model"),
259 compaction: Box::new(CompactionConfig::default()),
260 initial_routed_usage: Box::default(),
261 goal_objective: None,
262 goal_token_budget: None,
263 goal_status: crate::tools::goal::GoalStatus::Active,
264 reasoning_effort: None,
265 reasoning_effort_auto: false,
266 auto_model: false,
267 allow_shell: false,
268 trust_mode: false,
269 auto_approve: false,
270 approval_mode: ApprovalMode::Suggest,
271 translation_enabled: false,
272 allowed_tools: None,
273 dynamic_tools: Vec::new(),
274 hook_executor: None,
275 verbosity: None,
276 provenance: UserInputProvenance::ExternalUser,
277 submission_id: None,
278 }))
279 .await
280 .expect("send probe turn");
281
282 let mut published_objective = None;
283 loop {
284 let event = tokio::time::timeout(model_turn_event_timeout(), async {
285 handle.rx_event.write().await.recv().await
286 })
287 .await
288 .expect("probe event timeout")
289 .expect("probe event");
290 match event {
291 Event::GoalUpdated { snapshot } if snapshot.status == "active" => {
292 published_objective = snapshot.objective;
293 }
294 Event::TurnStarted { .. } => break,
295 _ => {}
296 }
297 }
298 tokio::time::timeout(model_turn_event_timeout(), request_entered.notified())
299 .await
300 .expect("provider request was never entered");
301 let active = goal_state.lock().expect("goal lock").is_active();
302 if active {
303 // Stop autonomous continuation after the one provider call.
304 handle
305 .send(Op::SetGoalStatus {
306 goal_id: None,
307 status: crate::tools::goal::GoalStatus::Paused,
308 clear: false,
309 })
310 .await
311 .expect("queue goal pause");
312 }
313 release_request.notify_one();
314 let snapshot = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot())
315 .await
316 .expect("probe turn did not settle")
317 .expect("probe session snapshot");
318 let contracts = snapshot
319 .messages
320 .iter()
321 .filter(|message| crate::runtime_handoff::is_operate_contract_message(message))
322 .count();
323 assert_eq!(model.calls.load(std::sync::atomic::Ordering::SeqCst), 1);
324
325 handle.send(Op::Shutdown).await.expect("shutdown engine");
326 run_task.await.expect("engine task");
327 (published_objective, active, contracts)
328 }
329
330 #[tokio::test]
331 async fn operate_never_promotes_wording_to_a_goal() {
332 let prompt =
333 "Migrate the settings loader to the new config crate and keep the old keys readable";
334
335 // The verb-list promotion is gone: an ordinary work prompt is an ordinary
336 // turn in every mode, and the model decides goals through `create_goal`
337 // (docs/design/TUI_DECONSTRUCTION.md, founder clarification 2026-09-09).
338 let (objective, active, contracts) = operate_goal_probe(AppMode::Operate, prompt).await;
339 assert_eq!(
340 objective, None,
341 "the host must not infer a goal from wording"
342 );
343 assert!(!active);
344 assert_eq!(contracts, 1, "Operate appends its contract exactly once");
345
346 let (objective, active, contracts) = operate_goal_probe(AppMode::Agent, prompt).await;
347 assert_eq!(objective, None, "Work must not promote a prompt to a goal");
348 assert!(!active);
349 assert_eq!(contracts, 0, "Work never sees the Operate contract");
350
351 // #6290 rework: even an explicit-looking declaration is ordinary
352 // prose now — the host never parses it, and the model decides goals
353 // through `create_goal`.
354 let (objective, active, contracts) =
355 operate_goal_probe(AppMode::Operate, "Please set /goal to ship the release").await;
356 assert_eq!(
357 objective, None,
358 "prose asking for a goal must not create one host-side"
359 );
360 assert!(!active);
361 assert_eq!(contracts, 1);
362 }
363
364 #[tokio::test]
365 async fn operate_leaves_followup_and_long_questions_as_ordinary_turns() {
366 let report = "what about like rust or docker builds or something";
367 let (objective, active, contracts) = operate_goal_probe(AppMode::Operate, report).await;
368 assert_eq!(
369 objective, None,
370 "conversational followup must remain ordinary turn in Operate"
371 );
372 assert!(!active);
373 assert_eq!(contracts, 1);
374
375 let long_q =
376 "why did the build fail on the last step when running under docker on macos with rust 1.80";
377 let (objective, active, contracts) = operate_goal_probe(AppMode::Operate, long_q).await;
378 assert_eq!(
379 objective, None,
380 "long question without punctuation must remain ordinary turn"
381 );
382 assert!(!active);
383 assert_eq!(contracts, 1);
384
385 let zh_followup = "那 rust 或者 docker 构建呢";
386 let (objective, active, contracts) = operate_goal_probe(AppMode::Operate, zh_followup).await;
387 assert_eq!(
388 objective, None,
389 "Chinese followup must remain ordinary turn in Operate"
390 );
391 assert!(!active);
392 assert_eq!(contracts, 1);
393 }
394
395 #[tokio::test]
396 async fn operate_does_not_create_a_goal_when_the_work_request_declines_one() {
397 let prompt = "Run one bounded cancellation check. Do not edit files, inspect other files, create a goal, spawn agents, or start any other tool.";
398 let (objective, active, contracts) = operate_goal_probe(AppMode::Operate, prompt).await;
399 assert_eq!(objective, None);
400 assert!(!active);
401 assert_eq!(
402 contracts, 1,
403 "the ordinary Operate turn still reaches the model"
404 );
405 }
406
407 #[tokio::test]
408 async fn operate_contract_is_appended_once_and_an_existing_goal_is_never_replaced() {
409 let first_entered = std::sync::Arc::new(tokio::sync::Notify::new());
410 let release_first = std::sync::Arc::new(tokio::sync::Notify::new());
411 let model = std::sync::Arc::new(IndexedGatedGoalModelClient {
412 calls: std::sync::atomic::AtomicUsize::new(0),
413 gates: HashMap::from([(
414 1,
415 (
416 std::sync::Arc::clone(&first_entered),
417 std::sync::Arc::clone(&release_first),
418 ),
419 )]),
420 max_calls: 2,
421 });
422 let config = goal_custom_route_config();
423 let client: crate::core::model_client::SharedModelClient = model.clone();
424 let (engine, handle) = Engine::new_with_model_client(
425 EngineConfig {
426 model: "local-model".to_string(),
427 max_steps: 1,
428 snapshots_enabled: false,
429 terminal_chrome_enabled: false,
430 ..EngineConfig::default()
431 },
432 &config,
433 client,
434 );
435 let goal_state = engine.config.goal_state.clone();
436 let run_task = tokio::spawn(engine.run());
437
438 let send = |content: &str, goal_objective: Option<String>, goal_status| {
439 Op::SendMessage(TurnSpec {
440 max_output_tokens: None,
441 content: content.to_string(),
442 images: Vec::new(),
443 mode: AppMode::Operate,
444 route: resolved_route_for_test(&config, "local-model"),
445 compaction: Box::new(CompactionConfig::default()),
446 initial_routed_usage: Box::default(),
447 goal_objective,
448 goal_token_budget: None,
449 goal_status,
450 reasoning_effort: None,
451 reasoning_effort_auto: false,
452 auto_model: false,
453 allow_shell: false,
454 trust_mode: false,
455 auto_approve: false,
456 approval_mode: ApprovalMode::Suggest,
457 translation_enabled: false,
458 allowed_tools: None,
459 dynamic_tools: Vec::new(),
460 hook_executor: None,
461 verbosity: None,
462 provenance: UserInputProvenance::ExternalUser,
463 submission_id: None,
464 })
465 };
466
467 let first_objective =
468 "Migrate the settings loader to the new config crate and keep the old keys readable";
469 // #6290 rework: prose no longer creates goals, so the unfinished goal
470 // this test needs is seeded directly — the same `GoalState::create` path
471 // the `/goal` command and the model's `create_goal` tool use.
472 goal_state
473 .lock()
474 .expect("goal lock")
475 .create(first_objective.to_string(), None)
476 .expect("seed unfinished goal");
477 handle
478 .send(send(
479 first_objective,
480 Some(first_objective.to_string()),
481 crate::tools::goal::GoalStatus::Active,
482 ))
483 .await
484 .expect("send first Operate turn");
485 tokio::time::timeout(model_turn_event_timeout(), first_entered.notified())
486 .await
487 .expect("first provider request was never entered");
488 assert_eq!(
489 goal_state.lock().expect("goal lock").objective(),
490 Some(first_objective)
491 );
492 handle
493 .send(Op::SetGoalStatus {
494 goal_id: None,
495 status: crate::tools::goal::GoalStatus::Paused,
496 clear: false,
497 })
498 .await
499 .expect("queue goal pause");
500 release_first.notify_one();
501 let _ = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot())
502 .await
503 .expect("first turn did not settle")
504 .expect("first session snapshot");
505
506 // Second Operate prompt while the (paused) goal is still unfinished: the
507 // host reports that goal, so the prompt is ordinary work under it.
508 handle
509 .send(send(
510 "Refactor the provider table so it survives a config reload",
511 Some(first_objective.to_string()),
512 crate::tools::goal::GoalStatus::Paused,
513 ))
514 .await
515 .expect("send second Operate turn");
516 let snapshot = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot())
517 .await
518 .expect("second turn did not settle")
519 .expect("second session snapshot");
520 assert_eq!(model.calls.load(std::sync::atomic::Ordering::SeqCst), 2);
521 let contracts = snapshot
522 .messages
523 .iter()
524 .filter(|message| crate::runtime_handoff::is_operate_contract_message(message))
525 .count();
526 assert_eq!(
527 contracts, 1,
528 "the contract must not repeat on later Operate turns"
529 );
530 let goal = goal_state.lock().expect("goal lock").snapshot();
531 assert_eq!(goal.objective.as_deref(), Some(first_objective));
532 assert_eq!(goal.status, "paused");
533
534 handle.send(Op::Shutdown).await.expect("shutdown engine");
535 run_task.await.expect("engine task");
536 }
537
538 fn without_named_custom_route(mut config: Config) -> Config {
539 config
540 .providers
541 .as_mut()
542 .expect("custom providers")
543 .custom
544 .clear();
545 config
546 }
547
548 #[tokio::test]
549 async fn exhausted_goal_reaches_route_failure_without_budget_pause() {
550 let config = goal_custom_route_config();
551 let engine_config = EngineConfig {
552 model: "local-model".to_string(),
553 snapshots_enabled: false,
554 terminal_chrome_enabled: false,
555 goal_objective: Some("stop at the budget".to_string()),
556 goal_token_budget: Some(10),
557 ..EngineConfig::default()
558 };
559 let (mut engine, handle) = Engine::new(engine_config, &config);
560 let goal_state = engine.config.goal_state.clone();
561 goal_state.lock().expect("goal lock").record_usage(11, 0);
562
563 let invalid_route_config = without_named_custom_route(config);
564 engine.authoritative_route_config =
565 Some(Arc::new(parking_lot::RwLock::new(invalid_route_config)));
566 assert!(
567 engine.current_runtime_route().is_err(),
568 "fixture must prove route resolution cannot succeed"
569 );
570 let run_task = tokio::spawn(engine.run());
571
572 handle
573 .send(Op::ContinueGoal {
574 dynamic_tools: Vec::new(),
575 engine_schedule_id: None,
576 })
577 .await
578 .expect("queue exhausted continuation");
579
580 let mut saw_route_error = false;
581 let mut saw_blocked_goal = false;
582 while !(saw_route_error && saw_blocked_goal) {
583 let event = tokio::time::timeout(model_turn_event_timeout(), async {
584 handle.rx_event.write().await.recv().await
585 })
586 .await
587 .expect("budget terminal event timeout")
588 .expect("budget terminal event");
589 match event {
590 Event::TurnStarted { .. } => {
591 panic!("an exhausted goal must not start a turn before route resolution")
592 }
593 Event::Error { envelope, .. } => {
594 assert!(
595 envelope.message.contains("route is no longer valid"),
596 "the exhausted goal must reach the invalid-route failure, got: {envelope:?}"
597 );
598 saw_route_error = true;
599 }
600 Event::GoalUpdated { snapshot } if snapshot.status == "paused" => {
601 panic!(
602 "budgets are telemetry-only in unbounded goal mode; the goal must not pause on budget (pause_reason={:?})",
603 snapshot.pause_reason
604 );
605 }
606 Event::GoalUpdated { snapshot } if snapshot.status == "blocked" => {
607 assert_eq!(snapshot.tokens_used, 11);
608 assert_eq!(snapshot.token_budget, Some(10));
609 assert_eq!(
610 snapshot.pause_reason, None,
611 "budget must never be the pause reason in unbounded goal mode"
612 );
613 saw_blocked_goal = true;
614 }
615 _ => {}
616 }
617 }
618
619 let snapshot = goal_state.lock().expect("goal lock").snapshot();
620 assert_eq!(snapshot.status, "blocked");
621 assert_eq!(snapshot.tokens_used, 11);
622 assert_eq!(snapshot.token_budget, Some(10));
623 assert_eq!(snapshot.pause_reason, None);
624
625 handle.send(Op::Shutdown).await.expect("shutdown engine");
626 run_task.await.expect("engine task");
627 }
628
629 #[tokio::test]
630 async fn continuation_circuit_breaker_pauses_with_run_limit_reason() {
631 // #5052: the backstop is configurable ([goal] max_continuations) and set
632 // deliberately past the retired hardcoded cap of 10 to prove an operate
633 // goal is no longer stopped there — only the configured backstop halts a
634 // pathological loop that never emits a terminal signal.
635 let backstop = 12u32;
636 let config = Config::default();
637 let (engine, handle) = Engine::new(
638 EngineConfig {
639 snapshots_enabled: false,
640 terminal_chrome_enabled: false,
641 goal_objective: Some("stop a runaway continuation loop".to_string()),
642 goal_max_continuations: backstop,
643 ..EngineConfig::default()
644 },
645 &config,
646 );
647 let goal_state = engine.config.goal_state.clone();
648 {
649 let mut goal = goal_state.lock().expect("goal lock");
650 for _ in 0..backstop {
651 goal.record_continuation();
652 }
653 }
654 let run_task = tokio::spawn(engine.run());
655
656 handle
657 .send(Op::ContinueGoal {
658 dynamic_tools: Vec::new(),
659 engine_schedule_id: None,
660 })
661 .await
662 .expect("queue capped continuation");
663
664 let mut saw_pause = false;
665 let mut saw_reason = false;
666 while !(saw_pause && saw_reason) {
667 let event = tokio::time::timeout(model_turn_event_timeout(), async {
668 handle.rx_event.write().await.recv().await
669 })
670 .await
671 .expect("continuation cap event timeout")
672 .expect("continuation cap event");
673 match event {
674 Event::TurnStarted { .. } => panic!("capped goal must not start another turn"),
675 Event::GoalUpdated { snapshot } if snapshot.status == "paused" => {
676 assert_eq!(
677 snapshot.pause_reason,
678 Some(crate::tools::goal::GoalPauseReason::Backoff)
679 );
680 saw_pause = true;
681 }
682 Event::Status { message } if message.contains("automatic continuations") => {
683 assert!(message.contains(&backstop.to_string()), "{message}");
684 assert!(message.contains("[goal] max_continuations"), "{message}");
685 saw_reason = true;
686 }
687 _ => {}
688 }
689 }
690
691 handle.send(Op::Shutdown).await.expect("shutdown engine");
692 run_task.await.expect("engine task");
693 }
694
695 #[tokio::test]
696 async fn goal_continues_past_legacy_ten_pass_cap_when_budget_remains() {
697 // #5052 regression: 10 automatic continuations used to be a terminal stop.
698 // With the default backstop and budget remaining, the loop must keep
699 // dispatching toward the completion gate.
700 let config = Config::default();
701 let (engine, _handle) = Engine::new(
702 EngineConfig {
703 snapshots_enabled: false,
704 terminal_chrome_enabled: false,
705 goal_objective: Some("run to the completion gate, not a pass count".to_string()),
706 goal_token_budget: Some(1_000_000),
707 ..EngineConfig::default()
708 },
709 &config,
710 );
711 {
712 let mut goal = engine.config.goal_state.lock().expect("goal lock");
713 for _ in 0..10 {
714 goal.record_continuation();
715 }
716 }
717
718 match engine.goal_continuation_if_active() {
719 GoalContinuationAction::Dispatch { snapshot, .. } => {
720 assert_eq!(snapshot.continuation_count, 11);
721 }
722 other => panic!("goal must continue past 10 passes, got {other:?}"),
723 }
724
725 // A backstop of 0 means unlimited-with-budget-stops: even a pathological
726 // pass count keeps continuing while budget remains.
727 let (engine, _handle) = Engine::new(
728 EngineConfig {
729 snapshots_enabled: false,
730 terminal_chrome_enabled: false,
731 goal_objective: Some("unlimited backstop".to_string()),
732 goal_token_budget: Some(1_000_000),
733 goal_max_continuations: 0,
734 ..EngineConfig::default()
735 },
736 &config,
737 );
738 {
739 let mut goal = engine.config.goal_state.lock().expect("goal lock");
740 for _ in 0..500 {
741 goal.record_continuation();
742 }
743 }
744 assert!(
745 matches!(
746 engine.goal_continuation_if_active(),
747 GoalContinuationAction::Dispatch { .. }
748 ),
749 "backstop 0 must not stop an in-budget goal"
750 );
751 }
752
753 /// T08-03: the cross-turn continuation gate stops at an *enforced* token
754 /// budget, the same stop the intra-turn gate and the host take, and keeps
755 /// the default advisory behavior when enforcement is off.
756 #[tokio::test]
757 async fn cross_turn_goal_continuation_stops_at_an_enforced_token_budget() {
758 let config = Config::default();
759 for enforce in [true, false] {
760 let (engine, _handle) = Engine::new(
761 EngineConfig {
762 snapshots_enabled: false,
763 terminal_chrome_enabled: false,
764 goal_objective: Some("stop at the enforced budget".to_string()),
765 goal_token_budget: Some(100),
766 goal_enforce_token_budget: enforce,
767 ..EngineConfig::default()
768 },
769 &config,
770 );
771 engine
772 .config
773 .goal_state
774 .lock()
775 .expect("goal lock")
776 .record_usage(100, 0);
777 let action = engine.goal_continuation_if_active();
778 if enforce {
779 assert!(
780 matches!(
781 action,
782 GoalContinuationAction::Stopped {
783 reason: crate::tools::goal::GoalPauseReason::BudgetLimit,
784 ..
785 }
786 ),
787 "an exhausted enforced budget must not dispatch another turn: {action:?}"
788 );
789 } else {
790 assert!(
791 matches!(action, GoalContinuationAction::Dispatch { .. }),
792 "an advisory budget keeps the goal running: {action:?}"
793 );
794 }
795 }
796 }
797
798 #[tokio::test]
799 async fn invalid_route_blocks_active_goal_and_refreshes_projections() {
800 let config = goal_custom_route_config();
801 let engine_config = EngineConfig {
802 model: "local-model".to_string(),
803 snapshots_enabled: false,
804 terminal_chrome_enabled: false,
805 goal_objective: Some("keep going across route drift".to_string()),
806 ..EngineConfig::default()
807 };
808 let (mut engine, handle) = Engine::new(engine_config, &config);
809 let goal_state = engine.config.goal_state.clone();
810 engine.authoritative_route_config = Some(Arc::new(parking_lot::RwLock::new(
811 without_named_custom_route(config),
812 )));
813 assert!(
814 engine.current_runtime_route().is_err(),
815 "fixture must fail before dispatch"
816 );
817 let run_task = tokio::spawn(engine.run());
818
819 handle
820 .send(Op::ContinueGoal {
821 dynamic_tools: Vec::new(),
822 engine_schedule_id: None,
823 })
824 .await
825 .expect("queue active continuation");
826 let session = handle
827 .get_session_snapshot()
828 .await
829 .expect("post-route-failure session snapshot");
830
831 let prompt = match session.system_prompt.expect("blocked system prompt") {
832 SystemPrompt::Text(text) => text,
833 SystemPrompt::Blocks(blocks) => blocks
834 .into_iter()
835 .map(|block| block.text)
836 .collect::<Vec<_>>()
837 .join("\n"),
838 };
839 assert!(!prompt.contains("<session_goal>"), "{prompt}");
840 let snapshot = goal_state.lock().expect("goal lock").snapshot();
841 assert_eq!(snapshot.status, "blocked");
842 assert!(
843 snapshot
844 .blocker
845 .as_deref()
846 .is_some_and(|blocker| blocker.contains("provider route is no longer valid")),
847 "{snapshot:?}"
848 );
849
850 let mut saw_route_error = false;
851 let mut saw_blocked_session = false;
852 let mut saw_blocked_goal = false;
853 let mut saw_blocked_status = false;
854 {
855 let mut events = handle.rx_event.write().await;
856 while let Ok(event) = events.try_recv() {
857 match event {
858 Event::TurnStarted { .. } => {
859 panic!("invalid route must not start a continuation turn")
860 }
861 Event::Error { envelope, .. } => {
862 assert!(format!("{envelope:?}").contains("provider route is no longer valid"));
863 saw_route_error = true;
864 }
865 Event::SessionUpdated {
866 system_prompt: Some(system_prompt),
867 ..
868 } => {
869 let prompt = match system_prompt {
870 SystemPrompt::Text(text) => text,
871 SystemPrompt::Blocks(blocks) => blocks
872 .into_iter()
873 .map(|block| block.text)
874 .collect::<Vec<_>>()
875 .join("\n"),
876 };
877 assert!(!prompt.contains("<session_goal>"), "{prompt}");
878 saw_blocked_session = true;
879 }
880 Event::GoalUpdated { snapshot } if snapshot.status == "blocked" => {
881 saw_blocked_goal = true;
882 }
883 Event::Status { message }
884 if message.contains("provider route is no longer valid") =>
885 {
886 assert!(message.contains("resume the goal"), "{message}");
887 saw_blocked_status = true;
888 }
889 _ => {}
890 }
891 }
892 }
893 assert!(saw_route_error, "route failure must remain visible");
894 assert!(
895 saw_blocked_session,
896 "session prompt projection must refresh"
897 );
898 assert!(saw_blocked_goal, "sidebar must receive blocked state");
899 assert!(saw_blocked_status, "blocked reason must remain visible");
900
901 handle.send(Op::Shutdown).await.expect("shutdown engine");
902 run_task.await.expect("engine task");
903 }
904
905 #[tokio::test]
906 async fn rejected_continuation_dispatch_blocks_goal_after_failed_turn() {
907 let config = goal_custom_route_config();
908 let engine_config = EngineConfig {
909 model: "local-model".to_string(),
910 snapshots_enabled: false,
911 terminal_chrome_enabled: false,
912 goal_objective: Some("keep going after dispatch".to_string()),
913 ..EngineConfig::default()
914 };
915 let (mut engine, handle) = Engine::new(engine_config, &config);
916 let goal_state = engine.config.goal_state.clone();
917 assert!(engine.current_runtime_route().is_ok());
918 // Exercise the continuation caller's `false` boundary deterministically:
919 // the route installs, but the injected-model authority has no client with
920 // which to start the request.
921 engine.model_client_injected = true;
922 engine.model_client = None;
923 let run_task = tokio::spawn(engine.run());
924
925 handle
926 .send(Op::ContinueGoal {
927 dynamic_tools: Vec::new(),
928 engine_schedule_id: None,
929 })
930 .await
931 .expect("queue rejected continuation");
932 let session = handle
933 .get_session_snapshot()
934 .await
935 .expect("post-rejection session snapshot");
936
937 let prompt = match session.system_prompt.expect("blocked system prompt") {
938 SystemPrompt::Text(text) => text,
939 SystemPrompt::Blocks(blocks) => blocks
940 .into_iter()
941 .map(|block| block.text)
942 .collect::<Vec<_>>()
943 .join("\n"),
944 };
945 assert!(!prompt.contains("<session_goal>"), "{prompt}");
946 let snapshot = goal_state.lock().expect("goal lock").snapshot();
947 assert_eq!(snapshot.status, "blocked");
948 assert!(
949 snapshot
950 .blocker
951 .as_deref()
952 .is_some_and(|blocker| blocker.contains("next model turn could not be started")),
953 "{snapshot:?}"
954 );
955
956 let mut starts = 0;
957 let mut saw_failed_turn = false;
958 let mut saw_dispatch_error = false;
959 let mut saw_blocked_session = false;
960 let mut saw_blocked_goal = false;
961 let mut saw_blocked_status = false;
962 {
963 let mut events = handle.rx_event.write().await;
964 while let Ok(event) = events.try_recv() {
965 match event {
966 Event::TurnStarted { .. } => starts += 1,
967 Event::TurnComplete { status, .. } => {
968 assert_eq!(status, TurnOutcomeStatus::Failed);
969 saw_failed_turn = true;
970 }
971 Event::Error { .. } => saw_dispatch_error = true,
972 Event::SessionUpdated {
973 system_prompt: Some(system_prompt),
974 ..
975 } => {
976 let prompt = match system_prompt {
977 SystemPrompt::Text(text) => text,
978 SystemPrompt::Blocks(blocks) => blocks
979 .into_iter()
980 .map(|block| block.text)
981 .collect::<Vec<_>>()
982 .join("\n"),
983 };
984 assert!(!prompt.contains("<session_goal>"), "{prompt}");
985 saw_blocked_session = true;
986 }
987 Event::GoalUpdated { snapshot } if snapshot.status == "blocked" => {
988 saw_blocked_goal = true;
989 }
990 Event::Status { message }
991 if message.contains("next model turn could not be started") =>
992 {
993 assert!(message.contains("resume the goal"), "{message}");
994 saw_blocked_status = true;
995 }
996 _ => {}
997 }
998 }
999 }
1000 assert_eq!(
1001 starts, 1,
1002 "dispatch reached exactly one engine turn boundary"
1003 );
1004 assert!(
1005 saw_failed_turn,
1006 "rejected dispatch must surface a failed turn"
1007 );
1008 assert!(
1009 saw_dispatch_error,
1010 "rejected dispatch must surface its error"
1011 );
1012 assert!(
1013 saw_blocked_session,
1014 "session prompt projection must refresh"
1015 );
1016 assert!(saw_blocked_goal, "sidebar must receive blocked state");
1017 assert!(saw_blocked_status, "blocked reason must remain visible");
1018
1019 handle.send(Op::Shutdown).await.expect("shutdown engine");
1020 run_task.await.expect("engine task");
1021 }
1022
1023 #[tokio::test]
1024 async fn started_nonretryable_continuation_failure_blocks_goal_with_bounded_reason() {
1025 let failure_marker = "HTTP 400 Bad Request: deterministic continuation failure";
1026 let leaked_secret = "sk-goal-secret-sentinel-123456";
1027 let failure_message = format!(
1028 "{failure_marker}: {leaked_secret} {}",
1029 "provider detail ".repeat(80)
1030 );
1031 let model = std::sync::Arc::new(FailingGoalModelClient {
1032 calls: std::sync::atomic::AtomicUsize::new(0),
1033 message: failure_message.clone(),
1034 });
1035 let config = goal_custom_route_config();
1036 let engine_config = EngineConfig {
1037 model: "local-model".to_string(),
1038 snapshots_enabled: false,
1039 terminal_chrome_enabled: false,
1040 goal_objective: Some("block a failed continuation truthfully".to_string()),
1041 ..EngineConfig::default()
1042 };
1043 let client: crate::core::model_client::SharedModelClient = model.clone();
1044 let (engine, handle) = Engine::new_with_model_client(engine_config, &config, client);
1045 let goal_state = engine.config.goal_state.clone();
1046 let run_task = tokio::spawn(engine.run());
1047
1048 handle
1049 .send(Op::ContinueGoal {
1050 dynamic_tools: Vec::new(),
1051 engine_schedule_id: None,
1052 })
1053 .await
1054 .expect("queue failing continuation");
1055 let session = tokio::time::timeout(model_turn_event_timeout(), handle.get_session_snapshot())
1056 .await
1057 .expect("failed continuation did not terminalize")
1058 .expect("post-failure session snapshot");
1059
1060 let prompt = match session.system_prompt.expect("blocked system prompt") {
1061 SystemPrompt::Text(text) => text,
1062 SystemPrompt::Blocks(blocks) => blocks
1063 .into_iter()
1064 .map(|block| block.text)
1065 .collect::<Vec<_>>()
1066 .join("\n"),
1067 };
1068 assert!(!prompt.contains("<session_goal>"), "{prompt}");
1069 let snapshot = goal_state.lock().expect("goal lock").snapshot();
1070 assert_eq!(snapshot.status, "blocked");
1071 let blocker = snapshot.blocker.as_deref().expect("failure blocker");
1072 assert!(blocker.contains(failure_marker), "{blocker}");
1073 assert!(blocker.contains("resume the goal"), "{blocker}");
1074 assert!(!blocker.contains(leaked_secret), "{blocker}");
1075 assert!(
1076 blocker.contains(codewhale_config::persistence::REDACTED),
1077 "{blocker}"
1078 );
1079 assert!(
1080 blocker.len() <= GOAL_CONTINUATION_FAILURE_DETAIL_MAX_BYTES + 160,
1081 "failure reason must remain bounded: {} bytes",
1082 blocker.len()
1083 );
1084 assert!(
1085 blocker.len() < failure_message.len(),
1086 "long provider detail must be truncated"
1087 );
1088 assert_eq!(
1089 model.calls.load(std::sync::atomic::Ordering::SeqCst),
1090 1,
1091 "a nonretryable started failure must not dispatch again"
1092 );
1093
1094 let mut starts = 0;
1095 let mut saw_failed_turn = false;
1096 let mut saw_provider_error = false;
1097 let mut saw_blocked_goal = false;
1098 let mut saw_blocked_status = false;
1099 {
1100 let mut events = handle.rx_event.write().await;
1101 while let Ok(event) = events.try_recv() {
1102 match event {
1103 Event::TurnStarted { .. } => starts += 1,
1104 Event::TurnComplete { status, error, .. } => {
1105 assert_eq!(status, TurnOutcomeStatus::Failed);
1106 assert!(
1107 error
1108 .as_deref()
1109 .is_some_and(|message| message.contains(failure_marker)),
1110 "{error:?}"
1111 );
1112 saw_failed_turn = true;
1113 }
1114 Event::Error { envelope, .. } => {
1115 if envelope.message.contains(failure_marker) {
1116 saw_provider_error = true;
1117 }
1118 }
1119 Event::GoalUpdated { snapshot } if snapshot.status == "blocked" => {
1120 saw_blocked_goal = true;
1121 }
1122 Event::Status { message } if message.contains(failure_marker) => {
1123 assert!(message.contains("resume the goal"), "{message}");
1124 saw_blocked_status = true;
1125 }
1126 _ => {}
1127 }
1128 }
1129 }
1130 assert_eq!(starts, 1, "exactly one continuation turn must start");
1131 assert!(saw_failed_turn, "failed turn receipt must remain visible");
1132 assert!(
1133 saw_provider_error,
1134 "provider error event must remain visible"
1135 );
1136 assert!(saw_blocked_goal, "goal must publish its blocked snapshot");
1137 assert!(saw_blocked_status, "bounded failure must remain visible");
1138
1139 handle.send(Op::Shutdown).await.expect("shutdown engine");
1140 run_task.await.expect("engine task");
1141 }
1142
1143 #[tokio::test]
1144 async fn headless_host_drains_existing_engine_completion_inbox_before_exit() {
1145 use crate::llm_client::mock::{MockLlmClient, canned};
1146
1147 let workspace = tempdir().unwrap();
1148 let config = goal_custom_route_config();
1149 let mock = Arc::new(MockLlmClient::new(vec![canned::simple_text_turn(
1150 "child evidence integrated by the existing Engine",
1151 )]));
1152 let (engine, handle) = Engine::new_with_model_client(
1153 EngineConfig {
1154 model: "local-model".into(),
1155 terminal_chrome_enabled: false,
1156 ..deterministic_engine_config(workspace.path())
1157 },
1158 &config,
1159 mock.clone(),
1160 );
1161 assert!(engine.subagent_settlement_snapshot().await.is_settled());
1162 // Reproduce the host boundary: the parent already ended, and a terminal
1163 // child's receipt is waiting for the Engine's normal idle fan-in path.
1164 engine
1165 .tx_event
1166 .send(Event::TurnComplete {
1167 usage: Usage::default(),
1168 parent_route_usage: Usage::default(),
1169 routed_usage_dropped_records: 0,
1170 status: TurnOutcomeStatus::Completed,
1171 error: None,
1172 tool_catalog: None,
1173 base_url: None,
1174 })
1175 .await
1176 .unwrap();
1177 engine
1178 .tx_subagent_completion
1179 .try_send(SubAgentCompletion {
1180 owner_session_id: engine.session.id.clone(),
1181 agent_id: "headless-settled-child".into(),
1182 payload: "bounded local fixture evidence".into(),
1183 })
1184 .unwrap();
1185 let pending = engine.subagent_settlement_snapshot().await;
1186 assert_eq!(pending.running_children, 0);
1187 assert_eq!(pending.pending_completions, 1);
1188 assert!(
1189 !pending.is_settled(),
1190 "terminal child alone cannot release the host"
1191 );
1192
1193 let run = tokio::spawn(engine.run());
1194 let mut events = crate::exec_agent::ExecAgentEvents::new(
1195 handle.clone(),
1196 Instant::now() + model_turn_event_timeout(),
1197 );
1198 let mut content = String::new();
1199 let mut starts = 0;
1200 tokio::time::timeout(model_turn_event_timeout(), async {
1201 loop {
1202 match events.next().await.expect("host event") {
1203 Event::TurnStarted { .. } => starts += 1,
1204 Event::MessageDelta { content: delta, .. } => content.push_str(&delta),
1205 Event::TurnComplete { status, .. } => {
1206 assert_eq!(status, TurnOutcomeStatus::Completed);
1207 break;
1208 }
1209 _ => {}
1210 }
1211 }
1212 })
1213 .await
1214 .expect("bounded headless settlement");
1215 assert_eq!(starts, 1, "fan-in uses exactly one existing Engine turn");
1216 assert_eq!(mock.call_count(), 1);
1217 assert!(content.contains("child evidence integrated"));
1218 assert!(!handle.is_cancelled());
1219 handle.send(Op::Shutdown).await.unwrap();
1220 tokio::time::timeout(Duration::from_secs(3), run)
1221 .await
1222 .unwrap()
1223 .unwrap();
1224 }
1225
1226 #[tokio::test]
1227 async fn host_managed_engine_does_not_self_dispatch_goal_continuation() {
1228 use crate::llm_client::mock::{MockLlmClient, canned};
1229
1230 let mut custom = HashMap::new();
1231 custom.insert(
1232 "custom-a".to_string(),
1233 crate::config::ProviderConfig {
1234 kind: Some("openai-compatible".to_string()),
1235 base_url: Some("http://127.0.0.1:18181/v1".to_string()),
1236 model: Some("local-model".to_string()),
1237 api_key: Some("local-test-key".to_string()),
1238 ..crate::config::ProviderConfig::default()
1239 },
1240 );
1241 let config = Config {
1242 provider: Some("custom-a".to_string()),
1243 providers: Some(crate::config::ProvidersConfig {
1244 custom,
1245 ..crate::config::ProvidersConfig::default()
1246 }),
1247 ..Config::default()
1248 };
1249 let runtime_services = crate::tools::spec::RuntimeToolServices {
1250 active_thread_id: Some("thr_host_managed".to_string()),
1251 ..crate::tools::spec::RuntimeToolServices::default()
1252 };
1253 let engine_config = EngineConfig {
1254 max_steps: 0,
1255 snapshots_enabled: false,
1256 terminal_chrome_enabled: false,
1257 goal_objective: Some("keep going".to_string()),
1258 runtime_services,
1259 ..EngineConfig::default()
1260 };
1261 let mock = Arc::new(MockLlmClient::new(vec![canned::simple_text_turn(
1262 "host-managed turn complete",
1263 )]));
1264 let client: crate::core::model_client::SharedModelClient = mock.clone();
1265 let (engine, handle) = Engine::new_with_model_client(engine_config, &config, client);
1266 let run_task = tokio::spawn(engine.run());
1267
1268 handle
1269 .send(Op::SendMessage(TurnSpec {
1270 max_output_tokens: None,
1271 content: "one host-owned turn".to_string(),
1272 images: Vec::new(),
1273 mode: AppMode::Agent,
1274 route: resolved_route_for_test(&config, "local-model"),
1275 compaction: Box::new(CompactionConfig::default()),
1276 initial_routed_usage: Box::default(),
1277 goal_objective: Some("keep going".to_string()),
1278 goal_token_budget: None,
1279 goal_status: crate::tools::goal::GoalStatus::Active,
1280 reasoning_effort: None,
1281 reasoning_effort_auto: false,
1282 auto_model: false,
1283 allow_shell: false,
1284 trust_mode: false,
1285 auto_approve: false,
1286 approval_mode: ApprovalMode::Suggest,
1287 translation_enabled: false,
1288 allowed_tools: None,
1289 dynamic_tools: Vec::new(),
1290 hook_executor: None,
1291 verbosity: None,
1292 provenance: UserInputProvenance::ExternalUser,
1293 submission_id: None,
1294 }))
1295 .await
1296 .expect("send host-owned goal turn");
1297
1298 let mut starts = 0;
1299 loop {
1300 let event = tokio::time::timeout(Duration::from_secs(3), async {
1301 handle.rx_event.write().await.recv().await
1302 })
1303 .await
1304 .expect("host engine event timeout")
1305 .expect("host engine event");
1306 match event {
1307 Event::TurnStarted { .. } => starts += 1,
1308 Event::TurnComplete { .. } => break,
1309 _ => {}
1310 }
1311 }
1312 assert_eq!(starts, 1);
1313 assert_eq!(
1314 mock.call_count(),
1315 1,
1316 "the host-owned turn runs exactly once"
1317 );
1318 assert!(
1319 tokio::time::timeout(Duration::from_millis(200), async {
1320 handle.rx_event.write().await.recv().await
1321 })
1322 .await
1323 .is_err(),
1324 "a hosted engine must wait for an explicit durable turn claim"
1325 );
1326
1327 handle.send(Op::Shutdown).await.expect("shutdown engine");
1328 run_task.await.expect("engine task");
1329 }
1330
1331 #[tokio::test]
1332 async fn cancellation_during_blocked_idle_handoff_survives_turn_admission() {
1333 use crate::llm_client::mock::{MockLlmClient, canned};
1334
1335 for child_completion in [true, false] {
1336 let workspace = tempdir().unwrap();
1337 let config = goal_custom_route_config();
1338 let mock = Arc::new(MockLlmClient::new(vec![canned::simple_text_turn(
1339 "must not dispatch after cancellation",
1340 )]));
1341 let (mut engine, mut handle) = Engine::new_with_model_client(
1342 EngineConfig {
1343 model: "local-model".into(),
1344 terminal_chrome_enabled: false,
1345 ..deterministic_engine_config(workspace.path())
1346 },
1347 &config,
1348 mock.clone(),
1349 );
1350 // Force the handoff to stop at its status send after its initial
1351 // cancellation check and before handle_send_message admits a turn.
1352 let (tx, rx) = tokio::sync::mpsc::channel(1);
1353 engine.tx_event = tx;
1354 handle.rx_event = Arc::new(RwLock::new(rx));
1355 engine.tx_event.send(Event::status("full")).await.unwrap();
1356 let completion = SubAgentCompletion {
1357 owner_session_id: engine.session.id.clone(),
1358 agent_id: "cancel-race-child".into(),
1359 payload: "retained-after-cancel-race".into(),
1360 };
1361 let mut wake: std::pin::Pin<Box<dyn std::future::Future<Output = ()> + '_>> =
1362 if child_completion {
1363 Box::pin(engine.handle_idle_subagent_completion(completion))
1364 } else {
1365 Box::pin(engine.handle_idle_shell_completion_wake())
1366 };
1367 assert!(
1368 tokio::time::timeout(Duration::from_millis(5), &mut wake)
1369 .await
1370 .is_err(),
1371 "the full event channel must hold the handoff before admission"
1372 );
1373 handle.cancel_with_reason(CancelReason::External);
1374 handle.rx_event.write().await.try_recv().unwrap();
1375 tokio::time::timeout(Duration::from_secs(2), &mut wake)
1376 .await
1377 .expect("cancelled handoff returns without a provider call");
1378 drop(wake);
1379 assert_eq!(mock.call_count(), 0);
1380 assert!(handle.is_cancelled());
1381 assert!(engine.delivered_subagent_completion_ids.is_empty());
1382 assert_eq!(
1383 engine.rx_subagent_completion.len(),
1384 usize::from(child_completion)
1385 );
1386 if child_completion {
1387 assert!(
1388 engine
1389 .rx_subagent_completion
1390 .try_recv()
1391 .unwrap()
1392 .payload
1393 .contains("retained-after-cancel-race")
1394 );
1395 }
1396 // A new explicit user action still receives a fresh turn control.
1397 let _turn = engine.begin_turn_control();
1398 assert!(!handle.is_cancelled());
1399 let existing_child = engine.cancel_token.child_token();
1400 drop(_turn);
1401 let _automatic =
1402 engine.begin_turn_control_for_provenance(UserInputProvenance::SubAgentHandoff);
1403 handle.cancel();
1404 assert!(
1405 existing_child.is_cancelled(),
1406 "stopping an automatic continuation must also stop earlier request siblings"
1407 );
1408 }
1409 }
1410
1411 #[tokio::test]
1412 async fn cancelled_parent_defers_child_receipts_until_an_explicit_turn() {
1413 use crate::llm_client::mock::{MockLlmClient, canned};
1414
1415 for reason in [CancelReason::User, CancelReason::External] {
1416 let workspace = tempdir().unwrap();
1417 let config = Config::default();
1418 let mock = Arc::new(MockLlmClient::new(vec![canned::simple_text_turn(
1419 "Explicit continuation completed.",
1420 )]));
1421 let (mut engine, handle) = Engine::new_with_model_client(
1422 deterministic_engine_config(workspace.path()),
1423 &config,
1424 mock.clone(),
1425 );
1426 handle.cancel_with_reason(reason);
1427 // Exercise a completion selected just before cancellation arrived.
1428 engine
1429 .handle_idle_subagent_completion(SubAgentCompletion {
1430 owner_session_id: engine.session.id.clone(),
1431 agent_id: "cancelled-worker".into(),
1432 payload: "parked-child-evidence".into(),
1433 })
1434 .await;
1435 assert_eq!(
1436 mock.call_count(),
1437 0,
1438 "cancellation must forbid a model wake"
1439 );
1440 assert!(engine.delivered_subagent_completion_ids.is_empty());
1441 assert_eq!(engine.rx_subagent_completion.len(), 1);
1442 assert!(
1443 tokio::time::timeout(Duration::from_millis(50), engine.next_run_input(false))
1444 .await
1445 .is_err(),
1446 "the idle loop must leave the receipt queued without spinning"
1447 );
1448
1449 let run = tokio::spawn(engine.run());
1450 handle
1451 .send(external_user_message_op(
1452 "Continue explicitly",
1453 AppMode::Agent,
1454 &config,
1455 ))
1456 .await
1457 .unwrap();
1458 tokio::time::timeout(Duration::from_secs(5), async {
1459 while let Some(event) = handle.rx_event.write().await.recv().await {
1460 if matches!(event, Event::TurnComplete { .. }) {
1461 break;
1462 }
1463 }
1464 })
1465 .await
1466 .expect("explicit turn completes");
1467 assert_eq!(mock.call_count(), 1);
1468 let snapshot = handle.get_session_snapshot().await.unwrap();
1469 assert!(
1470 serde_json::to_string(&snapshot.messages)
1471 .unwrap()
1472 .contains("parked-child-evidence"),
1473 "the next explicit turn must retain the completion receipt"
1474 );
1475 handle.send(Op::Shutdown).await.unwrap();
1476 run.await.unwrap();
1477 }
1478 }
1478 lines RUST