| 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 |