| 1 | //! Captured child authority projected onto the canonical Engine registry. |
| 2 | //! This module owns no model loop, tool registry, approval decision or history. |
| 3 | use super::*; |
| 4 | |
| 5 | /// Captured billing origin and the existing worker projection only. This |
| 6 | /// carries no ChildGrant, tool/config authority, job, kernel or new ledger. |
| 7 | #[derive(Clone)] |
| 8 | pub(crate) struct ChildAccountingProjection { |
| 9 | origin: SubAgentAccountingOrigin, |
| 10 | runtime_usage_lease: Option<crate::cost_status::RuntimeUsageLease>, |
| 11 | manager: SharedSubAgentManager, |
| 12 | mailbox: Option<Mailbox>, |
| 13 | owner: String, |
| 14 | } |
| 15 | impl ChildAccountingProjection { |
| 16 | pub(super) fn capture(runtime: &SubAgentRuntime, owner: &str) -> Self { |
| 17 | Self { |
| 18 | origin: runtime.accounting_origin.clone(), |
| 19 | runtime_usage_lease: runtime.runtime_usage_lease.clone(), |
| 20 | manager: runtime.manager.clone(), |
| 21 | mailbox: runtime.mailbox.clone(), |
| 22 | owner: owner.to_owned(), |
| 23 | } |
| 24 | } |
| 25 | pub(crate) fn publish( |
| 26 | &self, |
| 27 | source: &str, |
| 28 | route: &crate::cost_status::EffectiveRouteEnvelope, |
| 29 | usage: &Usage, |
| 30 | reason: crate::cost_status::RuntimeUsageMissingReason, |
| 31 | ) -> Option<u64> { |
| 32 | let known = usage_has_reported_data(usage); |
| 33 | let owner = self |
| 34 | .runtime_usage_lease |
| 35 | .as_ref() |
| 36 | .map(crate::cost_status::RuntimeUsageLease::owner); |
| 37 | if known { |
| 38 | match owner { |
| 39 | Some(owner) => crate::cost_status::report_effective_route_for_runtime( |
| 40 | self.origin.cost_scope, |
| 41 | Some(owner), |
| 42 | source, |
| 43 | route, |
| 44 | usage, |
| 45 | ), |
| 46 | None => self |
| 47 | .origin |
| 48 | .report_ownerless_usage(&self.owner, source, route, usage), |
| 49 | } |
| 50 | if let Some(mailbox) = &self.mailbox { |
| 51 | let _ = mailbox.send(MailboxMessage::token_usage( |
| 52 | &self.owner, |
| 53 | source, |
| 54 | route.clone(), |
| 55 | usage.clone(), |
| 56 | )); |
| 57 | } |
| 58 | } else { |
| 59 | match owner { |
| 60 | Some(owner) => crate::cost_status::report_missing_runtime_usage( |
| 61 | self.origin.cost_scope, |
| 62 | Some(owner), |
| 63 | source, |
| 64 | route, |
| 65 | reason, |
| 66 | ), |
| 67 | None => crate::cost_status::report_missing_usage_for_interactive_origin( |
| 68 | self.origin.cost_scope, |
| 69 | &self.origin.session_id, |
| 70 | &self.owner, |
| 71 | source, |
| 72 | route, |
| 73 | reason, |
| 74 | ), |
| 75 | } |
| 76 | } |
| 77 | known |
| 78 | .then(|| priced_usd_microusd(&route.audit(usage))) |
| 79 | .flatten() |
| 80 | } |
| 81 | /// Finish only the worker projection of an already-published receipt. |
| 82 | /// The caller supplies its held Engine scheduler; no runtime is created. |
| 83 | pub(crate) fn recover_settled( |
| 84 | &self, |
| 85 | scheduler: &tokio::runtime::Handle, |
| 86 | source: String, |
| 87 | route: crate::cost_status::EffectiveRouteEnvelope, |
| 88 | usage: Usage, |
| 89 | priced: Option<u64>, |
| 90 | reason: crate::cost_status::RuntimeUsageMissingReason, |
| 91 | ) { |
| 92 | let manager = self.manager.clone(); |
| 93 | let owner = self.owner.clone(); |
| 94 | scheduler.spawn(async move { |
| 95 | Self::project_worker_receipt(&manager, &owner, &source, &route, &usage, priced, reason) |
| 96 | .await; |
| 97 | }); |
| 98 | } |
| 99 | pub(crate) async fn project_settled( |
| 100 | &self, |
| 101 | source: &str, |
| 102 | route: &crate::cost_status::EffectiveRouteEnvelope, |
| 103 | usage: &Usage, |
| 104 | priced: Option<u64>, |
| 105 | reason: crate::cost_status::RuntimeUsageMissingReason, |
| 106 | ) { |
| 107 | Self::project_worker_receipt( |
| 108 | &self.manager, |
| 109 | &self.owner, |
| 110 | source, |
| 111 | route, |
| 112 | usage, |
| 113 | priced, |
| 114 | reason, |
| 115 | ) |
| 116 | .await; |
| 117 | } |
| 118 | async fn project_worker_receipt( |
| 119 | manager: &SharedSubAgentManager, |
| 120 | owner: &str, |
| 121 | source: &str, |
| 122 | route: &crate::cost_status::EffectiveRouteEnvelope, |
| 123 | usage: &Usage, |
| 124 | priced: Option<u64>, |
| 125 | reason: crate::cost_status::RuntimeUsageMissingReason, |
| 126 | ) { |
| 127 | let mut manager = manager.write().await; |
| 128 | if usage_has_reported_data(usage) { |
| 129 | manager.record_worker_routed_usage(owner, source, route, usage, priced); |
| 130 | } else { |
| 131 | manager.record_worker_missing_usage( |
| 132 | owner, |
| 133 | source, |
| 134 | crate::cost_status::MissingUsageCoverage::for_route(route, reason), |
| 135 | ); |
| 136 | } |
| 137 | } |
| 138 | } |
| 139 | |
| 140 | #[derive(Clone)] |
| 141 | pub(crate) struct ChildAuthority { |
| 142 | pub(crate) grant: crate::worker_profile::ChildGrant, |
| 143 | pub(super) disallowed_tools: Vec<String>, |
| 144 | pub(super) accept_edits: bool, |
| 145 | pub(crate) agent_type: FleetRole, |
| 146 | pub(crate) owner_agent_id: String, |
| 147 | pub(crate) owner_agent_name: String, |
| 148 | pub(super) coordination_manager: SharedSubAgentManager, |
| 149 | pub(super) enforce_write_claim: bool, |
| 150 | pub(crate) runtime: SubAgentRuntime, |
| 151 | pub(crate) person_wait: Arc<PersonWaitClock>, |
| 152 | } |
| 153 | impl std::fmt::Debug for ChildAuthority { |
| 154 | fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { |
| 155 | f.debug_struct("ChildAuthority") |
| 156 | .field("owner", &self.owner_agent_id) |
| 157 | .field("grant", &self.grant) |
| 158 | .finish_non_exhaustive() |
| 159 | } |
| 160 | } |
| 161 | impl ChildAuthority { |
| 162 | const ACTION_ALIASES: &'static [(&'static str, &'static str, &'static str)] = |
| 163 | CANONICAL_ACTION_ALIASES; |
| 164 | pub(crate) fn capture( |
| 165 | runtime: SubAgentRuntime, |
| 166 | agent_type: FleetRole, |
| 167 | owner_agent_id: String, |
| 168 | owner_agent_name: String, |
| 169 | explicit_allowed_tools: Option<Vec<String>>, |
| 170 | ) -> Arc<Self> { |
| 171 | let parent_shell = if tool_denied(Some(&runtime.context.disallowed_tools), "bash") { |
| 172 | ShellPolicy::None |
| 173 | } else { |
| 174 | ShellPolicy::from_legacy_allow_shell(runtime.allow_shell) |
| 175 | }; |
| 176 | let mut effective = runtime.worker_profile.clone(); |
| 177 | effective.shell = effective.shell.min_with(parent_shell); |
| 178 | let scope = intersect_explicit_tool_scope(&effective.tools, explicit_allowed_tools); |
| 179 | let grant = crate::worker_profile::ChildGrant::resolve( |
| 180 | &agent_type, |
| 181 | &effective, |
| 182 | scope, |
| 183 | !runtime.would_exceed_depth(), |
| 184 | ); |
| 185 | Arc::new(Self { |
| 186 | grant, |
| 187 | disallowed_tools: effective.denied_tools, |
| 188 | accept_edits: runtime.accept_edits, |
| 189 | agent_type, |
| 190 | owner_agent_id, |
| 191 | owner_agent_name, |
| 192 | coordination_manager: runtime.manager.clone(), |
| 193 | enforce_write_claim: true, |
| 194 | runtime, |
| 195 | person_wait: Arc::new(PersonWaitClock::default()), |
| 196 | }) |
| 197 | } |
| 198 | /// One child builder used by Core's actual turn registry. The supplied |
| 199 | /// runtime already carries the exact turn route and nested completion inbox. |
| 200 | pub(crate) fn tool_registry_builder( |
| 201 | &self, |
| 202 | runtime: SubAgentRuntime, |
| 203 | todos: SharedTodoList, |
| 204 | plan: SharedPlanState, |
| 205 | ) -> ToolRegistryBuilder { |
| 206 | let mut builder = ToolRegistryBuilder::new().with_full_agent_surface_options( |
| 207 | Some(runtime.client.clone()), |
| 208 | runtime.model.clone(), |
| 209 | runtime.manager.clone(), |
| 210 | runtime.clone(), |
| 211 | runtime.agent_tool_surface_options.clone(), |
| 212 | todos, |
| 213 | plan, |
| 214 | ); |
| 215 | if let Some(pool) = runtime.mcp_pool.as_ref() { |
| 216 | builder = builder.with_mcp_tools(pool.clone()); |
| 217 | } |
| 218 | builder |
| 219 | } |
| 220 | pub(crate) fn accounting_projection(&self) -> ChildAccountingProjection { |
| 221 | ChildAccountingProjection::capture(&self.runtime, &self.owner_agent_id) |
| 222 | } |
| 223 | pub(crate) fn accounting_origin(&self) -> (crate::cost_status::CostScopeToken, String, String) { |
| 224 | ( |
| 225 | self.runtime.accounting_origin.cost_scope, |
| 226 | self.runtime.accounting_origin.session_id.clone(), |
| 227 | self.owner_agent_id.clone(), |
| 228 | ) |
| 229 | } |
| 230 | pub(crate) async fn settle_response( |
| 231 | &self, |
| 232 | source_id: &str, |
| 233 | route: crate::cost_status::EffectiveRouteEnvelope, |
| 234 | usage: &Usage, |
| 235 | ) { |
| 236 | record_provider_response_usage( |
| 237 | &self.runtime, |
| 238 | &self.owner_agent_id, |
| 239 | source_id, |
| 240 | route, |
| 241 | usage, |
| 242 | ) |
| 243 | .await; |
| 244 | } |
| 245 | pub(crate) async fn project_settled_response( |
| 246 | &self, |
| 247 | source_id: &str, |
| 248 | route: &crate::cost_status::EffectiveRouteEnvelope, |
| 249 | usage: &Usage, |
| 250 | ) { |
| 251 | let priced = usage_has_reported_data(usage) |
| 252 | .then(|| priced_usd_microusd(&route.audit(usage))) |
| 253 | .flatten(); |
| 254 | self.runtime |
| 255 | .manager |
| 256 | .write() |
| 257 | .await |
| 258 | .record_worker_routed_usage(&self.owner_agent_id, source_id, route, usage, priced); |
| 259 | } |
| 260 | pub(crate) fn new_execution_id(&self) -> String { |
| 261 | new_child_execution_id(&self.owner_agent_id) |
| 262 | } |
| 263 | pub(crate) fn pause_person_wait(&self) -> PersonWaitPause { |
| 264 | self.person_wait.pause() |
| 265 | } |
| 266 | pub(crate) fn nested_runtime( |
| 267 | &self, |
| 268 | context: ToolContext, |
| 269 | completion: mpsc::Sender<SubAgentCompletion>, |
| 270 | fork: SubAgentForkContext, |
| 271 | ) -> SubAgentRuntime { |
| 272 | let mut runtime = self.runtime.clone(); |
| 273 | runtime.context = context; |
| 274 | runtime.parent_agent_id = Some(self.owner_agent_id.clone()); |
| 275 | runtime.worker_profile.shell = self.grant.shell_policy(); |
| 276 | runtime.agent_tool_surface_options.shell_policy = self.grant.shell_policy(); |
| 277 | runtime.parent_completion_tx = Some(completion); |
| 278 | runtime.fork_context = Some(fork); |
| 279 | runtime |
| 280 | } |
| 281 | pub(crate) fn typed_error(error: anyhow::Error) -> ToolError { |
| 282 | match error.downcast::<ToolError>() { |
| 283 | Ok(error) => error, |
| 284 | Err(error) => ToolError::permission_denied(error.to_string()), |
| 285 | } |
| 286 | } |
| 287 | pub(crate) fn validate_context( |
| 288 | &self, |
| 289 | context: &ToolContext, |
| 290 | ) -> std::result::Result<(), ToolError> { |
| 291 | if context.workspace != self.runtime.context.workspace |
| 292 | || context.owner_agent_id.as_deref() != Some(self.owner_agent_id.as_str()) |
| 293 | || context.state_namespace != self.runtime.context.state_namespace |
| 294 | { |
| 295 | return Err(ToolError::permission_denied( |
| 296 | "child caller identity or workspace changed before dispatch", |
| 297 | )); |
| 298 | } |
| 299 | if context |
| 300 | .cancel_token |
| 301 | .as_ref() |
| 302 | .is_some_and(CancellationToken::is_cancelled) |
| 303 | || self.runtime.cancel_token.is_cancelled() |
| 304 | { |
| 305 | return Err(ToolError::cancelled("child dispatch cancelled")); |
| 306 | } |
| 307 | Ok(()) |
| 308 | } |
| 309 | pub(crate) async fn run_tool_bounded<F, T>( |
| 310 | &self, |
| 311 | future: F, |
| 312 | ) -> std::result::Result<T, ToolError> |
| 313 | where |
| 314 | F: std::future::Future<Output = std::result::Result<T, ToolError>>, |
| 315 | { |
| 316 | let (deadline, _) = budget_handback::wall_deadlines(&self.runtime); |
| 317 | match run_tool_with_person_aware_timeout( |
| 318 | self.runtime.tool_timeout, |
| 319 | deadline, |
| 320 | &self.person_wait, |
| 321 | future, |
| 322 | ) |
| 323 | .await |
| 324 | { |
| 325 | Some(result) => result, |
| 326 | None => Err(ToolError::execution_failed( |
| 327 | "child tool or original work deadline exhausted", |
| 328 | )), |
| 329 | } |
| 330 | } |
| 331 | pub(crate) fn approval_receipt_store( |
| 332 | &self, |
| 333 | ) -> std::result::Result<crate::approval_log::ApprovalReceiptStore, String> { |
| 334 | self.runtime |
| 335 | .approval_receipt_store |
| 336 | .clone() |
| 337 | .unwrap_or_else(|| Err("child caller has no captured approval receipt store".into())) |
| 338 | } |
| 339 | pub(crate) fn context(self: &Arc<Self>) -> ToolContext { |
| 340 | let mut context = self |
| 341 | .runtime |
| 342 | .context |
| 343 | .clone() |
| 344 | .with_owner_agent(self.owner_agent_id.clone(), self.owner_agent_name.clone()) |
| 345 | .with_shell_policy(self.grant.shell_policy()); |
| 346 | context.disallowed_tools = self.disallowed_tools.clone(); |
| 347 | context.child_host = Some(self.clone()); |
| 348 | context |
| 349 | } |
| 350 | pub(crate) fn validate( |
| 351 | &self, |
| 352 | registry: &ToolRegistry, |
| 353 | name: &str, |
| 354 | input: &Value, |
| 355 | ) -> Result<()> { |
| 356 | if !self.grant.desktop && is_machine_control_tool(name) { |
| 357 | return Err(admission_denied(format!( |
| 358 | "[tool.family.denied] Desktop/computer-control tool `{name}` is not available to sub-agents. Run it in the parent session instead, or ask the user." |
| 359 | ))); |
| 360 | } |
| 361 | if self.grant_blocks_tool(registry, name) { |
| 362 | return Err(admission_denied(format!( |
| 363 | "Tool {name} is not available to this read-only worker because its process path does not share the hardened evidence boundary. Use read/search, classifier-bounded bash reads, or the verifier's bounded Run tool instead." |
| 364 | ))); |
| 365 | } |
| 366 | let action = input.get("action").and_then(Value::as_str); |
| 367 | if self.grant.surface == crate::worker_profile::ToolSurface::Evidence |
| 368 | && name == "Web" |
| 369 | && !matches!(action, Some("search" | "fetch")) |
| 370 | { |
| 371 | return Err(admission_denied( |
| 372 | "Tool Web is limited to search/fetch in the read-only evidence profile", |
| 373 | )); |
| 374 | } |
| 375 | // Catalog shaping is not authority. `agent` clears both name-keyed |
| 376 | // gates below by design, so the per-action gate has to be repeated |
| 377 | // here or a hand-written call would reach an action the role's own |
| 378 | // catalog withheld. |
| 379 | if name == "agent" |
| 380 | && matches!(parse_agent_tool_action(input), Ok(AgentToolAction::Claim)) |
| 381 | && !self.agent_action_permitted("claim") |
| 382 | { |
| 383 | return Err(admission_denied(format!( |
| 384 | "agent action=claim widens an enforced write scope, and the Fleet role `{role}` has no write authority to widen. Use an `implement` or `general` role.", |
| 385 | role = self.agent_type.as_str() |
| 386 | ))); |
| 387 | } |
| 388 | let family_action_allowed = if !Self::ACTION_ALIASES |
| 389 | .iter() |
| 390 | .any(|(family, _, _)| *family == name) |
| 391 | { |
| 392 | true |
| 393 | } else if let Some(action) = action { |
| 394 | self.is_action_allowed(name, action) |
| 395 | } else { |
| 396 | self.grant |
| 397 | .scope |
| 398 | .as_ref() |
| 399 | .is_none_or(|list| list.iter().any(|allowed| allowed == name)) |
| 400 | }; |
| 401 | if !self.is_tool_allowed(name) || !family_action_allowed { |
| 402 | return Err(admission_denied(format!( |
| 403 | "Tool {name} not allowed for this sub-agent; report the blocked probe to the parent instead of working around it" |
| 404 | ))); |
| 405 | } |
| 406 | // #3217: authoritative per-role posture — read-only roles cannot mutate |
| 407 | // and non-`Full`-shell roles cannot run shell, regardless of whether |
| 408 | // the parent session is auto-approved. This closes the auto-approve |
| 409 | // bypass where a read-only child could quietly write or shell out. |
| 410 | if !self.posture_permits_tool(registry, name, Some(input)) { |
| 411 | if self.allows_bounded_readonly_bash(name) { |
| 412 | // #6015: the same rule text and next steps as the durable |
| 413 | // authority and the executor, from the one classifier. |
| 414 | let lane = |
| 415 | crate::tools::shell::readonly_enforced_lane_available(registry.context()); |
| 416 | return Err(admission_denied( |
| 417 | match crate::tools::shell::agent_readonly_bash_verdict(input) { |
| 418 | Err(rejection) => format!( |
| 419 | "{} (tool {name}, Fleet role `{role}`)", |
| 420 | crate::tools::shell::readonly_refusal(&rejection, lane), |
| 421 | role = self.agent_type.as_str(), |
| 422 | ), |
| 423 | Ok(()) => format!( |
| 424 | "[shell.readonly.command] Tool {name} input did not match the bounded read-only shell grammar for Fleet role `{role}`. {guidance}", |
| 425 | role = self.agent_type.as_str(), |
| 426 | guidance = |
| 427 | codewhale_execpolicy::command_safety::readonly_command_help() |
| 428 | ), |
| 429 | }, |
| 430 | )); |
| 431 | } |
| 432 | return Err(admission_denied(format!( |
| 433 | "[role.posture.denied] Tool {name} is not permitted for the read-only Fleet role `{role}`. Use an `implement` or `general` role (or `custom` with an explicit allowed_tools list) to mutate the workspace or run shell commands.", |
| 434 | role = self.agent_type.as_str() |
| 435 | ))); |
| 436 | } |
| 437 | // Denied network capability cannot be expanded by answering a prompt. |
| 438 | if self.network_is_denied() { |
| 439 | reject_network_reaching_input(name, input).map_err(Self::typed_error)?; |
| 440 | } |
| 441 | reject_subagent_terminal_takeover(name, input).map_err(Self::typed_error)?; |
| 442 | if self.write_is_denied() { |
| 443 | reject_unbounded_verification(name, input, !self.shell_is_denied()) |
| 444 | .map_err(Self::typed_error)?; |
| 445 | } |
| 446 | // The centralized envelope check. Everything above is name- or |
| 447 | // shape-specific; this one is derived from the tool's real capabilities |
| 448 | // and this call's canonical action, so it also covers the tools no list |
| 449 | // in this file can name — repository plugins, runtime MCP server tools, |
| 450 | // and anything registered later. The bounded read-only bash carve-out |
| 451 | // (#5426/#5438) carries its proven-read-only evidence so the envelope |
| 452 | // classifies it Bounded instead of refusing it as Executes. |
| 453 | if let Some(spec) = registry.get(name) { |
| 454 | crate::tools::execution_envelope::enforce_execution_envelope( |
| 455 | name, |
| 456 | input, |
| 457 | spec.as_ref(), |
| 458 | self.execution_envelope(), |
| 459 | self.bounded_readonly_bash_evidence(name, input), |
| 460 | ) |
| 461 | .map_err(admission_denied)?; |
| 462 | } |
| 463 | Ok(()) |
| 464 | } |
| 465 | pub(crate) async fn validate_claim( |
| 466 | &self, |
| 467 | registry: &ToolRegistry, |
| 468 | name: &str, |
| 469 | input: &Value, |
| 470 | ) -> Result<Vec<String>> { |
| 471 | self.validate(registry, name, input)?; |
| 472 | let scope_aware_write = matches!( |
| 473 | name, |
| 474 | "write" | "edit" | "write_file" | "edit_file" | "apply_patch" | "fim_edit" |
| 475 | ) || (name == "File" |
| 476 | && input |
| 477 | .get("action") |
| 478 | .and_then(Value::as_str) |
| 479 | .is_some_and(|action| matches!(action, "write" | "edit" | "patch"))) |
| 480 | || (name == "pandoc_convert" && input.get("output_path").is_some()); |
| 481 | if scope_aware_write && self.enforce_write_claim { |
| 482 | let paths = mutation_paths(name, input)?; |
| 483 | if paths.is_empty() { |
| 484 | return Err(admission_denied(format!( |
| 485 | "Write tool {name} did not expose a bounded repo-relative target for coordination" |
| 486 | ))); |
| 487 | } |
| 488 | let held = self.coordination_manager.clone().read_owned().await; |
| 489 | let owner = self.owner_agent_id.clone(); |
| 490 | codewhale_app_server::daemon_socket::owner_work(move || { |
| 491 | held.validate_write_scope(&owner, &paths) |
| 492 | .map_err(anyhow::Error::msg) |
| 493 | }) |
| 494 | .await |
| 495 | .map_err(|error| admission_denied(error.to_string()))?; |
| 496 | } else if self.enforce_write_claim |
| 497 | // The typed read-only boundary above already rejected mutation. |
| 498 | && !self.write_is_denied() |
| 499 | && !is_internal_coordination_state_tool(name) |
| 500 | // A shell run the read-only classifier proves mutation-free cannot |
| 501 | // collide with the peer's writes no matter how contended the |
| 502 | // checkout is, so the gate below does not apply to it. |
| 503 | && !proven_readonly_shell_run(name, input) |
| 504 | && (is_unbounded_shell_run(name, input) |
| 505 | || registry.get(name).is_some_and(|spec| { |
| 506 | let canonical = canonical_action_alias(name, input); |
| 507 | let is_shell_control = matches!( |
| 508 | canonical, |
| 509 | "exec_shell_wait" | "exec_shell_interact" | "exec_shell_cancel" |
| 510 | ); |
| 511 | let capabilities = spec.capabilities(); |
| 512 | !is_shell_control |
| 513 | && (spec.approval_requirement_for(input) == ApprovalRequirement::Suggest |
| 514 | || (!spec.is_read_only_for(input) |
| 515 | && capabilities.iter().any(|capability| { |
| 516 | matches!( |
| 517 | capability, |
| 518 | ToolCapability::WritesFiles |
| 519 | | ToolCapability::ExecutesCode |
| 520 | | ToolCapability::Network |
| 521 | ) |
| 522 | }))) |
| 523 | })) |
| 524 | { |
| 525 | let held = self.coordination_manager.clone().read_owned().await; |
| 526 | let owner = self.owner_agent_id.clone(); |
| 527 | let blocking_peers = codewhale_app_server::daemon_socket::owner_work(move || { |
| 528 | // A peer-free checkout still requires the exact held origin. |
| 529 | // The file/root checks stay on the bounded owner worker. |
| 530 | if let Some(claim) = held.shared_write_claim(&owner) { |
| 531 | held.validate_coordination_claim_root(&owner, claim) |
| 532 | .map_err(anyhow::Error::msg)?; |
| 533 | Ok(held.live_peer_shared_write_claim_owners(&owner)) |
| 534 | } else { |
| 535 | Ok(Vec::new()) |
| 536 | } |
| 537 | }) |
| 538 | .await |
| 539 | .map_err(|error| admission_denied(error.to_string()))?; |
| 540 | if !blocking_peers.is_empty() { |
| 541 | return Err(admission_denied(format!( |
| 542 | "Tool {name} cannot prove a bounded file target or read-only execution while peers are writing in this shared checkout (blocking peers: {}). Use a bounded write tool, a proven read-only command, or bash with read_only=true for analysis under native enforcement. Executable work that needs writes requires worktree isolation. Disjoint write_roots alone do not constrain arbitrary code.", |
| 543 | blocking_peers.join(", ") |
| 544 | ))); |
| 545 | } |
| 546 | } |
| 547 | let observed_paths = if scope_aware_write { |
| 548 | mutation_paths(name, input)? |
| 549 | } else { |
| 550 | Vec::new() |
| 551 | }; |
| 552 | Ok(observed_paths) |
| 553 | } |
| 554 | pub(crate) async fn record_settled_writes(&self, paths: Vec<String>) { |
| 555 | if paths.is_empty() { |
| 556 | return; |
| 557 | } |
| 558 | let mut manager = self.coordination_manager.write().await; |
| 559 | if let Some(record) = manager.worker_records.get_mut(&self.owner_agent_id) { |
| 560 | record.delivery_evidence.observed_writes.extend(paths); |
| 561 | } |
| 562 | } |
| 563 | pub(crate) fn delegated_call( |
| 564 | &self, |
| 565 | registry: &ToolRegistry, |
| 566 | name: &str, |
| 567 | input: &Value, |
| 568 | ) -> bool { |
| 569 | if self.bounded_readonly_bash_evidence(name, input) { |
| 570 | return true; |
| 571 | } |
| 572 | registry |
| 573 | .get(name) |
| 574 | .is_some_and(|spec| match spec.approval_requirement_for(input) { |
| 575 | ApprovalRequirement::Auto => true, |
| 576 | ApprovalRequirement::Suggest => { |
| 577 | self.grant.files == crate::worker_profile::FileGrant::Write |
| 578 | && (self.accept_edits |
| 579 | || Self::role_can_delegate_writes(&self.agent_type) |
| 580 | || self.workspace_write_carve_out_permits(registry, name, input)) |
| 581 | } |
| 582 | ApprovalRequirement::Required => { |
| 583 | Self::is_delegated_builtin_verification(name, input) |
| 584 | } |
| 585 | }) |
| 586 | } |
| 587 | pub(crate) fn role_can_delegate_writes(agent_type: &FleetRole) -> bool { |
| 588 | // Builder is the named implementation role. Custom may write only when |
| 589 | // its profile, explicit scope, and execution envelope all agree. |
| 590 | matches!(agent_type, FleetRole::Builder | FleetRole::Custom) |
| 591 | } |
| 592 | |
| 593 | pub(crate) fn workspace_write_carve_out_permits( |
| 594 | &self, |
| 595 | registry: &ToolRegistry, |
| 596 | name: &str, |
| 597 | input: &Value, |
| 598 | ) -> bool { |
| 599 | // This is a bounded convenience for write-capable children, not an |
| 600 | // authority escalation: every target must resolve inside the workspace |
| 601 | // and still pass sensitive-path, repository-law, and claim checks. |
| 602 | if self.grant.files != crate::worker_profile::FileGrant::Write { |
| 603 | return false; |
| 604 | } |
| 605 | crate::core::authority::paths_within_workspace_write_carve_out( |
| 606 | ®istry.context().workspace, |
| 607 | &raw_mutation_target_paths(name, input), |
| 608 | ) |
| 609 | } |
| 610 | |
| 611 | pub(crate) fn is_delegated_builtin_verification(name: &str, input: &Value) -> bool { |
| 612 | use crate::tools::execution_envelope::{VerificationBound, classify_verification}; |
| 613 | |
| 614 | // Reuse the same classifier as the execution envelope. This prevents a |
| 615 | // second, looser notion of "test command" from growing in this module. |
| 616 | matches!( |
| 617 | classify_verification(canonical_action_alias(name, input), input), |
| 618 | Some(VerificationBound::Default | VerificationBound::Filter) |
| 619 | ) |
| 620 | } |
| 621 | |
| 622 | pub(crate) fn posture_permits_tool( |
| 623 | &self, |
| 624 | registry: &ToolRegistry, |
| 625 | name: &str, |
| 626 | input: Option<&Value>, |
| 627 | ) -> bool { |
| 628 | // Delegation depth governs `agent`; write posture must not accidentally |
| 629 | // suppress it or turn depth into a mutation permission. |
| 630 | if name == "agent" { |
| 631 | return true; |
| 632 | } |
| 633 | match registry.get(name) { |
| 634 | Some(spec) => match input.map_or_else( |
| 635 | || spec.approval_requirement(), |
| 636 | |input| spec.approval_requirement_for(input), |
| 637 | ) { |
| 638 | ApprovalRequirement::Auto => true, |
| 639 | ApprovalRequirement::Suggest => { |
| 640 | self.grant.files == crate::worker_profile::FileGrant::Write |
| 641 | } |
| 642 | ApprovalRequirement::Required => { |
| 643 | // An outbound read needs approval because its payload can |
| 644 | // disclose data. That hold does not grant shell/write |
| 645 | // authority; network/envelope and parent approval gates |
| 646 | // below independently decide whether this child may send it. |
| 647 | let capabilities = spec.capabilities(); |
| 648 | if capabilities.contains(&ToolCapability::ReadOnly) |
| 649 | && capabilities.contains(&ToolCapability::Network) |
| 650 | && !capabilities.contains(&ToolCapability::ExecutesCode) |
| 651 | && !capabilities.contains(&ToolCapability::WritesFiles) |
| 652 | { |
| 653 | return true; |
| 654 | } |
| 655 | |
| 656 | // #5426 acceptance point 1: the bounded read-only shell. |
| 657 | // `allows_bounded_readonly_bash` admits canonical `bash` |
| 658 | // to the inspection grant through the raw-shell deny |
| 659 | // list; a call the agent read-only classifier proves |
| 660 | // mutation-free is Auto-class evidence, not a held |
| 661 | // mutation, so the gate must not demand a full shell |
| 662 | // grant for it. Judged by the same predicate |
| 663 | // `BashTool::execute` enforces under |
| 664 | // `ShellPolicy::ReadOnly` (shell.rs), so this admission |
| 665 | // can never widen past the execute-time refusal — the |
| 666 | // first live dogfood against #5428 was denied all three |
| 667 | // canonical inspection commands here because the gate |
| 668 | // consulted only `Required` → `Full`. |
| 669 | if self.allows_bounded_readonly_bash(name) |
| 670 | && input.is_some_and(crate::tools::shell::agent_readonly_bash_input) |
| 671 | { |
| 672 | return true; |
| 673 | } |
| 674 | // `Verify` holds process-start authority for the bounded |
| 675 | // verification surface; the envelope below still refuses |
| 676 | // its ExecutesCode/WritesFiles calls by capability. |
| 677 | self.grant.shell >= crate::worker_profile::ShellGrant::Verify |
| 678 | } |
| 679 | }, |
| 680 | None => true, |
| 681 | } |
| 682 | } |
| 683 | |
| 684 | pub(crate) fn is_tool_denied(&self, name: &str) -> bool { |
| 685 | // The shared matcher canonicalizes legacy/lowercase spellings before |
| 686 | // applying exact or prefix rules. For example `exec_shell*` denies |
| 687 | // `bash`, and `write_file*` denies `write`, in roots and children alike. |
| 688 | tool_denied(Some(&self.disallowed_tools), name) |
| 689 | } |
| 690 | |
| 691 | pub(crate) fn allows_bounded_readonly_bash(&self, name: &str) -> bool { |
| 692 | name == "bash" && self.grant.shell == crate::worker_profile::ShellGrant::Inspect |
| 693 | } |
| 694 | |
| 695 | pub(crate) fn legacy_action_alias(family: &str, action: &str) -> Option<&'static str> { |
| 696 | Self::ACTION_ALIASES |
| 697 | .iter() |
| 698 | .find_map(|(candidate_family, candidate_action, alias)| { |
| 699 | (*candidate_family == family && *candidate_action == action).then_some(*alias) |
| 700 | }) |
| 701 | } |
| 702 | |
| 703 | pub(crate) fn is_action_allowed(&self, family: &str, action: &str) -> bool { |
| 704 | let alias = Self::legacy_action_alias(family, action); |
| 705 | // Read-only inspection keeps two deliberate evidence carve-outs: |
| 706 | // classifier-bounded Bash reads and Web search/fetch. They bypass only |
| 707 | // the coarse family sentinel; action, posture, and envelope checks still |
| 708 | // reject mutation, arbitrary shell, and non-evidence Web actions. |
| 709 | let bounded_readonly_bash = self.allows_bounded_readonly_bash(family) && action == "run"; |
| 710 | let web_readonly_action = family.eq_ignore_ascii_case("Web") |
| 711 | && matches!(action, "search" | "fetch") |
| 712 | && self.network_is_denied(); |
| 713 | if self.is_tool_denied(family) |
| 714 | || !bounded_readonly_bash |
| 715 | && !web_readonly_action |
| 716 | && alias.is_some_and(|name| self.is_tool_denied(name)) |
| 717 | { |
| 718 | return false; |
| 719 | } |
| 720 | match &self.grant.scope { |
| 721 | None => true, |
| 722 | Some(list) => { |
| 723 | list.iter().any(|name| name.eq_ignore_ascii_case(family)) |
| 724 | || alias.is_some_and(|alias| explicit_scope_permits(list, alias)) |
| 725 | } |
| 726 | } |
| 727 | } |
| 728 | |
| 729 | pub(crate) fn is_tool_allowed(&self, name: &str) -> bool { |
| 730 | if name == "agent" && !self.grant.spawn { |
| 731 | return false; |
| 732 | } |
| 733 | if self.is_tool_denied(name) && !self.allows_bounded_readonly_bash(name) { |
| 734 | return false; |
| 735 | } |
| 736 | match &self.grant.scope { |
| 737 | None => true, |
| 738 | Some(list) => { |
| 739 | explicit_scope_permits(list, name) |
| 740 | || Self::ACTION_ALIASES.iter().any(|(family, _, alias)| { |
| 741 | *family == name && list.iter().any(|allowed| allowed == alias) |
| 742 | }) |
| 743 | } |
| 744 | } |
| 745 | } |
| 746 | |
| 747 | pub(crate) fn grant_blocks_tool(&self, registry: &ToolRegistry, name: &str) -> bool { |
| 748 | let lower = name.to_ascii_lowercase(); |
| 749 | // Desktop / machine-control is never in a child's grant. |
| 750 | if !self.grant.desktop && is_machine_control_tool(name) { |
| 751 | return true; |
| 752 | } |
| 753 | // Core's virtual search returns this child's filtered catalog; it is |
| 754 | // not a registered process tool. Scope and deny checks still follow. |
| 755 | let evidence_tool = crate::core::engine::tool_catalog::is_tool_search_tool(name) |
| 756 | || crate::tools::registry::readonly_evidence_tool_name(name) |
| 757 | || registry |
| 758 | .get(name) |
| 759 | .is_some_and(|tool| crate::tools::registry::readonly_evidence_tool(tool.as_ref())); |
| 760 | // The Evidence surface admits only the hardened evidence tools and |
| 761 | // `agent` (delegation). |
| 762 | if self.grant.surface == crate::worker_profile::ToolSurface::Evidence |
| 763 | && lower != "agent" |
| 764 | && !evidence_tool |
| 765 | { |
| 766 | return true; |
| 767 | } |
| 768 | // Below a Full shell grant the raw process surface is gone — except |
| 769 | // canonical `bash` under Inspect, which the read-only classifier |
| 770 | // bounds at dispatch. |
| 771 | let raw_shell = lower == "bash" |
| 772 | || lower.starts_with("exec_shell") |
| 773 | || matches!( |
| 774 | lower.as_str(), |
| 775 | "exec_wait" | "exec_interact" | "task_shell_start" | "task_shell_wait" |
| 776 | ) |
| 777 | || lower.starts_with("terminal/"); |
| 778 | raw_shell |
| 779 | && self.grant.shell < crate::worker_profile::ShellGrant::Full |
| 780 | && !self.allows_bounded_readonly_bash(name) |
| 781 | } |
| 782 | |
| 783 | pub(crate) fn network_is_denied(&self) -> bool { |
| 784 | !self.grant.network |
| 785 | } |
| 786 | |
| 787 | pub(crate) fn write_is_denied(&self) -> bool { |
| 788 | self.grant.files != crate::worker_profile::FileGrant::Write |
| 789 | } |
| 790 | |
| 791 | pub(crate) fn shell_is_denied(&self) -> bool { |
| 792 | self.grant.shell < crate::worker_profile::ShellGrant::Verify |
| 793 | } |
| 794 | |
| 795 | pub(crate) fn execution_envelope(&self) -> crate::tools::execution_envelope::ExecutionEnvelope { |
| 796 | // The capability envelope is the grant's own projection: catalog and |
| 797 | // dispatch derive it from the same fields, never from a deny-list |
| 798 | // sentinel or a role re-mapping that could disagree with them. |
| 799 | crate::tools::execution_envelope::ExecutionEnvelope { |
| 800 | write: self.grant.files == crate::worker_profile::FileGrant::Write, |
| 801 | network: self.grant.network, |
| 802 | shell: !self.shell_is_denied(), |
| 803 | } |
| 804 | } |
| 805 | |
| 806 | pub(crate) fn envelope_permits( |
| 807 | &self, |
| 808 | registry: &ToolRegistry, |
| 809 | name: &str, |
| 810 | input: &Value, |
| 811 | ) -> bool { |
| 812 | let envelope = self.execution_envelope(); |
| 813 | if envelope.is_unrestricted() { |
| 814 | return true; |
| 815 | } |
| 816 | match registry.get(name) { |
| 817 | Some(spec) => crate::tools::execution_envelope::enforce_execution_envelope( |
| 818 | name, |
| 819 | input, |
| 820 | spec.as_ref(), |
| 821 | envelope, |
| 822 | self.bounded_readonly_bash_evidence(name, input), |
| 823 | ) |
| 824 | .is_ok(), |
| 825 | None => true, |
| 826 | } |
| 827 | } |
| 828 | |
| 829 | pub(crate) fn bounded_readonly_bash_evidence(&self, name: &str, input: &Value) -> bool { |
| 830 | self.allows_bounded_readonly_bash(name) |
| 831 | && crate::tools::shell::agent_readonly_bash_input(input) |
| 832 | } |
| 833 | |
| 834 | pub(crate) fn agent_action_permitted(&self, action: &str) -> bool { |
| 835 | // `release` (#5906) carried the same `agents/coordinate` authority as |
| 836 | // `claim` and is gated identically: a read-only role has no write |
| 837 | // scope, so no contention refusal to remediate. |
| 838 | if !matches!(action, "claim" | "release") { |
| 839 | return true; |
| 840 | } |
| 841 | self.grant.files == crate::worker_profile::FileGrant::Write |
| 842 | && !self.is_tool_denied("agents/coordinate") |
| 843 | } |
| 844 | |
| 845 | pub(crate) fn visibility_representative_input(&self, name: &str) -> Option<Value> { |
| 846 | // Visibility and dispatch consult the same capability guard. These |
| 847 | // representative calls let a read-only bash schema survive catalog |
| 848 | // shaping without treating an empty input as arbitrary shell authority. |
| 849 | if self.grant.shell != crate::worker_profile::ShellGrant::Inspect { |
| 850 | return None; |
| 851 | } |
| 852 | match name { |
| 853 | "bash" => Some(json!({"command": "pwd"})), |
| 854 | "Bash" => Some(json!({"action": "run", "command": "pwd"})), |
| 855 | _ => None, |
| 856 | } |
| 857 | } |
| 858 | |
| 859 | pub(crate) fn tools_for_model( |
| 860 | &self, |
| 861 | registry: &ToolRegistry, |
| 862 | agent_type: &FleetRole, |
| 863 | ) -> Vec<Tool> { |
| 864 | // Filter the full registry in deny-first order. These catalog filters |
| 865 | // reduce accidental exposure, but are never the authority boundary: |
| 866 | // execute() repeats role, scope, posture, envelope, and claim checks. |
| 867 | let _ = agent_type; |
| 868 | let api_tools = registry.to_api_tools(); |
| 869 | let filtered = match &self.grant.scope { |
| 870 | None => api_tools, |
| 871 | Some(list) => api_tools |
| 872 | .into_iter() |
| 873 | .filter(|tool| { |
| 874 | explicit_scope_permits(list, &tool.name) |
| 875 | || is_action_family(&tool.name) |
| 876 | && tool.input_schema["properties"]["action"]["enum"] |
| 877 | .as_array() |
| 878 | .is_some_and(|actions| { |
| 879 | actions.iter().any(|action| { |
| 880 | action.as_str().is_some_and(|action| { |
| 881 | Self::legacy_action_alias(&tool.name, action) |
| 882 | .is_some_and(|alias| { |
| 883 | list.iter().any(|n| n == alias) |
| 884 | }) |
| 885 | }) |
| 886 | }) |
| 887 | }) |
| 888 | }) |
| 889 | .collect::<Vec<_>>(), |
| 890 | }; |
| 891 | let mut tools = filtered |
| 892 | .into_iter() |
| 893 | .filter(|tool| tool.name != "agent" || self.grant.spawn) |
| 894 | .filter(|tool| { |
| 895 | !self.is_tool_denied(&tool.name) || self.allows_bounded_readonly_bash(&tool.name) |
| 896 | }) |
| 897 | .filter(|tool| !self.grant_blocks_tool(registry, &tool.name)) |
| 898 | .filter(|tool| { |
| 899 | let representative = self.visibility_representative_input(&tool.name); |
| 900 | tool.name == "File" |
| 901 | || self.posture_permits_tool(registry, &tool.name, representative.as_ref()) |
| 902 | }) |
| 903 | .filter(|tool| { |
| 904 | if is_action_family(&tool.name) { |
| 905 | return true; |
| 906 | } |
| 907 | let representative = self |
| 908 | .visibility_representative_input(&tool.name) |
| 909 | .unwrap_or_else(|| json!({})); |
| 910 | self.envelope_permits(registry, &tool.name, &representative) |
| 911 | }) |
| 912 | .collect::<Vec<_>>(); |
| 913 | |
| 914 | for tool in &mut tools { |
| 915 | if !is_action_family(&tool.name) { |
| 916 | continue; |
| 917 | } |
| 918 | // Indexing `["properties"]["action"]["enum"]` mutably would |
| 919 | // fabricate an `"action": {"enum": null}` property on schemas |
| 920 | // that have no action discriminator (the lowercase `bash` |
| 921 | // command/timeout shape) — a phantom node that fails Moonshot |
| 922 | // MFJS validation. Only shape enums that already exist. |
| 923 | let Some(actions) = tool |
| 924 | .input_schema |
| 925 | .pointer_mut("/properties/action/enum") |
| 926 | .and_then(serde_json::Value::as_array_mut) |
| 927 | else { |
| 928 | continue; |
| 929 | }; |
| 930 | actions.retain(|action| { |
| 931 | let Some(action) = action.as_str() else { |
| 932 | return false; |
| 933 | }; |
| 934 | let posture_allows = tool.name != "File" |
| 935 | || self.grant.files == crate::worker_profile::FileGrant::Write |
| 936 | || matches!(action, "read" | "list" | "search_name" | "search_content"); |
| 937 | let evidence_action = self.grant.surface |
| 938 | != crate::worker_profile::ToolSurface::Evidence |
| 939 | || tool.name != "Web" |
| 940 | || matches!(action, "search" | "fetch"); |
| 941 | let mut representative = self |
| 942 | .visibility_representative_input(&tool.name) |
| 943 | .unwrap_or_else(|| json!({})); |
| 944 | representative["action"] = json!(action); |
| 945 | posture_allows |
| 946 | && evidence_action |
| 947 | && self.is_action_allowed(&tool.name, action) |
| 948 | && self.envelope_permits(registry, &tool.name, &representative) |
| 949 | }); |
| 950 | } |
| 951 | // `agent` is not a `CANONICAL_ACTION_ALIASES` family, so the pruner |
| 952 | // above never reaches it — and it must not become one, because |
| 953 | // `canonical_action_alias` feeds `execution_envelope`, where `agent`'s |
| 954 | // `ExecutesCode` capability is deliberately reclassified `Bounded`. |
| 955 | // Shape its enum explicitly instead. |
| 956 | for tool in &mut tools { |
| 957 | if tool.name != "agent" { |
| 958 | continue; |
| 959 | } |
| 960 | let Some(actions) = tool |
| 961 | .input_schema |
| 962 | .pointer_mut("/properties/action/enum") |
| 963 | .and_then(serde_json::Value::as_array_mut) |
| 964 | else { |
| 965 | continue; |
| 966 | }; |
| 967 | actions.retain(|action| { |
| 968 | action |
| 969 | .as_str() |
| 970 | .is_some_and(|action| self.agent_action_permitted(action)) |
| 971 | }); |
| 972 | } |
| 973 | tools.retain(|tool| { |
| 974 | tool.input_schema["properties"]["action"]["enum"] |
| 975 | .as_array() |
| 976 | .is_none_or(|actions| !actions.is_empty()) |
| 977 | }); |
| 978 | tools |
| 979 | } |
| 980 | } |
| 981 | |
| 982 | /// Worker metadata and its existing transcript projection. Core Session remains |
| 983 | /// the request/history owner; these snapshots are delivery/resume receipts. |
| 984 | pub(crate) struct ChildJob { |
| 985 | pub(crate) authority: Arc<ChildAuthority>, |
| 986 | pub(crate) assignment: SubAgentAssignment, |
| 987 | pub(crate) started_at: Instant, |
| 988 | pub(crate) max_steps: u32, |
| 989 | pub(crate) work_max_steps: u32, |
| 990 | pub(crate) work_deadline: Option<Instant>, |
| 991 | pub(crate) hard_deadline: Option<Instant>, |
| 992 | pub(crate) fork_context: bool, |
| 993 | parking: Option<Arc<std::sync::atomic::AtomicBool>>, |
| 994 | artifact: tokio::sync::Mutex<Option<SubAgentTranscriptArtifactWriter>>, |
| 995 | scheduler: tokio::runtime::Handle, |
| 996 | pub(crate) requests: std::sync::atomic::AtomicU32, |
| 997 | logical_steps: std::sync::atomic::AtomicU32, |
| 998 | response_received: std::sync::atomic::AtomicBool, |
| 999 | selected_model: std::sync::Mutex<String>, |
| 1000 | pacing_sent: std::sync::atomic::AtomicBool, |
| 1001 | stop_reason: std::sync::Mutex<Option<String>>, |
| 1002 | } |
| 1003 | impl ChildJob { |
| 1004 | pub(crate) async fn admitted( |
| 1005 | authority: Arc<ChildAuthority>, |
| 1006 | assignment: SubAgentAssignment, |
| 1007 | started_at: Instant, |
| 1008 | max_steps: u32, |
| 1009 | fork_context: bool, |
| 1010 | parking: Option<Arc<std::sync::atomic::AtomicBool>>, |
| 1011 | ) -> Result<Arc<Self>> { |
| 1012 | let work_max_steps = if max_steps >= 2 { |
| 1013 | max_steps - 1 |
| 1014 | } else { |
| 1015 | max_steps |
| 1016 | }; |
| 1017 | let (work_deadline, hard_deadline) = budget_handback::wall_deadlines(&authority.runtime); |
| 1018 | let artifact = SubAgentTranscriptArtifactWriter::for_runtime( |
| 1019 | &authority.runtime, |
| 1020 | &authority.owner_agent_id, |
| 1021 | ) |
| 1022 | .await?; |
| 1023 | let selected_model = authority.runtime.model.clone(); |
| 1024 | Ok(Arc::new(Self { |
| 1025 | authority, |
| 1026 | assignment, |
| 1027 | started_at, |
| 1028 | max_steps, |
| 1029 | work_max_steps, |
| 1030 | work_deadline, |
| 1031 | hard_deadline, |
| 1032 | fork_context, |
| 1033 | parking, |
| 1034 | artifact: tokio::sync::Mutex::new(Some(artifact)), |
| 1035 | scheduler: tokio::runtime::Handle::try_current() |
| 1036 | .map_err(|_| anyhow!("child admission requires the existing Engine scheduler"))?, |
| 1037 | requests: std::sync::atomic::AtomicU32::new(0), |
| 1038 | logical_steps: std::sync::atomic::AtomicU32::new(0), |
| 1039 | response_received: std::sync::atomic::AtomicBool::new(false), |
| 1040 | selected_model: std::sync::Mutex::new(selected_model), |
| 1041 | pacing_sent: std::sync::atomic::AtomicBool::new(false), |
| 1042 | stop_reason: std::sync::Mutex::new(None), |
| 1043 | })) |
| 1044 | } |
| 1045 | pub(crate) async fn record_route_replacement( |
| 1046 | &self, |
| 1047 | current: &SubAgentRuntime, |
| 1048 | next: &SubAgentRuntime, |
| 1049 | source: SpawnRouteSource, |
| 1050 | note: String, |
| 1051 | error: &anyhow::Error, |
| 1052 | ) { |
| 1053 | let mut manager = self.authority.runtime.manager.write().await; |
| 1054 | if let Some(origin) = current.route_origin.as_deref() { |
| 1055 | manager.record_refused_route( |
| 1056 | &origin.route, |
| 1057 | ¤t.client.redact_model_bound_text(&error.to_string()), |
| 1058 | ); |
| 1059 | } |
| 1060 | manager.record_route_replacement(&self.authority.owner_agent_id, next, source, note); |
| 1061 | } |
| 1062 | pub(crate) fn can_replace_first_request(&self) -> bool { |
| 1063 | self.steps() == 1 |
| 1064 | && !self |
| 1065 | .response_received |
| 1066 | .load(std::sync::atomic::Ordering::Acquire) |
| 1067 | && !self.authority.runtime.cancel_token.is_cancelled() |
| 1068 | } |
| 1069 | pub(crate) fn installed_replacement(&self, model: &str) { |
| 1070 | self.logical_steps |
| 1071 | .store(0, std::sync::atomic::Ordering::Relaxed); |
| 1072 | *self |
| 1073 | .selected_model |
| 1074 | .lock() |
| 1075 | .unwrap_or_else(std::sync::PoisonError::into_inner) = model.to_owned(); |
| 1076 | } |
| 1077 | pub(crate) fn pacing_notice(&self) -> Option<String> { |
| 1078 | if self.pacing_sent.load(std::sync::atomic::Ordering::Relaxed) { |
| 1079 | return None; |
| 1080 | } |
| 1081 | let notice = child_budget_pacing_notice( |
| 1082 | self.started_at, |
| 1083 | self.work_deadline, |
| 1084 | self.steps(), |
| 1085 | self.work_max_steps, |
| 1086 | )?; |
| 1087 | (!self |
| 1088 | .pacing_sent |
| 1089 | .swap(true, std::sync::atomic::Ordering::Relaxed)) |
| 1090 | .then_some(notice) |
| 1091 | } |
| 1092 | pub(crate) async fn project(&self, messages: &[Message], steps: u32) -> Result<()> { |
| 1093 | let mut artifact = self.artifact.lock().await; |
| 1094 | if let Some(writer) = artifact.as_mut() { |
| 1095 | writer.sync_messages(messages, true)?; |
| 1096 | } |
| 1097 | let checkpoint = checkpoint_subagent_progress( |
| 1098 | &self.authority.runtime, |
| 1099 | &self.authority.owner_agent_id, |
| 1100 | "Core child Session checkpoint", |
| 1101 | messages, |
| 1102 | steps, |
| 1103 | true, |
| 1104 | ) |
| 1105 | .await; |
| 1106 | publish_live_subagent_transcript( |
| 1107 | &self.authority.runtime, |
| 1108 | &self.authority.owner_agent_id, |
| 1109 | &self.authority.agent_type, |
| 1110 | &self.assignment, |
| 1111 | None, |
| 1112 | Some(&checkpoint), |
| 1113 | artifact.as_mut(), |
| 1114 | messages, |
| 1115 | steps, |
| 1116 | self.started_at, |
| 1117 | self.fork_context, |
| 1118 | ) |
| 1119 | .await; |
| 1120 | Ok(()) |
| 1121 | } |
| 1122 | pub(crate) async fn before_replace(&self, old: &[Message], new: &[Message]) -> Result<()> { |
| 1123 | let mut artifact = self.artifact.lock().await; |
| 1124 | if let Some(writer) = artifact.as_mut() { |
| 1125 | writer.sync_messages(old, true)?; |
| 1126 | writer.record_compaction(old, new)?; |
| 1127 | writer.append_messages(&[], true)?; |
| 1128 | } |
| 1129 | Ok(()) |
| 1130 | } |
| 1131 | pub(crate) fn stop_for_budget(&self, reason: &str) { |
| 1132 | let mut slot = self |
| 1133 | .stop_reason |
| 1134 | .lock() |
| 1135 | .unwrap_or_else(std::sync::PoisonError::into_inner); |
| 1136 | slot.get_or_insert_with(|| reason.to_owned()); |
| 1137 | } |
| 1138 | pub(crate) fn budget_reason(&self) -> Option<String> { |
| 1139 | self.stop_reason |
| 1140 | .lock() |
| 1141 | .unwrap_or_else(std::sync::PoisonError::into_inner) |
| 1142 | .clone() |
| 1143 | } |
| 1144 | pub(crate) fn steps(&self) -> u32 { |
| 1145 | self.logical_steps |
| 1146 | .load(std::sync::atomic::Ordering::Relaxed) |
| 1147 | } |
| 1148 | pub(crate) fn provider_refused(&self, error: &anyhow::Error) { |
| 1149 | if matches!( |
| 1150 | error.downcast_ref::<LlmError>(), |
| 1151 | Some(LlmError::RateLimited { .. }) |
| 1152 | ) && let Some(governor) = self.authority.runtime.governor.as_ref() |
| 1153 | { |
| 1154 | governor.record_rate_limited(Instant::now()); |
| 1155 | } |
| 1156 | } |
| 1157 | pub(crate) fn dispatched( |
| 1158 | self: &Arc<Self>, |
| 1159 | logical_step: u32, |
| 1160 | source: String, |
| 1161 | route: crate::cost_status::EffectiveRouteEnvelope, |
| 1162 | ) -> ChildDispatchedRequest { |
| 1163 | if let Some(governor) = self.authority.runtime.governor.as_ref() { |
| 1164 | governor.record_attempt(Instant::now()); |
| 1165 | } |
| 1166 | self.logical_steps |
| 1167 | .fetch_max(logical_step, std::sync::atomic::Ordering::Relaxed); |
| 1168 | self.requests |
| 1169 | .fetch_add(1, std::sync::atomic::Ordering::Relaxed); |
| 1170 | ChildDispatchedRequest { |
| 1171 | job: self.clone(), |
| 1172 | source, |
| 1173 | route, |
| 1174 | settled: None, |
| 1175 | refused: false, |
| 1176 | projected: false, |
| 1177 | } |
| 1178 | } |
| 1179 | pub(crate) fn provider_responded(&self) { |
| 1180 | // Even an incomplete/failed decoder means the provider accepted this |
| 1181 | // request. It is never eligible for first-request route replacement. |
| 1182 | self.response_received |
| 1183 | .store(true, std::sync::atomic::Ordering::Release); |
| 1184 | } |
| 1185 | pub(crate) fn response_settled(&self, success: bool) { |
| 1186 | if success && let Some(governor) = self.authority.runtime.governor.as_ref() { |
| 1187 | governor.record_success(Instant::now()); |
| 1188 | } |
| 1189 | } |
| 1190 | pub(crate) async fn finish( |
| 1191 | &self, |
| 1192 | messages: &[Message], |
| 1193 | steps: u32, |
| 1194 | status: SubAgentStatus, |
| 1195 | text: Option<String>, |
| 1196 | reason: Option<&str>, |
| 1197 | ) -> Result<SubAgentResult> { |
| 1198 | self.project(messages, steps).await?; |
| 1199 | let mut artifact = self.artifact.lock().await; |
| 1200 | if self.authority.runtime.cancel_token.is_cancelled() { |
| 1201 | // `project` has durably synced this exact Core Session and updated |
| 1202 | // the captured worker. Return that checkpoint rather than discarding |
| 1203 | // it while projecting the cancellation result. |
| 1204 | let checkpoint = self |
| 1205 | .authority |
| 1206 | .runtime |
| 1207 | .manager |
| 1208 | .read() |
| 1209 | .await |
| 1210 | .get_result(&self.authority.owner_agent_id)? |
| 1211 | .checkpoint; |
| 1212 | return Ok(cancelled_subagent_result( |
| 1213 | &self.authority.runtime, |
| 1214 | &self.authority.owner_agent_id, |
| 1215 | &self.authority.agent_type, |
| 1216 | &self.assignment, |
| 1217 | messages, |
| 1218 | steps, |
| 1219 | self.max_steps, |
| 1220 | checkpoint.as_ref(), |
| 1221 | self.parking.as_ref(), |
| 1222 | &mut artifact, |
| 1223 | self.started_at, |
| 1224 | self.fork_context, |
| 1225 | " in Core child turn", |
| 1226 | ) |
| 1227 | .await); |
| 1228 | } |
| 1229 | let checkpoint = build_subagent_checkpoint( |
| 1230 | &self.authority.owner_agent_id, |
| 1231 | reason.unwrap_or_else(|| subagent_status_name(&status)), |
| 1232 | messages, |
| 1233 | steps, |
| 1234 | matches!(status, SubAgentStatus::Interrupted(_)), |
| 1235 | ); |
| 1236 | let duration_ms = u64::try_from(self.started_at.elapsed().as_millis()).unwrap_or(u64::MAX); |
| 1237 | insert_subagent_full_transcript_handle( |
| 1238 | &self.authority.runtime, |
| 1239 | &self.authority.owner_agent_id, |
| 1240 | &self.authority.agent_type, |
| 1241 | &self.assignment, |
| 1242 | &status, |
| 1243 | text.as_ref(), |
| 1244 | Some(&checkpoint), |
| 1245 | artifact.as_mut(), |
| 1246 | messages, |
| 1247 | steps, |
| 1248 | duration_ms, |
| 1249 | self.fork_context, |
| 1250 | ) |
| 1251 | .await; |
| 1252 | let mut result = self |
| 1253 | .authority |
| 1254 | .runtime |
| 1255 | .manager |
| 1256 | .read() |
| 1257 | .await |
| 1258 | .get_result(&self.authority.owner_agent_id)?; |
| 1259 | result.status = status; |
| 1260 | result.result = text; |
| 1261 | result.steps_taken = steps; |
| 1262 | result.checkpoint = Some(checkpoint); |
| 1263 | result.duration_ms = duration_ms; |
| 1264 | result.model = self |
| 1265 | .selected_model |
| 1266 | .lock() |
| 1267 | .unwrap_or_else(std::sync::PoisonError::into_inner) |
| 1268 | .clone(); |
| 1269 | result.needs_input = match &result.status { |
| 1270 | SubAgentStatus::Interrupted(reason) => result |
| 1271 | .checkpoint |
| 1272 | .as_ref() |
| 1273 | .map(|checkpoint| needs_input_for_interrupted_checkpoint(reason, checkpoint)), |
| 1274 | _ => None, |
| 1275 | }; |
| 1276 | result.usage = self |
| 1277 | .authority |
| 1278 | .runtime |
| 1279 | .manager |
| 1280 | .read() |
| 1281 | .await |
| 1282 | .worker_records |
| 1283 | .get(&self.authority.owner_agent_id) |
| 1284 | .map(|record| record.usage.clone()); |
| 1285 | if result.status == SubAgentStatus::BudgetExhausted { |
| 1286 | let cause = reason.unwrap_or("child work budget exhausted"); |
| 1287 | if let Some(note) = budget_work_preservation_note( |
| 1288 | &self.authority.runtime.manager, |
| 1289 | &self.authority.owner_agent_id, |
| 1290 | cause, |
| 1291 | ) |
| 1292 | .await |
| 1293 | { |
| 1294 | result |
| 1295 | .result |
| 1296 | .get_or_insert_with(String::new) |
| 1297 | .push_str(&format!("\n\n{note}")); |
| 1298 | } |
| 1299 | let artifact = if let Some(body) = result.result.clone() { |
| 1300 | budget_handback::write_digest_artifact( |
| 1301 | &self.authority.runtime, |
| 1302 | &self.authority.owner_agent_id, |
| 1303 | body, |
| 1304 | ) |
| 1305 | .await |
| 1306 | } else { |
| 1307 | None |
| 1308 | }; |
| 1309 | let mut handback_note = "The bounded Core hand-back preserves the recorded work; the assignment is not complete.".to_string(); |
| 1310 | if let Some(artifact) = artifact { |
| 1311 | handback_note.push_str(&format!( |
| 1312 | " Recorded work is saved as this child's deliverable: {}.", |
| 1313 | artifact.display() |
| 1314 | )); |
| 1315 | } |
| 1316 | result = budget_partial_result_with_note( |
| 1317 | result, |
| 1318 | reason.unwrap_or("child work budget exhausted"), |
| 1319 | &handback_note, |
| 1320 | ); |
| 1321 | } |
| 1322 | release_resident_leases_for(&self.authority.owner_agent_id); |
| 1323 | Ok(result) |
| 1324 | } |
| 1325 | } |
| 1326 | |
| 1327 | /// Exact primary dispatch ownership. Dropping an actually invoked request |
| 1328 | /// records ambiguity under its captured origin; an unopened/queued turn never |
| 1329 | /// constructs this guard. Billing is synchronous before worker projection. |
| 1330 | pub(crate) struct ChildDispatchedRequest { |
| 1331 | job: Arc<ChildJob>, |
| 1332 | source: String, |
| 1333 | route: crate::cost_status::EffectiveRouteEnvelope, |
| 1334 | settled: Option<( |
| 1335 | Usage, |
| 1336 | Option<u64>, |
| 1337 | crate::cost_status::RuntimeUsageMissingReason, |
| 1338 | )>, |
| 1339 | refused: bool, |
| 1340 | projected: bool, |
| 1341 | } |
| 1342 | /// Only a typed authoritative rejection proves the dispatched operation |
| 1343 | /// produced no billable response. Network/body/cancellation ambiguity does not. |
| 1344 | pub(crate) fn provider_request_refusal_proven(error: &anyhow::Error) -> bool { |
| 1345 | matches!( |
| 1346 | error.downcast_ref::<LlmError>(), |
| 1347 | Some( |
| 1348 | LlmError::RateLimited { .. } |
| 1349 | | LlmError::QuotaExhausted(_) |
| 1350 | | LlmError::AuthenticationError(_) |
| 1351 | | LlmError::AuthorizationError(_) |
| 1352 | | LlmError::InvalidRequest { .. } |
| 1353 | | LlmError::ModelError(_) |
| 1354 | | LlmError::ContentPolicyError(_) |
| 1355 | | LlmError::ContextLengthError(_) |
| 1356 | ) |
| 1357 | ) |
| 1358 | } |
| 1359 | impl ChildDispatchedRequest { |
| 1360 | pub(crate) async fn settle_open_error(&mut self, error: &anyhow::Error) { |
| 1361 | self.refused = provider_request_refusal_proven(error); |
| 1362 | if !self.refused { |
| 1363 | self.settle(&Usage::default(), false).await; |
| 1364 | } |
| 1365 | } |
| 1366 | pub(crate) async fn settle(&mut self, usage: &Usage, complete: bool) { |
| 1367 | let reason = if complete { |
| 1368 | crate::cost_status::RuntimeUsageMissingReason::SuccessWithoutUsage |
| 1369 | } else { |
| 1370 | crate::cost_status::RuntimeUsageMissingReason::RequestOutcomeUnknown |
| 1371 | }; |
| 1372 | let priced = report_provider_response_usage_origin( |
| 1373 | &self.job.authority.runtime, |
| 1374 | &self.job.authority.owner_agent_id, |
| 1375 | &self.source, |
| 1376 | &self.route, |
| 1377 | usage, |
| 1378 | reason, |
| 1379 | ); |
| 1380 | self.settled = Some((usage.clone(), priced, reason)); |
| 1381 | self.project(usage, priced, reason).await; |
| 1382 | self.projected = true; |
| 1383 | } |
| 1384 | async fn project( |
| 1385 | &self, |
| 1386 | usage: &Usage, |
| 1387 | priced: Option<u64>, |
| 1388 | reason: crate::cost_status::RuntimeUsageMissingReason, |
| 1389 | ) { |
| 1390 | self.job |
| 1391 | .authority |
| 1392 | .accounting_projection() |
| 1393 | .project_settled(&self.source, &self.route, usage, priced, reason) |
| 1394 | .await; |
| 1395 | } |
| 1396 | } |
| 1397 | impl Drop for ChildDispatchedRequest { |
| 1398 | fn drop(&mut self) { |
| 1399 | if self.refused || self.projected { |
| 1400 | return; |
| 1401 | } |
| 1402 | let (usage, priced, reason) = self.settled.clone().unwrap_or_else(|| { |
| 1403 | let reason = crate::cost_status::RuntimeUsageMissingReason::RequestOutcomeUnknown; |
| 1404 | report_provider_response_usage_origin( |
| 1405 | &self.job.authority.runtime, |
| 1406 | &self.job.authority.owner_agent_id, |
| 1407 | &self.source, |
| 1408 | &self.route, |
| 1409 | &Usage::default(), |
| 1410 | reason, |
| 1411 | ); |
| 1412 | (Usage::default(), None, reason) |
| 1413 | }); |
| 1414 | // If the actor was aborted while waiting for this projection, the |
| 1415 | // exact already-settled cost is not billed again. The existing held |
| 1416 | // scheduler retires only worker metadata under the same source id. |
| 1417 | self.job.authority.accounting_projection().recover_settled( |
| 1418 | &self.job.scheduler, |
| 1419 | self.source.clone(), |
| 1420 | self.route.clone(), |
| 1421 | usage, |
| 1422 | priced, |
| 1423 | reason, |
| 1424 | ); |
| 1425 | } |
| 1426 | } |
| 1427 | |
| 1428 | #[derive(Debug)] |
| 1429 | pub(super) struct UnsettledChildCancellation(pub(super) String); |
| 1430 | impl std::fmt::Display for UnsettledChildCancellation { |
| 1431 | fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { |
| 1432 | write!(f, "{}; approval receipts may remain pending", self.0) |
| 1433 | } |
| 1434 | } |
| 1435 | impl std::error::Error for UnsettledChildCancellation {} |
| 1436 | |
| 1437 | struct ChildActorGuard { |
| 1438 | handle: crate::core::engine::EngineHandle, |
| 1439 | actor: tokio::task::AbortHandle, |
| 1440 | joined: bool, |
| 1441 | } |
| 1442 | impl Drop for ChildActorGuard { |
| 1443 | fn drop(&mut self) { |
| 1444 | if self.joined { |
| 1445 | return; |
| 1446 | } |
| 1447 | self.handle.cancel(); |
| 1448 | self.actor.abort(); |
| 1449 | // Emergency abandonment cannot promise a durable decision. Preserve |
| 1450 | // unmatched Asked receipts for the existing protected recovery reader; |
| 1451 | // routine Stop stays in drive_child_actor through the owned join. |
| 1452 | } |
| 1453 | } |
| 1454 | |
| 1455 | /// One existing input's pending-count receipt. The actual Core outcome owns |
| 1456 | /// commit/drop; abandoning the transport also retires its UI pending counter. |
| 1457 | #[cfg(test)] |
| 1458 | pub(crate) fn send_test_child_input( |
| 1459 | tx: &mpsc::UnboundedSender<SubAgentInput>, |
| 1460 | text: &str, |
| 1461 | interrupt: bool, |
| 1462 | pending: Arc<std::sync::atomic::AtomicUsize>, |
| 1463 | ) { |
| 1464 | tx.send(SubAgentInput { |
| 1465 | text: text.into(), |
| 1466 | interrupt, |
| 1467 | pending: Some(pending), |
| 1468 | }) |
| 1469 | .unwrap(); |
| 1470 | } |
| 1471 | |
| 1472 | struct InputSettlement(SubAgentInput); |
| 1473 | impl Drop for InputSettlement { |
| 1474 | fn drop(&mut self) { |
| 1475 | self.0.mark_taken(); |
| 1476 | } |
| 1477 | } |
| 1478 | |
| 1479 | /// One transport adapter into the same Core actor and its run_turn pipeline. |
| 1480 | /// It owns no provider, tool planning, approval decision, or request history. |
| 1481 | #[allow(clippy::too_many_arguments)] |
| 1482 | pub(super) async fn run_child_agent( |
| 1483 | runtime: &SubAgentRuntime, |
| 1484 | agent_id: String, |
| 1485 | agent_type: FleetRole, |
| 1486 | prompt: String, |
| 1487 | assignment: SubAgentAssignment, |
| 1488 | allowed_tools: Option<Vec<String>>, |
| 1489 | fork_context: bool, |
| 1490 | started_at: Instant, |
| 1491 | max_steps: u32, |
| 1492 | parking: Option<Arc<std::sync::atomic::AtomicBool>>, |
| 1493 | input_rx: mpsc::UnboundedReceiver<SubAgentInput>, |
| 1494 | ) -> Result<SubAgentResult> { |
| 1495 | let (runtime, attachment) = if runtime.cancel_token.is_cancelled() { |
| 1496 | (runtime.clone(), None) |
| 1497 | } else { |
| 1498 | prepare_child_membership(runtime, &assignment, &agent_id).await? |
| 1499 | }; |
| 1500 | let authority = ChildAuthority::capture( |
| 1501 | runtime.clone(), |
| 1502 | agent_type.clone(), |
| 1503 | agent_id.clone(), |
| 1504 | assignment |
| 1505 | .role |
| 1506 | .clone() |
| 1507 | .unwrap_or_else(|| agent_type.as_str().to_owned()), |
| 1508 | allowed_tools.clone(), |
| 1509 | ); |
| 1510 | let system = |
| 1511 | build_subagent_system_prompt_with_skills(&agent_type, &assignment, &runtime.context); |
| 1512 | let work_max_steps = if max_steps >= 2 { |
| 1513 | max_steps - 1 |
| 1514 | } else { |
| 1515 | max_steps |
| 1516 | }; |
| 1517 | let prompt = format!( |
| 1518 | "{prompt}\n\n{}", |
| 1519 | child_runtime_budget_context(&runtime, max_steps, work_max_steps) |
| 1520 | ); |
| 1521 | let fork = if fork_context { |
| 1522 | match runtime.fork_context.as_ref() { |
| 1523 | Some(fork) => Some(fork.with_resolved_state_block().await), |
| 1524 | None => None, |
| 1525 | } |
| 1526 | } else { |
| 1527 | None |
| 1528 | }; |
| 1529 | let mut seed = build_initial_subagent_messages_with_system( |
| 1530 | &prompt, |
| 1531 | &assignment, |
| 1532 | &agent_type, |
| 1533 | &system, |
| 1534 | fork.as_ref(), |
| 1535 | ); |
| 1536 | let initial = seed |
| 1537 | .pop() |
| 1538 | .ok_or_else(|| anyhow!("child assignment has no task message"))?; |
| 1539 | let content = initial |
| 1540 | .content |
| 1541 | .iter() |
| 1542 | .filter_map(|block| match block { |
| 1543 | ContentBlock::Text { text, .. } => Some(text.as_str()), |
| 1544 | _ => None, |
| 1545 | }) |
| 1546 | .collect::<Vec<_>>() |
| 1547 | .join("\n"); |
| 1548 | let job = ChildJob::admitted( |
| 1549 | authority.clone(), |
| 1550 | assignment, |
| 1551 | started_at, |
| 1552 | max_steps, |
| 1553 | fork_context, |
| 1554 | parking, |
| 1555 | ) |
| 1556 | .await?; |
| 1557 | if runtime.cancel_token.is_cancelled() { |
| 1558 | // No tool/model is admitted. Preserve the same full seed and the |
| 1559 | // existing Park-vs-Cancel projection, including a step-zero checkpoint. |
| 1560 | seed.push(initial); |
| 1561 | job.project(&seed, 0).await?; |
| 1562 | let checkpoint = runtime |
| 1563 | .manager |
| 1564 | .read() |
| 1565 | .await |
| 1566 | .get_result(&agent_id) |
| 1567 | .ok() |
| 1568 | .and_then(|result| result.checkpoint); |
| 1569 | let mut artifact = job.artifact.lock().await; |
| 1570 | return Ok(cancelled_subagent_result( |
| 1571 | &runtime, |
| 1572 | &agent_id, |
| 1573 | &agent_type, |
| 1574 | &job.assignment, |
| 1575 | &seed, |
| 1576 | 0, |
| 1577 | max_steps, |
| 1578 | checkpoint.as_ref(), |
| 1579 | job.parking.as_ref(), |
| 1580 | &mut artifact, |
| 1581 | started_at, |
| 1582 | fork_context, |
| 1583 | " before Core child admission", |
| 1584 | ) |
| 1585 | .await); |
| 1586 | } |
| 1587 | let api = runtime.api_config.as_deref().ok_or_else(|| anyhow!( |
| 1588 | "child Engine requires its captured provider configuration; no ambient route is admitted" |
| 1589 | ))?; |
| 1590 | let config = crate::core::engine::EngineConfig { |
| 1591 | workspace: runtime.context.workspace.clone(), |
| 1592 | model: runtime.model.clone(), |
| 1593 | max_steps: crate::core::engine::turn_budget::resolve_max_model_steps(Some(work_max_steps)), |
| 1594 | compaction: runtime.compaction.clone(), |
| 1595 | subagent_api_timeout: runtime.step_api_timeout, |
| 1596 | auto_review_policy: runtime.auto_review_policy.as_ref().clone(), |
| 1597 | tools_always_load: allowed_tools.iter().flatten().cloned().collect(), |
| 1598 | locale_tag: runtime.locale_tag.clone(), |
| 1599 | ..Default::default() |
| 1600 | }; |
| 1601 | let (mut core, handle) = crate::core::engine::Engine::new_child_admitted( |
| 1602 | config, |
| 1603 | api, |
| 1604 | authority, |
| 1605 | subagent_request_system_prompt(&system), |
| 1606 | attachment, |
| 1607 | )?; |
| 1608 | core.install_child_job(job.clone(), seed)?; |
| 1609 | let spec = core.child_turn_spec(content, allowed_tools)?; |
| 1610 | handle.send(crate::core::ops::Op::SendMessage(spec)).await?; |
| 1611 | drive_child_actor(core, handle, job, input_rx).await |
| 1612 | } |
| 1613 | |
| 1614 | /// Project an already-decided Core observation under its captured owner. |
| 1615 | /// The existing audit/approval writers retain authority; capacity waits are |
| 1616 | /// bounded by the caller's cancellation and captured deadline/tool timeout. |
| 1617 | pub(crate) async fn forward_child_gate_observation( |
| 1618 | runtime: &SubAgentRuntime, |
| 1619 | owner: &str, |
| 1620 | mut event: Event, |
| 1621 | cancel: &CancellationToken, |
| 1622 | deadline: Option<tokio::time::Instant>, |
| 1623 | ) { |
| 1624 | let Event::ToolGateDecision { agent_id, .. } = &mut event else { |
| 1625 | return; |
| 1626 | }; |
| 1627 | if agent_id.is_none() { |
| 1628 | *agent_id = Some(owner.to_owned()); |
| 1629 | } |
| 1630 | let Some(tx) = runtime.event_tx.as_ref() else { |
| 1631 | return; |
| 1632 | }; |
| 1633 | let deadline = deadline |
| 1634 | .into_iter() |
| 1635 | .chain(runtime.context.turn_deadline) |
| 1636 | .chain(Some(tokio::time::Instant::now() + runtime.tool_timeout)) |
| 1637 | .min() |
| 1638 | .expect("captured tool timeout bounds observation delivery"); |
| 1639 | let permit = tokio::select! { |
| 1640 | biased; |
| 1641 | () = runtime.cancel_token.cancelled() => None, |
| 1642 | () = cancel.cancelled() => None, |
| 1643 | result = tokio::time::timeout_at(deadline, tx.reserve()) => result.ok().and_then(Result::ok), |
| 1644 | }; |
| 1645 | if let Some(permit) = permit { |
| 1646 | permit.send(event); |
| 1647 | } else { |
| 1648 | // Available capacity can still carry the final observation after |
| 1649 | // cancellation. A full/closed host cannot keep the actor alive. |
| 1650 | let _ = tx.try_send(event); |
| 1651 | } |
| 1652 | } |
| 1653 | |
| 1654 | /// Service the existing Core event/approval/input queues until its actor |
| 1655 | /// settles. Shared by the real launch and its exact Core transport tests. |
| 1656 | pub(crate) async fn drive_child_actor( |
| 1657 | core: crate::core::engine::Engine, |
| 1658 | handle: crate::core::engine::EngineHandle, |
| 1659 | job: Arc<ChildJob>, |
| 1660 | mut input_rx: mpsc::UnboundedReceiver<SubAgentInput>, |
| 1661 | ) -> Result<SubAgentResult> { |
| 1662 | let runtime = job.authority.runtime.clone(); |
| 1663 | let agent_id = job.authority.owner_agent_id.clone(); |
| 1664 | let mut actor = tokio::spawn(Box::pin(core.run_child())); |
| 1665 | let mut guard = ChildActorGuard { |
| 1666 | handle: handle.clone(), |
| 1667 | actor: actor.abort_handle(), |
| 1668 | joined: false, |
| 1669 | }; |
| 1670 | let mut approvals = futures_util::stream::FuturesUnordered::new(); |
| 1671 | let mut steers = futures_util::stream::FuturesUnordered::new(); |
| 1672 | let mut events = handle.rx_event.write().await; |
| 1673 | let mut input_open = true; |
| 1674 | let mut pending_input: Option<InputSettlement> = None; |
| 1675 | let mut cancellation_deadline = None; |
| 1676 | // The launch permit alone is not a started turn. Project only Core's |
| 1677 | // admitted lifecycle event, once across work, report and route retries. |
| 1678 | let mut started = false; |
| 1679 | let mut observe_start = |event: &Event| { |
| 1680 | if !started && matches!(event, Event::TurnStarted { .. }) { |
| 1681 | started = true; |
| 1682 | if let Some(mailbox) = runtime.mailbox.as_ref() { |
| 1683 | let _ = mailbox.send(MailboxMessage::started( |
| 1684 | &agent_id, |
| 1685 | job.authority.agent_type.clone(), |
| 1686 | )); |
| 1687 | } |
| 1688 | } |
| 1689 | }; |
| 1690 | let mut result = loop { |
| 1691 | if runtime.cancel_token.is_cancelled() && cancellation_deadline.is_none() { |
| 1692 | handle.cancel(); |
| 1693 | cancellation_deadline = Some(tokio::time::Instant::now() + CHILD_STOP_SETTLE_GRACE); |
| 1694 | } |
| 1695 | use futures_util::StreamExt; |
| 1696 | tokio::select! { |
| 1697 | biased; |
| 1698 | () = runtime.cancel_token.cancelled(), if !handle.is_cancelled() => handle.cancel(), |
| 1699 | joined = &mut actor => { |
| 1700 | guard.joined = true; |
| 1701 | break joined.map_err(|e| anyhow!(UnsettledChildCancellation(format!("child Core actor failed: {e}"))))?; |
| 1702 | }, |
| 1703 | () = async { |
| 1704 | match cancellation_deadline { |
| 1705 | Some(deadline) => tokio::time::sleep_until(deadline).await, |
| 1706 | None => std::future::pending::<()>().await, |
| 1707 | } |
| 1708 | } => { |
| 1709 | actor.abort(); |
| 1710 | let _ = (&mut actor).await; |
| 1711 | guard.joined = true; |
| 1712 | break Err(anyhow!(UnsettledChildCancellation("child cancellation settlement deadline elapsed".into()))); |
| 1713 | }, |
| 1714 | input = input_rx.recv(), if input_open && pending_input.is_none() => match input { |
| 1715 | Some(input) => pending_input = Some(InputSettlement(input)), |
| 1716 | None => input_open = false, |
| 1717 | }, |
| 1718 | reservation = handle.reserve_steer(), if pending_input.is_some() => { |
| 1719 | let input = pending_input.take().expect("selected pending input"); |
| 1720 | match reservation { |
| 1721 | Ok(permit) => { |
| 1722 | let outcome = if input.0.interrupt { permit.send_replacing_with_outcome(input.0.text.clone()) } |
| 1723 | else { permit.send_with_outcome(input.0.text.clone()) }; |
| 1724 | steers.push(async move { let result = outcome.await; (input, result) }); |
| 1725 | } |
| 1726 | Err(_) => drop(input), // Core closed before admission: explicitly dropped. |
| 1727 | } |
| 1728 | }, |
| 1729 | Some((input, _outcome)) = steers.next(), if !steers.is_empty() => { |
| 1730 | // Core's commit/drop outcome, not transport admission, retires it. |
| 1731 | drop(input); |
| 1732 | }, |
| 1733 | Some((id, decision)) = approvals.next(), if !approvals.is_empty() => { |
| 1734 | // Manager removed the slot before this exact Core inbox delivery. |
| 1735 | if runtime.cancel_token.is_cancelled() || handle.is_cancelled() { |
| 1736 | // Core observes the exact cancelled TurnControl before any |
| 1737 | // queued Allow; do not stall its join on another inbox send. |
| 1738 | handle.cancel(); |
| 1739 | } else { |
| 1740 | match decision { |
| 1741 | Ok(ChildApprovalOutcome::Approved) => handle.approve_tool_call(id).await?, |
| 1742 | Ok(ChildApprovalOutcome::Denied) => handle.deny_tool_call(id).await?, |
| 1743 | _ => handle.deny_tool_call_unavailable(id).await?, |
| 1744 | } |
| 1745 | } |
| 1746 | }, |
| 1747 | event = events.recv() => { |
| 1748 | let Some(event) = event else { break Err(anyhow!("child Core event stream closed before join")); }; |
| 1749 | observe_start(&event); |
| 1750 | match &event { |
| 1751 | Event::ApprovalRequired { id, tool_name, description, .. } => { |
| 1752 | let (_, answer) = runtime.manager.write().await.register_child_approval(&agent_id, id, tool_name, description)?; |
| 1753 | let key = id.clone(); approvals.push(async move { (key, answer.await) }); |
| 1754 | record_agent_progress(&runtime, &agent_id, AgentProgressEventMeta::new(AgentWorkerStatus::WaitingForUser) |
| 1755 | .with_tool(tool_name.clone()).with_approval_id(id.clone()), description.clone()); |
| 1756 | if let Some(tx) = runtime.event_tx.as_ref().filter(|_| runtime.parent_can_prompt) { |
| 1757 | let admitted_cancel = handle.captured_turn_cancel(); |
| 1758 | let sent = tokio::select! { |
| 1759 | biased; |
| 1760 | () = runtime.cancel_token.cancelled() => false, |
| 1761 | () = admitted_cancel.cancelled() => false, |
| 1762 | result = tx.send(event.clone()) => result.is_ok(), |
| 1763 | }; |
| 1764 | if !sent { |
| 1765 | if runtime.cancel_token.is_cancelled() || handle.is_cancelled() { |
| 1766 | handle.cancel(); |
| 1767 | } else { |
| 1768 | runtime.manager.write().await.cancel_child_approval(id); |
| 1769 | handle.deny_tool_call_unavailable(id).await?; |
| 1770 | } |
| 1771 | } |
| 1772 | } else if runtime.cancel_token.is_cancelled() || handle.is_cancelled() { |
| 1773 | handle.cancel(); |
| 1774 | } else { |
| 1775 | runtime.manager.write().await.cancel_child_approval(id); |
| 1776 | handle.deny_tool_call_unavailable(id).await?; |
| 1777 | } |
| 1778 | } |
| 1779 | Event::ApprovalWithdrawn { id } => { |
| 1780 | runtime.manager.write().await.cancel_child_approval(id); |
| 1781 | if let Some(tx) = runtime.event_tx.as_ref() { |
| 1782 | tokio::select! { |
| 1783 | biased; |
| 1784 | () = runtime.cancel_token.cancelled() => { let _ = tx.try_send(event.clone()); } |
| 1785 | result = tx.send(event.clone()) => { let _ = result; } |
| 1786 | } |
| 1787 | } |
| 1788 | } |
| 1789 | Event::ToolGateDecision { .. } => { |
| 1790 | forward_child_gate_observation( |
| 1791 | &runtime, |
| 1792 | &agent_id, |
| 1793 | event, |
| 1794 | &handle.captured_turn_cancel(), |
| 1795 | job.hard_deadline.map(Into::into), |
| 1796 | ).await; |
| 1797 | } |
| 1798 | Event::ToolExecutionStarted { id } => { |
| 1799 | if let Some(mailbox) = runtime.mailbox.as_ref() { |
| 1800 | let _ = mailbox.send(MailboxMessage::ToolCallStarted { agent_id: agent_id.clone(), tool_name: format!("Core tool {id}"), step: job.steps() }); |
| 1801 | } |
| 1802 | } |
| 1803 | Event::ToolCallComplete { id, name, result, .. } => { |
| 1804 | record_agent_progress(&runtime, &agent_id, AgentProgressEventMeta::new(AgentWorkerStatus::RunningTool).with_tool(name.clone()), |
| 1805 | format!("Core tool {id} {}", if result.is_ok() { "settled" } else { "refused or failed" })); |
| 1806 | } |
| 1807 | Event::Status { message, .. } => record_agent_progress(&runtime, &agent_id, |
| 1808 | AgentProgressEventMeta::new(AgentWorkerStatus::Running), message.clone()), |
| 1809 | _ => {} // Session is projected by Core; aggregate usage is never rebilled here. |
| 1810 | } |
| 1811 | }, |
| 1812 | } |
| 1813 | }; |
| 1814 | if !guard.joined { |
| 1815 | handle.cancel(); |
| 1816 | result = match tokio::time::timeout(CHILD_STOP_SETTLE_GRACE, &mut actor).await { |
| 1817 | Ok(joined) => joined.map_err(|error| { |
| 1818 | anyhow!(UnsettledChildCancellation(format!( |
| 1819 | "child Core actor failed: {error}" |
| 1820 | ))) |
| 1821 | })?, |
| 1822 | Err(_) => { |
| 1823 | actor.abort(); |
| 1824 | let _ = (&mut actor).await; |
| 1825 | Err(anyhow!(UnsettledChildCancellation( |
| 1826 | "child cancellation settlement deadline elapsed".into() |
| 1827 | ))) |
| 1828 | } |
| 1829 | }; |
| 1830 | guard.joined = true; |
| 1831 | } |
| 1832 | // The join can win over events already queued by Core. Retain its start |
| 1833 | // and final gate observations; this projection never asks or decides a call. |
| 1834 | while let Ok(event) = events.try_recv() { |
| 1835 | observe_start(&event); |
| 1836 | if let Event::ApprovalWithdrawn { id } = &event { |
| 1837 | runtime.manager.write().await.cancel_child_approval(id); |
| 1838 | if let Some(tx) = &runtime.event_tx { |
| 1839 | let _ = tx.try_send(event); |
| 1840 | } |
| 1841 | } else if matches!(event, Event::ToolGateDecision { .. }) { |
| 1842 | forward_child_gate_observation( |
| 1843 | &runtime, |
| 1844 | &agent_id, |
| 1845 | event, |
| 1846 | &handle.captured_turn_cancel(), |
| 1847 | job.hard_deadline.map(Into::into), |
| 1848 | ) |
| 1849 | .await; |
| 1850 | } |
| 1851 | } |
| 1852 | // Joining drops the Core receiver and every unsettled steer. Retire those |
| 1853 | // exact outcomes even when the join branch won before the outcome branch. |
| 1854 | use futures_util::StreamExt; |
| 1855 | while let Some((input, _outcome)) = steers.next().await { |
| 1856 | drop(input); |
| 1857 | } |
| 1858 | drop(pending_input); |
| 1859 | while let Ok(input) = input_rx.try_recv() { |
| 1860 | input.mark_taken(); |
| 1861 | } |
| 1862 | let pending = runtime |
| 1863 | .manager |
| 1864 | .read() |
| 1865 | .await |
| 1866 | .pending_requests_for_agent(&agent_id); |
| 1867 | for pending in pending { |
| 1868 | runtime |
| 1869 | .manager |
| 1870 | .write() |
| 1871 | .await |
| 1872 | .cancel_child_approval(&pending.approval_id); |
| 1873 | if let Some(tx) = runtime.event_tx.as_ref() { |
| 1874 | let event = Event::ApprovalWithdrawn { |
| 1875 | id: pending.approval_id, |
| 1876 | }; |
| 1877 | if runtime.cancel_token.is_cancelled() { |
| 1878 | let _ = tx.try_send(event); |
| 1879 | } else if let Some(deadline) = job.hard_deadline { |
| 1880 | let _ = tokio::time::timeout_at(deadline.into(), tx.send(event)).await; |
| 1881 | } else { |
| 1882 | let _ = tx.send(event).await; |
| 1883 | } |
| 1884 | } |
| 1885 | } |
| 1886 | result |
| 1887 | } |
| 1888 | |
| 1889 | async fn prepare_child_membership( |
| 1890 | runtime: &SubAgentRuntime, |
| 1891 | assignment: &SubAgentAssignment, |
| 1892 | agent_id: &str, |
| 1893 | ) -> Result<( |
| 1894 | SubAgentRuntime, |
| 1895 | Option<crate::extension_host::HostAttachment>, |
| 1896 | )> { |
| 1897 | let mut scoped_runtime = runtime.clone(); |
| 1898 | let extension_host = if crate::plugins::activation::extension_host_policy_enabled() { |
| 1899 | if let Some(plugins) = runtime.context.plugin_registry.as_ref() { |
| 1900 | let selected = if let Some(preset) = assignment.native_preset.as_ref() { |
| 1901 | let plugins = Arc::clone(plugins); |
| 1902 | let preset = preset.clone(); |
| 1903 | let policy = crate::plugins::activation::extension_host_policy_enabled(); |
| 1904 | #[cfg(test)] |
| 1905 | let env_scope = crate::test_support::env_scope_ticket(); |
| 1906 | tokio::task::spawn_blocking(move || { |
| 1907 | #[cfg(test)] |
| 1908 | let _env_scope = crate::test_support::join_env_scope(env_scope); |
| 1909 | let _policy = crate::plugins::activation::PolicyScope::propagate(policy); |
| 1910 | plugins.with_native_preset(preset) |
| 1911 | }) |
| 1912 | .await |
| 1913 | .map_err(|error| anyhow!("Native preset validation worker failed: {error}"))? |
| 1914 | .map_err(anyhow::Error::msg)? |
| 1915 | } else { |
| 1916 | plugins.as_ref().clone() |
| 1917 | }; |
| 1918 | let attachment = crate::extension_host::manager().attach(Arc::new(selected)); |
| 1919 | attachment.set_identity( |
| 1920 | runtime |
| 1921 | .context |
| 1922 | .execution |
| 1923 | .session_objects |
| 1924 | .as_ref() |
| 1925 | .map(|session| session.session_id.clone()), |
| 1926 | Some(agent_id.to_owned()), |
| 1927 | ); |
| 1928 | attachment.reconcile().await.map_err(anyhow::Error::msg)?; |
| 1929 | scoped_runtime.context = scoped_runtime |
| 1930 | .context |
| 1931 | .clone() |
| 1932 | .with_plugin_registry(attachment.plugin_view()); |
| 1933 | if let Some(parent_pool) = runtime.mcp_pool.as_ref() { |
| 1934 | let mut pool = parent_pool |
| 1935 | .lock() |
| 1936 | .await |
| 1937 | .fork_for_plugins(attachment.plugin_view())?; |
| 1938 | // Discover this caller's selected catalog through the existing |
| 1939 | // Core connection factory before the child registry snapshots it. |
| 1940 | // Optional connection failures retain Core's lazy retry behavior. |
| 1941 | for (name, error) in pool.connect_all().await { |
| 1942 | tracing::warn!("child MCP server {name} unavailable: {error}"); |
| 1943 | } |
| 1944 | scoped_runtime.mcp_pool = Some(Arc::new(tokio::sync::Mutex::new(pool))); |
| 1945 | } |
| 1946 | Some(attachment) |
| 1947 | } else if assignment.native_preset.is_some() { |
| 1948 | return Err(anyhow!( |
| 1949 | "Native composition requires the caller's reviewed plugin inventory" |
| 1950 | )); |
| 1951 | } else { |
| 1952 | None |
| 1953 | } |
| 1954 | } else if assignment.native_preset.is_some() { |
| 1955 | return Err(anyhow!( |
| 1956 | "Native composition requires Experimental extension_host" |
| 1957 | )); |
| 1958 | } else { |
| 1959 | None |
| 1960 | }; |
| 1961 | |
| 1962 | Ok((scoped_runtime, extension_host)) |
| 1963 | } |
| 1964 | |
| 1965 | /// Counters for one logical Engine model step. They do not dispatch or own |
| 1966 | /// a turn; the existing outer run_turn loop consumes this policy decision. |
| 1967 | #[derive(Default)] |
| 1968 | pub(crate) struct ChildRequestRetries { |
| 1969 | transient: u32, |
| 1970 | timeouts: u32, |
| 1971 | } |
| 1972 | pub(crate) enum ChildRequestRecovery { |
| 1973 | Retry { |
| 1974 | delay: Duration, |
| 1975 | note: String, |
| 1976 | }, |
| 1977 | Interrupted { |
| 1978 | checkpoint_reason: &'static str, |
| 1979 | message: String, |
| 1980 | }, |
| 1981 | } |
| 1982 | impl ChildRequestRetries { |
| 1983 | pub(crate) fn decide( |
| 1984 | &mut self, |
| 1985 | runtime: &SubAgentRuntime, |
| 1986 | error: &anyhow::Error, |
| 1987 | ) -> Option<ChildRequestRecovery> { |
| 1988 | if matches!(error.downcast_ref::<LlmError>(), Some(LlmError::Timeout(_))) { |
| 1989 | if self.timeouts >= SUBAGENT_API_TIMEOUT_MAX_RETRIES { |
| 1990 | return Some(ChildRequestRecovery::Interrupted { |
| 1991 | checkpoint_reason: "api_timeout", |
| 1992 | message: format!( |
| 1993 | "API call timed out after {}ms on {} API attempt(s); checkpoint preserved for continuation", |
| 1994 | runtime.step_api_timeout.as_millis(), |
| 1995 | self.timeouts.saturating_add(1) |
| 1996 | ), |
| 1997 | }); |
| 1998 | } |
| 1999 | self.timeouts = self.timeouts.saturating_add(1); |
| 2000 | let delay = subagent_api_timeout_retry_delay( |
| 2001 | self.timeouts, |
| 2002 | runtime.api_timeout_retry_base_backoff, |
| 2003 | ); |
| 2004 | return Some(ChildRequestRecovery::Retry { |
| 2005 | delay, |
| 2006 | note: format!( |
| 2007 | "API call timed out after {}ms; retrying API request {}/{} in {}ms", |
| 2008 | runtime.step_api_timeout.as_millis(), |
| 2009 | self.timeouts, |
| 2010 | SUBAGENT_API_TIMEOUT_MAX_RETRIES, |
| 2011 | delay.as_millis() |
| 2012 | ), |
| 2013 | }); |
| 2014 | } |
| 2015 | let retryable = |
| 2016 | retryable_subagent_provider_failure(error, self.transient.saturating_add(1))?; |
| 2017 | if self.transient >= SUBAGENT_TRANSIENT_PROVIDER_MAX_RETRIES { |
| 2018 | return Some(ChildRequestRecovery::Interrupted { |
| 2019 | checkpoint_reason: retryable.checkpoint_reason, |
| 2020 | message: format!( |
| 2021 | "{} after {} API attempt(s): {error:#}; checkpoint preserved for continuation", |
| 2022 | retryable.label, |
| 2023 | self.transient.saturating_add(1) |
| 2024 | ), |
| 2025 | }); |
| 2026 | } |
| 2027 | self.transient = self.transient.saturating_add(1); |
| 2028 | Some(ChildRequestRecovery::Retry { |
| 2029 | delay: retryable.delay, |
| 2030 | note: format!( |
| 2031 | "{}; retrying API request {}/{} in {}ms ({error:#})", |
| 2032 | retryable.label, |
| 2033 | self.transient, |
| 2034 | SUBAGENT_TRANSIENT_PROVIDER_MAX_RETRIES, |
| 2035 | retryable.delay.as_millis() |
| 2036 | ), |
| 2037 | }) |
| 2038 | } |
| 2039 | } |
| 2040 | |
| 2041 | /// Reuse the approved first-request route policy. The caller installs it in |
| 2042 | /// the existing Engine; this helper has no provider, planner or retry loop. |
| 2043 | pub(crate) fn approved_first_request_replacement( |
| 2044 | job: &ChildJob, |
| 2045 | current: &SubAgentRuntime, |
| 2046 | replacements_tried: &mut usize, |
| 2047 | pin_fallback_used: &mut bool, |
| 2048 | error: &anyhow::Error, |
| 2049 | ) -> Result<Option<(SubAgentRuntime, SpawnRouteSource, String)>> { |
| 2050 | if !job.can_replace_first_request() { |
| 2051 | return Ok(None); |
| 2052 | } |
| 2053 | let detail: String = current |
| 2054 | .client |
| 2055 | .redact_model_bound_text(&format!("{error}")) |
| 2056 | .chars() |
| 2057 | .take(160) |
| 2058 | .collect(); |
| 2059 | if !*pin_fallback_used |
| 2060 | && let Some(origin) = current.route_origin.as_deref() |
| 2061 | && let Some(parent) = origin.parent.as_ref() |
| 2062 | && let Some(why) = pin_refusal_reason(error) |
| 2063 | { |
| 2064 | *pin_fallback_used = true; |
| 2065 | let note: String = format!( |
| 2066 | "{}, which failed authorization ({why}: {detail}); ran on {} instead", |
| 2067 | origin.source, parent.label |
| 2068 | ) |
| 2069 | .chars() |
| 2070 | .take(480) |
| 2071 | .collect(); |
| 2072 | let mut next = current.clone(); |
| 2073 | parent.install(&mut next); |
| 2074 | next.route_origin = Some(Arc::new(spawn_route_origin( |
| 2075 | SpawnRouteSource::SessionFallback, |
| 2076 | current.worker_profile.role.as_str(), |
| 2077 | None, |
| 2078 | parent.label.clone(), |
| 2079 | Some(¬e), |
| 2080 | ))); |
| 2081 | return Ok(Some((next, SpawnRouteSource::SessionFallback, note))); |
| 2082 | } |
| 2083 | let Some(why) = route_replacement_reason(error) else { |
| 2084 | return Ok(None); |
| 2085 | }; |
| 2086 | let original = &job.authority.runtime; |
| 2087 | let mut skipped = Vec::new(); |
| 2088 | while let Some(route) = original.route_replacements.get(*replacements_tried) { |
| 2089 | *replacements_tried += 1; |
| 2090 | let to = format!( |
| 2091 | "{}/{}", |
| 2092 | route.provider.as_deref().unwrap_or_default(), |
| 2093 | route.model |
| 2094 | ); |
| 2095 | match replacement_route_runtime(original, route) { |
| 2096 | Ok(mut next) => { |
| 2097 | let from = format!( |
| 2098 | "{}/{}", |
| 2099 | current.client.api_provider().as_str(), |
| 2100 | current.model |
| 2101 | ); |
| 2102 | let mut note = format!( |
| 2103 | "{from} refused the first request before any work ({why}: {detail}); moved to approved replacement {to} (attempt {} of {})", |
| 2104 | *replacements_tried, |
| 2105 | original.route_replacements.len() |
| 2106 | ); |
| 2107 | if !skipped.is_empty() { |
| 2108 | note.push_str("; skipped "); |
| 2109 | note.push_str(&skipped.join("; ")); |
| 2110 | } |
| 2111 | let note: String = note.chars().take(480).collect(); |
| 2112 | next.route_origin = Some(Arc::new(spawn_route_origin( |
| 2113 | SpawnRouteSource::RoleReplacement, |
| 2114 | original.worker_profile.role.as_str(), |
| 2115 | None, |
| 2116 | runtime_route_label(&next), |
| 2117 | Some(¬e), |
| 2118 | ))); |
| 2119 | return Ok(Some((next, SpawnRouteSource::RoleReplacement, note))); |
| 2120 | } |
| 2121 | Err(unavailable) => skipped.push(format!( |
| 2122 | "{to} ({})", |
| 2123 | unavailable.chars().take(120).collect::<String>() |
| 2124 | )), |
| 2125 | } |
| 2126 | } |
| 2127 | if skipped.is_empty() { |
| 2128 | Ok(None) |
| 2129 | } else { |
| 2130 | Err(anyhow!( |
| 2131 | "no approved replacement route could take the task: {}", |
| 2132 | skipped.join("; ") |
| 2133 | )) |
| 2134 | } |
| 2135 | } |
| 2136 | |
| 2137 | #[cfg(test)] |
| 2138 | mod owner_origin_tests { |
| 2139 | use super::*; |
| 2140 | |
| 2141 | #[tokio::test] |
| 2142 | async fn canonical_child_unbounded_claim_checks_exact_held_origin_before_peer_projection() { |
| 2143 | let _environment = crate::test_support::lock_test_env(); |
| 2144 | let temporary = tempfile::tempdir().unwrap(); |
| 2145 | let root = temporary.path().canonicalize().unwrap(); |
| 2146 | let original = root.join("original"); |
| 2147 | let selected = root.join("selected"); |
| 2148 | fs::create_dir_all(original.join("src")).unwrap(); |
| 2149 | fs::create_dir_all(selected.join("src")).unwrap(); |
| 2150 | let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", root.join("home")); |
| 2151 | let manager = new_shared_subagent_manager(original.clone(), 2); |
| 2152 | { |
| 2153 | let mut held = manager.write().await; |
| 2154 | for lexical in [&original, &selected] { |
| 2155 | let (canonical, file) = |
| 2156 | crate::runtime_api::open_workspace_directory(lexical).unwrap(); |
| 2157 | held.admit_coordination_workspace(lexical.clone(), canonical, Arc::new(file)) |
| 2158 | .unwrap(); |
| 2159 | } |
| 2160 | let spec = super::super::owner_scoped_coordination_tests::writer( |
| 2161 | "selected-writer", |
| 2162 | &selected, |
| 2163 | "src", |
| 2164 | ); |
| 2165 | held.register_worker_with_coordination(spec).unwrap(); |
| 2166 | } |
| 2167 | let mut runtime = super::super::tests::stub_runtime(); |
| 2168 | runtime.manager = manager.clone(); |
| 2169 | runtime.context = ToolContext::new(selected.clone()); |
| 2170 | runtime.worker_profile = WorkerRuntimeProfile::for_role(FleetRole::Builder); |
| 2171 | runtime.allow_shell = true; |
| 2172 | let authority = ChildAuthority::capture( |
| 2173 | runtime, |
| 2174 | FleetRole::Builder, |
| 2175 | "selected-writer".into(), |
| 2176 | "writer".into(), |
| 2177 | None, |
| 2178 | ); |
| 2179 | let mut registry = ToolRegistry::new(ToolContext::new(selected.clone())); |
| 2180 | registry.register(Arc::new(crate::tools::shell::BashTool::new("Bash"))); |
| 2181 | let input = serde_json::json!({"command":"python3 -c 'print(1)'"}); |
| 2182 | // No peer exists: ordinary work is still admitted under its held root. |
| 2183 | assert!( |
| 2184 | authority |
| 2185 | .validate_claim(®istry, "bash", &input) |
| 2186 | .await |
| 2187 | .unwrap() |
| 2188 | .is_empty() |
| 2189 | ); |
| 2190 | { |
| 2191 | let mut held = manager.write().await; |
| 2192 | let foreign = held.admitted_coordination_roots[&original].receipt.clone(); |
| 2193 | held.worker_origins |
| 2194 | .get_mut("selected-writer") |
| 2195 | .unwrap() |
| 2196 | .scope = Some(foreign); |
| 2197 | } |
| 2198 | let refused = authority |
| 2199 | .validate_claim(®istry, "bash", &input) |
| 2200 | .await |
| 2201 | .unwrap_err(); |
| 2202 | assert!( |
| 2203 | refused.to_string().contains("owner re-admission"), |
| 2204 | "{refused}" |
| 2205 | ); |
| 2206 | assert!( |
| 2207 | manager |
| 2208 | .read() |
| 2209 | .await |
| 2210 | .live_peer_shared_write_claim_owners("selected-writer") |
| 2211 | .is_empty() |
| 2212 | ); |
| 2213 | assert!(!selected.join("effect").exists()); |
| 2214 | } |
| 2215 | } |
| 2216 |