返回 CodeWhale
codemode.rs
根目录 / crates / tui / src / tools / codemode.rs
1 //! `execute_tools` — code mode: run a model-provided JavaScript program that
2 //! composes tool calls through `tools.call(name, args)`.
3 //!
4 //! The program handles loops, branching, filtering, and data movement;
5 //! intermediate results stay in the VM and only the bounded return value plus
6 //! a host-owned receipt reach the model. The engine-side precedent is the
7 //! synthetic interpreter dispatch (`js_execution` / `code_execution`): this
8 //! tool is engine-injected, never registered, and dispatched in
9 //! `core::engine::tool_execution`.
10 //!
11 //! Authority stays entirely in Rust, and there is one gate. In a main-session
12 //! turn the engine hands the program a [`NestedCallGate`]; every nested call —
13 //! native, plugin, or MCP — is sent back to the turn loop and planned by the
14 //! same `plan_tool_calls` a direct call goes through (deny/allow lists,
15 //! preparation, hooks, ask-rules, Auto-Review, repo law, the worker authority
16 //! envelope, the K1 Computer Use refusal), with the source tagged as code mode
17 //! so a nested call never activates a deferred schema. A nested call that
18 //! needs approval suspends the program: the engine raises the normal
19 //! approval request (`request_tool_approval`, same receipt log, same card) and
20 //! the program resumes on allow; on deny only that nested call fails, as an
21 //! exception the program can catch. Approving the program itself therefore
22 //! grants nothing, and `execute_tools` is prepared as auto-approved.
23 //!
24 //! Without a gate (sub-agents, direct unit calls) the program keeps the
25 //! Phase-1 profile: read-only, auto-approved native calls only, no MCP.
26 //!
27 //! Receipts (#6509): the host records every nested call when it starts, so a
28 //! program that hits its deadline still reports what finished, what was
29 //! refused, and what was in flight. Oversized nested results keep the
30 //! `{content, metadata, truncated}` keys, and the cut is never silent:
31 //! `truncated` names the original size and the spillover file holding the
32 //! full output, and `content` becomes the leading text of the raw output
33 //! (a JSON result cannot stay parsed once cut), so a script checks
34 //! `truncated` before reading fields. A failed call's text is bounded and
35 //! spilled the same way.
36 //!
37 //! KV-cache effect: `[features] code_mode` (on by default) makes
38 //! `execute_tools` eager from the first request of a session; with the flag
39 //! off it is deferred like the interpreter tools. The flag is session config,
40 //! so the prefix is stable within a session either way. The definition text
41 //! is static. Program text, nested results, and the describe-only nested
42 //! `tool_search` (which returns schemas without activating them) live in
43 //! append-only turn history, never in the prefix.
44 //!
45 //! Known limitations:
46 //! - Nested calls do not take per-tool locks against sibling top-level calls;
47 //! the program runs under its own exclusive lock instead. Inside the
48 //! program, calls the gate marks non-parallel run one at a time.
49 //! - No nested `agent`, `workflow`, `rlm`, `request_user_input`, interpreter,
50 //! interactive shell, sandbox escalation, Computer Use consent/script, MCP
51 //! sign-in, or recursive `execute_tools`. Those stay direct calls.
52 //! - Rich content blocks (images) from nested results are dropped; text and
53 //! JSON payloads pass through bounded.
54 //! - Hidden from Plan mode and refused under a worker authority envelope,
55 //! like the other execution surfaces.
56
57 use std::sync::{Arc, Mutex};
58 use std::time::{Duration, Instant};
59
60 use async_trait::async_trait;
61 use serde_json::{Value, json};
62 use tokio::sync::{Mutex as AsyncMutex, RwLock, Semaphore, mpsc, oneshot};
63 use tokio_util::sync::CancellationToken;
64
65 use codewhale_models::Tool;
66 use codewhale_workflow_js::{
67 BudgetSnapshot, DriverError, ProgressEvent, SpawnedTask, TaskRequest, ToolCallRequest,
68 ToolCallResponse, ToolInvoker, WorkflowDriver, WorkflowRunCancel, WorkflowVm,
69 };
70
71 use crate::core::events::Event;
72 use crate::mcp::McpPool;
73 use crate::tools::registry::{ToolRegistry, enforce_tool_authority};
74 use crate::tools::spec::{
75 ApprovalRequirement, RichToolResult, ToolContext, ToolError, ToolResult, ToolSpec, required_str,
76 };
77
78 /// Tool name surfaced to the model. Dispatched alongside the synthetic
79 /// interpreter tools; see `core::engine::tool_execution`.
80 pub const EXECUTE_TOOLS_TOOL_NAME: &str = "execute_tools";
81
82 const EXECUTE_TOOLS_TOOL_TYPE: &str = "execute_tools_20260918";
83
84 /// Maximum program source accepted, in bytes.
85 const MAX_CODE_BYTES: usize = 64 * 1024;
86 /// Run deadline when no engine turn serves the program (sub-agents, direct
87 /// unit calls). Those callers bound the whole call themselves (the sub-agent
88 /// tool timeout), so this is an hour-long backstop rather than an invented
89 /// short cap. A gated run takes the remaining turn wall clock instead, which
90 /// is unbounded unless `[tui].turn_wall_clock_secs` is set.
91 const FALLBACK_RUN_DEADLINE: Duration = Duration::from_secs(3_600);
92 /// How often the run watchdog re-checks while the program is paused on the
93 /// gate (approval card, hook, review). Bounds the overrun after a pause.
94 const PAUSED_WATCHDOG_POLL: Duration = Duration::from_millis(200);
95 /// Maximum nested tool calls in flight at once, enforced host-side.
96 const MAX_CONCURRENT_CALLS: usize = 4;
97 /// Per nested-call result cap, in serialized bytes. The same threshold the
98 /// engine spills a direct tool result at, so a nested result is never cut
99 /// sooner than the same call made directly.
100 const PER_CALL_RESULT_CAP_BYTES: usize = crate::tools::truncate::SPILLOVER_THRESHOLD_BYTES;
101 /// Model-visible return cap, in serialized bytes.
102 const RETURN_CAP_BYTES: usize = 16 * 1024;
103
104 /// Names refused before the gate, with a message that names the supported
105 /// alternative. These either need the turn loop itself (a prompt, a
106 /// sub-agent, a schema activation) or would nest an execution surface.
107 const PROHIBITED_NESTED: &[&str] = &[
108 EXECUTE_TOOLS_TOOL_NAME,
109 "code_execution",
110 "js_execution",
111 "agent",
112 "workflow",
113 // Recursive RLM rounds are admitted by the turn loop serving the direct
114 // `rlm` call; a program has no such server for a nested one.
115 "rlm",
116 crate::core::engine::tool_catalog::REQUEST_USER_INPUT_NAME,
117 crate::core::engine::tool_catalog::MULTI_TOOL_PARALLEL_NAME,
118 ];
119
120 /// Model-facing definition. `defer_loading` is decided by the catalog
121 /// (eager under code mode, deferred otherwise); `allowed_callers` mirrors
122 /// the interpreter tools. The text is static so it never moves the prefix.
123 pub fn execute_tools_tool_definition() -> Tool {
124 Tool {
125 tool_type: Some(EXECUTE_TOOLS_TOOL_TYPE.to_string()),
126 name: EXECUTE_TOOLS_TOOL_NAME.to_string(),
127 description: "Run a JavaScript program that composes tool calls with \
128 `await tools.call(name, args)` and returns a bounded JSON result. Prefer it \
129 whenever you would make several dependent or repetitive calls, MCP and plugin \
130 tools included: intermediate results stay in the program and only what you \
131 return reaches the conversation. Every nested call passes the same permission \
132 checks as a direct call; a call that needs approval pauses the program until \
133 the user decides, and a denied or refused call throws inside the program (catch \
134 it to continue). Inside a program, tools.call('tool_search', {query}) returns \
135 matching tool names with their input schemas without loading them into the \
136 conversation. Each result is {content, metadata, truncated}. When a result is \
137 cut, truncated names its full size and saved copy and content is the leading \
138 text of the raw output instead of parsed JSON, so check truncated before \
139 reading fields. Not \
140 available inside programs: agent, workflow, request_user_input, nested \
141 execute_tools, interactive shells, sandbox escalation, Computer Use consent or \
142 scripts, and MCP sign-in. At most 50 nested calls, 4 concurrent; the return \
143 value is capped at 16 KiB."
144 .to_string(),
145 input_schema: json!({
146 "type": "object",
147 "properties": {
148 "code": {
149 "type": "string",
150 "description": "JavaScript program. The return value (or thrown error) becomes the result; use tools.call(name, argsObject) for tool calls."
151 }
152 },
153 "required": ["code"]
154 }),
155 allowed_callers: Some(vec!["direct".to_string()]),
156 defer_loading: Some(false),
157 input_examples: None,
158 strict: None,
159 cache_control: None,
160 }
161 }
162
163 // ---------------------------------------------------------------------------
164 // The engine-served gate
165 // ---------------------------------------------------------------------------
166
167 /// How the gate decided one nested call. Named in the receipt.
168 #[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)]
169 #[serde(rename_all = "snake_case")]
170 pub(crate) enum NestedDecision {
171 /// Admitted without a prompt: an auto-approved tier, a remembered grant,
172 /// or a posture that already covers it.
173 Auto,
174 /// A person approved this exact nested call.
175 Approved,
176 /// A person denied it.
177 Denied,
178 /// The approval card expired with no answer. The call did not run, and
179 /// the user did not deny it.
180 TimedOut,
181 /// A gate refused it before any prompt (policy, hook, authority, or a
182 /// name that stays a direct call).
183 Refused,
184 }
185
186 /// The gate's answer for one nested call.
187 pub(crate) enum NestedCallVerdict {
188 /// Run the call with the final (possibly hook-rewritten) name and input.
189 Run {
190 name: String,
191 input: Value,
192 supports_parallel: bool,
193 decision: NestedDecision,
194 /// `additionalContext` from tool_call_before hooks (#3026), recorded
195 /// on the call's receipt so it reaches the model.
196 hook_context: Option<String>,
197 },
198 /// The engine answered in place (describe-only `tool_search`, a guard).
199 Answered {
200 result: ToolResult,
201 hook_context: Option<String>,
202 },
203 /// Do not run it; the error is what the program sees.
204 Refused {
205 error: ToolError,
206 decision: NestedDecision,
207 },
208 }
209
210 /// One nested call waiting on the turn loop's gate.
211 pub(crate) struct NestedCallRequest {
212 pub(crate) name: String,
213 pub(crate) input: Value,
214 pub(crate) reply: oneshot::Sender<NestedCallVerdict>,
215 /// Fires when the caller no longer wants an answer (an extension's host
216 /// cancelled the request, its owner was revoked, the host exited, the
217 /// invocation ended). The server must not start work for a request whose
218 /// token has fired ([`Self::is_stale`]) and must stop waiting on a person
219 /// when it fires: an approval is withdrawn, never silently decided.
220 pub(crate) withdraw: Option<CancellationToken>,
221 }
222
223 impl NestedCallRequest {
224 /// Nobody is waiting for the answer any more.
225 pub(crate) fn is_stale(&self) -> bool {
226 self.reply.is_closed()
227 || self
228 .withdraw
229 .as_ref()
230 .is_some_and(CancellationToken::is_cancelled)
231 }
232 }
233
234 /// Who an extension's `core/call` is made for, as the core knows it: composed
235 /// by Rust from the extension tool's registration, never from anything the host
236 /// says.
237 #[derive(Debug, Clone, PartialEq, Eq)]
238 pub struct ExtensionCaller {
239 /// `extension:<plugin>`: named on the approval card and in audit records.
240 pub(crate) origin: String,
241 /// The extension tool the calls are made inside.
242 pub(crate) tool: String,
243 /// The approval-grant scope of that plugin build (`ext:<id>@<hash>`); the
244 /// approval keys of its calls are scoped by it, so a grant given to the
245 /// model never covers an extension's call and the reverse.
246 pub(crate) scope: String,
247 }
248
249 /// What an extension tool's gate carries beyond the channel: who it is for,
250 /// and the tool snapshot its calls run against.
251 struct ExtensionGate {
252 caller: ExtensionCaller,
253 specs: Vec<Arc<dyn ToolSpec>>,
254 }
255
256 /// Handle a running program uses to reach the engine turn that launched it.
257 /// Built per `execute_tools` call by the turn loop and carried on that call's
258 /// [`ToolContext`]; dropped with the call.
259 #[derive(Clone)]
260 pub(crate) struct NestedCallGate {
261 requests: mpsc::Sender<NestedCallRequest>,
262 mcp_pool: Option<Arc<AsyncMutex<McpPool>>>,
263 tx_event: mpsc::Sender<Event>,
264 deadline: Duration,
265 /// Set only on the gate of an extension tool's call
266 /// ([`Self::for_extension`]).
267 extension: Option<Arc<ExtensionGate>>,
268 }
269
270 impl NestedCallGate {
271 /// A gate plus the receiver the turn loop serves. `deadline` bounds the
272 /// program's own run time; time spent waiting on the gate (approval
273 /// cards, hooks, reviews) does not count against it.
274 pub(crate) fn new(
275 mcp_pool: Option<Arc<AsyncMutex<McpPool>>>,
276 tx_event: mpsc::Sender<Event>,
277 deadline: Duration,
278 ) -> (Self, mpsc::Receiver<NestedCallRequest>) {
279 let (requests, receiver) = mpsc::channel(MAX_CONCURRENT_CALLS);
280 (
281 Self {
282 requests,
283 mcp_pool,
284 tx_event,
285 deadline,
286 extension: None,
287 },
288 receiver,
289 )
290 }
291
292 /// This gate serves one extension tool's call: `caller` is who it is for
293 /// and `specs` the tool snapshot its core calls run against.
294 pub(crate) fn for_extension(
295 mut self,
296 caller: ExtensionCaller,
297 specs: Vec<Arc<dyn ToolSpec>>,
298 ) -> Self {
299 self.extension = Some(Arc::new(ExtensionGate { caller, specs }));
300 self
301 }
302
303 /// The caller and tool snapshot of an extension tool's gate; `None` for
304 /// the gate of an `execute_tools` program or an `rlm` call.
305 pub(crate) fn extension(&self) -> Option<(&ExtensionCaller, &[Arc<dyn ToolSpec>])> {
306 self.extension
307 .as_deref()
308 .map(|gate| (&gate.caller, gate.specs.as_slice()))
309 }
310
311 /// Ask the serving turn loop to decide one nested call. A gate nobody
312 /// serves any more refuses.
313 pub(crate) async fn ask(&self, name: String, input: Value) -> NestedCallVerdict {
314 self.ask_withdrawable(name, input, None).await
315 }
316
317 /// [`Self::ask`], withdrawn when `withdraw` fires: the request is dropped
318 /// if the server has not started it, and the server stops waiting on a
319 /// person if it has (the approval is recorded cancelled).
320 pub(crate) async fn ask_withdrawable(
321 &self,
322 name: String,
323 input: Value,
324 withdraw: Option<CancellationToken>,
325 ) -> NestedCallVerdict {
326 let unavailable = || NestedCallVerdict::Refused {
327 error: ToolError::not_available(
328 "the turn that launched this program is no longer serving its permission gate",
329 ),
330 decision: NestedDecision::Refused,
331 };
332 let cancelled = || NestedCallVerdict::Refused {
333 error: ToolError::cancelled("the call was withdrawn before the gate answered"),
334 decision: NestedDecision::Refused,
335 };
336 let (reply, answer) = oneshot::channel();
337 let request = NestedCallRequest {
338 name,
339 input,
340 reply,
341 withdraw: withdraw.clone(),
342 };
343 let wait = async {
344 if self.requests.send(request).await.is_err() {
345 return unavailable();
346 }
347 answer.await.unwrap_or_else(|_| unavailable())
348 };
349 match withdraw {
350 None => wait.await,
351 Some(withdraw) => tokio::select! {
352 biased;
353 () = withdraw.cancelled() => cancelled(),
354 verdict = wait => verdict,
355 },
356 }
357 }
358 }
359
360 #[cfg(test)]
361 impl NestedCallGate {
362 /// A gate whose server admits every call exactly as asked, for tests of
363 /// consumers that are not about admission. Needs a Tokio runtime.
364 pub(crate) fn admitting_for_test() -> Self {
365 Self::answering_for_test(|name, input| NestedCallVerdict::Run {
366 name: name.to_string(),
367 input: input.clone(),
368 supports_parallel: false,
369 decision: NestedDecision::Auto,
370 hook_context: None,
371 })
372 }
373
374 /// A gate whose server answers every call with `answer`.
375 pub(crate) fn answering_for_test(
376 answer: impl Fn(&str, &Value) -> NestedCallVerdict + Send + 'static,
377 ) -> Self {
378 let (tx_event, mut rx_event) = mpsc::channel(64);
379 tokio::spawn(async move { while rx_event.recv().await.is_some() {} });
380 let (gate, mut requests) = Self::new(None, tx_event, Duration::from_secs(60));
381 tokio::spawn(async move {
382 while let Some(request) = requests.recv().await {
383 let _ = request.reply.send(answer(&request.name, &request.input));
384 }
385 });
386 gate
387 }
388 }
389
390 /// Refusals decided from the request alone: calls that need the turn loop
391 /// itself or their own approval card stay direct. Checked on the name the
392 /// program sent, and again by the turn loop on the name planning resolved it
393 /// to (`Agent` resolves to `agent`) and the final, hook-rewritten input.
394 ///
395 /// Names compare ASCII case-insensitively: dispatch resolves `Agent` to
396 /// `agent`, so a case-sensitive list would be a bypass, not a policy.
397 pub(crate) fn refusal_before_gate(name: &str, input: &Value, gated: bool) -> Option<String> {
398 let lower = name.to_ascii_lowercase();
399 if PROHIBITED_NESTED.contains(&lower.as_str())
400 || (!gated && crate::core::engine::tool_catalog::is_tool_search_tool(&lower))
401 {
402 return Some(format!(
403 "`{name}` is not available inside execute_tools programs; call it directly (use workflow/task() for fan-out)"
404 ));
405 }
406 if matches!(lower.as_str(), "bash" | "exec_shell")
407 && input.get("interactive").and_then(Value::as_bool) == Some(true)
408 {
409 return Some(format!(
410 "`{name}` with interactive:true needs the terminal; call it directly"
411 ));
412 }
413 if input.get("sandbox_permissions").is_some() {
414 return Some(format!(
415 "`{name}` requests a sandbox escalation, which needs its own exact-call approval; call it directly"
416 ));
417 }
418 if crate::tools::approval_cache::computer_use_user_gate(name, input).is_some()
419 || crate::tools::approval_cache::computer_use_batch_hidden_gate(name, input).is_some()
420 {
421 return Some(format!(
422 "Computer Use call `{name}` grants consent or runs a script; it needs its own approval card, so call it directly and let the user decide"
423 ));
424 }
425 if McpPool::is_mcp_tool(name)
426 && name.ends_with(&format!("_{}", crate::mcp::AUTHENTICATE_TOOL_NAME))
427 {
428 return Some(format!("`{name}` starts an MCP sign-in; call it directly"));
429 }
430 None
431 }
432
433 // ---------------------------------------------------------------------------
434 // Receipts
435 // ---------------------------------------------------------------------------
436
437 /// What was cut from an oversized value, and where the whole value went.
438 #[derive(Debug, Clone, PartialEq, serde::Serialize)]
439 struct Truncation {
440 original_bytes: usize,
441 kept_bytes: usize,
442 spill_path: Option<String>,
443 }
444
445 #[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)]
446 #[serde(rename_all = "snake_case")]
447 enum CallStatus {
448 /// Started and not yet finished. Left in place when the run ends early,
449 /// so the receipt says the call may have partially run.
450 InFlight,
451 Ok,
452 Failed,
453 Refused,
454 }
455
456 /// One nested call, as recorded by the host — not the script. The receipt is
457 /// what makes "no failures found" distinguishable from "nothing ran".
458 #[derive(Debug, Clone, serde::Serialize)]
459 struct CallReceipt {
460 seq: usize,
461 tool: String,
462 decision: Option<NestedDecision>,
463 status: CallStatus,
464 ok: bool,
465 elapsed_ms: u64,
466 bytes: usize,
467 truncated: Option<Truncation>,
468 note: Option<String>,
469 /// `additionalContext` a tool_call_before hook attached to this call,
470 /// the same text a direct call appends to its result (#3026).
471 #[serde(skip_serializing_if = "Option::is_none")]
472 hook_context: Option<String>,
473 }
474
475 /// Run clock that stops while the program waits on the gate, so a person
476 /// taking a minute on an approval card does not spend the program's budget.
477 /// An extension tool's `tool/call` deadline runs on the same clock
478 /// (`extension_host::supervisor::HostProcess::call_with_clock`), paused while
479 /// one of its `core/call`s waits on the gate.
480 pub(crate) struct PauseClock {
481 started: Instant,
482 paused: Duration,
483 depth: usize,
484 since: Option<Instant>,
485 }
486
487 impl PauseClock {
488 pub(crate) fn new() -> Self {
489 Self {
490 started: Instant::now(),
491 paused: Duration::ZERO,
492 depth: 0,
493 since: None,
494 }
495 }
496
497 fn pause(&mut self) {
498 if self.depth == 0 {
499 self.since = Some(Instant::now());
500 }
501 self.depth += 1;
502 }
503
504 fn resume(&mut self) {
505 self.depth = self.depth.saturating_sub(1);
506 if self.depth == 0
507 && let Some(since) = self.since.take()
508 {
509 self.paused += since.elapsed();
510 }
511 }
512
513 fn active(&self) -> Duration {
514 let paused = self.paused + self.since.map_or(Duration::ZERO, |since| since.elapsed());
515 self.started.elapsed().saturating_sub(paused)
516 }
517
518 /// Time left before `deadline`, or `None` once it has passed. While the
519 /// clock is paused its budget is frozen, so the caller looks again shortly
520 /// ([`PAUSED_WATCHDOG_POLL`]): sleeping out the whole deadline would let it
521 /// overrun by that much once the pause ends.
522 pub(crate) fn remaining(&self, deadline: Duration) -> Option<Duration> {
523 if self.depth > 0 {
524 return Some(PAUSED_WATCHDOG_POLL.min(deadline));
525 }
526 deadline
527 .checked_sub(self.active())
528 .filter(|left| !left.is_zero())
529 }
530 }
531
532 struct PauseGuard<'a>(&'a Mutex<PauseClock>);
533
534 impl Drop for PauseGuard<'_> {
535 fn drop(&mut self) {
536 if let Ok(mut clock) = self.0.lock() {
537 clock.resume();
538 }
539 }
540 }
541
542 /// Abort a spawned nested call when the program stops waiting for it.
543 struct AbortOnDrop<T>(tokio::task::JoinHandle<T>);
544
545 impl<T> Drop for AbortOnDrop<T> {
546 fn drop(&mut self) {
547 self.0.abort();
548 }
549 }
550
551 // ---------------------------------------------------------------------------
552 // Invoker
553 // ---------------------------------------------------------------------------
554
555 /// [`ToolInvoker`] over a snapshot of the parent turn's registry.
556 ///
557 /// The snapshot (spec Arcs plus a cloned [`ToolContext`]) is taken at
558 /// dispatch so the invoker is `'static` for the VM thread. Nested calls run
559 /// on the engine's runtime (captured at dispatch), not the VM thread's.
560 pub(crate) struct CodemodeInvoker {
561 specs: Vec<Arc<dyn ToolSpec>>,
562 context: ToolContext,
563 gate: Option<NestedCallGate>,
564 runtime: Option<tokio::runtime::Handle>,
565 semaphore: Arc<Semaphore>,
566 /// Non-parallel nested calls take this exclusively; parallel-safe ones
567 /// share it, mirroring how the engine schedules direct calls.
568 order: RwLock<()>,
569 receipts: Mutex<Vec<CallReceipt>>,
570 clock: Arc<Mutex<PauseClock>>,
571 spill_prefix: String,
572 /// Refusals beyond code mode's own, decided from the request alone before
573 /// the gate is asked (an extension's list: `extension_host::core_call`).
574 extra_refusal: Option<ExtraRefusal>,
575 }
576
577 /// A refusal decided from the specs, the tool name and its input.
578 pub(crate) type ExtraRefusal = fn(&[Arc<dyn ToolSpec>], &str, &Value) -> Option<String>;
579
580 /// Why a gated call did not produce a result. The decision is what the gate
581 /// recorded; a caller that is not code mode maps it to its own error codes.
582 #[derive(Debug)]
583 pub(crate) enum NestedFailure {
584 /// Refused or denied: nothing ran.
585 Rejected {
586 decision: NestedDecision,
587 message: String,
588 },
589 /// The seam broke, the call timed out or was cancelled.
590 Unavailable(String),
591 }
592
593 impl From<NestedFailure> for DriverError {
594 fn from(failure: NestedFailure) -> Self {
595 match failure {
596 NestedFailure::Rejected { message, .. } => DriverError::Rejected(message),
597 NestedFailure::Unavailable(message) => DriverError::Unavailable(message),
598 }
599 }
600 }
601
602 impl CodemodeInvoker {
603 fn new(specs: Vec<Arc<dyn ToolSpec>>, context: ToolContext) -> Self {
604 let mut context = context;
605 let gate = context.execution.nested_call_gate.take();
606 let spill_prefix = context
607 .origin_tool_call_id
608 .clone()
609 .unwrap_or_else(|| EXECUTE_TOOLS_TOOL_NAME.to_string());
610 Self::with_gate(specs, context, gate, spill_prefix, None)
611 }
612
613 /// The machinery for an extension tool's `core/call`s: the same gate,
614 /// executor, receipts, concurrency cap and pausable clock as a program's
615 /// nested calls, with `extra_refusal` added to what is refused outright.
616 /// `context` is the extension tool's own; its gate is the one given here,
617 /// which the tool's calls then never see again.
618 pub(crate) fn for_extension(
619 specs: Vec<Arc<dyn ToolSpec>>,
620 mut context: ToolContext,
621 gate: NestedCallGate,
622 spill_prefix: String,
623 extra_refusal: ExtraRefusal,
624 ) -> Self {
625 context.execution.nested_call_gate = None;
626 Self::with_gate(
627 specs,
628 context,
629 Some(gate),
630 spill_prefix,
631 Some(extra_refusal),
632 )
633 }
634
635 fn with_gate(
636 specs: Vec<Arc<dyn ToolSpec>>,
637 context: ToolContext,
638 gate: Option<NestedCallGate>,
639 spill_prefix: String,
640 extra_refusal: Option<ExtraRefusal>,
641 ) -> Self {
642 Self {
643 specs,
644 context,
645 gate,
646 runtime: tokio::runtime::Handle::try_current().ok(),
647 semaphore: Arc::new(Semaphore::new(MAX_CONCURRENT_CALLS)),
648 order: RwLock::new(()),
649 receipts: Mutex::new(Vec::new()),
650 clock: Arc::new(Mutex::new(PauseClock::new())),
651 spill_prefix,
652 extra_refusal,
653 }
654 }
655
656 fn deadline(&self) -> Duration {
657 self.gate
658 .as_ref()
659 .map_or(FALLBACK_RUN_DEADLINE, |gate| gate.deadline)
660 }
661
662 fn pause(&self) -> PauseGuard<'_> {
663 if let Ok(mut clock) = self.clock.lock() {
664 clock.pause();
665 }
666 PauseGuard(&self.clock)
667 }
668
669 /// Time left before the deadline ([`PauseClock::remaining`]).
670 fn remaining(&self, deadline: Duration) -> Option<Duration> {
671 self.clock.lock().ok()?.remaining(deadline)
672 }
673
674 /// The clock the run's deadline is measured on, for a caller whose own
675 /// deadline must stop while this run waits on the gate.
676 pub(crate) fn clock(&self) -> Arc<Mutex<PauseClock>> {
677 Arc::clone(&self.clock)
678 }
679
680 fn begin(&self, tool: &str) -> usize {
681 let Ok(mut receipts) = self.receipts.lock() else {
682 return 0;
683 };
684 let seq = receipts.len() + 1;
685 receipts.push(CallReceipt {
686 seq,
687 tool: tool.to_string(),
688 decision: None,
689 status: CallStatus::InFlight,
690 ok: false,
691 elapsed_ms: 0,
692 bytes: 0,
693 truncated: None,
694 note: None,
695 hook_context: None,
696 });
697 seq
698 }
699
700 fn finish(&self, seq: usize, update: impl FnOnce(&mut CallReceipt)) {
701 if let Ok(mut receipts) = self.receipts.lock()
702 && let Some(receipt) = receipts.iter_mut().find(|receipt| receipt.seq == seq)
703 {
704 update(receipt);
705 }
706 }
707
708 fn drain(&self) -> Vec<CallReceipt> {
709 self.receipts
710 .lock()
711 .map(|receipts| receipts.clone())
712 .unwrap_or_default()
713 }
714
715 fn refused(
716 &self,
717 seq: usize,
718 started: Instant,
719 decision: NestedDecision,
720 note: String,
721 ) -> NestedFailure {
722 self.finish(seq, |receipt| {
723 receipt.decision = Some(decision);
724 receipt.status = CallStatus::Refused;
725 receipt.elapsed_ms = started.elapsed().as_millis() as u64;
726 receipt.note = Some(note.clone());
727 });
728 NestedFailure::Rejected {
729 decision,
730 message: note,
731 }
732 }
733
734 /// Phase-1 admission when no engine gate serves the program.
735 fn ungated_admission(&self, name: &str, input: &Value) -> Result<(), String> {
736 if McpPool::is_mcp_tool(name) {
737 return Err(format!(
738 "`{name}` is an MCP tool; this program has no session permission gate (sub-agent or host without a turn), so call it directly"
739 ));
740 }
741 let Some(spec) = self.specs.iter().find(|spec| spec.name() == name) else {
742 return Err(format!(
743 "unknown tool `{name}`; discover names with tool_search before writing the program"
744 ));
745 };
746 if !spec.is_read_only_for(input) {
747 return Err(format!(
748 "`{name}` can mutate; without a session permission gate code mode executes read-only calls only"
749 ));
750 }
751 if spec.approval_requirement_for(input) != ApprovalRequirement::Auto {
752 return Err(format!(
753 "`{name}` needs approval; without a session permission gate code mode executes only auto-approved calls"
754 ));
755 }
756 enforce_tool_authority(name, input, spec.as_ref(), &self.context)
757 .map_err(|err| err.to_string())
758 }
759
760 /// Execute an admitted call on the engine runtime: MCP through the
761 /// session pool (same dispatcher as a direct call), everything else
762 /// through its registry spec under the turn's authority envelope.
763 async fn execute(&self, name: &str, input: Value) -> Result<RichToolResult, ToolError> {
764 let future: std::pin::Pin<
765 Box<dyn std::future::Future<Output = Result<RichToolResult, ToolError>> + Send>,
766 > = if McpPool::is_mcp_tool(name) {
767 let Some(gate) = self.gate.as_ref() else {
768 return Err(ToolError::not_available(format!(
769 "MCP tool `{name}` needs the session MCP pool"
770 )));
771 };
772 let Some(pool) = gate.mcp_pool.clone() else {
773 return Err(ToolError::not_available(format!(
774 "MCP is not connected for this turn, so `{name}` cannot run"
775 )));
776 };
777 let tx_event = gate.tx_event.clone();
778 let disallowed = self.context.disallowed_tools.clone();
779 let name = name.to_string();
780 Box::pin(async move {
781 crate::core::engine::Engine::execute_mcp_tool_with_pool(
782 pool,
783 &tx_event,
784 &name,
785 input,
786 &disallowed,
787 // Calls a program makes never carry a person's decision.
788 None,
789 )
790 .await
791 })
792 } else {
793 let Some(spec) = self.specs.iter().find(|spec| spec.name() == name).cloned() else {
794 return Err(ToolError::not_available(format!(
795 "tool `{name}` is not registered"
796 )));
797 };
798 enforce_tool_authority(name, &input, spec.as_ref(), &self.context)?;
799 let context = self.context.clone();
800 Box::pin(async move {
801 crate::extension_host::validate_caller_plugins(context.plugin_registry.as_deref())
802 .map_err(ToolError::not_available)?;
803 spec.execute_rich(input, &context).await
804 })
805 };
806 match self.runtime.as_ref() {
807 Some(runtime) => {
808 let mut task = AbortOnDrop(runtime.spawn(future));
809 match (&mut task.0).await {
810 Ok(result) => result,
811 Err(err) => Err(ToolError::execution_failed(format!(
812 "nested call did not complete: {err}"
813 ))),
814 }
815 }
816 None => future.await,
817 }
818 }
819
820 /// Shape one finished call for the program and complete its receipt.
821 fn deliver(
822 &self,
823 seq: usize,
824 started: Instant,
825 name: &str,
826 decision: NestedDecision,
827 outcome: Result<RichToolResult, ToolError>,
828 ) -> Result<ToolCallResponse, NestedFailure> {
829 let elapsed_ms = started.elapsed().as_millis() as u64;
830 match outcome {
831 Ok(rich) => {
832 let result = rich.into_result();
833 let success = result.success;
834 let (payload, raw_len, truncated) = if success {
835 self.bounded_envelope(seq, name, &result)
836 } else {
837 let raw_len = result.content.len();
838 let (text, truncated) = bound_text(
839 result.content,
840 PER_CALL_RESULT_CAP_BYTES,
841 &self.spill_id(seq, name),
842 );
843 (text, raw_len, truncated)
844 };
845 self.finish(seq, |receipt| {
846 receipt.decision = Some(decision);
847 receipt.status = if success {
848 CallStatus::Ok
849 } else {
850 CallStatus::Failed
851 };
852 receipt.ok = success;
853 receipt.elapsed_ms = elapsed_ms;
854 receipt.bytes = raw_len;
855 receipt.truncated = truncated;
856 });
857 Ok(ToolCallResponse {
858 ok: success,
859 result: payload,
860 })
861 }
862 Err(err) => {
863 let message = err.to_string();
864 // Validation-shaped failures mean nothing ran (admission);
865 // execution failures ran and failed (agent kind via ok:false);
866 // seam breaks are unavailable.
867 match err {
868 ToolError::InvalidInput { .. }
869 | ToolError::MissingField { .. }
870 | ToolError::PathEscape { .. }
871 | ToolError::PermissionDenied { .. } => {
872 Err(self.refused(seq, started, decision, message))
873 }
874 ToolError::Timeout { .. }
875 | ToolError::Cancelled { .. }
876 | ToolError::NotAvailable { .. } => {
877 self.finish(seq, |receipt| {
878 receipt.decision = Some(decision);
879 receipt.status = CallStatus::Failed;
880 receipt.elapsed_ms = elapsed_ms;
881 receipt.note = Some(message.clone());
882 });
883 Err(NestedFailure::Unavailable(message))
884 }
885 ToolError::ExecutionFailed { .. } => {
886 let bytes = message.len();
887 let (text, truncated) = bound_text(
888 message,
889 PER_CALL_RESULT_CAP_BYTES,
890 &self.spill_id(seq, name),
891 );
892 self.finish(seq, |receipt| {
893 receipt.decision = Some(decision);
894 receipt.status = CallStatus::Failed;
895 receipt.elapsed_ms = elapsed_ms;
896 receipt.bytes = bytes;
897 receipt.truncated = truncated;
898 });
899 Ok(ToolCallResponse {
900 ok: false,
901 result: text,
902 })
903 }
904 }
905 }
906 }
907 }
908
909 /// Spillover id for one nested call's full output.
910 fn spill_id(&self, seq: usize, name: &str) -> String {
911 format!("{}-nested-{seq}-{name}", self.spill_prefix)
912 }
913
914 /// Stable envelope: the tool's text content (parsed as JSON when it is
915 /// JSON), its structured metadata, and `truncated` (null unless cut).
916 /// An oversized result keeps the same keys: `content` becomes the head
917 /// of the raw text (no longer parsed JSON) and `truncated` says how much
918 /// there was and where the whole output was saved.
919 fn bounded_envelope(
920 &self,
921 seq: usize,
922 name: &str,
923 result: &ToolResult,
924 ) -> (Value, usize, Option<Truncation>) {
925 let metadata = result.metadata.clone().unwrap_or(Value::Null);
926 let content = serde_json::from_str(&result.content)
927 .unwrap_or_else(|_| Value::String(result.content.clone()));
928 let payload = json!({ "content": content, "metadata": metadata, "truncated": Value::Null });
929 let raw = payload.to_string();
930 if raw.len() <= PER_CALL_RESULT_CAP_BYTES {
931 return (payload, raw.len(), None);
932 }
933 let spill_path = spill(&self.spill_id(seq, name), &raw);
934 let metadata_len = metadata.to_string().len();
935 let metadata = if metadata_len <= PER_CALL_RESULT_CAP_BYTES / 4 {
936 metadata
937 } else {
938 Value::Null
939 };
940 let head = char_prefix(&result.content, PER_CALL_RESULT_CAP_BYTES / 2);
941 let truncation = Truncation {
942 original_bytes: raw.len(),
943 kept_bytes: head.len(),
944 spill_path,
945 };
946 (
947 json!({
948 "content": head,
949 "metadata": metadata,
950 "truncated": truncation,
951 }),
952 raw.len(),
953 Some(truncation),
954 )
955 }
956 }
957
958 /// The receipts' bounds when they are attached to another tool's result.
959 const RECEIPT_NOTE_BYTES: usize = 256;
960
961 impl CodemodeInvoker {
962 /// The receipts of the calls so far as JSON, at most `max` of them with
963 /// each note cut short: what an extension tool's result records about the
964 /// core calls it made. The count says how many there were in all.
965 pub(crate) fn receipts_json(&self, max: usize) -> Value {
966 let receipts = self.drain();
967 let total = receipts.len();
968 let shown: Vec<CallReceipt> = receipts
969 .into_iter()
970 .take(max)
971 .map(|mut receipt| {
972 receipt.note = receipt
973 .note
974 .map(|note| char_prefix(&note, RECEIPT_NOTE_BYTES));
975 receipt.hook_context = None;
976 receipt
977 })
978 .collect();
979 json!({ "total": total, "calls": shown })
980 }
981
982 /// One gated call: refuse it outright if it is on a refusal list, ask the
983 /// engine's gate (the program is paused meanwhile), then run it and shape
984 /// the answer, recording a receipt throughout. `withdraw` stops the wait
985 /// for a slot, the gate's wait on a person, and the execution.
986 pub(crate) async fn call(
987 &self,
988 tool: String,
989 input: Value,
990 withdraw: Option<&CancellationToken>,
991 ) -> Result<ToolCallResponse, NestedFailure> {
992 let permit = self.semaphore.clone().acquire_owned();
993 let _permit = match withdraw {
994 None => permit.await,
995 Some(withdraw) => tokio::select! {
996 biased;
997 () = withdraw.cancelled() => {
998 return Err(NestedFailure::Unavailable("the call was withdrawn".to_string()));
999 }
1000 permit = permit => permit,
1001 },
1002 }
1003 .map_err(|_| NestedFailure::Unavailable("code-mode run shut down".to_string()))?;
1004 let started = Instant::now();
1005 let seq = self.begin(&tool);
1006
1007 if let Some(note) = refusal_before_gate(&tool, &input, self.gate.is_some()).or_else(|| {
1008 self.extra_refusal
1009 .and_then(|refuse| refuse(&self.specs, &tool, &input))
1010 }) {
1011 return Err(self.refused(seq, started, NestedDecision::Refused, note));
1012 }
1013
1014 let (name, input, supports_parallel, decision, hook_context) = match self.gate.as_ref() {
1015 Some(gate) => {
1016 let verdict = {
1017 let _paused = self.pause();
1018 gate.ask_withdrawable(tool.clone(), input, withdraw.cloned())
1019 .await
1020 };
1021 match verdict {
1022 NestedCallVerdict::Run {
1023 name,
1024 input,
1025 supports_parallel,
1026 decision,
1027 hook_context,
1028 } => (name, input, supports_parallel, decision, hook_context),
1029 NestedCallVerdict::Answered {
1030 result,
1031 hook_context,
1032 } => {
1033 self.finish(seq, |receipt| receipt.hook_context = hook_context);
1034 return self.deliver(
1035 seq,
1036 started,
1037 &tool,
1038 NestedDecision::Auto,
1039 Ok(RichToolResult::plain(result)),
1040 );
1041 }
1042 NestedCallVerdict::Refused { error, decision } => {
1043 return Err(self.refused(seq, started, decision, error.to_string()));
1044 }
1045 }
1046 }
1047 None => {
1048 if let Err(note) = self.ungated_admission(&tool, &input) {
1049 return Err(self.refused(seq, started, NestedDecision::Refused, note));
1050 }
1051 (tool.clone(), input, true, NestedDecision::Auto, None)
1052 }
1053 };
1054 self.finish(seq, |receipt| {
1055 if name != tool {
1056 receipt.tool = name.clone();
1057 }
1058 receipt.hook_context = hook_context;
1059 });
1060
1061 let run = async {
1062 if supports_parallel {
1063 let _shared = self.order.read().await;
1064 self.execute(&name, input).await
1065 } else {
1066 let _exclusive = self.order.write().await;
1067 self.execute(&name, input).await
1068 }
1069 };
1070 let outcome = match withdraw {
1071 None => run.await,
1072 Some(withdraw) => tokio::select! {
1073 biased;
1074 () = withdraw.cancelled() => {
1075 Err(ToolError::cancelled("the call was withdrawn while it ran"))
1076 }
1077 outcome = run => outcome,
1078 },
1079 };
1080 self.deliver(seq, started, &name, decision, outcome)
1081 }
1082 }
1083
1084 #[async_trait]
1085 impl ToolInvoker for CodemodeInvoker {
1086 async fn invoke(&self, request: ToolCallRequest) -> Result<ToolCallResponse, DriverError> {
1087 let ToolCallRequest { tool, input } = request;
1088 self.call(tool, input, None)
1089 .await
1090 .map_err(DriverError::from)
1091 }
1092 }
1093
1094 /// [`WorkflowDriver`] for code-mode runs: `task()` is refused (fan-out stays
1095 /// with `workflow`), the token budget is unconstrained, and progress events
1096 /// feed the run receipt.
1097 pub(crate) struct CodemodeDriver {
1098 events: Mutex<Vec<ProgressEvent>>,
1099 }
1100
1101 impl Default for CodemodeDriver {
1102 fn default() -> Self {
1103 Self {
1104 events: Mutex::new(Vec::new()),
1105 }
1106 }
1107 }
1108
1109 impl CodemodeDriver {
1110 fn log_lines(&self) -> Vec<String> {
1111 self.events
1112 .lock()
1113 .map(|events| {
1114 events
1115 .iter()
1116 .filter_map(|event| match event {
1117 ProgressEvent::Log { message } => Some(message.clone()),
1118 _ => None,
1119 })
1120 .collect()
1121 })
1122 .unwrap_or_default()
1123 }
1124 }
1125
1126 #[async_trait]
1127 impl WorkflowDriver for CodemodeDriver {
1128 async fn spawn_task(&self, _request: TaskRequest) -> Result<SpawnedTask, DriverError> {
1129 Err(DriverError::Rejected(
1130 "task() is unavailable in execute_tools programs; tools.call() composes direct tool calls"
1131 .to_string(),
1132 ))
1133 }
1134
1135 fn budget(&self) -> BudgetSnapshot {
1136 BudgetSnapshot {
1137 total: None,
1138 spent: 0,
1139 }
1140 }
1141
1142 fn progress(&self, event: ProgressEvent) {
1143 if let Ok(mut events) = self.events.lock() {
1144 events.push(event);
1145 }
1146 }
1147
1148 fn cancel_all(&self) {}
1149 }
1150
1151 /// Longest prefix of `text` within `max_bytes` that ends on a char boundary.
1152 fn char_prefix(text: &str, max_bytes: usize) -> String {
1153 let cut = max_bytes.min(text.len());
1154 let cut = (0..=cut)
1155 .rev()
1156 .find(|&index| text.is_char_boundary(index))
1157 .unwrap_or(0);
1158 text[..cut].to_string()
1159 }
1160
1161 /// Save a full value through the session spillover store; `None` when the
1162 /// store is unavailable (the truncation is still reported).
1163 fn spill(id: &str, content: &str) -> Option<String> {
1164 crate::tools::truncate::write_spillover(id, content)
1165 .map(|path| path.display().to_string())
1166 .map_err(|error| {
1167 tracing::warn!(target: "codemode", %error, "nested result spillover failed");
1168 })
1169 .ok()
1170 }
1171
1172 /// Bound the program's return value to `cap` serialized bytes. An oversized
1173 /// value becomes the head of its JSON text, and the truncation record says so.
1174 fn bound_json(value: Value, cap: usize, spill_id: &str) -> (Value, Option<Truncation>) {
1175 let raw = value.to_string();
1176 if raw.len() <= cap {
1177 return (value, None);
1178 }
1179 bound_text(raw, cap, spill_id)
1180 }
1181
1182 /// Bound `text` to `cap` bytes: an oversized text becomes its head, the
1183 /// whole text is saved, and the truncation record says so.
1184 fn bound_text(text: String, cap: usize, spill_id: &str) -> (Value, Option<Truncation>) {
1185 if text.len() <= cap {
1186 return (Value::String(text), None);
1187 }
1188 let head = char_prefix(&text, cap / 2);
1189 let truncation = Truncation {
1190 original_bytes: text.len(),
1191 kept_bytes: head.len(),
1192 spill_path: spill(spill_id, &text),
1193 };
1194 (Value::String(head), Some(truncation))
1195 }
1196
1197 fn receipt_payload(
1198 success: bool,
1199 body: Value,
1200 invoker: &CodemodeInvoker,
1201 driver: &CodemodeDriver,
1202 ) -> ToolResult {
1203 let calls = invoker.drain();
1204 let content = json!({
1205 "success": success,
1206 "body": body,
1207 "nested_calls": calls.len(),
1208 "calls": calls,
1209 "log": driver.log_lines(),
1210 })
1211 .to_string();
1212 ToolResult {
1213 content,
1214 success,
1215 metadata: None,
1216 }
1217 }
1218
1219 /// Execute one `execute_tools` call: validate, run the program under its
1220 /// deadline, and return the bounded program value plus the host-owned
1221 /// receipt. A script failure or a deadline is a `success: false` payload
1222 /// that still lists every nested call — only VM setup failures are `Err`.
1223 pub async fn execute_tools_tool(
1224 input: &Value,
1225 registry: &ToolRegistry,
1226 context: &ToolContext,
1227 ) -> Result<ToolResult, ToolError> {
1228 let code = required_str(input, "code")?;
1229 if code.trim().is_empty() {
1230 return Err(ToolError::missing_field("code"));
1231 }
1232 if code.len() > MAX_CODE_BYTES {
1233 return Err(ToolError::invalid_input(format!(
1234 "code exceeds {MAX_CODE_BYTES} bytes"
1235 )));
1236 }
1237 let invoker = Arc::new(CodemodeInvoker::new(registry.all(), context.clone()));
1238 let driver = Arc::new(CodemodeDriver::default());
1239 let deadline = invoker.deadline();
1240 let cancel = WorkflowRunCancel::new();
1241 let vm = WorkflowVm::new();
1242 let mut run = Box::pin(vm.run_tools_script(
1243 code,
1244 Value::Null,
1245 driver.clone(),
1246 invoker.clone(),
1247 cancel.clone(),
1248 ));
1249 let outcome = loop {
1250 let Some(remaining) = invoker.remaining(deadline) else {
1251 break None;
1252 };
1253 tokio::select! {
1254 result = &mut run => break Some(result),
1255 () = tokio::time::sleep(remaining) => {}
1256 }
1257 };
1258 let program_result = match outcome {
1259 None => {
1260 // Stop the VM (and with it every in-flight nested call) before
1261 // reading the receipt, so the listed state is final.
1262 cancel.cancel();
1263 drop(run);
1264 return Ok(receipt_payload(
1265 false,
1266 json!({
1267 "error": format!(
1268 "execute_tools stopped at its {}s run deadline (time waiting on approvals is not counted). Calls with status ok/failed/refused finished; calls still in_flight were cancelled and may have partially run.",
1269 deadline.as_secs()
1270 ),
1271 "timed_out": true,
1272 }),
1273 &invoker,
1274 &driver,
1275 ));
1276 }
1277 Some(Err(err)) => {
1278 return Ok(receipt_payload(
1279 false,
1280 json!({ "error": err.to_string() }),
1281 &invoker,
1282 &driver,
1283 ));
1284 }
1285 Some(Ok(value)) => value,
1286 };
1287 let spill_id = format!("{}-return", invoker.spill_prefix);
1288 let (bounded, truncated) = bound_json(program_result, RETURN_CAP_BYTES, &spill_id);
1289 Ok(receipt_payload(
1290 true,
1291 json!({ "return": bounded, "return_truncated": truncated }),
1292 &invoker,
1293 &driver,
1294 ))
1295 }
1296
1297 #[cfg(test)]
1298 mod tests {
1299 use super::*;
1300
1301 #[test]
1302 fn catalog_advertises_execute_tools_deferred_outside_plan() {
1303 use codewhale_config::AppMode;
1304 use std::collections::HashSet;
1305 let empty = HashSet::new();
1306 for mode in [AppMode::Agent, AppMode::Operate] {
1307 let mut catalog = Vec::new();
1308 crate::core::engine::tool_catalog::ensure_advanced_tooling(
1309 &mut catalog,
1310 mode,
1311 &empty,
1312 ToolMode::Direct,
1313 );
1314 let tool = catalog
1315 .iter()
1316 .find(|tool| tool.name == EXECUTE_TOOLS_TOOL_NAME)
1317 .unwrap_or_else(|| panic!("{mode:?} catalog must advertise execute_tools"));
1318 assert_eq!(tool.defer_loading, Some(true), "{mode:?} must defer it");
1319 }
1320 let mut catalog = Vec::new();
1321 crate::core::engine::tool_catalog::ensure_advanced_tooling(
1322 &mut catalog,
1323 AppMode::Plan,
1324 &empty,
1325 ToolMode::Direct,
1326 );
1327 assert!(
1328 catalog
1329 .iter()
1330 .all(|tool| tool.name != EXECUTE_TOOLS_TOOL_NAME)
1331 );
1332 }
1333
1334 #[test]
1335 fn definition_is_deferred_direct_only() {
1336 let tool = execute_tools_tool_definition();
1337 assert_eq!(tool.name, EXECUTE_TOOLS_TOOL_NAME);
1338 assert_eq!(tool.tool_type.as_deref(), Some(EXECUTE_TOOLS_TOOL_TYPE));
1339 assert!(tool.description.contains("tools.call"));
1340 // Catalog decides deferral; the name is absent from the eager set,
1341 // so injection marks it deferred like the interpreter tools.
1342 assert!(
1343 !crate::core::engine::tool_catalog::DEFAULT_ACTIVE_NATIVE_TOOLS
1344 .contains(&tool.name.as_str())
1345 );
1346 assert_eq!(tool.allowed_callers, Some(vec!["direct".to_string()]));
1347 }
1348
1349 #[test]
1350 fn bound_json_keeps_small_values_verbatim() {
1351 let (value, truncated) = bound_json(json!({"a": 1}), 1024, "bound-small");
1352 assert!(truncated.is_none());
1353 assert_eq!(value, json!({"a": 1}));
1354 }
1355
1356 #[test]
1357 fn bound_json_names_what_it_cut() {
1358 let _guard = crate::tools::truncate::TEST_SPILLOVER_GUARD
1359 .lock()
1360 .unwrap_or_else(|err| err.into_inner());
1361 let spill_root = tempfile::tempdir().unwrap();
1362 let previous =
1363 crate::tools::truncate::set_test_spillover_root(Some(spill_root.path().to_path_buf()));
1364 let big = "x".repeat(100);
1365 let (value, truncated) = bound_json(json!({ "blob": big }), 64, "bound-big");
1366 crate::tools::truncate::set_test_spillover_root(previous);
1367 let truncated = truncated.expect("oversized value is reported");
1368 assert_eq!(truncated.original_bytes, 111);
1369 assert_eq!(truncated.kept_bytes, 32);
1370 assert!(
1371 value
1372 .as_str()
1373 .is_some_and(|head| head.starts_with("{\"blob\""))
1374 );
1375 let saved = truncated.spill_path.expect("full value saved");
1376 assert_eq!(std::fs::read_to_string(saved).unwrap().len(), 111);
1377 }
1378
1379 use crate::core::engine::tool_catalog::ToolMode;
1380 use crate::tools::file_tool::{ReadTool, WriteTool};
1381 use crate::tools::registry::ToolRegistryBuilder;
1382
1383 fn workspace_with_note() -> (tempfile::TempDir, std::path::PathBuf) {
1384 let dir = tempfile::tempdir().unwrap();
1385 let note = dir.path().join("note.txt");
1386 std::fs::write(&note, "alpha\nbeta\n").unwrap();
1387 (dir, note)
1388 }
1389
1390 #[tokio::test]
1391 async fn nested_read_passes_gates_and_records_receipt() {
1392 let (_dir, note) = workspace_with_note();
1393 let context = ToolContext::new(note.parent().unwrap());
1394 let invoker = CodemodeInvoker::new(vec![Arc::new(ReadTool)], context);
1395 let response = invoker
1396 .invoke(ToolCallRequest {
1397 tool: "read".to_string(),
1398 input: json!({ "path": note.to_string_lossy() }),
1399 })
1400 .await
1401 .unwrap();
1402 assert!(response.ok);
1403 assert!(response.result.to_string().contains("alpha"));
1404 let receipts = invoker.drain();
1405 assert_eq!(receipts.len(), 1);
1406 assert_eq!(receipts[0].tool, "read");
1407 assert!(receipts[0].ok);
1408 }
1409
1410 #[tokio::test]
1411 async fn nested_write_is_refused_and_writes_nothing() {
1412 let (_dir, note) = workspace_with_note();
1413 let target = note.parent().unwrap().join("evil.txt");
1414 let context = ToolContext::new(note.parent().unwrap());
1415 let invoker = CodemodeInvoker::new(vec![Arc::new(WriteTool)], context);
1416 let err = invoker
1417 .invoke(ToolCallRequest {
1418 tool: "write".to_string(),
1419 input: json!({ "path": target.to_string_lossy(), "content": "x" }),
1420 })
1421 .await
1422 .unwrap_err();
1423 assert!(
1424 matches!(&err, DriverError::Rejected(message) if message.contains("read-only")),
1425 "unexpected: {err:?}"
1426 );
1427 assert!(!target.exists());
1428 assert_eq!(invoker.drain().len(), 1);
1429 }
1430
1431 #[tokio::test]
1432 async fn prohibited_nested_names_are_refused() {
1433 let dir = tempfile::tempdir().unwrap();
1434 let context = ToolContext::new(dir.path());
1435 let invoker = CodemodeInvoker::new(vec![], context);
1436 for name in [
1437 "agent",
1438 "workflow",
1439 "rlm",
1440 "execute_tools",
1441 "tool_search",
1442 "nope",
1443 ] {
1444 let err = invoker
1445 .invoke(ToolCallRequest {
1446 tool: name.to_string(),
1447 input: json!({}),
1448 })
1449 .await
1450 .unwrap_err();
1451 assert!(matches!(err, DriverError::Rejected(_)), "{name}: {err:?}");
1452 }
1453 }
1454
1455 #[tokio::test]
1456 async fn program_composes_nested_read_and_returns_receipt() {
1457 let (_dir, note) = workspace_with_note();
1458 let workspace = note.parent().unwrap().to_path_buf();
1459 let context = ToolContext::new(workspace.clone());
1460 let registry = ToolRegistryBuilder::new()
1461 .with_tool(Arc::new(ReadTool))
1462 .build(context.clone());
1463 let path = note.to_string_lossy().replace('\\', "\\\\");
1464 let code = format!(
1465 "const r = await tools.call('read', {{ path: '{path}' }}); return {{ hasAlpha: JSON.stringify(r).includes('alpha') }};"
1466 );
1467 let result = execute_tools_tool(&json!({ "code": code }), &registry, &context)
1468 .await
1469 .unwrap();
1470 assert!(result.success, "{}", result.content);
1471 let body: Value = serde_json::from_str(&result.content).unwrap();
1472 assert_eq!(body["nested_calls"], 1);
1473 assert_eq!(body["body"]["return"]["hasAlpha"], true);
1474 }
1475
1476 #[tokio::test]
1477 async fn program_loads_skills_at_runtime_through_load_skill() {
1478 // Skills-as-tools composes with code mode: `load_skill` is
1479 // read-only and auto-approved, so a program can list and load
1480 // skills at runtime without widening its authority.
1481 // A configured skills dir, not a project root: project skills load
1482 // only in a trusted workspace, which this composition test is not
1483 // about.
1484 let dir = tempfile::tempdir().unwrap();
1485 let workspace = dir.path().join("workspace");
1486 std::fs::create_dir_all(&workspace).unwrap();
1487 let skills_root = dir.path().join("configured-skills");
1488 let skill_dir = skills_root.join("greet");
1489 std::fs::create_dir_all(&skill_dir).unwrap();
1490 std::fs::write(
1491 skill_dir.join("SKILL.md"),
1492 "---\nname: greet\ndescription: Say hello\n---\n# Greet\nSay hello warmly.\n",
1493 )
1494 .unwrap();
1495 let context = ToolContext::new(&workspace)
1496 .with_skills_config(&skills_root, crate::skills::SkillDiscoveryMode::Compatible);
1497 let registry = ToolRegistryBuilder::new()
1498 .with_tool(Arc::new(crate::tools::skill::LoadSkillTool))
1499 .build(context.clone());
1500 let code = "const list = await tools.call('load_skill', { name: 'list' }); \
1501 const body = await tools.call('load_skill', { name: 'greet' }); \
1502 return { listed: JSON.stringify(list).includes('greet'), \
1503 loaded: JSON.stringify(body).includes('warmly') };";
1504 let result = execute_tools_tool(&json!({ "code": code }), &registry, &context)
1505 .await
1506 .unwrap();
1507 assert!(result.success, "{}", result.content);
1508 let body: Value = serde_json::from_str(&result.content).unwrap();
1509 assert_eq!(body["nested_calls"], 2);
1510 assert_eq!(body["body"]["return"]["listed"], true);
1511 assert_eq!(body["body"]["return"]["loaded"], true);
1512 }
1513
1514 // --- Gated (engine-served) programs -----------------------------------
1515
1516 use std::sync::atomic::{AtomicUsize, Ordering};
1517
1518 /// A stand-in for the turn loop's gate: answers every nested call with
1519 /// `answer` after `delay`, counting how often it was asked.
1520 fn gated_context(
1521 workspace: &std::path::Path,
1522 deadline: Duration,
1523 delay: Duration,
1524 mcp_pool: Option<Arc<AsyncMutex<McpPool>>>,
1525 answer: impl Fn(&str, &Value) -> NestedCallVerdict + Send + Sync + 'static,
1526 ) -> (ToolContext, Arc<AtomicUsize>) {
1527 let (tx_event, mut rx_event) = mpsc::channel(64);
1528 tokio::spawn(async move { while rx_event.recv().await.is_some() {} });
1529 let (gate, mut requests) = NestedCallGate::new(mcp_pool, tx_event, deadline);
1530 let asked = Arc::new(AtomicUsize::new(0));
1531 let counter = asked.clone();
1532 tokio::spawn(async move {
1533 while let Some(request) = requests.recv().await {
1534 counter.fetch_add(1, Ordering::SeqCst);
1535 tokio::time::sleep(delay).await;
1536 let verdict = answer(&request.name, &request.input);
1537 let _ = request.reply.send(verdict);
1538 }
1539 });
1540 let mut context = ToolContext::new(workspace);
1541 context.execution.nested_call_gate = Some(gate);
1542 (context, asked)
1543 }
1544
1545 fn run_as_asked(name: &str, input: &Value) -> NestedCallVerdict {
1546 NestedCallVerdict::Run {
1547 name: name.to_string(),
1548 input: input.clone(),
1549 supports_parallel: true,
1550 decision: NestedDecision::Auto,
1551 hook_context: None,
1552 }
1553 }
1554
1555 fn body(result: &ToolResult) -> Value {
1556 serde_json::from_str(&result.content).expect("receipt is JSON")
1557 }
1558
1559 /// A tool that takes longer than any test deadline.
1560 struct SlowTool;
1561
1562 #[async_trait]
1563 impl ToolSpec for SlowTool {
1564 fn name(&self) -> &str {
1565 "slow"
1566 }
1567 fn description(&self) -> &str {
1568 "sleeps"
1569 }
1570 fn input_schema(&self) -> Value {
1571 json!({"type": "object"})
1572 }
1573 fn capabilities(&self) -> Vec<crate::tools::spec::ToolCapability> {
1574 vec![crate::tools::spec::ToolCapability::ReadOnly]
1575 }
1576 async fn execute(
1577 &self,
1578 _input: Value,
1579 _context: &ToolContext,
1580 ) -> Result<ToolResult, ToolError> {
1581 tokio::time::sleep(Duration::from_secs(30)).await;
1582 Ok(ToolResult::success("late"))
1583 }
1584 }
1585
1586 /// A tool whose output is larger than the nested result cap.
1587 struct HugeTool;
1588
1589 #[async_trait]
1590 impl ToolSpec for HugeTool {
1591 fn name(&self) -> &str {
1592 "huge"
1593 }
1594 fn description(&self) -> &str {
1595 "big output"
1596 }
1597 fn input_schema(&self) -> Value {
1598 json!({"type": "object"})
1599 }
1600 fn capabilities(&self) -> Vec<crate::tools::spec::ToolCapability> {
1601 vec![crate::tools::spec::ToolCapability::ReadOnly]
1602 }
1603 async fn execute(
1604 &self,
1605 input: Value,
1606 _context: &ToolContext,
1607 ) -> Result<ToolResult, ToolError> {
1608 let text = "y".repeat(PER_CALL_RESULT_CAP_BYTES * 2);
1609 if input.get("fail").and_then(Value::as_bool) == Some(true) {
1610 return Ok(ToolResult::error(text));
1611 }
1612 Ok(ToolResult::success(text))
1613 }
1614 }
1615
1616 #[tokio::test]
1617 async fn gated_nested_mcp_call_runs_through_the_session_pool() {
1618 let dir = tempfile::tempdir().unwrap();
1619 let pool = Arc::new(AsyncMutex::new(McpPool::new(
1620 crate::mcp::McpConfig::default(),
1621 )));
1622 let (context, asked) = gated_context(
1623 dir.path(),
1624 Duration::from_secs(30),
1625 Duration::ZERO,
1626 Some(pool),
1627 run_as_asked,
1628 );
1629 let registry = ToolRegistryBuilder::new().build(context.clone());
1630 let code = "const r = await tools.call('list_mcp_resources', {}); \
1631 return { hasContent: r.content !== undefined, truncated: r.truncated };";
1632 let result = execute_tools_tool(&json!({ "code": code }), &registry, &context)
1633 .await
1634 .unwrap();
1635 assert!(result.success, "{}", result.content);
1636 let body = body(&result);
1637 assert_eq!(body["body"]["return"]["hasContent"], true);
1638 assert_eq!(body["body"]["return"]["truncated"], Value::Null);
1639 assert_eq!(body["calls"][0]["tool"], "list_mcp_resources");
1640 assert_eq!(body["calls"][0]["decision"], "auto");
1641 assert_eq!(body["calls"][0]["status"], "ok");
1642 assert_eq!(
1643 asked.load(Ordering::SeqCst),
1644 1,
1645 "the MCP call went through the gate"
1646 );
1647 }
1648
1649 #[tokio::test]
1650 async fn gated_denial_fails_only_that_nested_call() {
1651 let (_dir, note) = workspace_with_note();
1652 let workspace = note.parent().unwrap().to_path_buf();
1653 let target = workspace.join("denied.txt");
1654 let (context, _asked) = gated_context(
1655 &workspace,
1656 Duration::from_secs(30),
1657 Duration::ZERO,
1658 None,
1659 |name, input| {
1660 if name == "write" {
1661 NestedCallVerdict::Refused {
1662 error: ToolError::permission_denied("Tool 'write' denied by user"),
1663 decision: NestedDecision::Denied,
1664 }
1665 } else {
1666 run_as_asked(name, input)
1667 }
1668 },
1669 );
1670 let registry = ToolRegistryBuilder::new()
1671 .with_tool(Arc::new(ReadTool))
1672 .with_tool(Arc::new(WriteTool))
1673 .build(context.clone());
1674 let note_path = note.to_string_lossy().replace('\\', "\\\\");
1675 let target_path = target.to_string_lossy().replace('\\', "\\\\");
1676 let code = format!(
1677 "let denied = null; \
1678 try {{ await tools.call('write', {{ path: '{target_path}', content: 'x' }}); }} \
1679 catch (e) {{ denied = String(e.message || e); }} \
1680 const r = await tools.call('read', {{ path: '{note_path}' }}); \
1681 return {{ denied, read: JSON.stringify(r).includes('alpha') }};"
1682 );
1683 let result = execute_tools_tool(&json!({ "code": code }), &registry, &context)
1684 .await
1685 .unwrap();
1686 assert!(result.success, "{}", result.content);
1687 let body = body(&result);
1688 assert!(
1689 body["body"]["return"]["denied"]
1690 .as_str()
1691 .is_some_and(|message| message.contains("denied by user")),
1692 "{body}"
1693 );
1694 assert_eq!(body["body"]["return"]["read"], true);
1695 assert_eq!(body["calls"][0]["decision"], "denied");
1696 assert_eq!(body["calls"][0]["status"], "refused");
1697 assert_eq!(body["calls"][1]["status"], "ok");
1698 assert!(!target.exists(), "a denied write never runs");
1699 }
1700
1701 #[tokio::test]
1702 async fn computer_use_consent_is_refused_inside_a_program_before_the_gate() {
1703 let dir = tempfile::tempdir().unwrap();
1704 let (context, asked) = gated_context(
1705 dir.path(),
1706 Duration::from_secs(30),
1707 Duration::ZERO,
1708 None,
1709 run_as_asked,
1710 );
1711 let invoker = CodemodeInvoker::new(Vec::new(), context);
1712 for (name, input) in [
1713 (
1714 "mcp_plugin-12-computer-use-computer_consent",
1715 json!({"action": "allow", "app": "Safari", "bundle_id": "com.apple.Safari"}),
1716 ),
1717 (
1718 "mcp_plugin-12-computer-use-computer_app_script",
1719 json!({"script": "do shell script \"id\""}),
1720 ),
1721 ("mcp_github_authenticate", json!({})),
1722 ("exec_shell", json!({"command": "ls", "interactive": true})),
1723 ("request_user_input", json!({})),
1724 // Dispatch resolves names case-insensitively, so the refusals do.
1725 ("Agent", json!({})),
1726 ("WORKFLOW", json!({})),
1727 ("Execute_Tools", json!({"code": "return 1;"})),
1728 ("BASH", json!({"command": "ls", "interactive": true})),
1729 ("Exec_Shell", json!({"command": "ls", "interactive": true})),
1730 ] {
1731 let err = invoker
1732 .invoke(ToolCallRequest {
1733 tool: name.to_string(),
1734 input,
1735 })
1736 .await
1737 .unwrap_err();
1738 assert!(matches!(err, DriverError::Rejected(_)), "{name}: {err:?}");
1739 }
1740 assert_eq!(asked.load(Ordering::SeqCst), 0, "nothing reached the gate");
1741 assert!(
1742 invoker
1743 .drain()
1744 .iter()
1745 .all(|receipt| receipt.status == CallStatus::Refused)
1746 );
1747 }
1748
1749 #[tokio::test]
1750 async fn deadline_keeps_finished_receipts_and_names_in_flight_calls() {
1751 let (_dir, note) = workspace_with_note();
1752 let workspace = note.parent().unwrap().to_path_buf();
1753 let (context, _asked) = gated_context(
1754 &workspace,
1755 Duration::from_millis(400),
1756 Duration::ZERO,
1757 None,
1758 run_as_asked,
1759 );
1760 let registry = ToolRegistryBuilder::new()
1761 .with_tool(Arc::new(ReadTool))
1762 .with_tool(Arc::new(SlowTool))
1763 .build(context.clone());
1764 let note_path = note.to_string_lossy().replace('\\', "\\\\");
1765 let code = format!(
1766 "await tools.call('read', {{ path: '{note_path}' }}); \
1767 await tools.call('slow', {{}}); return 'unreachable';"
1768 );
1769 let result = tokio::time::timeout(
1770 Duration::from_secs(10),
1771 execute_tools_tool(&json!({ "code": code }), &registry, &context),
1772 )
1773 .await
1774 .expect("the run deadline ends the program")
1775 .expect("a deadline is a receipt, not a host error");
1776 assert!(!result.success);
1777 let body = body(&result);
1778 assert_eq!(body["body"]["timed_out"], true, "{body}");
1779 assert_eq!(body["nested_calls"], 2);
1780 assert_eq!(body["calls"][0]["tool"], "read");
1781 assert_eq!(body["calls"][0]["status"], "ok");
1782 assert_eq!(body["calls"][1]["tool"], "slow");
1783 assert_eq!(body["calls"][1]["status"], "in_flight");
1784 }
1785
1786 #[tokio::test]
1787 async fn time_waiting_on_the_gate_does_not_spend_the_deadline() {
1788 let (_dir, note) = workspace_with_note();
1789 let workspace = note.parent().unwrap().to_path_buf();
1790 // Each gate answer takes longer than the whole run deadline, as a
1791 // person deciding an approval would.
1792 let (context, _asked) = gated_context(
1793 &workspace,
1794 Duration::from_millis(300),
1795 Duration::from_millis(500),
1796 None,
1797 run_as_asked,
1798 );
1799 let registry = ToolRegistryBuilder::new()
1800 .with_tool(Arc::new(ReadTool))
1801 .build(context.clone());
1802 let note_path = note.to_string_lossy().replace('\\', "\\\\");
1803 let code = format!(
1804 "const a = await tools.call('read', {{ path: '{note_path}' }}); \
1805 const b = await tools.call('read', {{ path: '{note_path}' }}); \
1806 return JSON.stringify([a, b]).includes('alpha');"
1807 );
1808 let result = execute_tools_tool(&json!({ "code": code }), &registry, &context)
1809 .await
1810 .unwrap();
1811 assert!(result.success, "{}", result.content);
1812 assert_eq!(body(&result)["body"]["return"], true);
1813 }
1814
1815 #[test]
1816 fn watchdog_rechecks_promptly_while_paused_on_the_gate() {
1817 let dir = tempfile::tempdir().unwrap();
1818 let invoker = CodemodeInvoker::new(Vec::new(), ToolContext::new(dir.path()));
1819 let deadline = Duration::from_secs(600);
1820 let paused = invoker.pause();
1821 // A paused program never times out, but the watchdog must not sleep
1822 // out the whole deadline or the run would overrun by that much once
1823 // the gate answers.
1824 assert_eq!(invoker.remaining(deadline), Some(PAUSED_WATCHDOG_POLL));
1825 drop(paused);
1826 assert!(
1827 invoker
1828 .remaining(deadline)
1829 .is_some_and(|left| left > PAUSED_WATCHDOG_POLL)
1830 );
1831 }
1832
1833 // This test deliberately serializes access to process-global spillover
1834 // state while awaiting the program.
1835 #[allow(clippy::await_holding_lock)]
1836 #[tokio::test]
1837 async fn oversized_nested_result_keeps_its_envelope_and_says_what_was_cut() {
1838 let _guard = crate::tools::truncate::TEST_SPILLOVER_GUARD
1839 .lock()
1840 .unwrap_or_else(|err| err.into_inner());
1841 let spill_root = tempfile::tempdir().unwrap();
1842 let previous =
1843 crate::tools::truncate::set_test_spillover_root(Some(spill_root.path().to_path_buf()));
1844 let dir = tempfile::tempdir().unwrap();
1845 let (context, _asked) = gated_context(
1846 dir.path(),
1847 Duration::from_secs(30),
1848 Duration::ZERO,
1849 None,
1850 run_as_asked,
1851 );
1852 let registry = ToolRegistryBuilder::new()
1853 .with_tool(Arc::new(HugeTool))
1854 .build(context.clone());
1855 let code = "const r = await tools.call('huge', {}); \
1856 let failed = null; \
1857 try { await tools.call('huge', { fail: true }); } \
1858 catch (e) { failed = String(e.message || e).length; } \
1859 return { keys: Object.keys(r).sort(), cut: r.truncated, \
1860 head: r.content.length, kind: typeof r.content, failed };";
1861 let result = execute_tools_tool(&json!({ "code": code }), &registry, &context)
1862 .await
1863 .unwrap();
1864 crate::tools::truncate::set_test_spillover_root(previous);
1865 assert!(result.success, "{}", result.content);
1866 let body = body(&result);
1867 let returned = &body["body"]["return"];
1868 assert_eq!(
1869 returned["keys"],
1870 json!(["content", "metadata", "truncated"])
1871 );
1872 assert!(
1873 returned["cut"]["original_bytes"].as_u64().unwrap() > PER_CALL_RESULT_CAP_BYTES as u64
1874 );
1875 assert_eq!(returned["head"], json!(PER_CALL_RESULT_CAP_BYTES / 2));
1876 assert_eq!(returned["kind"], "string", "a cut result is its raw head");
1877 assert!(returned["cut"]["spill_path"].as_str().is_some(), "{body}");
1878 assert!(body["calls"][0]["truncated"]["original_bytes"].is_u64());
1879 // A failed call's text is bounded and spilled the same way.
1880 assert_eq!(returned["failed"], json!(PER_CALL_RESULT_CAP_BYTES / 2));
1881 assert_eq!(body["calls"][1]["status"], "failed");
1882 assert_eq!(
1883 body["calls"][1]["truncated"]["original_bytes"],
1884 json!(PER_CALL_RESULT_CAP_BYTES * 2)
1885 );
1886 assert!(
1887 body["calls"][1]["truncated"]["spill_path"]
1888 .as_str()
1889 .is_some(),
1890 "{body}"
1891 );
1892 }
1893 }
1894
1894 lines RUST