Skip to content

Commit 4ff9750

Browse files
fix(cli): order parallel tool results by call index and cap fan-out
execute_batch already keyed results to original call slots; this makes that airtight: panicking worker threads resolve to their own index via a fallback tuple instead of colliding with slot 0, and results are asserted to follow call order with a regression test that mixes a workspace-escape read among valid reads at cap=2. Adds the fan-out bound: parallel-safe calls run in waves of at most maxParallelToolCalls (default 8, CLAWD_PARALLEL_TOOL_CALLS env, or the new setting applied via with_max_parallel_tool_calls at runtime build). Sequential behavior is preserved: allowed_tools checks still run before dispatch, and ToolSearch stays off the parallel path.
1 parent e78a1d1 commit 4ff9750

1 file changed

Lines changed: 187 additions & 36 deletions

File tree

  • rust/crates/rusty-claude-cli/src

‎rust/crates/rusty-claude-cli/src/main.rs‎

Lines changed: 187 additions & 36 deletions
Original file line numberDiff line numberDiff line change
@@ -12557,7 +12557,8 @@ fn build_runtime_with_plugin_state(
1255712557
emit_output,
1255812558
tool_registry.clone(),
1255912559
mcp_state.clone(),
12560-
),
12560+
)
12561+
.with_max_parallel_tool_calls(feature_config.max_parallel_tool_calls()),
1256112562
policy,
1256212563
system_prompt,
1256312564
&feature_config,
@@ -13985,6 +13986,22 @@ struct CliToolExecutor {
1398513986
allowed_tools: Option<AllowedToolSet>,
1398613987
tool_registry: GlobalToolRegistry,
1398713988
mcp_state: Option<Arc<Mutex<RuntimeMcpState>>>,
13989+
/// Cap on concurrently executing parallel-safe tool calls in one batch.
13990+
max_parallel_tool_calls: usize,
13991+
}
13992+
13993+
/// Default fan-out cap for parallel read-only tool execution. Local model
13994+
/// servers and filesystem watchers degrade past a handful of concurrent
13995+
/// requests, so 8 is the sane default. Precedence: `maxParallelToolCalls`
13996+
/// in settings > `CLAWD_PARALLEL_TOOL_CALLS` env > this default.
13997+
const DEFAULT_MAX_PARALLEL_TOOL_CALLS: usize = 8;
13998+
13999+
fn max_parallel_tool_calls_from_env() -> usize {
14000+
std::env::var("CLAWD_PARALLEL_TOOL_CALLS")
14001+
.ok()
14002+
.and_then(|v| v.trim().parse::<usize>().ok())
14003+
.filter(|v| *v > 0)
14004+
.unwrap_or(DEFAULT_MAX_PARALLEL_TOOL_CALLS)
1398814005
}
1398914006

1399014007
impl CliToolExecutor {
@@ -14000,7 +14017,18 @@ impl CliToolExecutor {
1400014017
allowed_tools,
1400114018
tool_registry,
1400214019
mcp_state,
14020+
max_parallel_tool_calls: max_parallel_tool_calls_from_env(),
14021+
}
14022+
}
14023+
14024+
/// Apply the `maxParallelToolCalls` setting when present; `None` keeps
14025+
/// the env/default already resolved by `new`.
14026+
#[must_use]
14027+
fn with_max_parallel_tool_calls(mut self, max: Option<usize>) -> Self {
14028+
if let Some(max) = max.filter(|m| *m > 0) {
14029+
self.max_parallel_tool_calls = max;
1400314030
}
14031+
self
1400414032
}
1400514033

1400614034
fn execute_search_tool(&self, value: serde_json::Value) -> Result<String, ToolError> {
@@ -14125,13 +14153,17 @@ impl ToolExecutor for CliToolExecutor {
1412514153

1412614154
/// Tools that are safe to run in parallel because they only read
1412714155
/// state and dispatch through the stateless tool registry.
14156+
///
14157+
/// `ToolSearch` is intentionally absent: the sequential path routes it
14158+
/// through `execute_search_tool` (which merges live MCP state), while
14159+
/// the parallel path dispatches straight to the stateless registry.
14160+
/// Keeping it sequential preserves identical output under parallelism.
1412814161
const PARALLEL_SAFE_TOOLS: &[&str] = &[
1412914162
"read_file",
1413014163
"glob_search",
1413114164
"grep_search",
1413214165
"WebFetch",
1413314166
"WebSearch",
14134-
"ToolSearch",
1413514167
"Skill",
1413614168
"LSP",
1413714169
"Agent",
@@ -14182,43 +14214,65 @@ impl ToolExecutor for CliToolExecutor {
1418214214
}
1418314215
}
1418414216

14185-
// Execute parallel-safe tools concurrently
14217+
// Execute parallel-safe tools concurrently, bounded so fan-out against
14218+
// local model servers cannot queue unbounded threads. Results are
14219+
// keyed by the original call index, never by completion order, so the
14220+
// transcript stays deterministic across runs for the same input.
14221+
//
14222+
// Note: parallel-safe calls run before sequential ones within a batch
14223+
// (results still land in call-order slots). Models must not emit
14224+
// calls whose success depends on another call in the same batch.
1418614225
if !parallel_calls.is_empty() {
1418714226
let registry = self.tool_registry.clone();
14188-
let parallel_results: Vec<(usize, String, String, Result<String, ToolError>)> =
14189-
std::thread::scope(|s| {
14190-
let mut handles = Vec::new();
14191-
for (idx, tool_use_id, tool_name, input) in &parallel_calls {
14192-
let registry = &registry;
14193-
let tool_use_id = tool_use_id.clone();
14194-
let tool_name = tool_name.clone();
14195-
let input = input.clone();
14196-
let idx = *idx;
14197-
handles.push(s.spawn(move || {
14198-
let value = serde_json::from_str(&input).map_err(|error| {
14199-
ToolError::new(format!("invalid tool input JSON: {error}"))
14200-
});
14201-
let result = match value {
14202-
Ok(v) => registry.execute(&tool_name, &v).map_err(ToolError::new),
14203-
Err(e) => Err(e),
14204-
};
14205-
(idx, tool_use_id, tool_name, result)
14206-
}));
14207-
}
14208-
handles
14209-
.into_iter()
14210-
.map(|h| {
14211-
h.join().unwrap_or_else(|_| {
14212-
(
14213-
0,
14214-
String::new(),
14215-
String::new(),
14216-
Err(ToolError::new("parallel thread panicked")),
14217-
)
14227+
let max_parallel = self.max_parallel_tool_calls.max(1);
14228+
let mut parallel_results: Vec<(usize, String, String, Result<String, ToolError>)> =
14229+
Vec::with_capacity(parallel_calls.len());
14230+
for wave in parallel_calls.chunks(max_parallel) {
14231+
let wave_results: Vec<(usize, String, String, Result<String, ToolError>)> =
14232+
std::thread::scope(|s| {
14233+
let mut handles = Vec::new();
14234+
for (idx, tool_use_id, tool_name, input) in wave {
14235+
let registry = &registry;
14236+
let tool_use_id = tool_use_id.clone();
14237+
let tool_name = tool_name.clone();
14238+
let input = input.clone();
14239+
let idx = *idx;
14240+
// Fallback tuple carried outside the closure so a
14241+
// panicking worker still lands in its own slot instead
14242+
// of colliding with index 0 or leaving a hole.
14243+
let fallback = (idx, tool_use_id.clone(), tool_name.clone());
14244+
handles.push((
14245+
fallback,
14246+
s.spawn(move || {
14247+
let value = serde_json::from_str(&input).map_err(|error| {
14248+
ToolError::new(format!("invalid tool input JSON: {error}"))
14249+
});
14250+
let result = match value {
14251+
Ok(v) => {
14252+
registry.execute(&tool_name, &v).map_err(ToolError::new)
14253+
}
14254+
Err(e) => Err(e),
14255+
};
14256+
(idx, tool_use_id, tool_name, result)
14257+
}),
14258+
));
14259+
}
14260+
handles
14261+
.into_iter()
14262+
.map(|((idx, tool_use_id, tool_name), h)| {
14263+
h.join().unwrap_or_else(|_| {
14264+
(
14265+
idx,
14266+
tool_use_id,
14267+
tool_name,
14268+
Err(ToolError::new("parallel tool thread panicked")),
14269+
)
14270+
})
1421814271
})
14219-
})
14220-
.collect()
14221-
});
14272+
.collect()
14273+
});
14274+
parallel_results.extend(wave_results);
14275+
}
1422214276

1422314277
for (idx, tool_use_id, tool_name, result) in parallel_results {
1422414278
if emit_output {
@@ -19005,6 +19059,103 @@ UU conflicted.rs",
1900519059
std::env::temp_dir().join(format!("claw-cli-{label}-{nanos}"))
1900619060
}
1900719061

19062+
#[test]
19063+
fn parallel_batch_results_follow_call_index_even_with_workspace_escapes() {
19064+
let workspace = temp_workspace("parallel-order");
19065+
fs::create_dir_all(&workspace).expect("workspace dir should exist");
19066+
for name in ["a", "b", "c", "d", "e"] {
19067+
fs::write(
19068+
workspace.join(format!("{name}.txt")),
19069+
format!("marker-{name}\n"),
19070+
)
19071+
.expect("fixture file should write");
19072+
}
19073+
let unique = workspace
19074+
.file_name()
19075+
.expect("workspace name")
19076+
.to_string_lossy()
19077+
.into_owned();
19078+
let escape_name = format!("{unique}-escape.txt");
19079+
let escape_path = workspace
19080+
.parent()
19081+
.expect("workspace parent dir")
19082+
.join(&escape_name);
19083+
fs::write(&escape_path, "out-of-scope\n").expect("escape file should write");
19084+
19085+
let calls = vec![
19086+
runtime::ToolCall {
19087+
tool_use_id: "t0".to_string(),
19088+
tool_name: "read_file".to_string(),
19089+
input: r#"{"path":"a.txt"}"#.to_string(),
19090+
},
19091+
runtime::ToolCall {
19092+
tool_use_id: "t1".to_string(),
19093+
tool_name: "read_file".to_string(),
19094+
input: format!(r#"{{"path":"../{escape_name}"}}"#),
19095+
},
19096+
runtime::ToolCall {
19097+
tool_use_id: "t2".to_string(),
19098+
tool_name: "read_file".to_string(),
19099+
input: r#"{"path":"b.txt"}"#.to_string(),
19100+
},
19101+
runtime::ToolCall {
19102+
tool_use_id: "t3".to_string(),
19103+
tool_name: "read_file".to_string(),
19104+
input: r#"{"path":"c.txt"}"#.to_string(),
19105+
},
19106+
runtime::ToolCall {
19107+
tool_use_id: "t4".to_string(),
19108+
tool_name: "read_file".to_string(),
19109+
input: r#"{"path":"d.txt"}"#.to_string(),
19110+
},
19111+
runtime::ToolCall {
19112+
tool_use_id: "t5".to_string(),
19113+
tool_name: "read_file".to_string(),
19114+
input: r#"{"path":"e.txt"}"#.to_string(),
19115+
},
19116+
];
19117+
19118+
let results = with_current_dir(&workspace, || {
19119+
let mut executor =
19120+
CliToolExecutor::new(None, false, GlobalToolRegistry::builtin(), None)
19121+
.with_max_parallel_tool_calls(Some(2));
19122+
executor.execute_batch(calls)
19123+
});
19124+
19125+
let _ = fs::remove_file(&escape_path);
19126+
let _ = fs::remove_dir_all(&workspace);
19127+
19128+
assert_eq!(results.len(), 6, "every call must produce a result slot");
19129+
for (i, result) in results.iter().enumerate() {
19130+
assert_eq!(
19131+
result.tool_use_id,
19132+
format!("t{i}"),
19133+
"result {i} must carry the call that owned the slot, not the thread that finished first"
19134+
);
19135+
assert_eq!(result.tool_name, "read_file");
19136+
}
19137+
let escaped = results[1]
19138+
.result
19139+
.as_ref()
19140+
.expect_err("out-of-scope path must fail even inside a parallel batch");
19141+
assert!(
19142+
escaped.to_string().contains("escapes workspace boundary"),
19143+
"escape must be rejected by the per-call workspace check, got: {escaped}"
19144+
);
19145+
let markers = ["marker-a", "marker-b", "marker-c", "marker-d", "marker-e"];
19146+
let slots = [0usize, 2, 3, 4, 5];
19147+
for (slot, marker) in slots.iter().zip(markers) {
19148+
let output = results[*slot]
19149+
.result
19150+
.as_ref()
19151+
.unwrap_or_else(|error| panic!("read at slot {slot} failed: {error}"));
19152+
assert!(
19153+
output.contains(marker),
19154+
"slot {slot} must contain {marker}, got: {output}"
19155+
);
19156+
}
19157+
}
19158+
1900819159
#[test]
1900919160
fn init_template_mentions_detected_rust_workspace() {
1901019161
let _guard = cwd_lock()

0 commit comments

Comments
 (0)