返回 CodeWhale
dispatch.rs
根目录 / crates / tui / src / core / engine / dispatch.rs
1 //! Tool dispatch — plan/execute helpers for the per-turn tool batch.
2 //!
3 //! Extracted from `core/engine.rs` (P1.3). The high-level ordering still
4 //! lives in `Engine::run_turn`; this module owns:
5 //!
6 //! * Streaming-buffer parsing into a finalized `serde_json::Value` tool input
7 //! (`final_tool_input`, `parse_tool_input`, fenced/JSON segment helpers).
8 //! * Policy predicates the turn loop consults — when a batch can run in
9 //! parallel and the small set of read-only MCP tools that are safe to run
10 //! in parallel.
11 //! * The tool execution plan/outcome types the batch driver passes around.
12 //!
13 //! All items are `pub(super)`-only: the public engine surface (Op/Event,
14 //! `EngineHandle`, `spawn_engine`) stays in `core/engine.rs`.
15
16 use serde_json::json;
17
18 use crate::tools::spec::{
19 ResourceClaim, ToolError, ToolExecutionOutcome, ToolResult, ToolResultContentBlock,
20 schedule_non_conflicting,
21 };
22 use codewhale_models::{Tool, ToolCaller};
23
24 use super::ToolUseState;
25
26 const MAX_SCHEMA_CONTAINER_REPAIR_BYTES: usize = 64 * 1024;
27
28 // === Types ============================================================
29
30 #[allow(dead_code)] // `index` mirrors batch order for diagnostic ergonomics.
31 pub(super) struct ToolExecOutcome {
32 pub(super) index: usize,
33 pub(super) id: String,
34 pub(super) model_call: Option<crate::core::events::ModelToolCall>,
35 pub(super) name: String,
36 pub(super) input: serde_json::Value,
37 pub(super) started_at: std::time::Instant,
38 pub(super) terminal: ToolExecutionOutcome,
39 pub(super) content_blocks: Vec<ToolResultContentBlock>,
40 /// Read-result bytes before spillover adds call-specific artifact paths.
41 pub(super) original_content_digest: Option<[u8; 32]>,
42 }
43
44 /// Notice appended as a user-role message when the guard first asks the worker
45 /// to change strategy after repeated no-progress denials (#6015).
46 pub(crate) const FLEET_STRATEGY_SWITCH_NOTICE: &str = "Fleet strategy switch required: repeated permission denials produced no new evidence. The rejected action is held. Use another permitted tool from the current catalog to make progress, or report completed work and the blocker. Do not work around permissions or request the same approval again.";
47
48 /// Notice appended when denials continue past the strategy switch: the next
49 /// response is report-only and its tool calls are admission-held (#6015).
50 pub(crate) const FLEET_FINAL_REPORT_NOTICE: &str = "Fleet no-progress final report: permission denials continued after the strategy switch without new evidence. Your next response is report-only; no tools will execute. Report what you completed, exact evidence, the permission blocker and remaining work. This is the last response unless the user changes direction or authority.";
51
52 /// Terminal reason once the report-only response has been recorded (#6015).
53 pub(crate) const FLEET_NO_PROGRESS_STOP: &str = "Fleet worker stopped after repeated permission denials without new evidence. Work and tool results are retained in the transcript; review the blocker before resuming.";
54
55 /// Progress observations for one provider response, independent of tool finish
56 /// order. Only typed permission denials contribute to the retry guard (#6015).
57 #[derive(Default)]
58 pub(crate) struct FleetDenialBatch {
59 denied: std::collections::HashSet<String>,
60 made_progress: bool,
61 }
62
63 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
64 pub(crate) enum FleetDenialAction {
65 Continue,
66 SwitchStrategy,
67 FinalReport,
68 }
69
70 /// Turn-local guard for an engine with a Fleet authority envelope. This is an
71 /// admission predicate and result accumulator, not another execution loop.
72 /// Three responses give the model two opportunities to use denial feedback;
73 /// after one strategy notice, three more denied responses request a report.
74 /// Fleet sub-agent workers run the same guard in their own loop (#6015).
75 #[derive(Default)]
76 pub(crate) struct FleetDenialGuard {
77 denied_rounds: std::collections::HashMap<String, u8>,
78 switch_requested: bool,
79 recovery_denied_rounds: u8,
80 denial_rounds_without_progress: u32,
81 report_only: bool,
82 // Last observed bytes per read request: paths alone cannot distinguish a
83 // changed file, and an unchanged read must not repeatedly reset denials.
84 // Coverage is bounded; an evicted observation is treated conservatively
85 // as new evidence. No raw arguments/bytes are kept.
86 reads: std::collections::VecDeque<([u8; 32], [u8; 32])>,
87 }
88
89 impl FleetDenialGuard {
90 const REPEATED_DENIAL_ROUNDS: u8 = 3;
91 const MAX_OBSERVATIONS: usize = 32;
92
93 pub(crate) fn reset(&mut self) {
94 *self = Self::default();
95 }
96
97 pub(crate) fn report_only(&self) -> bool {
98 self.report_only
99 }
100
101 pub(crate) fn awaiting_strategy_change(&self) -> bool {
102 self.switch_requested
103 }
104
105 pub(super) fn denial_rounds_without_progress(&self) -> u32 {
106 self.denial_rounds_without_progress
107 }
108
109 pub(crate) fn original_content_digest(
110 name: &str,
111 input: &serde_json::Value,
112 output: &ToolResult,
113 ) -> Option<[u8; 32]> {
114 use sha2::{Digest, Sha256};
115 let action = crate::tools::canonical_action::canonical_action_alias(name, input);
116 (output.success
117 && matches!(
118 action,
119 "read_file" | "list_dir" | "file_search" | "grep_files"
120 ))
121 .then(|| Sha256::digest(output.content.as_bytes()).into())
122 }
123
124 pub(crate) fn admission_error(
125 &self,
126 name: &str,
127 input: &serde_json::Value,
128 ) -> Option<ToolError> {
129 let action = crate::tools::canonical_action::canonical_action_alias(name, input);
130 if self.report_only {
131 Some(ToolError::permission_denied(
132 "Fleet no-progress final report: no tools may execute in this response. Report completed work, evidence and the remaining blocker; do not change permission mode or retry tools.",
133 ))
134 } else if self.switch_requested && self.denied_rounds.contains_key(action) {
135 Some(ToolError::permission_denied(
136 "Fleet permission-denial loop: this action is held until useful permitted work or an explicit authority change. Use another permitted tool or report the blocker; do not change permission mode or request permission again.",
137 ))
138 } else {
139 None
140 }
141 }
142
143 pub(crate) fn observe(
144 &mut self,
145 batch: &mut FleetDenialBatch,
146 name: &str,
147 input: &serde_json::Value,
148 status: crate::tools::spec::ToolTerminalStatus,
149 result: &Result<ToolResult, ToolError>,
150 original_content_digest: Option<[u8; 32]>,
151 ) {
152 use crate::tools::spec::ToolTerminalStatus;
153
154 let action = crate::tools::canonical_action::canonical_action_alias(name, input);
155 if status == ToolTerminalStatus::Denied
156 && matches!(result, Err(ToolError::PermissionDenied { .. }))
157 {
158 batch.denied.insert(action.to_owned());
159 return;
160 }
161 let Ok(output) = result else { return };
162 if status != ToolTerminalStatus::Succeeded
163 || !output.success
164 || output.metadata.as_ref().is_some_and(|metadata| {
165 metadata
166 .get("executed")
167 .and_then(serde_json::Value::as_bool)
168 == Some(false)
169 || metadata
170 .get("cancelled")
171 .and_then(serde_json::Value::as_bool)
172 == Some(true)
173 })
174 {
175 return;
176 }
177 // Waiting is useful coordination, but its repeated success receipt is
178 // neither new evidence nor a failure. Its own timeouts still govern it.
179 if matches!(
180 action,
181 "exec_shell_wait" | "exec_wait" | "terminal/wait" | "wait_for_dev_server" | "sleep"
182 ) || name == "agent"
183 && input
184 .get("action")
185 .and_then(serde_json::Value::as_str)
186 .is_some_and(|action| matches!(action, "wait" | "status" | "list"))
187 {
188 return;
189 }
190 if matches!(
191 action,
192 "read_file" | "list_dir" | "file_search" | "grep_files"
193 ) {
194 use sha2::{Digest, Sha256};
195
196 let mut semantic_input = input.clone();
197 if action != name
198 && let Some(object) = semantic_input.as_object_mut()
199 {
200 object.remove("action");
201 }
202 let mut hasher = Sha256::new();
203 hasher.update(action.as_bytes());
204 hasher.update([0]);
205 // Tool JSON preserves insertion order; reordered equivalent keys
206 // must not manufacture a new read request.
207 hasher.update(crate::client::canonical_json(&semantic_input).as_bytes());
208 let key: [u8; 32] = hasher.finalize().into();
209 let contents = original_content_digest
210 .unwrap_or_else(|| Sha256::digest(output.content.as_bytes()).into());
211 let previous = self
212 .reads
213 .iter()
214 .position(|(old_key, _)| *old_key == key)
215 .and_then(|index| self.reads.remove(index));
216 batch.made_progress |=
217 previous.is_none_or(|(_, old_contents)| old_contents != contents);
218 self.reads.push_back((key, contents));
219 if self.reads.len() > Self::MAX_OBSERVATIONS {
220 self.reads.pop_front();
221 }
222 } else {
223 // A successful mutation or unfamiliar tool is useful work. Do not
224 // terminate it based on guesses about its content or side effects.
225 batch.made_progress = true;
226 }
227 }
228
229 pub(crate) fn finish_batch(&mut self, batch: FleetDenialBatch) -> FleetDenialAction {
230 if self.report_only {
231 return FleetDenialAction::Continue;
232 }
233 if batch.made_progress {
234 self.denied_rounds.clear();
235 self.switch_requested = false;
236 self.recovery_denied_rounds = 0;
237 self.denial_rounds_without_progress = 0;
238 return FleetDenialAction::Continue;
239 }
240 if batch.denied.is_empty() {
241 return FleetDenialAction::Continue;
242 }
243 self.denial_rounds_without_progress = self.denial_rounds_without_progress.saturating_add(1);
244 if self.switch_requested {
245 self.recovery_denied_rounds = self.recovery_denied_rounds.saturating_add(1);
246 if self.recovery_denied_rounds >= Self::REPEATED_DENIAL_ROUNDS {
247 self.report_only = true;
248 return FleetDenialAction::FinalReport;
249 }
250 return FleetDenialAction::Continue;
251 }
252 for action in batch.denied {
253 // The guard never retains payloads. Unknown families beyond this
254 // bounded window do not evict an already observed denial streak.
255 if self.denied_rounds.contains_key(&action)
256 || self.denied_rounds.len() < Self::MAX_OBSERVATIONS
257 {
258 let count = self.denied_rounds.entry(action).or_default();
259 *count = count.saturating_add(1);
260 }
261 }
262 if self
263 .denied_rounds
264 .values()
265 .any(|count| *count >= Self::REPEATED_DENIAL_ROUNDS)
266 {
267 self.switch_requested = true;
268 FleetDenialAction::SwitchStrategy
269 } else {
270 FleetDenialAction::Continue
271 }
272 }
273 }
274
275 #[derive(Debug, Clone)]
276 pub(super) struct ToolExecutionPlan {
277 pub(super) index: usize,
278 pub(super) id: String,
279 pub(super) model_call: Option<crate::core::events::ModelToolCall>,
280 pub(super) name: String,
281 pub(super) input: serde_json::Value,
282 pub(super) caller: Option<ToolCaller>,
283 pub(super) interactive: bool,
284 pub(super) approval_required: bool,
285 pub(super) approval_description: String,
286 pub(super) approval_force_prompt: bool,
287 pub(super) supports_parallel: bool,
288 pub(super) read_only: bool,
289 pub(super) detached_start: bool,
290 pub(super) resources: Vec<ResourceClaim>,
291 pub(super) blocked_error: Option<ToolError>,
292 pub(super) guard_result: Option<ToolResult>,
293 }
294
295 pub(super) enum ToolExecutionBatch {
296 Parallel(Vec<ToolExecutionPlan>),
297 Serial(Box<ToolExecutionPlan>),
298 }
299
300 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
301 pub(super) enum ToolApprovalStamp {
302 ApprovedByUser,
303 ApprovedWithPolicy,
304 }
305
306 impl ToolApprovalStamp {
307 const ALL: [Self; 2] = [Self::ApprovedByUser, Self::ApprovedWithPolicy];
308
309 fn decision(self) -> &'static str {
310 match self {
311 Self::ApprovedByUser => "approved_by_user",
312 Self::ApprovedWithPolicy => "approved_with_policy",
313 }
314 }
315
316 fn model_visible_note(self) -> &'static str {
317 match self {
318 Self::ApprovedByUser => {
319 "[approval] This tool call required approval and was approved by the user before execution."
320 }
321 Self::ApprovedWithPolicy => {
322 "[approval] This tool call required approval and was approved by the user with an adjusted execution policy before execution."
323 }
324 }
325 }
326 }
327
328 pub(super) fn stamp_tool_result_approval(result: &mut ToolResult, approval: ToolApprovalStamp) {
329 let approval_metadata = json!({
330 "required": true,
331 "decision": approval.decision(),
332 "model_visible": true,
333 });
334 let metadata = result.metadata.get_or_insert_with(|| json!({}));
335 if let Some(object) = metadata.as_object_mut() {
336 object.insert("approval".to_string(), approval_metadata);
337 } else {
338 let prior = std::mem::replace(metadata, json!({}));
339 if let Some(object) = metadata.as_object_mut() {
340 object.insert("_prior".to_string(), prior);
341 object.insert("approval".to_string(), approval_metadata);
342 }
343 }
344
345 let note = approval.model_visible_note();
346 if result.content.starts_with("[approval] ") {
347 return;
348 }
349 if result.content.is_empty() {
350 result.content = note.to_string();
351 } else {
352 result.content = format!("{note}\n\n{}", result.content);
353 }
354 }
355
356 /// The tool output a person reads: the result without the note
357 /// [`stamp_tool_result_approval`] put in front of it for the model (#6566).
358 ///
359 /// Only a result the engine stamped (`metadata.approval.model_visible`) loses
360 /// a note, and only that decision's exact note text followed by a blank line
361 /// (or nothing). Output that merely begins with "[approval] " — from a
362 /// command, a file, or anything else a tool read — is shown whole, so
363 /// injected text cannot hide tool output from the person.
364 pub(crate) fn content_without_approval_note(result: &ToolResult) -> &str {
365 let content = result.content.as_str();
366 let approval = result
367 .metadata
368 .as_ref()
369 .and_then(|metadata| metadata.get("approval"));
370 let stamped = approval
371 .and_then(|approval| approval.get("model_visible"))
372 .and_then(serde_json::Value::as_bool)
373 == Some(true);
374 let Some(stamp) = approval
375 .and_then(|approval| approval.get("decision"))
376 .and_then(serde_json::Value::as_str)
377 .and_then(|decision| {
378 ToolApprovalStamp::ALL
379 .into_iter()
380 .find(|stamp| stamp.decision() == decision)
381 })
382 .filter(|_| stamped)
383 else {
384 return content;
385 };
386 match content.strip_prefix(stamp.model_visible_note()) {
387 Some("") => "",
388 Some(rest) => rest.strip_prefix("\n\n").unwrap_or(content),
389 None => content,
390 }
391 }
392
393 // Hold the lock guard for the duration of a tool execution.
394 // The inner guards are held for RAII purposes (dropped when the guard is dropped).
395 pub(super) enum ToolExecGuard<'a> {
396 Read(#[allow(dead_code)] tokio::sync::RwLockReadGuard<'a, ()>),
397 Write(#[allow(dead_code)] tokio::sync::RwLockWriteGuard<'a, ()>),
398 }
399
400 // === Caller policy and errors ========================================
401
402 pub(super) fn caller_type_for_tool_use(caller: Option<&ToolCaller>) -> &str {
403 caller.map_or("direct", |c| c.caller_type.as_str())
404 }
405
406 pub(super) fn caller_allowed_for_tool(
407 caller: Option<&ToolCaller>,
408 tool_def: Option<&Tool>,
409 ) -> bool {
410 let requested = caller_type_for_tool_use(caller);
411 if let Some(def) = tool_def
412 && let Some(allowed) = &def.allowed_callers
413 {
414 if allowed.is_empty() {
415 return requested == "direct";
416 }
417 return allowed.iter().any(|item| item == requested);
418 }
419 requested == "direct"
420 }
421
422 /// Whole-word check for "mode"/"modes" — a plain `contains("mode")` also
423 /// matched "model", letting provider model errors skip the actionable-hint
424 /// suffix (#3020).
425 fn mentions_mode_word(lower: &str) -> bool {
426 lower
427 .split(|ch: char| !ch.is_ascii_alphanumeric())
428 .any(|word| word == "mode" || word == "modes")
429 }
430
431 #[cfg(test)]
432 pub(super) fn format_tool_error(err: &ToolError, tool_name: &str) -> String {
433 format_tool_error_with_schema(err, tool_name, None)
434 }
435
436 pub(super) fn format_tool_error_with_schema(
437 err: &ToolError,
438 tool_name: &str,
439 input_schema: Option<&serde_json::Value>,
440 ) -> String {
441 let message = match err {
442 ToolError::InvalidInput { message } => {
443 format!("Invalid input for tool '{tool_name}': {message}")
444 }
445 ToolError::MissingField { field } => {
446 format!("Tool '{tool_name}' is missing required field '{field}'")
447 }
448 ToolError::PathEscape { path } => format!(
449 "Path escapes workspace: {}. Use a workspace-relative path or enable trust mode.",
450 path.display()
451 ),
452 ToolError::ExecutionFailed { message, .. } => message.clone(),
453 ToolError::Timeout { seconds } => format!(
454 "Tool '{tool_name}' timed out after {seconds}s. Try a narrower scope or a longer timeout."
455 ),
456 ToolError::Cancelled { message } => message.clone(),
457 ToolError::NotAvailable { message } => {
458 let lower = message.to_ascii_lowercase();
459 // #3020: Pass through self-explanatory messages that already name the
460 // cause (mode switch, allow_shell, feature flag). Avoids appending a
461 // conflicting "Check mode, feature flags" suffix on top of
462 // "switch to Act mode" which already gives the recovery path.
463 if lower.contains("current tool catalog")
464 || lower.contains("did you mean:")
465 || mentions_mode_word(&lower)
466 || lower.contains("allow_shell")
467 || lower.contains("feature flag")
468 {
469 message.clone()
470 } else {
471 format!(
472 "Tool '{tool_name}' is not available: {message}. Check mode, feature flags, or tool name."
473 )
474 }
475 }
476 ToolError::PermissionDenied { message } => {
477 let lower = message.to_ascii_lowercase();
478 // #3020: messages that already name the denial cause get no
479 // conflicting "Adjust approval mode" suffix. They keep the
480 // `Tool '…' was denied:` lead, which is how a receipt tells a
481 // call Codewhale blocked from one that ran and failed; an
482 // approval denial already starts `Tool '…' denied by user`.
483 if lower.contains("denied by user") {
484 message.clone()
485 } else if mentions_mode_word(&lower) || lower.contains("allow_shell") {
486 format!("Tool '{tool_name}' was denied: {message}")
487 } else {
488 format!(
489 "Tool '{tool_name}' was denied: {message}. Adjust approval mode or request permission."
490 )
491 }
492 }
493 };
494
495 let (category, bad_field) = match err {
496 ToolError::InvalidInput { .. } => ("invalid_input", None),
497 ToolError::MissingField { field } => ("missing_field", Some(field.as_str())),
498 ToolError::PathEscape { .. } => ("path_escape", Some("path")),
499 ToolError::NotAvailable { .. } => ("tool_not_available", Some("tool_name")),
500 _ => return message,
501 };
502 let valid_shape = input_schema.cloned().unwrap_or_else(|| {
503 serde_json::json!({
504 "type": "object",
505 "guidance": format!("Use the advertised input schema for '{tool_name}'")
506 })
507 });
508 let feedback = serde_json::json!({
509 "category": category,
510 "bad_field": bad_field,
511 "valid_shape": valid_shape,
512 "retryable": true,
513 "side_effect_status": "not_started"
514 });
515 format!("{message}\nTool validation feedback: {feedback}")
516 }
517
518 // === Streaming-buffer parsing =========================================
519
520 /// Promote a streaming `ToolUseState` to a finalized JSON input.
521 ///
522 /// Order of preference:
523 ///
524 /// 1. `input_buffer` (the raw streamed delta concatenation) — parsed as
525 /// JSON. This is the most authoritative because it's what the model
526 /// actually emitted.
527 /// 2. `input` (the per-delta best-effort parse mirror) — used when the
528 /// buffer is empty (pre-streaming tool calls take this path).
529 /// 3. `input_buffer` non-empty but unparseable → fall back to `input`
530 /// (the per-delta parser has already mirrored the most recent valid
531 /// partial parse into `tool_state.input`).
532 pub(super) fn final_tool_input(state: &ToolUseState) -> serde_json::Value {
533 if state.input_parse_error.is_some() {
534 return malformed_tool_arguments_input(&state.input_buffer);
535 }
536 if !state.input_buffer.trim().is_empty()
537 && let Some(parsed) = parse_tool_input(&state.input_buffer)
538 {
539 // Structure was synthesized to make this parse, so the argument text
540 // was cut off. Route it to the same malformed-arguments path as an
541 // outright parse failure rather than dispatching a completed guess.
542 if parsed.structure_synthesized {
543 return malformed_tool_arguments_input(&state.input_buffer);
544 }
545 return parsed.value;
546 }
547 state.input.clone()
548 }
549
550 /// A parsed tool-argument buffer, plus whether the parse only succeeded
551 /// because the repair ladder synthesized structure (see
552 /// `crate::tools::arg_repair`). Mid-stream callers mirroring partial state
553 /// may ignore the flag; the caller making the final dispatch decision must
554 /// not, because synthesized structure means the argument text was cut off.
555 pub(super) struct ParsedToolInput {
556 pub(super) value: serde_json::Value,
557 pub(super) structure_synthesized: bool,
558 }
559
560 pub(super) fn parse_tool_input(buffer: &str) -> Option<ParsedToolInput> {
561 let trimmed = buffer.trim();
562 if trimmed.is_empty() {
563 return None;
564 }
565 // Try the deterministic arg-repair ladder first (handles trailing commas,
566 // unclosed braces, embedded control chars, etc.)
567 if let Ok(repaired) = crate::tools::arg_repair::repair(trimmed) {
568 return Some(ParsedToolInput {
569 value: repaired.value,
570 structure_synthesized: repaired.structure_synthesized,
571 });
572 }
573 // Fall back to existing strategies for code-fenced, double-encoded, and
574 // segment-extraction patterns that the repair ladder doesn't cover.
575 if let Some(stripped) = strip_code_fences(trimmed)
576 && let Ok(value) = serde_json::from_str::<serde_json::Value>(&stripped)
577 {
578 return Some(ParsedToolInput {
579 value,
580 structure_synthesized: false,
581 });
582 }
583 if let Ok(serde_json::Value::String(inner)) = serde_json::from_str::<serde_json::Value>(trimmed)
584 && let Ok(value) = serde_json::from_str::<serde_json::Value>(&inner)
585 {
586 return Some(ParsedToolInput {
587 value,
588 structure_synthesized: false,
589 });
590 }
591 extract_json_segment(trimmed)
592 .and_then(|segment| serde_json::from_str::<serde_json::Value>(&segment).ok())
593 .map(|value| ParsedToolInput {
594 value,
595 structure_synthesized: false,
596 })
597 }
598
599 /// Decode a JSON container that a provider encoded as a string when the tool
600 /// schema explicitly requires an object or array.
601 ///
602 /// This intentionally avoids general argument coercion: the string must be
603 /// bounded, parse as strict JSON, and decode to the declared container type.
604 /// Primitive strings are never coerced.
605 pub(super) fn normalize_schema_json_containers(
606 value: &mut serde_json::Value,
607 schema: &serde_json::Value,
608 ) -> usize {
609 let expected_container = if schema_declares_type(schema, "object") {
610 Some("object")
611 } else if schema_declares_type(schema, "array") {
612 Some("array")
613 } else {
614 None
615 };
616
617 if let (Some(expected), serde_json::Value::String(encoded)) = (expected_container, &*value)
618 && encoded.len() <= MAX_SCHEMA_CONTAINER_REPAIR_BYTES
619 && let Ok(decoded) = serde_json::from_str::<serde_json::Value>(encoded)
620 && ((expected == "object" && decoded.is_object())
621 || (expected == "array" && decoded.is_array()))
622 {
623 *value = decoded;
624 return 1 + normalize_schema_json_containers(value, schema);
625 }
626
627 match value {
628 serde_json::Value::Object(object) => {
629 let properties = schema
630 .get("properties")
631 .and_then(serde_json::Value::as_object);
632 object
633 .iter_mut()
634 .map(|(key, child)| {
635 properties
636 .and_then(|items| items.get(key))
637 .map(|child_schema| normalize_schema_json_containers(child, child_schema))
638 .unwrap_or(0)
639 })
640 .sum()
641 }
642 serde_json::Value::Array(items) => schema
643 .get("items")
644 .map(|item_schema| {
645 items
646 .iter_mut()
647 .map(|item| normalize_schema_json_containers(item, item_schema))
648 .sum()
649 })
650 .unwrap_or(0),
651 _ => 0,
652 }
653 }
654
655 fn schema_declares_type(schema: &serde_json::Value, expected: &str) -> bool {
656 match schema.get("type") {
657 Some(serde_json::Value::String(value)) => value == expected,
658 Some(serde_json::Value::Array(values)) => values.iter().any(|value| value == expected),
659 _ => false,
660 }
661 }
662
663 pub(super) fn malformed_tool_arguments_input(buffer: &str) -> serde_json::Value {
664 json!({ "raw_arguments": buffer })
665 }
666
667 pub(super) fn malformed_tool_arguments_error(buffer: &str) -> String {
668 format!("malformed tool arguments from model: expected valid JSON, got {buffer:?}")
669 }
670
671 fn strip_code_fences(text: &str) -> Option<String> {
672 if !text.contains("```") {
673 return None;
674 }
675 let line_count = text.lines().count();
676 let mut lines = Vec::with_capacity(line_count);
677 for line in text.lines() {
678 if line.trim_start().starts_with("```") {
679 continue;
680 }
681 lines.push(line);
682 }
683 let stripped = lines.join("\n");
684 let stripped = stripped.trim();
685 if stripped.is_empty() {
686 None
687 } else {
688 Some(stripped.to_string())
689 }
690 }
691
692 fn extract_json_segment(text: &str) -> Option<String> {
693 extract_balanced_segment(text, '{', '}').or_else(|| extract_balanced_segment(text, '[', ']'))
694 }
695
696 fn extract_balanced_segment(text: &str, open: char, close: char) -> Option<String> {
697 let start = text.find(open)?;
698 let mut depth = 0i32;
699 let mut end = None;
700 for (offset, ch) in text[start..].char_indices() {
701 if ch == open {
702 depth += 1;
703 } else if ch == close {
704 depth -= 1;
705 if depth == 0 {
706 end = Some(start + offset + ch.len_utf8());
707 break;
708 }
709 }
710 }
711 end.map(|end_idx| text[start..end_idx].to_string())
712 }
713
714 // === Dispatch policy ==================================================
715
716 #[cfg(test)]
717 pub(super) fn should_parallelize_tool_batch(plans: &[ToolExecutionPlan]) -> bool {
718 if plans.is_empty() || !plans.iter().all(tool_plan_can_join_parallel_batch) {
719 return false;
720 }
721 schedule_non_conflicting(
722 plans
723 .iter()
724 .map(|plan| ((), plan.resources.clone()))
725 .collect(),
726 )
727 .len()
728 == 1
729 }
730
731 pub(super) fn tool_plan_is_parallel_safe(plan: &ToolExecutionPlan) -> bool {
732 plan.read_only && plan.supports_parallel && !plan.approval_required && !plan.interactive
733 }
734
735 pub(super) fn tool_plan_can_join_parallel_batch(plan: &ToolExecutionPlan) -> bool {
736 plan.blocked_error.is_none()
737 && (tool_plan_is_parallel_safe(plan)
738 || (plan.detached_start && !plan.approval_required && !plan.interactive))
739 }
740
741 pub(super) fn plan_tool_execution_batches(
742 plans: Vec<ToolExecutionPlan>,
743 ) -> Vec<ToolExecutionBatch> {
744 let mut batches = Vec::new();
745 let mut parallel_candidates = Vec::new();
746
747 let flush_parallel = |parallel_candidates: &mut Vec<_>,
748 batches: &mut Vec<ToolExecutionBatch>| {
749 for chunk in schedule_non_conflicting(std::mem::take(parallel_candidates)) {
750 batches.push(ToolExecutionBatch::Parallel(chunk));
751 }
752 };
753
754 for plan in plans {
755 if tool_plan_can_join_parallel_batch(&plan) {
756 let resources = plan.resources.clone();
757 parallel_candidates.push((plan, resources));
758 continue;
759 }
760
761 flush_parallel(&mut parallel_candidates, &mut batches);
762 batches.push(ToolExecutionBatch::Serial(Box::new(plan)));
763 }
764
765 flush_parallel(&mut parallel_candidates, &mut batches);
766
767 batches
768 }
769
770 pub(super) fn mcp_tool_is_parallel_safe(name: &str) -> bool {
771 matches!(
772 name,
773 "list_mcp_resources"
774 | "list_mcp_resource_templates"
775 | "mcp_read_resource"
776 | "read_mcp_resource"
777 | "mcp_get_prompt"
778 )
779 }
780
781 pub(super) fn mcp_tool_is_read_only(name: &str) -> bool {
782 matches!(
783 name,
784 "list_mcp_resources"
785 | "list_mcp_resource_templates"
786 | "mcp_read_resource"
787 | "read_mcp_resource"
788 | "mcp_get_prompt"
789 )
790 }
791
792 pub(super) fn mcp_tool_approval_description(name: &str, input: &serde_json::Value) -> String {
793 use crate::tools::approval_cache::{ComputerUseUserGate, computer_use_user_gate};
794
795 // K1/K2: a Computer Use consent or script card names exactly what the
796 // person is granting. Generic "may have side effects" text is how a
797 // model-issued consent used to read as routine.
798 match computer_use_user_gate(name, input) {
799 Some(ComputerUseUserGate::Consent {
800 action,
801 app,
802 bundle_id,
803 scope,
804 remember,
805 confirm,
806 }) => {
807 if confirm {
808 // The plugin paused on an action that cannot be taken back
809 // and handed the model a token; approving this card is the
810 // person's confirmation of that one action.
811 return "Computer Use confirmation requested by the model: allow the irreversible action (pay, buy, send, transfer or delete) the plugin just paused on. Approve only if you asked for exactly that action.".to_string();
812 }
813 let target = match scope {
814 "foreground" => {
815 "shared-desktop foreground control (take the pointer and focus)".to_string()
816 }
817 _ => {
818 let app = app.as_deref().unwrap_or("<unnamed app>");
819 match bundle_id.as_deref() {
820 Some(bundle) => format!("app '{app}' (bundle id {bundle})"),
821 None => format!("app '{app}' (bundle id not given)"),
822 }
823 }
824 };
825 let lifetime = if action == "revoke" {
826 "clears session and persisted decisions, including a saved deny"
827 } else if remember {
828 "persisted until revoked"
829 } else {
830 "this session"
831 };
832 let verb = if action == "revoke" {
833 "revoke recorded decisions for"
834 } else {
835 "allow"
836 };
837 return format!(
838 "Computer Use consent requested by the model: {verb} {target}; scope: {scope}; {lifetime}. Approve only if you want this."
839 );
840 }
841 Some(ComputerUseUserGate::AppScript {
842 language,
843 script_sha256,
844 first_line,
845 line_count,
846 }) => {
847 let shown = if line_count > 1 {
848 format!("first of {line_count} lines")
849 } else {
850 "1 line".to_string()
851 };
852 return format!(
853 "Computer Use app_script: run this exact {language} script outside the sandbox (sha256 {}, {shown}): {first_line}",
854 &script_sha256[..16]
855 );
856 }
857 Some(ComputerUseUserGate::Computer {
858 action,
859 transport,
860 destination,
861 }) => {
862 let transport = transport.as_deref().unwrap_or("default");
863 let destination = destination.as_deref().unwrap_or("<not given>");
864 return format!(
865 "Computer Use: {action} a computer the model will then drive (transport {transport}, {destination}). Approve only if you want this machine controlled."
866 );
867 }
868 None => {}
869 }
870 match crate::mcp::mcp_tool_approval_hint(name) {
871 _ if mcp_tool_is_read_only(name) => format!("Read-only MCP tool '{name}'"),
872 Some(crate::mcp::McpToolApprovalHint::TrustedReadOnly) => {
873 format!("Read-only MCP tool '{name}' (declared by a reviewed plugin)")
874 }
875 Some(crate::mcp::McpToolApprovalHint::Destructive) => {
876 format!("MCP tool '{name}' is marked destructive by its server")
877 }
878 None => format!("MCP tool '{name}' may have side effects"),
879 }
880 }
881
882 #[cfg(test)]
883 mod schema_json_container_tests {
884 use super::*;
885 use crate::tools::spec::ToolSpec;
886 use serde_json::json;
887
888 #[test]
889 fn decodes_nested_containers_and_passes_tool_validation() {
890 let schema = crate::tools::user_input::RequestUserInputTool::default().input_schema();
891 let encoded_options = serde_json::to_string(&json!([
892 { "label": "Repository", "description": "Inspect the current repository" },
893 { "label": "Workspace", "description": "Inspect the whole workspace" }
894 ]))
895 .expect("encode options");
896 let encoded_questions = serde_json::to_string(&json!([{
897 "header": "Scope",
898 "id": "scope",
899 "question": "Which scope should be inspected?",
900 "options": encoded_options
901 }]))
902 .expect("encode questions");
903 let mut input = json!({ "questions": encoded_questions });
904
905 assert_eq!(normalize_schema_json_containers(&mut input, &schema), 2);
906 assert!(input["questions"].is_array());
907 assert!(input["questions"][0]["options"].is_array());
908 crate::tools::user_input::UserInputRequest::from_value(&input)
909 .expect("normalized input must still pass tool-specific validation");
910 }
911
912 #[test]
913 fn leaves_primitives_wrong_types_and_unbounded_strings_unchanged() {
914 let schema = json!({
915 "type": "object",
916 "properties": {
917 "text": { "type": "string" },
918 "count": { "type": "integer" },
919 "items": { "type": "array" },
920 "oversized": { "type": "array" }
921 }
922 });
923 let oversized = format!("[\"{}\"]", "x".repeat(MAX_SCHEMA_CONTAINER_REPAIR_BYTES));
924 let mut input = json!({
925 "text": "[\"still text\"]",
926 "count": "10",
927 "items": "{\"wrong\":\"container\"}",
928 "oversized": oversized
929 });
930 let before = input.clone();
931
932 assert_eq!(normalize_schema_json_containers(&mut input, &schema), 0);
933 assert_eq!(input, before);
934 }
935 }
936
936 lines RUST