返回 CodeWhale
test_cases_01.rs
根目录 / crates / tui / src / core / engine / tests / test_cases_01.rs
1 const WORKING_SET_SUMMARY_MARKER: &str = "## Repo Working Set";
2
3 #[tokio::test]
4 async fn event_capacity_cancel_before_admission_has_no_started_turn_and_keeps_classifier_cost() {
5 use crate::llm_client::mock::{MockLlmClient, canned};
6 let _cost = crate::cost_status::test_scope();
7 let workspace = tempdir().unwrap();
8 let config = Config::default();
9 let mock = Arc::new(MockLlmClient::new(vec![canned::simple_text_turn(
10 "must not run",
11 )]));
12 let mut engine_config = deterministic_engine_config(workspace.path());
13 engine_config.features.disable(Feature::Mcp);
14 let (mut engine, handle) = Engine::new_with_model_client(engine_config, &config, mock.clone());
15 let (tx, rx) = mpsc::channel(1);
16 engine.tx_event = tx;
17 let mut handle = handle;
18 handle.rx_event = Arc::new(RwLock::new(rx));
19 engine
20 .tx_event
21 .try_send(Event::status("existing idle receipt"))
22 .unwrap();
23 let mut op = external_user_message_op("not admitted", AppMode::Agent, &config);
24 let Op::SendMessage(spec) = &mut op else {
25 unreachable!()
26 };
27 spec.initial_routed_usage
28 .records
29 .push(crate::cost_status::RuntimeUsageRecord {
30 source_id: "event-capacity:classifier".into(),
31 usage: crate::cost_status::EffectiveRouteUsage {
32 route: crate::cost_status::EffectiveRouteEnvelope::capture(
33 None,
34 ProviderKind::Openai,
35 "openai",
36 "classifier",
37 None,
38 chrono::Utc::now(),
39 ),
40 usage: Usage {
41 input_tokens: 7,
42 output_tokens: 3,
43 ..Usage::default()
44 },
45 },
46 });
47 let controls = Arc::clone(&engine.turn_controls);
48 handle.send(op).await.unwrap();
49 let task = tokio::spawn(engine.run());
50 tokio::time::timeout(Duration::from_secs(2), async {
51 loop {
52 if controls.lock().unwrap().active.is_some() {
53 break;
54 }
55 tokio::task::yield_now().await;
56 }
57 })
58 .await
59 .expect("queued op enters the existing control scope");
60 handle.cancel();
61 // The oneshot snapshot can settle while the event receiver remains full.
62 let snapshot = tokio::time::timeout(Duration::from_secs(2), handle.get_session_snapshot())
63 .await
64 .expect("cancel releases admission reservation")
65 .unwrap();
66 assert!(
67 snapshot.messages.is_empty(),
68 "no session mutation before admission"
69 );
70 assert_eq!(mock.call_count(), 0);
71 assert!(controls.lock().unwrap().active.is_none());
72 let cost = crate::cost_status::drain();
73 assert!(
74 cost.usage_source_fingerprints
75 .contains(&crate::cost_status::usage_source_fingerprint(
76 "event-capacity:classifier"
77 ),),
78 "already billed classifier work remains accounted"
79 );
80 let mut events = handle.rx_event.write().await;
81 assert!(matches!(events.try_recv(), Ok(Event::Status { .. })));
82 assert!(
83 events.try_recv().is_err(),
84 "unadmitted work has no fabricated terminal event"
85 );
86 drop(events);
87 handle.send(Op::Shutdown).await.unwrap();
88 tokio::time::timeout(Duration::from_secs(2), task)
89 .await
90 .unwrap()
91 .unwrap();
92 }
93
94 #[tokio::test]
95 async fn event_capacity_admitted_full_channel_settles_once_with_partial_usage_without_drain() {
96 use crate::llm_client::mock::{MockLlmClient, canned};
97 let _cost = crate::cost_status::test_scope();
98 let workspace = tempdir().unwrap();
99 let config = Config::default();
100 let mock = Arc::new(MockLlmClient::new(Vec::new()));
101 let mut engine_config = deterministic_engine_config(workspace.path());
102 engine_config.features.disable(Feature::Mcp);
103 let (engine, handle) = Engine::new_with_model_client(engine_config, &config, mock.clone());
104 let tx = engine.tx_event.clone();
105 let entered = Arc::new(tokio::sync::Notify::new());
106 let filled = Arc::clone(&entered);
107 mock.push_factory(move |_| {
108 while tx.try_send(Event::status("backpressure receipt")).is_ok() {}
109 filled.notify_one();
110 vec![
111 canned::message_start("billed-before-cancel"),
112 canned::message_delta(
113 "end_turn",
114 Some(Usage {
115 input_tokens: 11,
116 output_tokens: 5,
117 ..Usage::default()
118 }),
119 ),
120 canned::text_block_start(0),
121 canned::text_delta(0, "must not render after cancellation"),
122 canned::message_stop(),
123 ]
124 });
125 let controls = Arc::clone(&engine.turn_controls);
126 let task = tokio::spawn(engine.run());
127 handle
128 .send(external_user_message_op(
129 "admitted",
130 AppMode::Agent,
131 &config,
132 ))
133 .await
134 .unwrap();
135 tokio::time::timeout(Duration::from_secs(2), entered.notified())
136 .await
137 .unwrap();
138 handle.cancel();
139 let snapshot = tokio::time::timeout(Duration::from_secs(2), handle.get_session_snapshot())
140 .await
141 .expect("admitted cancellation settles before any event drain")
142 .unwrap();
143 assert_eq!(mock.call_count(), 1);
144 assert_eq!(
145 snapshot.total_tokens, 16,
146 "known partial provider usage survives cancellation"
147 );
148 assert!(controls.lock().unwrap().active.is_none());
149 let mut events = handle.rx_event.write().await;
150 let mut started = 0;
151 let mut completed = 0;
152 let mut terminal_was_last = false;
153 while let Ok(event) = events.try_recv() {
154 assert!(
155 !terminal_was_last,
156 "completion stays after prior accepted observations"
157 );
158 match event {
159 Event::TurnStarted { .. } => started += 1,
160 Event::TurnComplete {
161 status,
162 usage,
163 parent_route_usage,
164 error,
165 ..
166 } => {
167 completed += 1;
168 terminal_was_last = true;
169 assert_eq!(status, TurnOutcomeStatus::Interrupted);
170 assert!(
171 error.is_none(),
172 "cancellation is not a provider failure: {error:?}"
173 );
174 assert_eq!((usage.input_tokens, usage.output_tokens), (11, 5));
175 assert_eq!(usage, parent_route_usage);
176 }
177 Event::MessageDelta { content, .. } => assert!(!content.contains("must not render")),
178 _ => {}
179 }
180 }
181 assert_eq!((started, completed), (1, 1));
182 assert!(terminal_was_last);
183 drop(events);
184 handle.send(Op::Shutdown).await.unwrap();
185 tokio::time::timeout(Duration::from_secs(2), task)
186 .await
187 .unwrap()
188 .unwrap();
189 }
190
191 #[tokio::test]
192 async fn event_capacity_invalid_images_release_queued_control_and_keep_classifier_receipt() {
193 use crate::llm_client::mock::MockLlmClient;
194 for full in [false, true] {
195 let _cost = crate::cost_status::test_scope();
196 let workspace = tempdir().unwrap();
197 let config = Config::default();
198 let mock = Arc::new(MockLlmClient::new(Vec::new()));
199 let mut engine_config = deterministic_engine_config(workspace.path());
200 engine_config.features.disable(Feature::Mcp);
201 let (engine, handle) = Engine::new_with_model_client(engine_config, &config, mock.clone());
202 if full {
203 while engine.tx_event.try_send(Event::status("occupied")).is_ok() {}
204 }
205 let controls = Arc::clone(&engine.turn_controls);
206 let mut op = external_user_message_op("invalid attachment", AppMode::Agent, &config);
207 let Op::SendMessage(spec) = &mut op else {
208 unreachable!()
209 };
210 spec.images
211 .push(codewhale_protocol::runtime::RuntimeImageInput {
212 mime: "image/png".into(),
213 data_base64: "invalid@base64".into(),
214 });
215 spec.initial_routed_usage
216 .records
217 .push(crate::cost_status::RuntimeUsageRecord {
218 source_id: "event-capacity:invalid-image-classifier".into(),
219 usage: crate::cost_status::EffectiveRouteUsage {
220 route: crate::cost_status::EffectiveRouteEnvelope::capture(
221 None,
222 ProviderKind::Openai,
223 "openai",
224 "classifier",
225 None,
226 chrono::Utc::now(),
227 ),
228 usage: Usage {
229 input_tokens: 7,
230 output_tokens: 3,
231 ..Usage::default()
232 },
233 },
234 });
235 handle.send(op).await.unwrap();
236 let task = tokio::spawn(engine.run());
237 if full {
238 tokio::time::timeout(Duration::from_secs(2), async {
239 loop {
240 if controls.lock().unwrap().active.is_some() {
241 break;
242 }
243 tokio::task::yield_now().await;
244 }
245 })
246 .await
247 .unwrap();
248 handle.cancel();
249 }
250 let snapshot = tokio::time::timeout(Duration::from_secs(2), handle.get_session_snapshot())
251 .await
252 .expect("rejected input releases the current control")
253 .unwrap();
254 assert!(snapshot.messages.is_empty());
255 assert_eq!(mock.call_count(), 0);
256 assert!(controls.lock().unwrap().active.is_none());
257 let cost = crate::cost_status::drain();
258 assert!(cost.usage_source_fingerprints.contains(
259 &crate::cost_status::usage_source_fingerprint(
260 "event-capacity:invalid-image-classifier"
261 ),
262 ));
263 let mut events = handle.rx_event.write().await;
264 let mut invalid_errors = 0;
265 while let Ok(event) = events.try_recv() {
266 assert!(!matches!(
267 event,
268 Event::TurnStarted { .. } | Event::TurnComplete { .. }
269 ));
270 if let Event::Error { envelope, .. } = event {
271 assert_eq!(envelope.code, "image_input_invalid");
272 invalid_errors += 1;
273 }
274 }
275 assert_eq!(
276 invalid_errors,
277 usize::from(!full),
278 "a cancelled full-channel rejection is not fabricated as delivered"
279 );
280 drop(events);
281 handle.send(Op::Shutdown).await.unwrap();
282 tokio::time::timeout(Duration::from_secs(2), task)
283 .await
284 .unwrap()
285 .unwrap();
286 }
287 }
288
289 #[tokio::test]
290 async fn event_capacity_releases_all_admitted_senders_and_keeps_idle_receipts_lossless() {
291 let workspace = tempdir().unwrap();
292 let (mut engine, _handle) = Engine::new(
293 deterministic_engine_config(workspace.path()),
294 &Config::default(),
295 );
296 let (tx, mut rx) = mpsc::channel(1);
297 engine.tx_event = tx;
298 engine.tx_event.try_send(Event::status("occupied")).unwrap();
299 let turn = engine.begin_turn_control();
300 let mut senders =
301 Box::pin(futures_util::future::join_all((0..8).map(|_| {
302 engine.send_event(Event::status("admitted observation"))
303 })));
304 assert!(
305 tokio::time::timeout(Duration::from_millis(20), &mut senders)
306 .await
307 .is_err()
308 );
309 engine.cancel_token.cancel();
310 let outcomes = tokio::time::timeout(Duration::from_secs(1), &mut senders)
311 .await
312 .unwrap();
313 assert_eq!(outcomes, vec![Err(streaming::EventSendError::Cancelled); 8]);
314 drop(senders);
315 assert!(matches!(rx.try_recv(), Ok(Event::Status { .. })));
316 // A cancelled turn still delivers its usage/status receipts when they
317 // fit immediately. The stream suffix remains strictly cancelled.
318 engine
319 .send_event(Event::status(
320 "cancelled turn receipt with available capacity",
321 ))
322 .await
323 .unwrap();
324 assert!(
325 matches!(rx.try_recv(), Ok(Event::Status { message, .. }) if message == "cancelled turn receipt with available capacity")
326 );
327 assert!(
328 !engine
329 .send_stream_event(Event::status("forbidden stream suffix"))
330 .await
331 );
332 assert!(rx.try_recv().is_err());
333 engine
334 .tx_event
335 .try_send(Event::status("occupied before idle receipt"))
336 .unwrap();
337 drop(turn);
338 let mut idle = Box::pin(engine.send_event(Event::status("idle receipt after cancelled turn")));
339 assert!(
340 tokio::time::timeout(Duration::from_millis(20), &mut idle)
341 .await
342 .is_err()
343 );
344 assert!(matches!(rx.try_recv(), Ok(Event::Status { .. })));
345 tokio::time::timeout(Duration::from_secs(1), &mut idle)
346 .await
347 .unwrap()
348 .unwrap();
349 drop(idle);
350 assert!(
351 matches!(rx.try_recv(), Ok(Event::Status { message, .. }) if message == "idle receipt after cancelled turn")
352 );
353 assert!(rx.try_recv().is_err());
354 }
355
356 #[tokio::test]
357 async fn event_capacity_nested_vm_fanout_is_bounded_per_program_not_per_turn() {
358 struct Counter(std::sync::atomic::AtomicUsize);
359 #[async_trait::async_trait]
360 impl codewhale_workflow_js::ToolInvoker for Counter {
361 async fn invoke(
362 &self,
363 _: codewhale_workflow_js::ToolCallRequest,
364 ) -> Result<codewhale_workflow_js::ToolCallResponse, codewhale_workflow_js::DriverError>
365 {
366 self.0.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
367 tokio::task::yield_now().await;
368 Ok(codewhale_workflow_js::ToolCallResponse {
369 ok: true,
370 result: json!(null),
371 })
372 }
373 }
374 let invoker = Arc::new(Counter(std::sync::atomic::AtomicUsize::new(0)));
375 for program in 1..=2 {
376 let result = codewhale_workflow_js::WorkflowVm::new().run_tools_script(
377 r#"const calls = await Promise.allSettled(Array.from({length:256}, () => tools.call('read', {})));
378 return {ok:calls.filter(x=>x.status==='fulfilled').length,
379 rejected:calls.filter(x=>x.status==='rejected' && x.reason.kind==='admission').length};"#,
380 json!(null),
381 Arc::new(codewhale_workflow_js::testing::FakeDriver::new()),
382 invoker.clone(), codewhale_workflow_js::WorkflowRunCancel::new(),
383 ).await.unwrap();
384 assert_eq!(result, json!({"ok":50,"rejected":206}));
385 assert_eq!(
386 invoker.0.load(std::sync::atomic::Ordering::SeqCst),
387 50 * program
388 );
389 }
390 }
391
392 #[tokio::test]
393 async fn event_capacity_admitted_user_shell_cancels_before_side_effect_and_settles_without_drain() {
394 let workspace = tempdir().unwrap();
395 let marker = workspace.path().join("shell-must-not-start.txt");
396 let mut engine_config = deterministic_engine_config(workspace.path());
397 engine_config.features.disable(Feature::Mcp);
398 let (engine, handle) = Engine::new(engine_config, &Config::default());
399 // Leave precisely the two lifecycle slots free. TurnStarted fills one;
400 // the reserved terminal owns the other. The next tool observation must
401 // wait before the human-provenance command can reach its executor.
402 for _ in 0..engine.tx_event.max_capacity() - 2 {
403 engine
404 .tx_event
405 .try_send(Event::status("prior receipt"))
406 .unwrap();
407 }
408 let controls = Arc::clone(&engine.turn_controls);
409 handle
410 .send(Op::RunShellCommand {
411 command: format!("echo must-not-run > \"{}\"", marker.display()),
412 mode: AppMode::Agent,
413 allow_shell: true,
414 trust_mode: true,
415 auto_approve: true,
416 approval_mode: ApprovalMode::Bypass,
417 })
418 .await
419 .unwrap();
420 let task = tokio::spawn(engine.run());
421 tokio::time::timeout(Duration::from_secs(2), async {
422 loop {
423 if controls.lock().unwrap().active.is_some() {
424 break;
425 }
426 tokio::task::yield_now().await;
427 }
428 })
429 .await
430 .unwrap();
431 handle.cancel();
432 tokio::time::timeout(Duration::from_secs(2), handle.get_session_snapshot())
433 .await
434 .expect("shell cancellation settles before any event drain")
435 .unwrap();
436 assert!(
437 !marker.exists(),
438 "cancellation preserves the execution gate"
439 );
440 assert!(controls.lock().unwrap().active.is_none());
441 let mut events = handle.rx_event.write().await;
442 let mut started = 0;
443 let mut completed = 0;
444 let mut terminal_was_last = false;
445 while let Ok(event) = events.try_recv() {
446 assert!(!terminal_was_last);
447 match event {
448 Event::TurnStarted { .. } => started += 1,
449 Event::TurnComplete { status, usage, .. } => {
450 completed += 1;
451 terminal_was_last = true;
452 assert_eq!(status, TurnOutcomeStatus::Interrupted);
453 assert_eq!(usage, Usage::default());
454 }
455 _ => {}
456 }
457 }
458 assert_eq!((started, completed), (1, 1));
459 assert!(terminal_was_last);
460 drop(events);
461 handle.send(Op::Shutdown).await.unwrap();
462 tokio::time::timeout(Duration::from_secs(2), task)
463 .await
464 .unwrap()
465 .unwrap();
466 }
467
468 #[tokio::test]
469 async fn event_capacity_cancelled_repl_child_keeps_unknown_cost_and_discards_kernel() {
470 use crate::llm_client::mock::{MockLlmClient, canned};
471 struct ReplClient {
472 inner: MockLlmClient,
473 stream_requests: std::sync::atomic::AtomicUsize,
474 child_entered: Arc<tokio::sync::Notify>,
475 child_dropped: Arc<std::sync::atomic::AtomicBool>,
476 }
477 #[async_trait::async_trait]
478 impl crate::core::model_client::ModelClient for ReplClient {
479 fn provider_name(&self) -> &str {
480 "backpressure-repl-fixture"
481 }
482 fn model(&self) -> &str {
483 "mock-model"
484 }
485 async fn create_message(
486 &self,
487 _: codewhale_models::MessageRequest,
488 ) -> anyhow::Result<codewhale_models::MessageResponse> {
489 anyhow::bail!("fixture expects canonical streaming requests")
490 }
491 async fn create_message_stream(
492 &self,
493 request: codewhale_models::MessageRequest,
494 ) -> anyhow::Result<crate::llm_client::StreamEventBox> {
495 if self
496 .stream_requests
497 .fetch_add(1, std::sync::atomic::Ordering::SeqCst)
498 == 0
499 {
500 return crate::core::model_client::ModelClient::create_message_stream(
501 &self.inner,
502 request,
503 )
504 .await;
505 }
506 let _drop = DropSignal(Arc::clone(&self.child_dropped));
507 self.child_entered.notify_one();
508 std::future::pending().await
509 }
510 async fn health_check(&self) -> anyhow::Result<bool> {
511 Ok(true)
512 }
513 }
514 let _cost = crate::cost_status::test_scope();
515 let workspace = tempdir().unwrap();
516 let entered = Arc::new(tokio::sync::Notify::new());
517 let dropped = Arc::new(std::sync::atomic::AtomicBool::new(false));
518 let mut response = canned::simple_text_turn(
519 "```repl\nchild = sub_query('hang until cancelled')\nfinalize(child)\n```",
520 );
521 for event in &mut response {
522 if let codewhale_models::StreamEvent::MessageDelta { usage, .. } = event {
523 *usage = Some(Usage {
524 input_tokens: 11,
525 output_tokens: 5,
526 ..Usage::default()
527 });
528 }
529 }
530 let client = Arc::new(ReplClient {
531 inner: MockLlmClient::new(vec![response]),
532 stream_requests: std::sync::atomic::AtomicUsize::new(0),
533 child_entered: Arc::clone(&entered),
534 child_dropped: Arc::clone(&dropped),
535 });
536 let api_config = rlm_host::fixture_config("mock-model");
537 let (mut engine, handle) = Engine::new_with_model_client(
538 EngineConfig {
539 model: "mock-model".into(),
540 ..deterministic_engine_config(workspace.path())
541 },
542 &api_config,
543 client.clone(),
544 );
545 rlm_host::install_fixture_route(&mut engine);
546 engine.session.auto_approve = true;
547 engine.session.add_message(Message {
548 role: Role::User,
549 content: vec![ContentBlock::Text {
550 text: "run cancellable REPL".into(),
551 cache_control: None,
552 }],
553 });
554 let turn_guard = engine.begin_turn_control();
555 let mut turn = TurnContext::new(4);
556 let registry = crate::tools::ToolRegistry::new(rlm_host::admitted_context(&engine, &turn.id));
557 let policy = test_tool_surface(
558 &engine,
559 registry,
560 Some(vec![catalog_tool(CODE_EXECUTION_TOOL_NAME)]),
561 AppMode::Agent,
562 );
563 let tx = engine.tx_event.clone();
564 let mut run = Box::pin(engine.run_turn(&mut turn, policy, None, None));
565 tokio::time::timeout(Duration::from_secs(10), async {
566 tokio::select! {
567 result = &mut run => panic!("REPL did not enter the child provider request: {result:?}"),
568 () = entered.notified() => {},
569 }
570 }).await.expect("actual Python kernel dispatches the pending child request");
571 while tx
572 .try_send(Event::status("full during nested REPL"))
573 .is_ok()
574 {}
575 handle.cancel();
576 let (status, error) = tokio::time::timeout(Duration::from_secs(2), &mut run)
577 .await
578 .expect("cancel drops the pending REPL round even with a full queue");
579 drop(run);
580 assert_eq!(status, TurnOutcomeStatus::Interrupted, "{error:?}");
581 assert!(error.is_none());
582 assert!(
583 dropped.load(std::sync::atomic::Ordering::SeqCst),
584 "nested provider future is released"
585 );
586 assert!(
587 engine.repl_kernel.is_none(),
588 "a cancelled round cannot preserve an executing process"
589 );
590 assert_eq!((turn.usage.input_tokens, turn.usage.output_tokens), (11, 5));
591 let cost = crate::cost_status::drain();
592 // The child provider future never returned a response; its canonical
593 // dispatch guard records an unknown outcome, not success without usage.
594 assert!(
595 cost.unpriced_reasons
596 .contains(crate::cost_status::RuntimeUsageMissingReason::RequestOutcomeUnknown.label()),
597 "pending child usage stays unknown, never zero"
598 );
599 assert_eq!(cost.priced_turns, 0);
600 assert_eq!(cost.unpriced_turns, 1);
601 assert_eq!(cost.missing_usage_sources.len(), 1);
602 assert!(cost.missing_usage_sources.values().all(|coverage| {
603 coverage.reason == crate::cost_status::RuntimeUsageMissingReason::RequestOutcomeUnknown
604 && coverage.money_metered
605 }));
606 assert!(cost.resolved_missing_usage_sources.is_empty());
607 assert_eq!(
608 client.inner.call_count(),
609 1,
610 "no next root provider request after cancellation"
611 );
612 drop(turn_guard);
613 }
614 #[tokio::test]
615 async fn event_capacity_cancelled_parallel_tool_keeps_completed_span_and_call_when_available() {
616 use crate::llm_client::mock::{MockLlmClient, canned};
617 use crate::tools::spec::{
618 ApprovalRequirement, PreparedToolCall, ResourceClaim, ToolCapability, ToolSpec,
619 };
620 use codewhale_protocol::engine_owner::OwnerOperationOutcome;
621 use std::sync::atomic::{AtomicUsize, Ordering};
622
623 struct CompleteThenCancel {
624 cancel: tokio_util::sync::CancellationToken,
625 tx: mpsc::Sender<Event>,
626 fill_queue: bool,
627 executed: Arc<AtomicUsize>,
628 }
629 #[async_trait::async_trait]
630 impl ToolSpec for CompleteThenCancel {
631 // Registered under the canonical read identity so the engine's
632 // central resource authority (not this fixture) grants the disjoint
633 // ReadPath claims that form one real parallel chunk.
634 fn name(&self) -> &str {
635 "read_file"
636 }
637 fn description(&self) -> &str {
638 "Finish an observed operation before firing its turn cancellation token."
639 }
640 fn input_schema(&self) -> Value {
641 json!({"type": "object"})
642 }
643 fn capabilities(&self) -> Vec<ToolCapability> {
644 vec![ToolCapability::ReadOnly]
645 }
646 fn supports_parallel(&self) -> bool {
647 true
648 }
649 fn prepare(
650 &self,
651 input: Value,
652 context: &ToolContext,
653 ) -> Result<PreparedToolCall, ToolError> {
654 let path = input["path"].as_str().expect("fixture path");
655 Ok(PreparedToolCall {
656 name: self.name().to_string(),
657 description: self.description().to_string(),
658 read_only: true,
659 supports_parallel: true,
660 starts_detached: false,
661 approval: ApprovalRequirement::Auto,
662 // Replaced by `registered_resource_claims`; kept honest anyway.
663 resources: vec![ResourceClaim::ReadPath(context.workspace.join(path))],
664 input,
665 })
666 }
667 async fn execute(&self, _: Value, _: &ToolContext) -> Result<ToolResult, ToolError> {
668 self.executed.fetch_add(1, Ordering::SeqCst);
669 if self.fill_queue {
670 while self
671 .tx
672 .try_send(Event::status("full at tool completion"))
673 .is_ok()
674 {}
675 }
676 self.cancel.cancel();
677 Ok(ToolResult::success("completed before cancellation"))
678 }
679 }
680
681 for fill_queue in [false, true] {
682 let workspace = tempdir().unwrap();
683 for path in ["one", "two"] {
684 std::fs::write(workspace.path().join(path), path).unwrap();
685 }
686 let mock = Arc::new(MockLlmClient::new(vec![
687 tool_batch_turn(&[
688 ("one", "read_file", r#"{"path":"one"}"#),
689 ("two", "read_file", r#"{"path":"two"}"#),
690 ]),
691 canned::simple_text_turn("must not run after cancellation"),
692 ]));
693 let (mut engine, handle) = Engine::new_with_model_client(
694 deterministic_engine_config(workspace.path()),
695 &Config::default(),
696 mock.clone(),
697 );
698 let turn_guard = engine.begin_turn_control();
699 let executed = Arc::new(AtomicUsize::new(0));
700 let mut registry = crate::tools::ToolRegistry::new(ToolContext::new(workspace.path()));
701 registry.register(Arc::new(CompleteThenCancel {
702 cancel: engine.cancel_token.clone(),
703 tx: engine.tx_event.clone(),
704 fill_queue,
705 executed: executed.clone(),
706 }));
707 let tools = Some(registry.to_api_tools_with_cache(true));
708 let policy = test_tool_surface(&engine, registry, tools, AppMode::Agent);
709 let mut turn = TurnContext::new(4);
710 let (status, error) = tokio::time::timeout(
711 Duration::from_secs(2),
712 engine.run_turn(&mut turn, policy, None, None),
713 )
714 .await
715 .expect("actual parallel tool completion cannot park on a full cancelled queue");
716 assert_eq!(status, TurnOutcomeStatus::Interrupted, "{error:?}");
717 assert!(error.is_none());
718 assert_eq!(
719 executed.load(Ordering::SeqCst),
720 1,
721 "the peer never executes"
722 );
723 assert_eq!(
724 mock.call_count(),
725 1,
726 "no provider continuation after cancellation"
727 );
728 let mut rx = handle.rx_event.write().await;
729 let events: Vec<_> = std::iter::from_fn(|| rx.try_recv().ok()).collect();
730 assert!(events.iter().any(|event| matches!(event,
731 Event::Status { message } if message == "Executing 2 read-only tools in 1 parallel chunk(s)"
732 )), "fixture must exercise one actual two-tool parallel chunk");
733 let starts: Vec<_> = events
734 .iter()
735 .filter_map(|event| match event {
736 Event::OperationActivityStarted {
737 span_id,
738 activity_kind,
739 } => Some((span_id, activity_kind)),
740 _ => None,
741 })
742 .collect();
743 assert_eq!(starts.len(), 1, "one actual tool activity was admitted");
744 let completed: Vec<_> = events
745 .iter()
746 .filter_map(|event| match event {
747 Event::OperationActivityCompleted {
748 span_id,
749 activity_kind,
750 outcome,
751 } => Some((span_id, activity_kind, outcome)),
752 _ => None,
753 })
754 .collect();
755 let calls: Vec<_> = events
756 .iter()
757 .filter_map(|event| match event {
758 Event::ToolCallComplete {
759 id,
760 model_call,
761 result,
762 ..
763 } => Some((id, model_call, result)),
764 _ => None,
765 })
766 .collect();
767 if fill_queue {
768 assert!(
769 completed.is_empty() && calls.is_empty(),
770 "cancelled full waits release without inventing delivery"
771 );
772 } else {
773 assert_eq!(
774 completed.len(),
775 1,
776 "preserve the completed activity after cancellation"
777 );
778 assert_eq!(completed[0].0, starts[0].0, "retain the span relationship");
779 assert_eq!(completed[0].1, starts[0].1);
780 assert_eq!(*completed[0].2, OwnerOperationOutcome::Succeeded);
781 assert_eq!(
782 calls.len(),
783 2,
784 "completed call and cancelled peer both settle exactly once"
785 );
786 let succeeded: Vec<_> = calls
787 .iter()
788 .filter(|(_, _, result)| result.as_ref().is_ok_and(|result| result.success))
789 .collect();
790 assert_eq!(succeeded.len(), 1);
791 let cancelled_peer = calls
792 .iter()
793 .find(|(_, _, result)| !result.as_ref().is_ok_and(|result| result.success))
794 .expect("cancelled peer has its own completion");
795 let peer = cancelled_peer
796 .2
797 .as_ref()
798 .expect("legacy cancelled peer result");
799 assert_eq!(peer.metadata.as_ref().unwrap()["cancelled"], true);
800 assert_eq!(peer.metadata.as_ref().unwrap()["cleanup_confirmed"], false);
801 assert_eq!(
802 succeeded[0].2.as_ref().unwrap().content,
803 "completed before cancellation"
804 );
805 assert!(
806 starts[0].0.starts_with(&format!("{}#", succeeded[0].0)),
807 "span retains its completed execution identity"
808 );
809 let provider_ids: HashSet<_> = calls
810 .iter()
811 .map(|(_, model_call, _)| model_call.as_ref().unwrap().provider_id.as_str())
812 .collect();
813 assert_eq!(provider_ids, HashSet::from(["one", "two"]));
814 let call_starts: HashSet<_> = events
815 .iter()
816 .filter_map(|event| match event {
817 Event::ToolCallStarted { id, .. } => Some(id),
818 _ => None,
819 })
820 .collect();
821 assert!(calls.iter().all(|(id, _, _)| call_starts.contains(id)));
822 }
823 drop(rx);
824 drop(turn_guard);
825 }
826 }
827
828 const REPRESENTATIVE_FIXTURE_ID: &str = "representative-v1";
829 const REPRESENTATIVE_PROJECT_AUTHORITY: &str = "REPRESENTATIVE_PROJECT_AUTHORITY";
830 const REPRESENTATIVE_PROJECT_AUTHORITY_BODY: &str = concat!(
831 "# Representative Project Authority\n\n",
832 "REPRESENTATIVE_PROJECT_AUTHORITY\n\n",
833 "- Keep all work local to the isolated fixture workspace.\n",
834 "- Treat the checked-in repository instructions as the authority for edits.\n",
835 "- Preserve unrelated files and report unsupported checks as unrun.\n",
836 "- Prefer one owner for each runtime fact and delete duplicated derivations.\n",
837 "- Use deterministic provider-free tests before claiming a behavior is verified.\n",
838 "- Keep durable state atomic, recoverable, and explicit about unavailable facts.\n",
839 "- Do not contact remotes, providers, registries, or production services.\n",
840 "- Record exact measurements and distinguish source proof from installed proof.\n",
841 );
842
843 #[test]
844 fn snapshot_notice_precedes_first_provider_call_and_is_owned_by_session() {
845 use crate::llm_client::mock::{MockLlmClient, canned};
846 let _env = lock_test_env();
847 let root = tempdir().unwrap();
848 let _home = EnvVarGuard::set("CODEWHALE_HOME", root.path());
849 let _user_home = EnvVarGuard::set("HOME", root.path());
850 let _user_profile = EnvVarGuard::set("USERPROFILE", root.path());
851 let workspace = root.path().join("workspace");
852 fs::create_dir(&workspace).unwrap();
853 fs::write(workspace.join("large.txt"), vec![b'x'; 4096]).unwrap();
854 let runtime = tokio::runtime::Builder::new_current_thread()
855 .enable_all()
856 .build()
857 .unwrap();
858 runtime.block_on(async {
859 // A resumed Engine for session-a must not warn again. Session-b in
860 // the same process/workspace must receive its own first-turn notice.
861 for (session_id, expected_notices) in [("session-a", 1), ("session-b", 1), ("session-a", 0)]
862 {
863 let config = Config::default();
864 let client = std::sync::Arc::new(MockLlmClient::new(Vec::new()));
865 let (engine, handle) = Engine::new_with_model_client(
866 EngineConfig {
867 session_id: Some(session_id.into()),
868 snapshots_enabled: true,
869 snapshots_max_workspace_bytes: 1024,
870 ..deterministic_engine_config(&workspace)
871 },
872 &config,
873 client.clone(),
874 );
875 let events = std::sync::Arc::clone(&handle.rx_event);
876 let observations = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
877 let observed = std::sync::Arc::clone(&observations);
878 client.push_factory(move |_| {
879 let mut events = events.try_write().expect("fixture owns the event receiver");
880 let mut notices = Vec::new();
881 while let Ok(event) = events.try_recv() {
882 if let Event::SnapshotsDisabled { reason, .. } = event {
883 notices.push(reason);
884 }
885 }
886 observed.lock().unwrap().push(notices);
887 canned::simple_text_turn("snapshot fixture done")
888 });
889 let run = tokio::spawn(engine.run());
890 handle
891 .send(external_user_message_op(
892 "check snapshots",
893 AppMode::Agent,
894 &config,
895 ))
896 .await
897 .unwrap();
898 let snapshot =
899 tokio::time::timeout(Duration::from_secs(10), handle.get_session_snapshot())
900 .await
901 .unwrap()
902 .unwrap();
903 // The Engine catches provider panics, so assertions inside the
904 // factory are not a test oracle. Inspect its observations here.
905 {
906 let observed = observations.lock().unwrap();
907 assert_eq!(
908 observed.len(),
909 1,
910 "factory must have recorded an observation"
911 );
912 assert_eq!(
913 observed[0].len(),
914 expected_notices,
915 "notice must precede provider dispatch for this session"
916 );
917 assert!(observed[0].iter().all(|reason| {
918 // One rendered line: consequence, cause, and the remedy
919 // that lifts this gate, each stated once.
920 reason.lines().count() == 1
921 && reason.contains("Snapshots and /undo are off")
922 && reason
923 .matches(crate::core::turn::SNAPSHOTS_CAP_CONFIG_KEY)
924 .count()
925 == 1
926 }));
927 }
928 assert_eq!(client.call_count(), 1);
929 assert!(
930 serde_json::to_string(&snapshot.messages)
931 .unwrap()
932 .contains("snapshot fixture done")
933 );
934 handle.send(Op::Shutdown).await.unwrap();
935 tokio::time::timeout(Duration::from_secs(10), run)
936 .await
937 .unwrap()
938 .unwrap();
939 }
940 });
941 // Await the owned blocking post-turn snapshots before restoring test home.
942 drop(runtime);
943 }
944
945 /// A recording host (`record_restore_points`) receives every workspace
946 /// snapshot receipt of a turn before its `TurnComplete`: the pre-turn restore
947 /// point, a `tool` snapshot naming the file-mutating call and the paths it
948 /// declared, the `post_tool` snapshot closing it, and the post-turn state. A
949 /// read-only call takes none. The pre/post-turn trees bracket exactly the
950 /// turn's write, so a host derives the turn's workspace delta from them.
951 #[test]
952 fn recorded_snapshot_receipts_bracket_the_turn_and_its_file_writes() {
953 use crate::llm_client::mock::{MockLlmClient, canned};
954 use crate::snapshot::WorkspaceSnapshotKind;
955 let _env = lock_test_env();
956 let root = tempdir().unwrap();
957 let _home = EnvVarGuard::set("CODEWHALE_HOME", root.path());
958 let _user_home = EnvVarGuard::set("HOME", root.path());
959 let _user_profile = EnvVarGuard::set("USERPROFILE", root.path());
960 let workspace = root.path().join("workspace");
961 fs::create_dir(&workspace).unwrap();
962 fs::write(workspace.join("README.md"), "fixture\n").unwrap();
963 let runtime = tokio::runtime::Builder::new_current_thread()
964 .enable_all()
965 .build()
966 .unwrap();
967 runtime.block_on(async {
968 let config = Config::default();
969 let client = std::sync::Arc::new(MockLlmClient::new(vec![
970 canned::tool_call_turn(
971 "call-write",
972 "File",
973 r#"{"action":"write","path":"out.md","content":"out\n"}"#,
974 ),
975 canned::tool_call_turn(
976 "call-read",
977 "File",
978 r#"{"action":"read","path":"README.md"}"#,
979 ),
980 canned::simple_text_turn("done"),
981 ]));
982 let (engine, handle) = Engine::new_with_model_client(
983 EngineConfig {
984 session_id: Some("session-restore".into()),
985 snapshots_enabled: true,
986 snapshots_max_workspace_bytes: 0,
987 record_restore_points: true,
988 ..deterministic_engine_config(&workspace)
989 },
990 &config,
991 client,
992 );
993 let run = tokio::spawn(engine.run());
994 let Op::SendMessage(mut spec) =
995 external_user_message_op("write out.md", AppMode::Agent, &config)
996 else {
997 unreachable!("external_user_message_op builds a SendMessage");
998 };
999 spec.auto_approve = true;
1000 spec.trust_mode = true;
1001 spec.approval_mode = ApprovalMode::Bypass;
1002 handle.send(Op::SendMessage(spec)).await.unwrap();
1003
1004 let mut completions = HashMap::new();
1005 let mut local_ids = HashMap::new();
1006 let mut receipts = Vec::new();
1007 let mut rx = handle.rx_event.write().await;
1008 while let Some(event) = tokio::time::timeout(model_turn_event_timeout(), rx.recv())
1009 .await
1010 .expect("turn events")
1011 {
1012 match event {
1013 Event::ToolCallStarted {
1014 id,
1015 model_call: Some(model_call),
1016 ..
1017 } => {
1018 uuid::Uuid::parse_str(&id).expect("host execution id");
1019 assert_ne!(id, model_call.provider_id);
1020 assert!(local_ids.insert(model_call.provider_id, id).is_none());
1021 }
1022 Event::ToolCallComplete {
1023 id,
1024 result,
1025 model_call: Some(model_call),
1026 ..
1027 } => {
1028 assert_eq!(local_ids.get(&model_call.provider_id), Some(&id));
1029 completions.insert(model_call.provider_id, result.expect("tool result"));
1030 }
1031 Event::WorkspaceSnapshotTaken { snapshot } => receipts.push(snapshot),
1032 Event::TurnComplete { status, error, .. } => {
1033 assert_eq!(status, TurnOutcomeStatus::Completed, "{error:?}");
1034 break;
1035 }
1036 _ => {}
1037 }
1038 }
1039 drop(rx);
1040
1041 assert!(completions.get("call-write").expect("write ran").success);
1042 assert!(completions.get("call-read").expect("read ran").success);
1043 assert_eq!(local_ids.len(), 2);
1044 assert_ne!(local_ids["call-write"], local_ids["call-read"]);
1045 // Every receipt, post-turn included, arrived before TurnComplete.
1046 assert_eq!(
1047 receipts
1048 .iter()
1049 .map(|receipt| (receipt.kind, receipt.tool_call_id.as_deref()))
1050 .collect::<Vec<_>>(),
1051 [
1052 (WorkspaceSnapshotKind::PreTurn, None),
1053 (
1054 WorkspaceSnapshotKind::Tool,
1055 Some(local_ids["call-write"].as_str())
1056 ),
1057 (
1058 WorkspaceSnapshotKind::PostTool,
1059 Some(local_ids["call-write"].as_str())
1060 ),
1061 (WorkspaceSnapshotKind::PostTurn, None),
1062 ],
1063 "a read-only call takes no restore point: {receipts:?}"
1064 );
1065 assert!(
1066 receipts
1067 .iter()
1068 .all(|receipt| receipt.session_id == "session-restore")
1069 );
1070 assert_eq!(
1071 receipts[1].write_paths.as_deref(),
1072 Some(&["out.md".to_string()][..])
1073 );
1074 assert_eq!(
1075 receipts[2].changed_paths.as_deref(),
1076 Some(&["out.md".to_string()][..])
1077 );
1078
1079 let repo = crate::snapshot::SnapshotRepo::open_existing(&workspace)
1080 .unwrap()
1081 .expect("snapshot repo");
1082 let listed = repo.list(usize::MAX).unwrap();
1083 assert!(
1084 listed.iter().any(|snapshot| receipts[1].matches(snapshot)
1085 && snapshot.label == format!("tool:{}", local_ids["call-write"])),
1086 "the tool receipt names a live snapshot"
1087 );
1088 let delta = repo
1089 .diff_snapshots(
1090 &crate::snapshot::SnapshotId::parse(&receipts[0].tree_id).unwrap(),
1091 &crate::snapshot::SnapshotId::parse(&receipts[3].tree_id).unwrap(),
1092 100,
1093 )
1094 .unwrap();
1095 assert_eq!(
1096 delta
1097 .entries
1098 .iter()
1099 .map(|entry| entry.path.as_str())
1100 .collect::<Vec<_>>(),
1101 ["out.md"],
1102 "the pre/post-turn trees bracket exactly this turn's write"
1103 );
1104
1105 handle.send(Op::Shutdown).await.unwrap();
1106 tokio::time::timeout(Duration::from_secs(10), run)
1107 .await
1108 .unwrap()
1109 .unwrap();
1110 });
1111 // Await the owned blocking post-turn snapshots before restoring test home.
1112 drop(runtime);
1113 }
1114
1115 #[test]
1116 fn preview_request_error_preserves_non_semantic_context_chain() {
1117 let error = anyhow::Error::msg("root cause").context("request preparation failed");
1118 assert_eq!(
1119 preview_request_error_user_message("en", &error),
1120 "request preparation failed: root cause"
1121 );
1122 assert_eq!(
1123 initial_stream_error_user_message("en", &error),
1124 "request preparation failed: root cause"
1125 );
1126 }
1127
1128 #[test]
1129 fn initial_stream_failure_preserves_sanitized_context_and_typed_category() {
1130 let error = anyhow::Error::new(crate::llm_client::LlmError::InvalidRequest {
1131 status: 400,
1132 message: "image input is unsupported; api_key=fixture-credential-value".to_string(),
1133 })
1134 .context("Responses API request failed");
1135 let display = initial_stream_error_user_message("en", &error);
1136 assert!(
1137 display.contains("Responses API request failed"),
1138 "{display}"
1139 );
1140 assert!(display.contains("Invalid request (400)"), "{display}");
1141 assert!(display.contains("image input is unsupported"), "{display}");
1142 assert!(!display.contains("fixture-credential-value"), "{display}");
1143 assert!(display.contains("[redacted]"), "{display}");
1144
1145 // The real boundary classifies the original error independently of its
1146 // expanded display text. Preserve the typed terminal invalid-input result.
1147 let policy_message = error.to_string();
1148 let mut envelope = crate::error_taxonomy::envelope_for_llm_error(error, policy_message);
1149 envelope.message = display;
1150 assert_eq!(
1151 envelope.category,
1152 crate::error_taxonomy::ErrorCategory::InvalidInput
1153 );
1154 assert!(!envelope.recoverable);
1155 assert_eq!(envelope.code, "llm_invalid_request");
1156 }
1157 const REPRESENTATIVE_INLINE_INSTRUCTIONS: &str = "REPRESENTATIVE_INLINE_INSTRUCTIONS";
1158 const REPRESENTATIVE_SKILL_DESCRIPTION: &str = "REPRESENTATIVE_SKILL_DESCRIPTION";
1159 const REPRESENTATIVE_MEMORY_CHECKPOINT: &str = "REPRESENTATIVE_MEMORY_CHECKPOINT";
1160 const REPRESENTATIVE_GOAL_OBJECTIVE: &str = "REPRESENTATIVE_GOAL_OBJECTIVE";
1161
1162 #[test]
1163 fn cancellation_wins_at_the_terminal_child_settlement_seam() {
1164 assert_eq!(
1165 terminal_turn_status_at_settlement(TurnOutcomeStatus::Completed, true),
1166 TurnOutcomeStatus::Interrupted
1167 );
1168 assert_eq!(
1169 terminal_turn_status_at_settlement(TurnOutcomeStatus::Completed, false),
1170 TurnOutcomeStatus::Completed
1171 );
1172 assert_eq!(
1173 terminal_turn_status_at_settlement(TurnOutcomeStatus::Failed, true),
1174 TurnOutcomeStatus::Failed
1175 );
1176 }
1177
1178 #[tokio::test]
1179 async fn terminal_barrier_keeps_healthy_child_and_late_completion_alive() {
1180 use std::sync::atomic::Ordering;
1181 let turn_token = CancellationToken::new();
1182 let (mailbox, _receiver) = Mailbox::new(turn_token.clone());
1183 let children = Arc::new(ForegroundChildRegistry::new());
1184 let child_token = turn_token.child_token();
1185 let registration = children
1186 .register(child_token.clone(), "agent_healthy")
1187 .unwrap();
1188 let parking = registration.parking_signal();
1189 let (complete_tx, mut complete_rx) = tokio::sync::mpsc::unbounded_channel();
1190 let (release_tx, release_rx) = tokio::sync::oneshot::channel();
1191 let child = tokio::spawn(async move {
1192 release_rx.await.unwrap();
1193 assert!(!child_token.is_cancelled());
1194 complete_tx.send("existing completion inbox").unwrap();
1195 drop(registration);
1196 });
1197 let (flush_tx, flush_rx) = tokio::sync::oneshot::channel();
1198 let drain_handle = tokio::spawn(async {
1199 let _ = flush_rx.await;
1200 });
1201 let barrier = TurnMailboxBarrier {
1202 mailbox,
1203 cancel_token: turn_token.clone(),
1204 foreground_children: Arc::clone(&children),
1205 flush_tx,
1206 drain_handle,
1207 settle_grace: Duration::from_secs(1),
1208 };
1209 tokio::time::timeout(Duration::from_secs(1), barrier.continue_and_flush())
1210 .await
1211 .unwrap();
1212 assert_eq!(children.active_count(), 1);
1213 assert!(!turn_token.is_cancelled());
1214 assert!(!parking.load(Ordering::Acquire));
1215 release_tx.send(()).unwrap();
1216 assert_eq!(complete_rx.recv().await, Some("existing completion inbox"));
1217 child.await.unwrap();
1218 assert_eq!(children.active_count(), 0);
1219 }
1220
1221 #[tokio::test]
1222 async fn terminal_barrier_explicit_cancel_still_joins_owned_child() {
1223 let turn_token = CancellationToken::new();
1224 let (mailbox, _receiver) = Mailbox::new(turn_token.clone());
1225 let children = Arc::new(ForegroundChildRegistry::new());
1226 let child_token = turn_token.child_token();
1227 let registration = children
1228 .register(child_token.clone(), "agent_owned")
1229 .unwrap();
1230 let child = tokio::spawn(async move {
1231 child_token.cancelled().await;
1232 drop(registration);
1233 });
1234 let (flush_tx, flush_rx) = tokio::sync::oneshot::channel();
1235 let drain_handle = tokio::spawn(async {
1236 let _ = flush_rx.await;
1237 });
1238 let barrier = TurnMailboxBarrier {
1239 mailbox,
1240 cancel_token: turn_token,
1241 foreground_children: Arc::clone(&children),
1242 flush_tx,
1243 drain_handle,
1244 settle_grace: Duration::from_secs(1),
1245 };
1246 let unsettled = tokio::time::timeout(Duration::from_secs(1), barrier.cancel_and_flush())
1247 .await
1248 .unwrap();
1249 assert!(unsettled.is_empty(), "cooperative children join cleanly");
1250 child.await.unwrap();
1251 assert_eq!(children.active_count(), 0);
1252 }
1253
1254 /// Regression for #6184: a foreground child parked on an await that never
1255 /// observes its cancel token must not withhold the terminal turn event. The
1256 /// join gives up at `settle_grace` and names the child it left behind.
1257 #[tokio::test]
1258 async fn terminal_barrier_cancel_names_child_that_ignores_cancellation() {
1259 let turn_token = CancellationToken::new();
1260 let (mailbox, _receiver) = Mailbox::new(turn_token.clone());
1261 let children = Arc::new(ForegroundChildRegistry::new());
1262 let child_token = turn_token.child_token();
1263 // The fake child keeps its registration for the whole test — it never
1264 // observes the cancel token, like a task parked on a blocking await.
1265 let registration = children
1266 .register(child_token.clone(), "agent_stuck")
1267 .unwrap();
1268 let (flush_tx, flush_rx) = tokio::sync::oneshot::channel();
1269 let drain_handle = tokio::spawn(async {
1270 let _ = flush_rx.await;
1271 });
1272 let barrier = TurnMailboxBarrier {
1273 mailbox,
1274 cancel_token: turn_token.clone(),
1275 foreground_children: Arc::clone(&children),
1276 flush_tx,
1277 drain_handle,
1278 settle_grace: Duration::from_millis(50),
1279 };
1280 // Esc latches the turn token before the barrier runs, so the grace
1281 // window — not the token — is the bound under test.
1282 turn_token.cancel();
1283 let unsettled = tokio::time::timeout(Duration::from_secs(1), barrier.cancel_and_flush())
1284 .await
1285 .expect("the bounded join must not wait on a parked child");
1286 assert_eq!(unsettled, vec!["agent_stuck".to_string()]);
1287 assert!(
1288 child_token.is_cancelled(),
1289 "the bounded join still cancels the child's token"
1290 );
1291 assert_eq!(
1292 children.active_count(),
1293 1,
1294 "the stuck child is left registered — leaked, not awaited"
1295 );
1296 drop(registration);
1297 assert_eq!(children.active_count(), 0);
1298 }
1299
1300 /// Regression for #6184: the mailbox drainer parked in an untimed
1301 /// `tx_event.send().await` (event channel full, UI not draining) must not
1302 /// withhold the terminal turn event. `flush` gives up at `settle_grace` and
1303 /// aborts the drainer, so the turn settles instead of waiting hours.
1304 #[tokio::test]
1305 async fn terminal_barrier_flush_bounds_a_drainer_parked_on_a_full_event_channel() {
1306 let turn_token = CancellationToken::new();
1307 let (mailbox, _receiver) = Mailbox::new(turn_token.clone());
1308 let children = Arc::new(ForegroundChildRegistry::new());
1309 // The wedged shape from the report: a one-slot event channel, already
1310 // full, with nobody draining — so the drainer's forward parks in `send`.
1311 let (event_tx, mut event_rx) = tokio::sync::mpsc::channel::<()>(1);
1312 event_tx.send(()).await.unwrap();
1313 let drain_handle = tokio::spawn(async move {
1314 // Parked forever: the channel is full and the receiver never drains.
1315 let _ = event_tx.send(()).await;
1316 });
1317 // Give the drainer a moment to park before the flush races it.
1318 tokio::task::yield_now().await;
1319 let (flush_tx, _flush_rx) = tokio::sync::oneshot::channel();
1320 let barrier = TurnMailboxBarrier {
1321 mailbox,
1322 cancel_token: turn_token,
1323 foreground_children: Arc::clone(&children),
1324 flush_tx,
1325 drain_handle,
1326 settle_grace: Duration::from_millis(50),
1327 };
1328 let started = Instant::now();
1329 tokio::time::timeout(Duration::from_secs(2), barrier.continue_and_flush())
1330 .await
1331 .expect("flush must give up at its grace, not park with the drainer");
1332 assert!(
1333 started.elapsed() < Duration::from_secs(2),
1334 "flush returned in {:?}",
1335 started.elapsed()
1336 );
1337 // Abort receipt: an aborted drainer drops its sender half, so the
1338 // buffered item drains and the channel then reads closed.
1339 assert!(event_rx.recv().await.is_some());
1340 assert!(
1341 event_rx.recv().await.is_none(),
1342 "aborted drainer must release the event channel"
1343 );
1344 }
1345
1346 /// A turn that fails without user cancellation runs the same barrier with a
1347 /// live turn token: Esc during the join must still break it rather than sit
1348 /// out the whole grace period (#6184).
1349 #[tokio::test]
1350 async fn terminal_barrier_cancel_join_breaks_on_fresh_esc() {
1351 let turn_token = CancellationToken::new();
1352 let (mailbox, _receiver) = Mailbox::new(turn_token.clone());
1353 let children = Arc::new(ForegroundChildRegistry::new());
1354 let registration = children
1355 .register(turn_token.child_token(), "agent_stuck")
1356 .unwrap();
1357 let (flush_tx, flush_rx) = tokio::sync::oneshot::channel();
1358 let drain_handle = tokio::spawn(async {
1359 let _ = flush_rx.await;
1360 });
1361 let barrier = TurnMailboxBarrier {
1362 mailbox,
1363 cancel_token: turn_token.clone(),
1364 foreground_children: Arc::clone(&children),
1365 flush_tx,
1366 drain_handle,
1367 settle_grace: Duration::from_secs(30),
1368 };
1369 let esc = turn_token.clone();
1370 let esc_task = tokio::spawn(async move {
1371 tokio::time::sleep(Duration::from_millis(20)).await;
1372 esc.cancel();
1373 });
1374 let unsettled = tokio::time::timeout(Duration::from_secs(1), barrier.cancel_and_flush())
1375 .await
1376 .expect("a fresh Esc must break the join before its grace expires");
1377 assert_eq!(unsettled, vec!["agent_stuck".to_string()]);
1378 esc_task.await.unwrap();
1379 drop(registration);
1380 }
1381
1381 lines RUST