返回 CodeWhale
engine.rs
根目录 / crates / tui / src / tools / subagent / engine.rs
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 &registry.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 &current.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(&note),
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(&note),
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(&registry, "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(&registry, "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
2216 lines RUST