use super::caller::Caller;
use super::{
Captured, ExecError, INPUT_FILE, OUTPUT_FILE, SpawnArgs, SpawnTool, ToolCall, ToolInputRecord,
ToolOutcome, ToolOutputRecord, atomic_write_json, bound, envelope, killed_by_signal,
spawn_and_capture, tool_call_dir,
};
use crate::config::ToolOutputBound;
use crate::prompt::tool::inject::{RoutedCall, RoutedCapture, ToolInjection};
use std::ffi::OsString;
use std::path::PathBuf;
use std::sync::atomic::AtomicBool;
use std::time::Duration;
pub(super) type Answered = Result<(Prepared, RoutedCapture), ExecError>;
pub(super) struct Prepared {
dir: PathBuf,
caller: Caller,
binary: OsString,
args: Vec<OsString>,
stdin: Vec<u8>,
extra_env: Vec<(&'static str, OsString)>,
name: String,
}
impl Prepared {
pub(super) fn spawn_args<'x>(
&'x self,
stop: &'x AtomicBool,
deadline: Duration,
etxtbsy_budget: u32,
) -> SpawnArgs<'x> {
SpawnArgs {
binary: &self.binary,
args: &self.args,
stdin_bytes: &self.stdin,
extra_env: &self.extra_env,
cwd: &self.caller.cwd,
stop,
deadline,
etxtbsy_budget,
tool_name: &self.name,
}
}
}
impl<'a> SpawnTool<'a> {
pub(super) fn prepare(
&self,
call: ToolCall<'_>,
step_dir: &std::path::Path,
) -> Result<Prepared, ExecError> {
let caller =
Caller::resolve(step_dir, &*self.git).ok_or_else(|| ExecError::NoWorktree {
name: call.name.to_string(),
step_dir: step_dir.to_path_buf(),
})?;
let dir = tool_call_dir(step_dir, call.id);
std::fs::create_dir_all(&dir).map_err(|source| ExecError::Io {
dir: dir.clone(),
source,
})?;
let input_record = ToolInputRecord {
id: call.id.to_string(),
name: call.name.to_string(),
input: call.input.clone(),
};
atomic_write_json(&dir, INPUT_FILE, &input_record)?;
let (binary, args) = self.resolve(call.name);
let extra_env = caller.env();
Ok(Prepared {
dir,
caller,
binary,
args,
stdin: serde_json::to_vec(call.input).expect("Value is always serializable"),
extra_env,
name: call.name.to_string(),
})
}
pub(super) fn spawn_one(
&self,
prepared: &Prepared,
stop: &AtomicBool,
) -> Result<RoutedCapture, ExecError> {
let captured =
spawn_and_capture(&prepared.spawn_args(stop, self.deadline, self.etxtbsy_budget))?;
classify(prepared, captured)
}
pub(super) fn spawn_fan(
&self,
prepared: Vec<Result<Prepared, ExecError>>,
stop: &AtomicBool,
) -> Vec<Answered> {
let (deadline, etxtbsy_budget) = (self.deadline, self.etxtbsy_budget);
let captured: Vec<Result<Captured, ExecError>> = std::thread::scope(|scope| {
let handles: Vec<_> = prepared
.iter()
.filter_map(|p| p.as_ref().ok())
.map(|p| {
scope.spawn(move || {
spawn_and_capture(&p.spawn_args(stop, deadline, etxtbsy_budget))
})
})
.collect();
handles
.into_iter()
.map(|h| h.join().unwrap_or_else(|p| std::panic::resume_unwind(p)))
.collect()
});
let mut captured = captured.into_iter();
prepared
.into_iter()
.map(|prepared| {
let prepared = prepared?;
let captured = captured.next().expect("one capture per spawned call")?;
let captured = classify(&prepared, captured)?;
Ok((prepared, captured))
})
.collect()
}
pub(super) fn route(
&self,
injection: &dyn ToolInjection,
prepared: &Prepared,
call: ToolCall<'_>,
stop: &AtomicBool,
) -> RoutedCapture {
injection.route(RoutedCall {
id: call.id,
name: call.name,
input: call.input,
workspace: &prepared.caller.workspace,
agent: &prepared.caller.agent_id,
cwd: &prepared.caller.cwd,
stop,
})
}
pub(super) fn route_fan(
&self,
prepared: Vec<Result<Prepared, ExecError>>,
calls: &[ToolCall<'_>],
injection: &dyn ToolInjection,
stop: &AtomicBool,
) -> Vec<Answered> {
prepared
.into_iter()
.zip(calls)
.map(|(prepared, call)| {
let prepared = prepared?;
let captured = self.route(injection, &prepared, *call, stop);
Ok((prepared, captured))
})
.collect()
}
pub(super) fn land(
&self,
prepared: &Prepared,
captured: &RoutedCapture,
output_bound: Option<ToolOutputBound>,
started_at: &str,
ended_at: &str,
) -> Result<ToolOutcome, ExecError> {
let exit_code = captured.exit_code;
let record = prepared.caller.record_rel(&prepared.dir).join(OUTPUT_FILE);
let stdout = bound::apply(&captured.stdout, "stdout", output_bound, &record);
let stderr = bound::apply(&captured.stderr, "stderr", output_bound, &record);
let content = envelope::render(exit_code, &stdout, &stderr);
let output_record = ToolOutputRecord {
stdout: String::from_utf8_lossy(&captured.stdout).into_owned(),
stderr: String::from_utf8_lossy(&captured.stderr).into_owned(),
exit_code,
started_at: started_at.to_string(),
ended_at: ended_at.to_string(),
};
atomic_write_json(&prepared.dir, OUTPUT_FILE, &output_record)?;
Ok(ToolOutcome {
content,
is_error: exit_code != 0,
})
}
}
fn classify(prepared: &Prepared, captured: Captured) -> Result<RoutedCapture, ExecError> {
match captured.status.code() {
Some(exit_code) => Ok(RoutedCapture {
stdout: captured.stdout,
stderr: captured.stderr,
exit_code,
}),
None => Err(killed_by_signal(&prepared.name, &captured.status)),
}
}