返回 CodeWhale
child_host.rs
根目录 / crates / tui / src / core / engine / child_host.rs
1 //! Private captured child configuration for the canonical Engine.
2 use super::*;
3 use crate::tools::subagent::engine::ChildAuthority;
4 use anyhow::anyhow;
5
6 pub(super) struct ChildHostSetup {
7 pub authority: Arc<ChildAuthority>,
8 pub system_prompt: SystemPrompt,
9 pub attachment: Option<crate::extension_host::HostAttachment>,
10 }
11 pub(super) struct ChildHostState {
12 pub authority: Arc<ChildAuthority>,
13 pub system_prompt: SystemPrompt,
14 job: Option<Arc<crate::tools::subagent::engine::ChildJob>>,
15 handle: Option<EngineHandle>,
16 pub(super) report: Option<crate::tools::subagent::budget_handback::ReportAdmission>,
17 route_runtime: crate::tools::subagent::SubAgentRuntime,
18 replacements_tried: usize,
19 pin_fallback_used: bool,
20 pending_route: Option<(
21 crate::tools::subagent::SubAgentRuntime,
22 ValidatedRuntimeRoute,
23 )>,
24 pub(super) request_stop_reason: Option<&'static str>,
25 }
26 impl From<ChildHostSetup> for ChildHostState {
27 fn from(setup: ChildHostSetup) -> Self {
28 let route_runtime = setup.authority.runtime.clone();
29 Self {
30 route_runtime,
31 replacements_tried: 0,
32 pin_fallback_used: false,
33 pending_route: None,
34 request_stop_reason: None,
35 authority: setup.authority,
36 system_prompt: setup.system_prompt,
37 job: None,
38 handle: None,
39 report: None,
40 }
41 }
42 }
43 impl Engine {
44 pub(super) fn new_tool_execution_id(&self) -> String {
45 self.child_host.as_ref().map_or_else(
46 || uuid::Uuid::new_v4().to_string(),
47 |child| child.authority.new_execution_id(),
48 )
49 }
50
51 pub(crate) fn new_child_admitted(
52 mut config: EngineConfig,
53 api_config: &Config,
54 authority: Arc<ChildAuthority>,
55 system_prompt: SystemPrompt,
56 attachment: Option<crate::extension_host::HostAttachment>,
57 ) -> Result<(Self, EngineHandle)> {
58 anyhow::ensure!(
59 config.workspace == authority.runtime.context.workspace,
60 "child Engine workspace differs from its captured caller"
61 );
62 anyhow::ensure!(
63 !authority.runtime.context.state_namespace.trim().is_empty(),
64 "child Engine needs its captured root session namespace"
65 );
66 config.session_id = Some(authority.runtime.context.state_namespace.clone());
67 config.plugin_registry = authority.runtime.context.plugin_registry.clone();
68 config.runtime_services = authority.runtime.context.runtime.clone();
69 config.todos = authority.runtime.todos.clone();
70 config.max_spawn_depth = authority.runtime.max_spawn_depth;
71 config.allow_shell = authority.runtime.allow_shell;
72 config.trust_mode = authority.runtime.context.trust_mode;
73 config.network_policy = authority.runtime.context.network_policy.clone();
74 if let Some(skills_dir) = authority.runtime.context.skills_dir.as_ref() {
75 config.skills_dir = skills_dir.clone();
76 }
77 config.skills_discovery_mode = authority.runtime.context.skills_discovery_mode;
78 config.features = authority.runtime.context.features.clone();
79 config.snapshots_enabled = false;
80 config.memory_enabled = false;
81 config.terminal_chrome_enabled = false;
82 config.advisor_config = crate::tools::subagent::AdvisorConfig::disabled();
83 config.goal_objective = None;
84 config.goal_status = GoalStatus::Paused;
85 config.exec_policy_engine = api_config.exec_policy_engine.clone();
86 let (mut engine, handle) = Self::new_admitted(
87 config,
88 api_config,
89 Some(EngineHostSetup::Child(ChildHostSetup {
90 authority,
91 system_prompt,
92 attachment,
93 })),
94 );
95 engine
96 .child_host
97 .as_mut()
98 .expect("captured child setup")
99 .handle = Some(handle.clone());
100 // The actual caller supplied this frozen concrete client. Route
101 // identity is still resolved and installed by the same Core seam.
102 engine.model_client_injected = true;
103 let posture = engine.runtime_authority_snapshot();
104 engine.apply_runtime_mode_policy(&TurnAuthority::from_effective_fields(
105 posture.mode,
106 posture.allow_shell,
107 posture.trust_mode,
108 posture.auto_approve,
109 posture.approval_mode,
110 ));
111 Ok((engine, handle))
112 }
113 }
114
115 impl Engine {
116 pub(crate) fn install_child_job(
117 &mut self,
118 job: Arc<crate::tools::subagent::engine::ChildJob>,
119 seed: Vec<Message>,
120 ) -> Result<()> {
121 let state = self
122 .child_host
123 .as_mut()
124 .ok_or_else(|| anyhow!("not a captured child Engine"))?;
125 anyhow::ensure!(
126 Arc::ptr_eq(&state.authority, &job.authority),
127 "child job changed captured authority"
128 );
129 anyhow::ensure!(state.job.is_none(), "child assignment is already installed");
130 let system = state.system_prompt.clone();
131 state.job = Some(job);
132 self.restore_session_history(seed, Some(system), true);
133 Ok(())
134 }
135
136 pub(crate) fn child_turn_spec(
137 &self,
138 content: String,
139 allowed_tools: Option<Vec<String>>,
140 ) -> Result<TurnSpec> {
141 let state = self
142 .child_host
143 .as_ref()
144 .ok_or_else(|| anyhow!("not a captured child Engine"))?;
145 state
146 .authority
147 .validate_context(&state.authority.context())?;
148 let posture = self.runtime_authority_snapshot();
149 Ok(TurnSpec {
150 content,
151 images: Vec::new(),
152 mode: posture.mode,
153 route: Box::new(self.current_runtime_route().map_err(anyhow::Error::msg)?),
154 compaction: Box::new(self.config.compaction.clone()),
155 initial_routed_usage: Box::default(),
156 goal_objective: None,
157 goal_token_budget: None,
158 goal_status: GoalStatus::Paused,
159 reasoning_effort: state.route_runtime.reasoning_effort.clone(),
160 reasoning_effort_auto: false,
161 auto_model: false,
162 allow_shell: state.authority.runtime.allow_shell && posture.allow_shell,
163 trust_mode: state.authority.runtime.context.trust_mode,
164 auto_approve: posture.auto_approve,
165 approval_mode: posture.approval_mode,
166 translation_enabled: false,
167 allowed_tools,
168 dynamic_tools: Vec::new(),
169 hook_executor: self.config.hook_executor.clone(),
170 verbosity: None,
171 provenance: UserInputProvenance::Runtime,
172 submission_id: None,
173 max_output_tokens: None,
174 })
175 }
176
177 pub(super) fn child_job(&self) -> Option<Arc<crate::tools::subagent::engine::ChildJob>> {
178 self.child_host.as_ref().and_then(|state| state.job.clone())
179 }
180 pub(super) fn child_request_protocol(&self) -> Option<codewhale_config::provider::WireFormat> {
181 self.child_host
182 .as_ref()
183 .map(|state| state.route_runtime.client.wire_format())
184 }
185 pub(super) fn child_report(&self) -> bool {
186 self.child_host
187 .as_ref()
188 .is_some_and(|state| state.report.is_some())
189 }
190
191 /// Executes a queued assignment or its single bounded reporting turn.
192 /// Both go through handle_send_message -> run_turn, the same Session and
193 /// Core approval inbox. The next report is reserved on the existing FIFO.
194 pub(super) async fn handle_child_send_message(
195 &mut self,
196 spec: TurnSpec,
197 ) -> Result<Option<crate::tools::subagent::SubAgentResult>> {
198 use crate::tools::subagent::SubAgentStatus;
199 let job = self
200 .child_job()
201 .ok_or_else(|| anyhow!("child Engine has no admitted assignment"))?;
202 self.child_host
203 .as_mut()
204 .expect("captured child")
205 .request_stop_reason = None;
206 let reporting = self.child_report();
207 let replay_content = spec.content.clone();
208 let replay_scope = spec.allowed_tools.clone();
209 let deadline = if reporting {
210 self.child_host
211 .as_ref()
212 .and_then(|state| state.report.as_ref())
213 .map(|report| report.deadline)
214 } else {
215 job.work_deadline
216 };
217 let admitted = self
218 .admitted_turn_control
219 .as_ref()
220 .map(|control| control.cancel.clone());
221 // This is the exact queued turn token, never the parent's or a later
222 // report's token. Dropping the timer aborts it after settlement.
223 let timer = deadline.zip(admitted).map(|(deadline, cancel)| {
224 let job = job.clone();
225 turn_heartbeat::AbortOnDrop(tokio::spawn(async move {
226 tokio::time::sleep_until(deadline.into()).await;
227 if !reporting {
228 job.stop_for_budget("child wall-time work budget exhausted");
229 }
230 cancel.cancel();
231 }))
232 });
233 let outcome = Box::pin(self.handle_send_message(spec)).await;
234 // Capture expiry at the actual Core outcome. Later checkpoint/ledger
235 // contention cannot relabel an already-settled provider refusal.
236 let report_deadline_expired =
237 reporting && deadline.is_some_and(|deadline| Instant::now() >= deadline);
238 drop(timer);
239 job.project(&self.session.messages, job.steps()).await?;
240 if job.authority.runtime.cancel_token.is_cancelled() {
241 return job
242 .finish(
243 &self.session.messages,
244 job.steps(),
245 SubAgentStatus::Cancelled,
246 None,
247 None,
248 )
249 .await
250 .map(Some);
251 }
252 let pending_route = self
253 .child_host
254 .as_mut()
255 .and_then(|state| state.pending_route.take());
256 if let Some((runtime, route)) = pending_route {
257 // The old request failed before a response. A fresh admitted turn
258 // reuses the exact original task on this same Engine/FIFO; it rebuilds
259 // tools, Native prompt sections and Core dispatch receipts together.
260 let client = runtime.client.clone();
261 self.install_validated_runtime_route(route);
262 self.model_client = Some(Arc::new(client));
263 self.session.reasoning_effort = runtime.reasoning_effort.clone();
264 self.session.reasoning_effort_auto = runtime.reasoning_effort_auto;
265 job.installed_replacement(&runtime.model);
266 self.child_host
267 .as_mut()
268 .expect("captured child")
269 .route_runtime = runtime;
270 let spec = self.child_turn_spec(replay_content, replay_scope)?;
271 self.child_host
272 .as_ref()
273 .expect("captured child")
274 .handle
275 .as_ref()
276 .expect("same Engine handle")
277 .send(Op::SendMessage(spec))
278 .await?;
279 return Ok(None);
280 }
281 let (status, error) = match outcome {
282 SendMessageOutcome::Finished { status, error } => (status, error),
283 SendMessageOutcome::NotStarted { error } => (TurnOutcomeStatus::Failed, error),
284 };
285 let text = self
286 .session
287 .messages
288 .iter()
289 .rev()
290 .find(|message| {
291 message.role == codewhale_models::Role::Assistant
292 || (matches!(
293 status,
294 TurnOutcomeStatus::Interrupted | TurnOutcomeStatus::Failed
295 ) && message.role == codewhale_models::Role::InterruptedAssistant)
296 })
297 .map(|message| {
298 message
299 .content
300 .iter()
301 .filter_map(|block| match block {
302 ContentBlock::Text { text, .. } if !text.trim().is_empty() => {
303 Some(text.as_str())
304 }
305 _ => None,
306 })
307 .collect::<Vec<_>>()
308 .join("\n")
309 })
310 .filter(|text| !text.trim().is_empty());
311 if let Some(cause) = job.budget_reason() {
312 if !reporting {
313 match crate::tools::subagent::budget_handback::admit_report(
314 &job,
315 &self.session.messages,
316 &cause,
317 self.model_client
318 .as_ref()
319 .expect("installed selected client")
320 .effective_max_output_tokens(&self.session.model),
321 )
322 .await
323 {
324 Ok(report) => {
325 let content = report
326 .messages
327 .iter()
328 .flat_map(|message| &message.content)
329 .filter_map(|block| match block {
330 ContentBlock::Text { text, .. } => Some(text.as_str()),
331 _ => None,
332 })
333 .collect::<Vec<_>>()
334 .join("\n");
335 let mut spec = self.child_turn_spec(content, Some(Vec::new()))?;
336 spec.allow_shell = false;
337 spec.max_output_tokens = std::num::NonZeroU32::new(report.output_tokens);
338 self.config.max_steps = 1;
339 let state = self.child_host.as_mut().expect("captured child");
340 state.report = Some(report);
341 // A fresh queued scope cannot erase parent cancellation:
342 // it is rechecked by the transport and every dispatch.
343 state
344 .handle
345 .as_ref()
346 .expect("same Core handle")
347 .send(Op::SendMessage(spec))
348 .await?;
349 return Ok(None);
350 }
351 Err(reason) => {
352 let partial =
353 crate::tools::subagent::budget_handback::fallback_partial_text(
354 &self.session.messages,
355 );
356 return job
357 .finish(
358 &self.session.messages,
359 job.steps(),
360 SubAgentStatus::BudgetExhausted,
361 Some(format!("{partial}\n\n{reason}")),
362 Some(&cause),
363 )
364 .await
365 .map(Some);
366 }
367 }
368 }
369 let report = if status == TurnOutcomeStatus::Completed && error.is_none() {
370 text
371 } else {
372 None
373 };
374 let mut partial = report.unwrap_or_else(|| {
375 crate::tools::subagent::budget_handback::fallback_partial_text(
376 &self.session.messages,
377 )
378 });
379 let mut receipt = if status == TurnOutcomeStatus::Completed && error.is_none() {
380 "Host budget hand-back receipt: one completed bounded report; assignment remains incomplete.".to_string()
381 } else if report_deadline_expired {
382 "Host budget hand-back receipt: report deadline expired; recorded work is preserved.".to_string()
383 } else {
384 format!(
385 "Host budget hand-back receipt: provider call failed or report did not finish: {}. Recorded work is preserved.",
386 error.as_deref().unwrap_or("Core report interrupted")
387 )
388 };
389 if job
390 .authority
391 .runtime
392 .manager
393 .read()
394 .await
395 .get_worker_record(&job.authority.owner_agent_id)
396 .is_some_and(|record| record.has_unreported_usage)
397 {
398 receipt.push_str(" Measured usage is only a subtotal, not a zero-cost report; unreported provider usage remains unknown.");
399 }
400 partial.push_str(&format!("\n\n{receipt}"));
401 // Record the host's actual outcome alongside this same canonical
402 // Session, before finish saves/projects the result checkpoint.
403 self.add_session_message(
404 self.runtime_text_message_with_turn_metadata(receipt, UserInputProvenance::Runtime),
405 )
406 .await;
407 return job
408 .finish(
409 &self.session.messages,
410 job.steps(),
411 SubAgentStatus::BudgetExhausted,
412 Some(partial),
413 Some(&cause),
414 )
415 .await
416 .map(Some);
417 }
418 let child_status = match (status, error.as_ref(), text.as_ref()) {
419 (TurnOutcomeStatus::Completed, None, Some(_)) => SubAgentStatus::Completed,
420 (TurnOutcomeStatus::Interrupted, _, _) => SubAgentStatus::Interrupted(
421 error
422 .clone()
423 .unwrap_or_else(|| "Core child turn interrupted".into()),
424 ),
425 (TurnOutcomeStatus::Failed, Some(reason), _)
426 if reason == super::dispatch::FLEET_NO_PROGRESS_STOP =>
427 {
428 // Core's terminal no-progress report stays Failed even after
429 // earlier steps; checkpoint preservation cannot erase its cause.
430 SubAgentStatus::Failed(reason.clone())
431 }
432 (TurnOutcomeStatus::Failed, _, _) if job.steps() > 1 => {
433 SubAgentStatus::Interrupted(error.clone().unwrap_or_else(|| {
434 "child stopped after prior work; checkpoint preserved".into()
435 }))
436 }
437 _ => SubAgentStatus::Failed(
438 error
439 .clone()
440 .unwrap_or_else(|| "child stopped without a final summary".into()),
441 ),
442 };
443 job.finish(
444 &self.session.messages,
445 job.steps(),
446 child_status,
447 text,
448 self.child_host
449 .as_ref()
450 .and_then(|state| state.request_stop_reason)
451 .or(error.as_deref()),
452 )
453 .await
454 .map(Some)
455 }
456
457 pub(super) fn clear_child_pending_route(&mut self) {
458 if let Some(state) = self.child_host.as_mut() {
459 state.pending_route = None;
460 }
461 }
462
463 pub(super) async fn retract_child_unsent_message(
464 &mut self,
465 mark: crate::core::turn::UnansweredUserMessage,
466 ) -> Result<bool> {
467 if !self.can_retract_unanswered_user_message(mark) {
468 return Ok(false);
469 }
470 if let Some(job) = self.child_job() {
471 let old = self.session.messages.to_vec();
472 // Preserve durable dispatch/history evidence before changing the
473 // live projection. The one actor owns both sides of this await.
474 if let Err(error) = job.before_replace(&old, &old[..mark.len - 1]).await {
475 self.child_host
476 .as_mut()
477 .expect("captured child")
478 .pending_route = None;
479 return Err(error);
480 }
481 }
482 Ok(self.retract_unanswered_user_message(mark))
483 }
484
485 pub(super) async fn admit_child_first_request_replacement(
486 &mut self,
487 error: &anyhow::Error,
488 ) -> Result<bool> {
489 let Some(job) = self.child_job() else {
490 return Ok(false);
491 };
492 if self.child_report() || !job.can_replace_first_request() {
493 return Ok(false);
494 }
495 let state = self.child_host.as_mut().expect("captured child");
496 let Some((next, source, note)) =
497 crate::tools::subagent::engine::approved_first_request_replacement(
498 &job,
499 &state.route_runtime,
500 &mut state.replacements_tried,
501 &mut state.pin_fallback_used,
502 error,
503 )?
504 else {
505 return Ok(false);
506 };
507 let api = next
508 .api_config
509 .as_deref()
510 .ok_or_else(|| anyhow!("approved child replacement has no captured config"))?;
511 let resolved = crate::route_runtime::resolve_runtime_route_for_identity(
512 api,
513 next.client.admitted_provider_identity(),
514 Some(&next.model),
515 )
516 .map_err(anyhow::Error::msg)?;
517 let mut route = resolved.validate().map_err(anyhow::Error::msg)?;
518 anyhow::ensure!(
519 next.client.admitted_provider_identity() == &route.identity
520 && route.candidate.endpoint().base_url == next.client.base_url(),
521 "approved child replacement client differs from its exact route"
522 );
523 route.client = next.client.clone();
524 job.authority.validate_context(&job.authority.context())?;
525 job.record_route_replacement(&state.route_runtime, &next, source, note.clone(), error)
526 .await;
527 state.pending_route = Some((next, route));
528 let _ = self.send_event(Event::status(note)).await;
529 Ok(true)
530 }
531
532 pub(crate) async fn run_child(self) -> Result<crate::tools::subagent::SubAgentResult> {
533 self.run_owned().await.ok_or_else(|| {
534 anyhow!("child Core actor stopped without a terminal assignment receipt")
535 })?
536 }
537 }
538
539 /// Test observer of Core's actual surface policy and activation cache. It has
540 /// no executor, permission policy or request producer of its own.
541 #[cfg(test)]
542 pub(crate) struct ChildSurfaceProbe {
543 pub(super) policy: ToolSurfacePolicy,
544 pub(super) cache: crate::core::session::ToolActivationCache,
545 }
546 #[cfg(test)]
547 impl ChildSurfaceProbe {
548 pub(crate) fn new(catalog: Vec<codewhale_models::Tool>, warm: &[String]) -> Self {
549 let policy = ToolSurfacePolicy::new(
550 crate::tools::ToolRegistry::new(ToolContext::for_empty_registry()),
551 Some(catalog),
552 AppMode::Agent,
553 &HashSet::new(),
554 &[],
555 false,
556 None,
557 None,
558 None,
559 tool_catalog::ToolMode::Direct,
560 );
561 let mut probe = Self {
562 policy,
563 cache: Default::default(),
564 };
565 let activated = probe.cache.activate(&probe.policy.catalog, warm);
566 probe.policy.active_names.extend(activated.admitted);
567 probe
568 }
569 pub(crate) fn catalog(&self) -> &[codewhale_models::Tool] {
570 &self.policy.catalog
571 }
572 pub(crate) fn catalog_mut(&mut self) -> &mut Vec<codewhale_models::Tool> {
573 &mut self.policy.catalog
574 }
575 pub(crate) fn active_names(&self) -> &HashSet<String> {
576 &self.policy.active_names
577 }
578 pub(crate) fn request_tools(
579 &mut self,
580 catalog: Vec<codewhale_models::Tool>,
581 strict: bool,
582 ) -> Vec<codewhale_models::Tool> {
583 self.policy.catalog = catalog;
584 self.cache.revalidate(&self.policy.catalog);
585 self.policy.active_names = tool_catalog::initial_active_tools(&self.policy.catalog);
586 self.policy
587 .active_names
588 .extend(self.cache.names().map(str::to_owned));
589 tool_catalog::active_tools_for_request(
590 &self.policy.catalog,
591 &self.policy.active_names,
592 strict,
593 )
594 .unwrap_or_default()
595 }
596 /// Pure cache admission probe; execution still runs Core's planner below.
597 pub(crate) fn hydrate(
598 &mut self,
599 name: &str,
600 input: &serde_json::Value,
601 ) -> Result<Option<String>> {
602 let activation = self
603 .cache
604 .activate(&self.policy.catalog, &[name.to_owned()]);
605 tool_catalog::remove_evicted_cache_activations(
606 &self.policy.catalog,
607 &mut self.policy.active_names,
608 activation.evicted,
609 );
610 self.policy
611 .active_names
612 .extend(activation.admitted.iter().cloned());
613 anyhow::ensure!(
614 activation.admitted.iter().any(|admitted| admitted == name),
615 "tool was not admitted by Core's bounded activation cache"
616 );
617 let definition = self
618 .policy
619 .catalog
620 .iter()
621 .find(|tool| tool.name == name)
622 .ok_or_else(|| anyhow!("tool left the Core catalog"))?;
623 Ok(
624 (!tool_catalog::deferred_first_call_matches_schema(definition, input)).then(|| {
625 tool_catalog::deferred_tool_schema_hydration_result(definition, input).content
626 }),
627 )
628 }
629 }
630 #[cfg(test)]
631 pub(crate) struct ChildProbeCall {
632 pub(crate) id: String,
633 pub(crate) execution_id: String,
634 pub(crate) name: String,
635 pub(crate) input: serde_json::Value,
636 }
637
638 #[cfg(test)]
639 impl Engine {
640 /// Direct tests use the real admitted Core planner/executor and approval
641 /// inbox. The loop below only observes/relays those Core events.
642 pub(crate) async fn probe_child_call(
643 authority: Arc<ChildAuthority>,
644 registry: crate::tools::ToolRegistry,
645 surface: &mut ChildSurfaceProbe,
646 call: ChildProbeCall,
647 ) -> Result<crate::tools::spec::RichToolResult> {
648 let runtime = authority.runtime.clone();
649 let owner = authority.owner_agent_id.clone();
650 let api = runtime
651 .api_config
652 .as_deref()
653 .ok_or_else(|| anyhow!("fixture needs its explicit captured Config"))?;
654 let (mut engine, handle) = Self::new_child_admitted(
655 EngineConfig {
656 workspace: runtime.context.workspace.clone(),
657 model: runtime.model.clone(),
658 max_steps: 1,
659 auto_review_policy: runtime.auto_review_policy.as_ref().clone(),
660 ..Default::default()
661 },
662 api,
663 authority,
664 SystemPrompt::Text("Core child batch probe".into()),
665 None,
666 )?;
667 surface.policy.registry = registry;
668 let mut events = handle.rx_event.write().await;
669 let execute = engine.probe_child_tool_batch(surface, call);
670 tokio::pin!(execute);
671 let mut approvals = futures_util::stream::FuturesUnordered::new();
672 use futures_util::StreamExt;
673 let answer = loop {
674 tokio::select! {
675 biased;
676 result = &mut execute => break result,
677 () = runtime.cancel_token.cancelled(), if !handle.is_cancelled() => handle.cancel(),
678 Some((id, decision)) = approvals.next(), if !approvals.is_empty() => {
679 if runtime.cancel_token.is_cancelled() || handle.is_cancelled() {
680 handle.cancel();
681 } else {
682 match decision {
683 Ok(crate::tools::subagent::ChildApprovalOutcome::Approved) => handle.approve_tool_call(id).await?,
684 Ok(crate::tools::subagent::ChildApprovalOutcome::Denied) => handle.deny_tool_call(id).await?,
685 _ => handle.deny_tool_call_unavailable(id).await?,
686 }
687 }
688 }
689 event = events.recv() => {
690 let Some(event) = event else { break Err(anyhow!("Core probe event stream closed")) };
691 match &event {
692 Event::ApprovalRequired { id, tool_name, description, .. } => {
693 let (_, receiver) = runtime.manager.write().await.register_child_approval(&owner, id, tool_name, description)?;
694 let key = id.clone(); approvals.push(async move { (key, receiver.await) });
695 if let Some(tx) = runtime.event_tx.as_ref().filter(|_| runtime.parent_can_prompt) {
696 let sent = tokio::select! {
697 biased;
698 () = runtime.cancel_token.cancelled() => false,
699 result = tx.send(event.clone()) => result.is_ok(),
700 };
701 if !sent {
702 runtime.manager.write().await.cancel_child_approval(id);
703 handle.deny_tool_call_unavailable(id).await?;
704 }
705 } else {
706 runtime.manager.write().await.cancel_child_approval(id);
707 handle.deny_tool_call_unavailable(id).await?;
708 }
709 }
710 Event::ApprovalWithdrawn { id } => {
711 runtime.manager.write().await.cancel_child_approval(id);
712 if let Some(tx) = &runtime.event_tx {
713 tokio::select! {
714 biased;
715 () = runtime.cancel_token.cancelled() => { let _ = tx.try_send(event.clone()); },
716 _ = tx.send(event.clone()) => {},
717 }
718 }
719 }
720 Event::ToolGateDecision { .. } => {
721 crate::tools::subagent::engine::forward_child_gate_observation(
722 &runtime,
723 &owner,
724 event,
725 &handle.captured_turn_cancel(),
726 runtime.context.turn_deadline,
727 ).await;
728 }
729 _ => {}
730 }
731 }
732 }
733 };
734 // The canonical batch may settle immediately after enqueueing its
735 // last receipt. Observe those exact queued decisions before handback.
736 while let Ok(event) = events.try_recv() {
737 if let Event::ApprovalWithdrawn { id } = &event {
738 runtime.manager.write().await.cancel_child_approval(id);
739 if let Some(tx) = &runtime.event_tx {
740 let _ = tx.try_send(event);
741 }
742 } else if matches!(event, Event::ToolGateDecision { .. }) {
743 crate::tools::subagent::engine::forward_child_gate_observation(
744 &runtime,
745 &owner,
746 event,
747 &handle.captured_turn_cancel(),
748 runtime.context.turn_deadline,
749 )
750 .await;
751 }
752 }
753 let pending = runtime
754 .manager
755 .read()
756 .await
757 .pending_requests_for_agent(&owner);
758 for pending in pending {
759 runtime
760 .manager
761 .write()
762 .await
763 .cancel_child_approval(&pending.approval_id);
764 if let Some(tx) = &runtime.event_tx {
765 let _ = tx.try_send(Event::ApprovalWithdrawn {
766 id: pending.approval_id,
767 });
768 }
769 }
770 answer
771 }
772 }
773
774 #[cfg(test)]
775 impl Engine {
776 pub(crate) fn probe_child_catalog(
777 authority: &ChildAuthority,
778 registry: &crate::tools::ToolRegistry,
779 ) -> Vec<codewhale_models::Tool> {
780 let catalog = tool_catalog::build_model_tool_catalog(
781 authority.tools_for_model(registry, &authority.agent_type),
782 Vec::new(),
783 AppMode::Agent,
784 &HashSet::new(),
785 );
786 let mut copied = crate::tools::ToolRegistry::new(ToolContext::for_empty_registry());
787 copied.register_all(registry.all());
788 let mut policy = ToolSurfacePolicy::new(
789 copied,
790 Some(catalog),
791 AppMode::Agent,
792 &HashSet::new(),
793 &[],
794 false,
795 None,
796 None,
797 None,
798 tool_catalog::ToolMode::Direct,
799 );
800 Self::narrow_child_surface(authority, &mut policy);
801 policy.catalog
802 }
803 }
804 impl Engine {
805 pub(super) fn narrow_child_surface(
806 authority: &ChildAuthority,
807 surface: &mut ToolSurfacePolicy,
808 ) {
809 let discoverable = !authority.grant.scope.as_ref().is_some_and(Vec::is_empty);
810 surface.catalog.retain(|tool| {
811 surface.registry.contains(&tool.name)
812 || (discoverable && tool_catalog::is_tool_search_tool(&tool.name))
813 });
814 surface
815 .active_names
816 .retain(|name| surface.catalog.iter().any(|tool| &tool.name == name));
817 surface.active = tool_catalog::active_tools_for_request(
818 &surface.catalog,
819 &surface.active_names,
820 surface.strict_tool_mode,
821 );
822 }
823 }
824
824 lines RUST