返回 CodeWhale
child_host.rs
根目录 / crates / tui / src / core / engine / tests / child_host.rs
1 //! Actual Engine admission/registry contracts for captured child authority.
2 use super::*;
3 use crate::tools::subagent::engine::ChildAuthority;
4 use anyhow::anyhow;
5
6 fn fixture(
7 workspace: &Path,
8 scope: Option<Vec<String>>,
9 ) -> (Engine, EngineHandle, Arc<ChildAuthority>, TurnRouteContext) {
10 let api = Config {
11 allow_shell: Some(true),
12 ..Default::default()
13 }
14 .with_legacy_root(Some("child-local-fixture".into()), None);
15 let (mut parent, _parent_handle) = Engine::new(
16 EngineConfig {
17 workspace: workspace.into(),
18 session_id: Some("origin-session".into()),
19 allow_shell: true,
20 snapshots_enabled: false,
21 memory_enabled: false,
22 ..Default::default()
23 },
24 &api,
25 );
26 parent.config.features.disable(Feature::Mcp);
27 let route = TurnRouteContext {
28 provider: ProviderKind::Deepseek,
29 model: DEFAULT_TEXT_MODEL.into(),
30 capabilities: Default::default(),
31 limits: None,
32 client: parent.codewhale_client.clone(),
33 api_config: Box::new(api.clone()),
34 locale_tag: parent.config.locale_tag.clone(),
35 role_models: HashMap::new(),
36 auto_model: false,
37 reasoning_effort: None,
38 reasoning_effort_auto: false,
39 };
40 let policy = TurnAuthority::from_effective_fields(
41 AppMode::Agent,
42 true,
43 false,
44 false,
45 ApprovalMode::Suggest,
46 );
47 let context = parent.build_tool_context_for_turn(&policy, &route);
48 let runtime = SubAgentRuntime::new(
49 parent.codewhale_client.clone().unwrap(),
50 DEFAULT_TEXT_MODEL.into(),
51 context,
52 true,
53 None,
54 parent.subagent_manager.clone(),
55 )
56 .with_api_config(api.clone())
57 .with_approval_receipt_store(parent.approval_receipt_store.clone());
58 let authority = ChildAuthority::capture(
59 runtime,
60 FleetRole::Scout,
61 "child-a".into(),
62 "inspect".into(),
63 scope,
64 );
65 let (engine, handle) = Engine::new_child_admitted(
66 EngineConfig {
67 workspace: workspace.into(),
68 model: DEFAULT_TEXT_MODEL.into(),
69 ..Default::default()
70 },
71 &api,
72 authority.clone(),
73 SystemPrompt::Text("captured child role".into()),
74 None,
75 )
76 .unwrap();
77 (engine, handle, authority, route)
78 }
79
80 #[tokio::test(flavor = "current_thread")]
81 async fn child_constructor_reuses_captured_manager_services_namespace_and_live_authority() {
82 let dir = tempdir().unwrap();
83 let _home = crate::test_support::SealedHome::at(dir.path());
84 let (engine, _handle, authority, route) = fixture(dir.path(), None);
85 assert!(Arc::ptr_eq(
86 &engine.subagent_manager,
87 &authority.runtime.manager
88 ));
89 assert!(Arc::ptr_eq(
90 &engine.shell_manager,
91 &authority.runtime.context.shell_manager
92 ));
93 assert!(Arc::ptr_eq(
94 &engine.file_read_tracker,
95 &authority.runtime.context.file_read_tracker
96 ));
97 assert_eq!(engine.session.id, "origin-session");
98 assert_eq!(engine.host_profile, EngineHostProfile::Child);
99 assert!(!Arc::ptr_eq(
100 &engine.live_runtime_authority,
101 &authority
102 .runtime
103 .context
104 .live_posture
105 .as_ref()
106 .unwrap()
107 .state
108 ));
109 let context = engine.build_tool_context_for_turn(
110 &TurnAuthority::from_effective_fields(
111 AppMode::Agent,
112 true,
113 false,
114 true,
115 ApprovalMode::Auto,
116 ),
117 &route,
118 );
119 assert!(Arc::ptr_eq(
120 &context.live_posture.as_ref().unwrap().state,
121 &authority
122 .runtime
123 .context
124 .live_posture
125 .as_ref()
126 .unwrap()
127 .state,
128 ));
129 assert_eq!(context.owner_agent_id.as_deref(), Some("child-a"));
130 assert_eq!(context.state_namespace, "origin-session");
131 assert_eq!(context.shell_policy, authority.grant.shell_policy());
132 let id = engine.new_tool_execution_id();
133 assert!(uuid::Uuid::parse_str(id.strip_prefix("agent:child-a:approval:").unwrap()).is_ok());
134 }
135
136 #[tokio::test(flavor = "current_thread")]
137 async fn child_ceiling_filters_acp_foreground_catalog_without_optional_allow_list() {
138 let dir = tempdir().unwrap();
139 let _home = crate::test_support::SealedHome::at(dir.path());
140 let (engine, _handle, _authority, route) = fixture(dir.path(), Some(vec!["read".into()]));
141 let build = engine.acp_tool_build(
142 &TurnAuthority::from_effective_fields(
143 AppMode::Agent,
144 true,
145 false,
146 true,
147 ApprovalMode::Auto,
148 ),
149 &route,
150 None,
151 );
152 assert!(build.surface.catalog.iter().any(|tool| tool.name == "read"));
153 assert!(
154 !build
155 .surface
156 .catalog
157 .iter()
158 .any(|tool| tool.name == "write")
159 );
160 assert!(
161 !build
162 .surface
163 .catalog
164 .iter()
165 .any(|tool| tool.name == "agent")
166 );
167 let error = build
168 .surface
169 .registry
170 .execute_full(
171 "File",
172 json!({"action":"write","path":"blocked.txt","content":"must not write"}),
173 )
174 .await
175 .unwrap_err();
176 assert!(matches!(error, ToolError::PermissionDenied { .. }));
177 assert!(!dir.path().join("blocked.txt").exists());
178 }
179
180 #[tokio::test(flavor = "current_thread")]
181 async fn child_registry_rejects_context_override_identity_before_tool_effect() {
182 let dir = tempdir().unwrap();
183 let _home = crate::test_support::SealedHome::at(dir.path());
184 let (engine, _handle, _authority, route) = fixture(dir.path(), None);
185 let build = engine.acp_tool_build(
186 &TurnAuthority::from_effective_fields(
187 AppMode::Agent,
188 true,
189 false,
190 true,
191 ApprovalMode::Auto,
192 ),
193 &route,
194 None,
195 );
196 let override_context = build
197 .surface
198 .registry
199 .context()
200 .clone()
201 .with_owner_agent("different-child", "different-child");
202 let error = build
203 .surface
204 .registry
205 .execute_rich_full_with_context(
206 "File",
207 json!({"action":"write","path":"blocked.txt","content":"must not write"}),
208 Some(&override_context),
209 )
210 .await
211 .unwrap_err();
212 assert!(error.to_string().contains("identity"));
213 assert!(!dir.path().join("blocked.txt").exists());
214 }
215
216 #[tokio::test(flavor = "current_thread")]
217 async fn child_shell_ceiling_survives_every_context_policy_setter() {
218 let dir = tempdir().unwrap();
219 let _home = crate::test_support::SealedHome::at(dir.path());
220 let (_engine, _handle, authority, _route) = fixture(dir.path(), None);
221 let ceiling = authority.grant.shell_policy();
222 assert_ne!(ceiling, crate::worker_profile::ShellPolicy::Full);
223 let mut context = authority
224 .context()
225 .with_shell_policy(crate::worker_profile::ShellPolicy::Full);
226 assert_eq!(context.shell_policy, ceiling);
227 context.set_shell_policy(crate::worker_profile::ShellPolicy::None);
228 assert_eq!(
229 context.shell_policy,
230 crate::worker_profile::ShellPolicy::None
231 );
232 context.set_shell_policy(crate::worker_profile::ShellPolicy::Full);
233 assert_eq!(context.shell_policy, ceiling);
234 }
235
236 #[tokio::test(flavor = "current_thread")]
237 async fn child_role_prompt_is_retained_by_the_existing_route_refresh_composer() {
238 let dir = tempdir().unwrap();
239 let _home = crate::test_support::SealedHome::at(dir.path());
240 let (mut engine, _handle, _authority, _route) = fixture(dir.path(), None);
241 engine.current_mode = AppMode::Plan;
242 engine.refresh_system_prompt_with_reason("mode");
243 match engine.session.system_prompt.as_ref().unwrap() {
244 SystemPrompt::Text(text) => assert_eq!(text, "captured child role"),
245 _ => panic!("fixture expected the captured role text"),
246 }
247 assert!(engine.host_managed_turns());
248 }
249
250 fn activate_fixture_control(engine: &mut Engine) -> u64 {
251 let mut controls = engine.turn_controls.lock().unwrap();
252 let control = controls.fresh();
253 let id = control.id;
254 controls.active = Some(control);
255 id
256 }
257
258 #[tokio::test(flavor = "current_thread")]
259 async fn replacing_child_steer_drops_only_uncommitted_same_control_inputs() {
260 let dir = tempdir().unwrap();
261 let _home = crate::test_support::SealedHome::at(dir.path());
262 let (mut engine, handle, _authority, _route) = fixture(dir.path(), None);
263 let active = activate_fixture_control(&mut engine);
264 let old = handle
265 .reserve_steer()
266 .await
267 .unwrap()
268 .send_with_outcome("superseded".into());
269 // A future turn is retained, including its own exact original outcome.
270 let future = {
271 let mut controls = engine.turn_controls.lock().unwrap();
272 let control = controls.fresh();
273 let id = control.id;
274 controls.pending.push_back(control);
275 id
276 };
277 let (future_tx, mut future_rx) = tokio::sync::oneshot::channel();
278 handle
279 .tx_steer
280 .send(handle::SteerInput {
281 turn_id: Some(future),
282 replace_pending: false,
283 content: "future input".into(),
284 outcome: Some(future_tx),
285 })
286 .await
287 .unwrap();
288 let replacement = handle
289 .reserve_steer()
290 .await
291 .unwrap()
292 .send_replacing_with_outcome("current replacement".into());
293 let pending = engine.next_turn_steer().unwrap();
294 assert!(pending.replace_pending);
295 assert_eq!(old.await.unwrap(), handle::SteerOutcome::Dropped);
296 assert!(matches!(
297 future_rx.try_recv(),
298 Err(tokio::sync::oneshot::error::TryRecvError::Empty)
299 ));
300 assert_eq!(pending.commit(), "current replacement");
301 assert_eq!(replacement.await.unwrap(), handle::SteerOutcome::Accepted);
302 assert!(engine.next_turn_steer().is_none());
303 assert_eq!(
304 engine
305 .turn_controls
306 .lock()
307 .unwrap()
308 .active
309 .as_ref()
310 .unwrap()
311 .id,
312 active
313 );
314 {
315 let mut controls = engine.turn_controls.lock().unwrap();
316 controls.active = controls.pending.pop_front();
317 }
318 assert_eq!(engine.next_turn_steer().unwrap().commit(), "future input");
319 assert_eq!(future_rx.await.unwrap(), handle::SteerOutcome::Accepted);
320 }
321
322 #[tokio::test(flavor = "current_thread")]
323 async fn replacing_child_steer_preserves_committed_history_and_drops_on_cancelled_claim() {
324 let dir = tempdir().unwrap();
325 let _home = crate::test_support::SealedHome::at(dir.path());
326 let (mut engine, handle, _authority, _route) = fixture(dir.path(), None);
327 activate_fixture_control(&mut engine);
328 let original = handle
329 .reserve_steer()
330 .await
331 .unwrap()
332 .send_with_outcome("already committed".into());
333 let text = engine.next_turn_steer().unwrap().commit();
334 engine
335 .add_session_message(engine.user_text_message_with_turn_metadata(text))
336 .await;
337 assert_eq!(original.await.unwrap(), handle::SteerOutcome::Accepted);
338 let replaced = handle
339 .reserve_steer()
340 .await
341 .unwrap()
342 .send_replacing_with_outcome("new instruction".into());
343 let claimed = engine.next_turn_steer().unwrap();
344 handle.cancel();
345 drop(claimed);
346 assert_eq!(replaced.await.unwrap(), handle::SteerOutcome::Dropped);
347 assert!(
348 engine
349 .session
350 .messages
351 .iter()
352 .flat_map(|m| &m.content)
353 .any(|b| matches!(b, ContentBlock::Text { text, .. } if text == "already committed"))
354 );
355 }
356
357 #[tokio::test(flavor = "current_thread")]
358 async fn future_control_lookahead_stays_bounded_and_never_retargets_input() {
359 let dir = tempdir().unwrap();
360 let _home = crate::test_support::SealedHome::at(dir.path());
361 let (mut engine, handle, _authority, _route) = fixture(dir.path(), None);
362 activate_fixture_control(&mut engine);
363 let future = {
364 let mut controls = engine.turn_controls.lock().unwrap();
365 let control = controls.fresh();
366 let id = control.id;
367 controls.pending.push_back(control);
368 id
369 };
370 let cap = handle.tx_steer.max_capacity();
371 for _ in 0..cap {
372 handle
373 .tx_steer
374 .send(handle::SteerInput {
375 turn_id: Some(future),
376 replace_pending: false,
377 content: "future only".into(),
378 outcome: None,
379 })
380 .await
381 .unwrap();
382 }
383 assert!(engine.next_turn_steer().is_none());
384 assert_eq!(engine.queued_steers.len(), cap);
385 for _ in 0..cap {
386 handle
387 .tx_steer
388 .send(handle::SteerInput {
389 turn_id: Some(future),
390 replace_pending: false,
391 content: "also future".into(),
392 outcome: None,
393 })
394 .await
395 .unwrap();
396 }
397 for _ in 0..3 {
398 assert!(engine.next_turn_steer().is_none());
399 }
400 assert_eq!(engine.queued_steers.len(), cap);
401 assert_eq!(engine.rx_steer.len(), cap);
402 assert!(handle.tx_steer.try_reserve().is_err());
403 handle.cancel();
404 assert!(engine.cancel_token.is_cancelled());
405 assert!(engine.session.messages.is_empty());
406 }
407
408 async fn actor_fixture(
409 workspace: &Path,
410 scripts: Vec<Vec<StreamEvent>>,
411 max_steps: u32,
412 ) -> (
413 Engine,
414 EngineHandle,
415 Arc<crate::tools::subagent::engine::ChildJob>,
416 Arc<crate::llm_client::mock::MockLlmClient>,
417 ) {
418 let (_old_engine, _old_handle, authority, _route) =
419 fixture(workspace, Some(vec!["read".into()]));
420 actor_fixture_from_authority(workspace, authority, scripts, max_steps).await
421 }
422
423 async fn actor_fixture_from_authority(
424 workspace: &Path,
425 authority: Arc<ChildAuthority>,
426 scripts: Vec<Vec<StreamEvent>>,
427 max_steps: u32,
428 ) -> (
429 Engine,
430 EngineHandle,
431 Arc<crate::tools::subagent::engine::ChildJob>,
432 Arc<crate::llm_client::mock::MockLlmClient>,
433 ) {
434 let mut authority = authority.as_ref().clone();
435 {
436 let mut manager = authority.runtime.manager.write().await;
437 authority.owner_agent_id = manager.insert_test_running_agent("core-child", workspace);
438 manager.assign_test_session_owner(&authority.owner_agent_id, "origin-session");
439 }
440 let api = authority.runtime.api_config.as_deref().unwrap().clone();
441 let authority = Arc::new(authority);
442 let job = crate::tools::subagent::engine::ChildJob::admitted(
443 authority.clone(),
444 crate::tools::subagent::SubAgentAssignment {
445 objective: "inspect evidence".into(),
446 role: None,
447 native_preset: None,
448 },
449 Instant::now(),
450 max_steps,
451 false,
452 None,
453 )
454 .await
455 .unwrap();
456 let (mut engine, handle) = Engine::new_child_admitted(
457 EngineConfig {
458 workspace: workspace.into(),
459 model: DEFAULT_TEXT_MODEL.into(),
460 max_steps: job.work_max_steps,
461 ..Default::default()
462 },
463 &api,
464 authority,
465 SystemPrompt::Text("inspect evidence".into()),
466 None,
467 )
468 .unwrap();
469 let mock = Arc::new(
470 crate::llm_client::mock::MockLlmClient::new(scripts).with_model(DEFAULT_TEXT_MODEL),
471 );
472 engine.model_client = Some(mock.clone());
473 engine.install_child_job(job.clone(), Vec::new()).unwrap();
474 let spec = engine
475 .child_turn_spec("inspect evidence".into(), job.authority.grant.scope.clone())
476 .unwrap();
477 handle.send(Op::SendMessage(spec)).await.unwrap();
478 (engine, handle, job, mock)
479 }
480
481 #[tokio::test(flavor = "current_thread")]
482 async fn child_work_and_one_bounded_report_share_actual_engine_session_and_dispatch() {
483 use crate::llm_client::mock::canned;
484 let dir = tempdir().unwrap();
485 let _home = crate::test_support::SealedHome::at(dir.path());
486 std::fs::write(dir.path().join("evidence.txt"), "exact evidence bytes").unwrap();
487 let (engine, handle, job, mock) = actor_fixture(
488 dir.path(),
489 vec![
490 canned::tool_call_turn("read-evidence", "read", r#"{"path":"evidence.txt"}"#),
491 canned::simple_text_turn("Read evidence.txt; work remains."),
492 ],
493 2,
494 )
495 .await;
496 let collect = async {
497 let mut events = handle.rx_event.write().await;
498 let mut seen = Vec::new();
499 while let Some(event) = events.recv().await {
500 seen.push(event);
501 }
502 seen
503 };
504 let (result, events) = tokio::time::timeout(Duration::from_secs(5), async {
505 tokio::join!(Box::pin(engine.run_child()), collect)
506 })
507 .await
508 .expect("same actor must settle work and reporting");
509 let result = result.unwrap();
510 assert_eq!(
511 result.status,
512 crate::tools::subagent::SubAgentStatus::BudgetExhausted
513 );
514 assert_eq!(mock.call_count(), 2);
515 let requests = mock.captured_requests();
516 assert!(
517 requests[0]
518 .tools
519 .as_ref()
520 .unwrap()
521 .iter()
522 .any(|tool| tool.name == "read")
523 );
524 assert!(requests[1].tools.is_none());
525 assert!(requests[1].tool_choice.is_none());
526 assert!(requests[1].max_tokens <= 1024);
527 assert_eq!(job.steps(), 2);
528 assert!(requests[1].messages.iter().flat_map(|m| &m.content).any(
529 |b| matches!(b, ContentBlock::Text { text, .. } if text.contains("exact evidence bytes"))
530 ));
531 assert_eq!(
532 events
533 .iter()
534 .filter(|e| matches!(e, Event::ToolCallComplete { .. }))
535 .count(),
536 1
537 );
538 assert!(result.checkpoint.is_some());
539 }
540
541 #[tokio::test(flavor = "current_thread")]
542 async fn queued_child_cancel_preserves_checkpoint_without_dispatching_a_fresh_request() {
543 use crate::llm_client::mock::canned;
544 let dir = tempdir().unwrap();
545 let _home = crate::test_support::SealedHome::at(dir.path());
546 let (engine, handle, job, mock) = actor_fixture(
547 dir.path(),
548 vec![canned::simple_text_turn("must not dispatch")],
549 2,
550 )
551 .await;
552 job.authority.runtime.cancel_token.cancel();
553 handle.cancel();
554 let collect = async {
555 let mut events = handle.rx_event.write().await;
556 while events.recv().await.is_some() {}
557 };
558 let (result, ()) = tokio::time::timeout(Duration::from_secs(5), async {
559 tokio::join!(Box::pin(engine.run_child()), collect)
560 })
561 .await
562 .unwrap();
563 let result = result.unwrap();
564 assert_eq!(
565 result.status,
566 crate::tools::subagent::SubAgentStatus::Cancelled
567 );
568 assert_eq!(mock.call_count(), 0);
569 assert!(result.checkpoint.is_some());
570 let checkpoint = result.checkpoint.as_ref().unwrap();
571 let durable = crate::tools::subagent::load_subagent_transcript_artifact(
572 dir.path(),
573 &job.authority.owner_agent_id,
574 )
575 .expect("cancelled Core Session is durably saved under its captured worker");
576 assert_eq!(checkpoint.agent_id, job.authority.owner_agent_id);
577 assert_eq!(checkpoint.message_count, durable.len());
578 assert_eq!(
579 serde_json::to_value(&checkpoint.messages).unwrap(),
580 serde_json::to_value(&durable).unwrap()
581 );
582 }
583
584 #[tokio::test(flavor = "current_thread")]
585 async fn actual_child_transport_services_approval_and_cancel_with_saturated_steer_queue() {
586 use crate::llm_client::mock::canned;
587 use crate::tools::subagent::engine::{drive_child_actor, send_test_child_input};
588 use crate::tools::subagent::{ChildApprovalOutcome, SubAgentStatus};
589 for cancel in [false, true] {
590 let dir = tempdir().unwrap();
591 let _home = crate::test_support::SealedHome::at(dir.path());
592 let (_old, _old_handle, captured, _) = fixture(dir.path(), Some(vec!["bash".into()]));
593 let mut runtime = captured.runtime.clone();
594 runtime.worker_profile =
595 crate::worker_profile::WorkerRuntimeProfile::for_role(FleetRole::Worker);
596 runtime.parent_can_prompt = true;
597 runtime.context.auto_approve = false;
598 runtime.context.approval_mode = ApprovalMode::Suggest;
599 let (parent_tx, mut parent_rx) = tokio::sync::mpsc::channel(16);
600 runtime.event_tx = Some(parent_tx);
601 let authority = ChildAuthority::capture(
602 runtime,
603 FleetRole::Worker,
604 "temporary".into(),
605 "worker".into(),
606 Some(vec!["bash".into()]),
607 );
608 let (core, handle, job, mock) = actor_fixture_from_authority(
609 dir.path(),
610 authority,
611 vec![
612 canned::tool_call_turn(
613 "held-effect",
614 "bash",
615 r#"{"command":"printf child > effect.txt"}"#,
616 ),
617 canned::simple_text_turn("Effect settled."),
618 ],
619 4,
620 )
621 .await;
622 let manager = job.authority.runtime.manager.clone();
623 let (input_tx, input_rx) = tokio::sync::mpsc::unbounded_channel();
624 let pending = Arc::new(std::sync::atomic::AtomicUsize::new(2));
625 let completed = CancellationToken::new();
626 let control = async {
627 let approval_id = loop {
628 if let Event::ApprovalRequired { id, .. } = parent_rx
629 .recv()
630 .await
631 .expect("actual parent approval stream")
632 {
633 break id;
634 }
635 };
636 let cap = handle.tx_steer.max_capacity();
637 let mut outcomes = Vec::new();
638 for index in 0..cap {
639 outcomes.push(
640 handle
641 .reserve_steer()
642 .await
643 .unwrap()
644 .send_with_outcome(format!("queued {index}")),
645 );
646 }
647 assert_eq!(handle.tx_steer.capacity(), 0);
648 send_test_child_input(
649 &input_tx,
650 "replace old queued inputs",
651 true,
652 pending.clone(),
653 );
654 send_test_child_input(&input_tx, "latest instruction", true, pending.clone());
655 // The transport is live with an exact pending Core card and a
656 // full existing steer queue; reserving input cannot own the poll.
657 for _ in 0..32 {
658 tokio::task::yield_now().await;
659 }
660 assert!(
661 !manager
662 .read()
663 .await
664 .pending_requests_for_agent(&job.authority.owner_agent_id)
665 .is_empty()
666 );
667 if cancel {
668 job.authority.runtime.cancel_token.cancel();
669 } else {
670 assert!(
671 manager
672 .write()
673 .await
674 .resolve_child_approval(&approval_id, ChildApprovalOutcome::Approved)
675 );
676 }
677 drop(input_tx);
678 loop {
679 tokio::select! {
680 () = completed.cancelled() => break,
681 event = parent_rx.recv() => if event.is_none() { break; },
682 }
683 }
684 for outcome in outcomes {
685 let settled = outcome.await.unwrap();
686 if cancel {
687 assert_eq!(settled, handle::SteerOutcome::Dropped);
688 }
689 }
690 };
691 let (result, ()) = tokio::time::timeout(Duration::from_secs(5), async {
692 tokio::join!(
693 async {
694 let result =
695 drive_child_actor(core, handle.clone(), job.clone(), input_rx).await;
696 completed.cancel();
697 result
698 },
699 control
700 )
701 })
702 .await
703 .expect("approval/cancellation must stay serviceable with a full steer queue");
704 let result = result.unwrap();
705 assert_eq!(
706 pending.load(std::sync::atomic::Ordering::Acquire),
707 0,
708 "every input settles after actor join"
709 );
710 assert!(
711 manager
712 .read()
713 .await
714 .pending_requests_for_agent(&job.authority.owner_agent_id)
715 .is_empty()
716 );
717 if cancel {
718 assert_eq!(result.status, SubAgentStatus::Cancelled);
719 assert!(!dir.path().join("effect.txt").exists());
720 assert_eq!(mock.call_count(), 1);
721 } else {
722 assert_eq!(result.status, SubAgentStatus::Completed);
723 assert_eq!(
724 std::fs::read_to_string(dir.path().join("effect.txt")).unwrap(),
725 "child"
726 );
727 assert_eq!(mock.call_count(), 2);
728 }
729 assert!(result.checkpoint.is_some());
730 let checkpoint = result.checkpoint.as_ref().unwrap();
731 let durable = crate::tools::subagent::load_subagent_transcript_artifact(
732 dir.path(),
733 &job.authority.owner_agent_id,
734 )
735 .expect("approval/cancellation returns the saved canonical Core Session");
736 assert!(!durable.is_empty());
737 assert_eq!(checkpoint.agent_id, job.authority.owner_agent_id);
738 assert_eq!(checkpoint.message_count, durable.len());
739 assert_eq!(
740 serde_json::to_value(&checkpoint.messages).unwrap(),
741 serde_json::to_value(&durable).unwrap()
742 );
743 }
744 }
745
746 struct ChildTimeoutClient {
747 successful: Arc<crate::llm_client::mock::MockLlmClient>,
748 attempts: std::sync::atomic::AtomicUsize,
749 held_attempts: usize,
750 body: bool,
751 }
752 impl crate::llm_client::LlmClient for ChildTimeoutClient {
753 fn provider_name(&self) -> &'static str {
754 "mock"
755 }
756 fn model(&self) -> &str {
757 DEFAULT_TEXT_MODEL
758 }
759 async fn create_message(
760 &self,
761 _request: codewhale_models::MessageRequest,
762 ) -> Result<codewhale_models::MessageResponse> {
763 Err(anyhow!("child must use the canonical streaming producer"))
764 }
765 async fn create_message_stream(
766 &self,
767 request: codewhale_models::MessageRequest,
768 ) -> Result<crate::llm_client::StreamEventBox> {
769 use futures_util::StreamExt;
770 let attempt = self
771 .attempts
772 .fetch_add(1, std::sync::atomic::Ordering::AcqRel);
773 if attempt < self.held_attempts {
774 if !self.body {
775 return std::future::pending().await;
776 }
777 let mut start = crate::llm_client::mock::canned::message_start("timed-out-body");
778 if let StreamEvent::MessageStart { message } = &mut start {
779 message.usage = Usage {
780 input_tokens: 17,
781 output_tokens: 3,
782 ..Default::default()
783 };
784 }
785 return Ok(Box::pin(
786 futures_util::stream::once(async { Ok(start) })
787 .chain(futures_util::stream::pending()),
788 ));
789 }
790 crate::llm_client::LlmClient::create_message_stream(self.successful.as_ref(), request).await
791 }
792 }
793
794 #[tokio::test(flavor = "current_thread")]
795 async fn child_step_timeout_bounds_open_and_body_with_one_outer_retry_and_exact_usage_sources() {
796 use crate::llm_client::mock::canned;
797 for body in [false, true] {
798 let dir = tempdir().unwrap();
799 let _home = crate::test_support::SealedHome::at(dir.path());
800 let (_old, _old_handle, captured, _) = fixture(dir.path(), Some(Vec::new()));
801 let mut authority = captured.as_ref().clone();
802 authority.runtime.step_api_timeout = Duration::from_millis(25);
803 authority.runtime.api_timeout_retry_base_backoff = Duration::from_millis(1);
804 let authority = Arc::new(authority);
805 let mut success = canned::simple_text_turn("Recovered without restarting the step.");
806 for event in &mut success {
807 if matches!(event, StreamEvent::MessageDelta { .. }) {
808 *event = canned::message_delta(
809 "end_turn",
810 Some(Usage {
811 input_tokens: 10,
812 output_tokens: 6,
813 ..Default::default()
814 }),
815 );
816 }
817 }
818 let (mut core, handle, job, successful) =
819 actor_fixture_from_authority(dir.path(), authority, vec![success], 3).await;
820 let timeout_client = Arc::new(ChildTimeoutClient {
821 successful: successful.clone(),
822 attempts: std::sync::atomic::AtomicUsize::new(0),
823 held_attempts: 2,
824 body,
825 });
826 core.model_client = Some(timeout_client.clone());
827 let (_input_tx, input_rx) = tokio::sync::mpsc::unbounded_channel();
828 let result = tokio::time::timeout(
829 Duration::from_secs(2),
830 crate::tools::subagent::engine::drive_child_actor(core, handle, job.clone(), input_rx),
831 )
832 .await
833 .expect("per-step timeout includes both stream opening and reading")
834 .unwrap();
835 assert_eq!(
836 result.status,
837 crate::tools::subagent::SubAgentStatus::Completed
838 );
839 assert_eq!(
840 timeout_client
841 .attempts
842 .load(std::sync::atomic::Ordering::Acquire),
843 3
844 );
845 assert_eq!(
846 job.steps(),
847 1,
848 "retries do not consume a fresh logical work step"
849 );
850 assert_eq!(successful.call_count(), 1);
851 let record = job
852 .authority
853 .runtime
854 .manager
855 .read()
856 .await
857 .get_worker_record(&job.authority.owner_agent_id)
858 .unwrap();
859 if body {
860 assert_eq!(
861 record.usage_source_fingerprints.len(),
862 3,
863 "each opened response settles once before its retry"
864 );
865 assert_eq!(record.usage.input_tokens, Some(44));
866 assert_eq!(record.usage.output_tokens, Some(12));
867 } else {
868 assert_eq!(record.usage_source_fingerprints.len(), 3);
869 assert_eq!(record.missing_usage_sources.len(), 2);
870 assert_eq!(record.usage.input_tokens, Some(10));
871 assert_eq!(record.usage.output_tokens, Some(6));
872 assert!(record.missing_usage_sources.values().all(|coverage| {
873 coverage.reason
874 == crate::cost_status::RuntimeUsageMissingReason::RequestOutcomeUnknown
875 }));
876 assert!(record.has_unreported_usage);
877 }
878 assert!(
879 !job.can_replace_first_request(),
880 "opened/accepted response forbids route replacement"
881 );
882 assert!(result.checkpoint.is_some());
883 }
884 }
885
886 #[tokio::test(flavor = "current_thread")]
887 async fn child_timeout_exhaustion_interrupts_with_checkpoint_and_never_dispatches_report_as_retry()
888 {
889 use crate::llm_client::mock::canned;
890 let dir = tempdir().unwrap();
891 let _home = crate::test_support::SealedHome::at(dir.path());
892 let (_old, _old_handle, captured, _) = fixture(dir.path(), Some(Vec::new()));
893 let mut authority = captured.as_ref().clone();
894 authority.runtime.step_api_timeout = Duration::from_millis(10);
895 authority.runtime.api_timeout_retry_base_backoff = Duration::from_millis(1);
896 let (mut core, handle, job, successful) = actor_fixture_from_authority(
897 dir.path(),
898 Arc::new(authority),
899 vec![canned::simple_text_turn("must not execute")],
900 3,
901 )
902 .await;
903 let timeout_client = Arc::new(ChildTimeoutClient {
904 successful: successful.clone(),
905 attempts: std::sync::atomic::AtomicUsize::new(0),
906 held_attempts: usize::MAX,
907 body: false,
908 });
909 core.model_client = Some(timeout_client.clone());
910 let (_input_tx, input_rx) = tokio::sync::mpsc::unbounded_channel();
911 let result = tokio::time::timeout(
912 Duration::from_secs(2),
913 crate::tools::subagent::engine::drive_child_actor(core, handle, job.clone(), input_rx),
914 )
915 .await
916 .unwrap()
917 .unwrap();
918 assert!(
919 matches!(result.status, crate::tools::subagent::SubAgentStatus::Interrupted(ref reason)
920 if reason.contains("6 API attempt(s)"))
921 );
922 assert_eq!(
923 timeout_client
924 .attempts
925 .load(std::sync::atomic::Ordering::Acquire),
926 6
927 );
928 assert_eq!(successful.call_count(), 0);
929 assert_eq!(job.steps(), 1);
930 assert!(result.checkpoint.is_some());
931 assert!(result.needs_input.is_some());
932 }
933
934 #[tokio::test(flavor = "current_thread")]
935 async fn child_cancelled_dispatched_open_keeps_unknown_origin_while_queued_cancel_has_none() {
936 use crate::llm_client::mock::canned;
937 let dir = tempdir().unwrap();
938 let _home = crate::test_support::SealedHome::at(dir.path());
939 let _cost = crate::cost_status::test_scope();
940 let (_old, _handle, captured, _) = fixture(dir.path(), Some(Vec::new()));
941 let (mut core, handle, job, successful) = actor_fixture_from_authority(
942 dir.path(),
943 captured,
944 vec![canned::simple_text_turn("never received")],
945 3,
946 )
947 .await;
948 let held = Arc::new(ChildTimeoutClient {
949 successful: successful.clone(),
950 attempts: std::sync::atomic::AtomicUsize::new(0),
951 held_attempts: usize::MAX,
952 body: false,
953 });
954 core.model_client = Some(held.clone());
955 let (_tx, rx) = tokio::sync::mpsc::unbounded_channel();
956 let transport =
957 crate::tools::subagent::engine::drive_child_actor(core, handle.clone(), job.clone(), rx);
958 let cancel = async {
959 while held.attempts.load(std::sync::atomic::Ordering::Acquire) == 0 {
960 tokio::task::yield_now().await;
961 }
962 handle.cancel();
963 };
964 let (result, ()) = tokio::time::timeout(Duration::from_secs(2), async {
965 tokio::join!(transport, cancel)
966 })
967 .await
968 .expect("cancellation stays serviced while stream opening is pending");
969 let result = result.unwrap();
970 assert!(matches!(
971 result.status,
972 crate::tools::subagent::SubAgentStatus::Interrupted(_)
973 ));
974 assert_eq!(successful.call_count(), 0);
975 let record = tokio::time::timeout(Duration::from_secs(2), async {
976 loop {
977 let record = job
978 .authority
979 .runtime
980 .manager
981 .read()
982 .await
983 .get_worker_record(&job.authority.owner_agent_id)
984 .unwrap();
985 if !record.missing_usage_sources.is_empty() {
986 break record;
987 }
988 tokio::task::yield_now().await;
989 }
990 })
991 .await
992 .unwrap();
993 assert_eq!(record.missing_usage_sources.len(), 1);
994 assert!(
995 record
996 .missing_usage_sources
997 .values()
998 .all(|coverage| coverage.reason
999 == crate::cost_status::RuntimeUsageMissingReason::RequestOutcomeUnknown)
1000 );
1001 assert!(record.usage.total_tokens.is_none());
1002 assert!(record.has_unreported_usage);
1003 let pending = crate::cost_status::drain();
1004 assert_eq!(pending.priced_turns, 0);
1005 assert_eq!(pending.unpriced_turns, 1);
1006 assert!(pending.unpriced_reasons.contains("request_outcome_unknown"));
1007 // The separate queued-cancellation case above asserts zero dispatches.
1008 // This case requires the actual model invocation to begin before cancel.
1009 }
1010
1011 #[tokio::test(flavor = "current_thread")]
1012 async fn child_completion_and_cancel_do_not_consume_parent_pending_posture_revision() {
1013 use crate::llm_client::mock::canned;
1014 for cancelled in [false, true] {
1015 let dir = tempdir().unwrap();
1016 let _home = crate::test_support::SealedHome::at(dir.path());
1017 let (engine, handle, job, mock) =
1018 actor_fixture(dir.path(), vec![canned::simple_text_turn("done")], 2).await;
1019 let parent = job
1020 .authority
1021 .runtime
1022 .context
1023 .live_posture
1024 .as_ref()
1025 .unwrap()
1026 .state
1027 .clone();
1028 let prior_applied = {
1029 let mut state = parent.lock().unwrap();
1030 state.authority.approval_mode = ApprovalMode::Never;
1031 state.revision += 1;
1032 state.applied_revision
1033 };
1034 if cancelled {
1035 job.authority.runtime.cancel_token.cancel();
1036 handle.cancel();
1037 }
1038 let collect = async {
1039 let mut events = handle.rx_event.write().await;
1040 while events.recv().await.is_some() {}
1041 };
1042 let (result, ()) = tokio::time::timeout(Duration::from_secs(5), async {
1043 tokio::join!(Box::pin(engine.run_child()), collect)
1044 })
1045 .await
1046 .unwrap();
1047 assert_eq!(
1048 result.unwrap().status,
1049 if cancelled {
1050 crate::tools::subagent::SubAgentStatus::Cancelled
1051 } else {
1052 crate::tools::subagent::SubAgentStatus::Completed
1053 }
1054 );
1055 assert_eq!(mock.call_count(), usize::from(!cancelled));
1056 let state = parent.lock().unwrap();
1057 assert_eq!(state.applied_revision, prior_applied);
1058 assert!(state.revision > state.applied_revision);
1059 assert_eq!(state.authority.approval_mode, ApprovalMode::Never);
1060 }
1061 }
1062
1063 #[tokio::test(flavor = "current_thread")]
1064 async fn actual_child_tool_output_cap_reaches_core_fanout_with_full_artifact() {
1065 use crate::llm_client::mock::canned;
1066 let dir = tempdir().unwrap();
1067 let _home = crate::test_support::SealedHome::at(dir.path());
1068 let raw = "owned-cap-evidence".repeat(32);
1069 std::fs::write(dir.path().join("evidence.txt"), &raw).unwrap();
1070 let (_old, _old_handle, authority, _route) = fixture(dir.path(), Some(vec!["read".into()]));
1071 let mut authority = authority.as_ref().clone();
1072 authority.runtime.max_output_tokens = Some(std::num::NonZeroU32::new(2).unwrap());
1073 let (engine, handle, _job, mock) = actor_fixture_from_authority(
1074 dir.path(),
1075 Arc::new(authority),
1076 vec![
1077 canned::tool_call_turn("read-capped", "read", r#"{"path":"evidence.txt"}"#),
1078 canned::simple_text_turn("Bounded report."),
1079 ],
1080 2,
1081 )
1082 .await;
1083 let collect = async {
1084 let mut rx = handle.rx_event.write().await;
1085 let mut seen = Vec::new();
1086 while let Some(event) = rx.recv().await {
1087 seen.push(event);
1088 }
1089 seen
1090 };
1091 let (result, events) = tokio::time::timeout(Duration::from_secs(5), async {
1092 tokio::join!(Box::pin(engine.run_child()), collect)
1093 })
1094 .await
1095 .expect("actual child output projection must settle");
1096 result.unwrap();
1097 assert_eq!(mock.call_count(), 2);
1098 let completed = events
1099 .iter()
1100 .filter_map(|event| match event {
1101 Event::ToolCallComplete {
1102 result: Ok(output), ..
1103 } => Some(output),
1104 _ => None,
1105 })
1106 .collect::<Vec<_>>();
1107 assert_eq!(completed.len(), 1);
1108 let output = completed[0];
1109 assert!(output.success);
1110 assert!(output.content.len() <= 6 + "\n[truncated: true]".len());
1111 assert!(output.content.ends_with("\n[truncated: true]"));
1112 let metadata = output
1113 .metadata
1114 .as_ref()
1115 .expect("small explicit cap still saves raw bytes");
1116 let path = metadata["artifact_path"].as_str().unwrap();
1117 let saved = std::fs::read_to_string(path).unwrap();
1118 assert!(
1119 saved.contains(&raw),
1120 "artifact must contain full original tool evidence"
1121 );
1122 assert!(saved.len() > output.content.len());
1123 assert_eq!(
1124 metadata["content_digest"],
1125 format!("sha256:{}", crate::hashing::sha256_hex(saved.as_bytes()))
1126 );
1127 assert_eq!(metadata["truncated"], true);
1128 }
1129
1130 #[tokio::test(flavor = "current_thread")]
1131 async fn captured_child_output_caps_preserve_utf8_error_kinds_metadata_and_raw_bytes() {
1132 use super::super::turn_loop::preserve_tool_output_before_fanout;
1133 use crate::tools::spec::ToolSpec as _;
1134 use base64::Engine as _;
1135 let dir = tempdir().unwrap();
1136 let _home = crate::test_support::SealedHome::at(dir.path());
1137 let (default_child, _default_handle, authority, _route) = fixture(dir.path(), None);
1138 assert_eq!(
1139 default_child.child_tool_result_token_cap().unwrap().get(),
1140 10_000
1141 );
1142 let default_raw = "好".repeat(15_000);
1143 let output = preserve_tool_output_before_fanout(
1144 Ok(RichToolResult::plain(ToolResult::success(
1145 default_raw.clone(),
1146 ))),
1147 ProviderKind::Deepseek,
1148 DEFAULT_TEXT_MODEL,
1149 None,
1150 "origin-session",
1151 ("default-cap", "read"),
1152 default_child.child_tool_result_token_cap(),
1153 )
1154 .await
1155 .unwrap()
1156 .into_result();
1157 assert_eq!(output.content.len(), 30_000 + "\n[truncated: true]".len());
1158 assert!(output.content.ends_with("\n[truncated: true]"));
1159 let default_path = output.metadata.as_ref().unwrap()["artifact_path"]
1160 .as_str()
1161 .unwrap();
1162 assert_eq!(std::fs::read_to_string(default_path).unwrap(), default_raw);
1163
1164 let mut authority = authority.as_ref().clone();
1165 authority.runtime.max_output_tokens = Some(std::num::NonZeroU32::new(2).unwrap());
1166 let retrieval_context = authority.runtime.context.clone();
1167 let api = authority.runtime.api_config.as_deref().unwrap().clone();
1168 let (child, _child_handle) = Engine::new_child_admitted(
1169 EngineConfig {
1170 workspace: dir.path().into(),
1171 model: DEFAULT_TEXT_MODEL.into(),
1172 ..Default::default()
1173 },
1174 &api,
1175 Arc::new(authority),
1176 SystemPrompt::Text("captured child".into()),
1177 None,
1178 )
1179 .unwrap();
1180 let raw = "好好好1234567".to_string();
1181 assert_eq!(raw.len(), 16);
1182 let cases = [
1183 ("denied-cap", ToolError::permission_denied(raw.clone())),
1184 ("cancelled-cap", ToolError::cancelled(raw.clone())),
1185 (
1186 "failed-cap",
1187 ToolError::execution_failed_with_metadata(
1188 raw.clone(),
1189 json!({"exit_code": 42, "status": "failed"}),
1190 ),
1191 ),
1192 (
1193 "failed-scalar-cap",
1194 ToolError::execution_failed_with_metadata(raw.clone(), json!(["original", 42])),
1195 ),
1196 ];
1197 for (id, original) in cases {
1198 let error = preserve_tool_output_before_fanout(
1199 Err(original),
1200 ProviderKind::Deepseek,
1201 DEFAULT_TEXT_MODEL,
1202 None,
1203 "origin-session",
1204 (id, "read"),
1205 child.child_tool_result_token_cap(),
1206 )
1207 .await
1208 .unwrap_err();
1209 let content = match &error {
1210 ToolError::PermissionDenied { message } if id == "denied-cap" => message,
1211 ToolError::Cancelled { message } if id == "cancelled-cap" => message,
1212 ToolError::ExecutionFailed { message, metadata } if id == "failed-cap" => {
1213 assert_eq!(metadata.as_ref().unwrap()["exit_code"], 42);
1214 assert_eq!(metadata.as_ref().unwrap()["status"], "failed");
1215 assert!(metadata.as_ref().unwrap()["artifact_id"].is_string());
1216 message
1217 }
1218 ToolError::ExecutionFailed { message, metadata } if id == "failed-scalar-cap" => {
1219 assert_eq!(metadata.as_ref(), Some(&json!(["original", 42])));
1220 message
1221 }
1222 _ => panic!("output projection changed typed error classification: {error:?}"),
1223 };
1224 assert!(content.starts_with("好好\n[truncated: true]\n"));
1225 let artifact = crate::artifacts::artifact_id_for_tool_call(id);
1226 assert!(content.ends_with(&format!(
1227 "omitted range recovery: retrieve_tool_result ref=\"{artifact}\""
1228 )));
1229 assert!(content.len() <= 6 + "\n[truncated: true]".len() + 80);
1230 let relative = crate::artifacts::session_artifact_relative_path(&artifact);
1231 let path =
1232 crate::artifacts::session_artifact_absolute_path("origin-session", &relative).unwrap();
1233 assert_eq!(std::fs::read_to_string(path).unwrap(), raw);
1234 let retrieved = crate::tools::tool_result_retrieval::RetrieveToolResultTool
1235 .execute(
1236 json!({"ref": artifact, "mode": "bytes"}),
1237 &retrieval_context,
1238 )
1239 .await
1240 .unwrap();
1241 let payload: serde_json::Value = serde_json::from_str(&retrieved.content).unwrap();
1242 assert_eq!(payload["total_bytes"], raw.len());
1243 assert_eq!(
1244 base64::engine::general_purpose::STANDARD
1245 .decode(payload["data"].as_str().unwrap())
1246 .unwrap(),
1247 raw.as_bytes(),
1248 );
1249 }
1250 // An immutable-ID conflict is a real save failure, even for a tiny cap.
1251 // The existing evidence must survive and no false retrieval receipt may
1252 // describe the different, unsaved bytes.
1253 let artifact = crate::artifacts::artifact_id_for_tool_call("save-failed-cap");
1254 let (path, _) = crate::artifacts::write_session_artifact_immutable(
1255 "origin-session",
1256 &artifact,
1257 b"prior immutable bytes",
1258 )
1259 .unwrap();
1260 let error = preserve_tool_output_before_fanout(
1261 Err(ToolError::execution_failed_with_metadata(
1262 raw.clone(),
1263 json!({"exit_code": 13, "status": "refused"}),
1264 )),
1265 ProviderKind::Deepseek,
1266 DEFAULT_TEXT_MODEL,
1267 None,
1268 "origin-session",
1269 ("save-failed-cap", "read"),
1270 child.child_tool_result_token_cap(),
1271 )
1272 .await
1273 .unwrap_err();
1274 let ToolError::ExecutionFailed { message, metadata } = error else {
1275 panic!("output preservation changed the typed error");
1276 };
1277 assert!(message.starts_with("好好\n[truncated: true]\n"));
1278 assert!(message.contains("full output could not be saved"));
1279 assert!(message.ends_with("omitted range recovery: re-run with narrower output"));
1280 assert!(!message.contains("retrieve_tool_result"));
1281 assert!(!message.contains(&artifact));
1282 assert!(message.len() <= 6 + "\n[truncated: true]".len() + 100);
1283 let metadata = metadata.unwrap();
1284 assert_eq!(metadata["exit_code"], 13);
1285 assert_eq!(metadata["status"], "refused");
1286 assert_eq!(metadata["output_persistence_failed"], true);
1287 assert_eq!(metadata["truncated"], true);
1288 assert!(metadata.get("artifact_id").is_none());
1289 assert_eq!(std::fs::read(path).unwrap(), b"prior immutable bytes");
1290 // The shared Normal/RLM lane remains byte-identical with no Child cap.
1291 let ordinary = preserve_tool_output_before_fanout(
1292 Err(ToolError::permission_denied(raw.clone())),
1293 ProviderKind::Deepseek,
1294 DEFAULT_TEXT_MODEL,
1295 None,
1296 "origin-session",
1297 ("ordinary", "read"),
1298 None,
1299 )
1300 .await
1301 .unwrap_err();
1302 assert!(matches!(ordinary, ToolError::PermissionDenied { message } if message == raw));
1303 }
1304
1305 #[tokio::test(flavor = "current_thread")]
1306 async fn detached_child_catastrophic_shell_floor_uses_captured_origin_in_every_posture() {
1307 use crate::llm_client::mock::canned;
1308 for approval_mode in [
1309 ApprovalMode::Bypass,
1310 ApprovalMode::Auto,
1311 ApprovalMode::Suggest,
1312 ApprovalMode::Never,
1313 ] {
1314 let dir = tempdir().unwrap();
1315 let _home = crate::test_support::SealedHome::at(dir.path());
1316 let (_old, _handle, captured, _) = fixture(dir.path(), Some(vec!["bash".into()]));
1317 let mut runtime = captured.runtime.clone();
1318 runtime.worker_profile =
1319 crate::worker_profile::WorkerRuntimeProfile::for_role(FleetRole::Worker);
1320 runtime.context.live_posture = None;
1321 runtime.context.approval_mode = approval_mode;
1322 runtime.context.auto_approve = approval_mode == ApprovalMode::Bypass;
1323 assert!(!runtime.has_foreground_ownership());
1324 runtime.parent_can_prompt = false;
1325 let authority = ChildAuthority::capture(
1326 runtime,
1327 FleetRole::Worker,
1328 "temporary".into(),
1329 "worker".into(),
1330 Some(vec!["bash".into()]),
1331 );
1332 let (core, handle, _job, mock) = actor_fixture_from_authority(
1333 dir.path(),
1334 authority,
1335 vec![
1336 canned::tool_call_turn(
1337 "detached-floor",
1338 "bash",
1339 r#"{"command":"dd if=/dev/zero of=/dev/null count=0"}"#,
1340 ),
1341 canned::simple_text_turn("Held the destructive background call."),
1342 ],
1343 4,
1344 )
1345 .await;
1346 let collect = async {
1347 let mut receiver = handle.rx_event.write().await;
1348 let mut events = Vec::new();
1349 while let Some(event) = receiver.recv().await {
1350 events.push(event);
1351 }
1352 events
1353 };
1354 let (result, events) = tokio::time::timeout(Duration::from_secs(5), async {
1355 tokio::join!(Box::pin(core.run_child()), collect)
1356 })
1357 .await
1358 .expect("detached safety floor never waits for a missing human host");
1359 result.unwrap();
1360 assert!(mock.call_count() >= 1);
1361 assert!(
1362 events.iter().any(|event| matches!(
1363 event,
1364 Event::ToolCallComplete { result: Err(error), .. }
1365 if error.to_string().contains("destructive background")
1366 || error.to_string().contains("no host that can answer this approval")
1367 )),
1368 "{approval_mode:?}: the canonical child planner must hold the call: {events:?}"
1369 );
1370 assert!(
1371 !events.iter().any(|event| matches!(
1372 event,
1373 Event::ToolGateDecision {
1374 gate: crate::core::events::ToolGate::AutoReviewGuardian,
1375 ..
1376 }
1377 )),
1378 "the deterministic catastrophic-action floor cannot become model self-approval"
1379 );
1380 assert!(
1381 !events.iter().any(|event| matches!(
1382 event,
1383 Event::ToolCallComplete { result: Ok(result), .. }
1384 if result.content.contains("records out")
1385 )),
1386 "the harmless dd fixture must never reach the shell"
1387 );
1388 }
1389 }
1390
1391 #[tokio::test(flavor = "current_thread")]
1392 async fn child_gate_observation_is_owned_once_and_full_channel_respects_original_cancellation() {
1393 use crate::core::events::{ToolGate, ToolGateVerdict};
1394 use crate::tools::subagent::engine::forward_child_gate_observation;
1395 for cancel_parent in [false, true] {
1396 let dir = tempdir().unwrap();
1397 let _home = crate::test_support::SealedHome::at(dir.path());
1398 let (_engine, _handle, authority, _) = fixture(dir.path(), None);
1399 let mut runtime = authority.runtime.clone();
1400 let (sender, mut receiver) = tokio::sync::mpsc::channel(1);
1401 runtime.event_tx = Some(sender.clone());
1402 runtime.tool_timeout = Duration::from_secs(60);
1403 let turn_cancel = tokio_util::sync::CancellationToken::new();
1404 let receipt = Event::ToolGateDecision {
1405 agent_id: None,
1406 tool_id: "canonical-call".into(),
1407 tool_name: "bash".into(),
1408 gate: ToolGate::AutoReviewGuardian,
1409 decision: ToolGateVerdict::Denied,
1410 risk: Some("high".into()),
1411 reason: "actual reviewer held the call".into(),
1412 };
1413 sender
1414 .send(Event::status("preexisting caller observation"))
1415 .await
1416 .unwrap();
1417 let forward = forward_child_gate_observation(
1418 &runtime,
1419 &authority.owner_agent_id,
1420 receipt.clone(),
1421 &turn_cancel,
1422 None,
1423 );
1424 tokio::pin!(forward);
1425 tokio::select! {
1426 () = &mut forward => panic!("full channel must exercise the capacity wait"),
1427 () = tokio::time::sleep(Duration::from_millis(10)) => {},
1428 }
1429 if cancel_parent {
1430 runtime.cancel_token.cancel();
1431 } else {
1432 turn_cancel.cancel();
1433 }
1434 tokio::time::timeout(Duration::from_millis(250), &mut forward)
1435 .await
1436 .expect("original cancellation bounds a full observational channel");
1437 assert!(
1438 matches!(receiver.recv().await, Some(Event::Status { message, .. })
1439 if message == "preexisting caller observation")
1440 );
1441 assert!(
1442 receiver.try_recv().is_err(),
1443 "a full cancelled queue cannot duplicate a receipt"
1444 );
1445 // Available capacity may carry the already-decided terminal receipt,
1446 // even after cancellation, without making another gate decision.
1447 forward_child_gate_observation(
1448 &runtime,
1449 &authority.owner_agent_id,
1450 receipt,
1451 &turn_cancel,
1452 None,
1453 )
1454 .await;
1455 assert!(
1456 matches!(receiver.recv().await, Some(Event::ToolGateDecision {
1457 agent_id: Some(owner), tool_id, gate: ToolGate::AutoReviewGuardian,
1458 decision: ToolGateVerdict::Denied, risk: Some(risk), reason, ..
1459 }) if owner == authority.owner_agent_id && tool_id == "canonical-call"
1460 && risk == "high" && reason == "actual reviewer held the call")
1461 );
1462 assert!(
1463 receiver.try_recv().is_err(),
1464 "exactly one canonical observation was projected"
1465 );
1466 }
1467 }
1468
1468 lines RUST