mod batch;
mod caller;
pub(super) mod lookup;
use batch::Prepared;
pub use lookup::{EnvPath, PathLookup};
use super::inject::{InjectedTool, RoutedCapture, ToolInjection};
use super::subprocess::{Captured, SpawnArgs, spawn_and_capture};
use super::{
ExecError, IN_PROCESS_SUBCOMMAND, INPUT_FILE, OUTPUT_FILE, ToolCall, ToolExecutor,
ToolInputRecord, ToolOutcome, ToolOutputRecord, atomic_write_json, bound, envelope,
tool_call_dir,
};
use crate::config::ToolOutputBound;
use crate::prompt::Clock;
use crate::template::{GitRunner, RealGit};
use std::ffi::OsString;
use std::os::unix::process::ExitStatusExt;
use std::path::Path;
use std::sync::atomic::AtomicBool;
use std::time::Duration;
pub struct SpawnTool<'a> {
data_root: &'a Path,
clock: &'a dyn Clock,
driver_target: &'a Path,
deadline: Duration,
etxtbsy_budget: u32,
path_lookup: Box<dyn PathLookup + 'a>,
git: Box<dyn GitRunner + 'a>,
injection: Option<&'a dyn ToolInjection>,
}
impl<'a> SpawnTool<'a> {
pub fn new(data_root: &'a Path, clock: &'a dyn Clock, driver_target: &'a Path) -> Self {
Self {
data_root,
clock,
driver_target,
deadline: super::DEFAULT_TOOL_DEADLINE,
etxtbsy_budget: super::subprocess::ETXTBSY_RETRY_ATTEMPTS,
path_lookup: Box::new(EnvPath),
git: Box::new(RealGit::new()),
injection: None,
}
}
pub fn with_injection(mut self, injection: Option<&'a dyn ToolInjection>) -> Self {
self.injection = injection;
self
}
#[cfg(test)] pub fn with_etxtbsy_budget(mut self, attempts: u32) -> Self {
self.etxtbsy_budget = attempts;
self
}
#[cfg(test)] pub fn with_deadline(mut self, d: Duration) -> Self {
self.deadline = d;
self
}
#[cfg(test)] pub fn with_path_lookup(mut self, l: Box<dyn PathLookup + 'a>) -> Self {
self.path_lookup = l;
self
}
#[cfg(test)] pub fn with_git(mut self, g: Box<dyn GitRunner + 'a>) -> Self {
self.git = g;
self
}
fn resolve(&self, name: &str) -> (OsString, Vec<OsString>) {
let external_name = format!("{}{}", super::EXTERNAL_PREFIX, name);
let harness_path = self.data_root.join(super::TOOLS_DIR).join(&external_name);
if harness_path.is_file() {
return (harness_path.into_os_string(), Vec::new());
}
if let Some(p) = self.path_lookup.which_on_path(&external_name) {
return (p.into_os_string(), Vec::new());
}
let args = vec![OsString::from(IN_PROCESS_SUBCOMMAND), OsString::from(name)];
(self.driver_target.as_os_str().to_owned(), args)
}
}
impl<'a> ToolExecutor for SpawnTool<'a> {
fn execute(
&self,
call: ToolCall<'_>,
step_dir: &Path,
stop: &AtomicBool,
output_bound: Option<ToolOutputBound>,
) -> Result<ToolOutcome, ExecError> {
let prepared = self.prepare(call, step_dir)?;
let started_at = self.clock.now_iso8601();
if let Some(routed) = self.route(&prepared, call, stop) {
let ended_at = self.clock.now_iso8601();
return self.land(&prepared, &routed, output_bound, &started_at, &ended_at);
}
let captured =
spawn_and_capture(&prepared.spawn_args(stop, self.deadline, self.etxtbsy_budget))?;
let ended_at = self.clock.now_iso8601();
self.finish(&prepared, captured, output_bound, &started_at, &ended_at)
}
fn injected(&self) -> Vec<InjectedTool> {
self.injection.map(ToolInjection::tools).unwrap_or_default()
}
fn execute_all(
&self,
calls: &[ToolCall<'_>],
step_dir: &Path,
stop: &AtomicBool,
output_bound: Option<ToolOutputBound>,
) -> Vec<Result<ToolOutcome, ExecError>> {
let prepared: Vec<Result<Prepared, ExecError>> = calls
.iter()
.map(|call| self.prepare(*call, step_dir))
.collect();
let started_at = self.clock.now_iso8601();
let routed: Vec<Option<RoutedCapture>> = prepared
.iter()
.zip(calls)
.map(|(p, call)| p.as_ref().ok().and_then(|p| self.route(p, *call, stop)))
.collect();
let (deadline, etxtbsy_budget) = (self.deadline, self.etxtbsy_budget);
let captured: Vec<Result<Captured, ExecError>> = std::thread::scope(|scope| {
let handles: Vec<_> = prepared
.iter()
.zip(&routed)
.filter_map(|(p, routed)| p.as_ref().ok().filter(|_| routed.is_none()))
.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 ended_at = self.clock.now_iso8601();
let mut captured = captured.into_iter();
prepared
.into_iter()
.zip(routed)
.map(|(prepared, routed)| {
let prepared = prepared?;
match routed {
Some(routed) => {
self.land(&prepared, &routed, output_bound, &started_at, &ended_at)
}
None => self.finish(
&prepared,
captured.next().expect("one capture per spawned call")?,
output_bound,
&started_at,
&ended_at,
),
}
})
.collect()
}
}
pub(super) fn killed_by_signal(name: &str, status: &std::process::ExitStatus) -> ExecError {
let signal = status.signal().unwrap_or(0);
ExecError::KilledBySignal {
name: name.to_string(),
signal,
}
}