use super::subprocess::{SpawnArgs, spawn_and_capture};
use super::{
ExecError, IN_PROCESS_SUBCOMMAND, INPUT_FILE, OUTPUT_FILE, ToolCall, ToolExecutor,
ToolInputRecord, ToolOutcome, ToolOutputRecord, atomic_write_json, tool_call_dir,
};
use crate::prompt::Clock;
use crate::prompt::step::STEPS_DIR;
use crate::workspace;
use std::ffi::{OsStr, OsString};
use std::os::unix::process::ExitStatusExt;
use std::path::{Path, PathBuf};
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>,
}
pub trait PathLookup {
fn which_on_path(&self, prefixed_name: &str) -> Option<PathBuf>;
}
pub struct EnvPath;
impl PathLookup for EnvPath {
fn which_on_path(&self, prefixed_name: &str) -> Option<PathBuf> {
which_in_path(prefixed_name)
}
}
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),
}
}
#[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
}
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,
) -> Result<ToolOutcome, ExecError> {
let caller = Caller::from_step_dir(step_dir).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 stdin = serde_json::to_vec(call.input).expect("Value is always serializable");
let extra_env = caller.env();
let binary_ref = &binary;
let req = SpawnArgs {
binary: binary_ref,
args: &args,
stdin_bytes: &stdin,
extra_env: &extra_env,
cwd: &caller.worktree,
stop,
deadline: self.deadline,
etxtbsy_budget: self.etxtbsy_budget,
tool_name: call.name,
};
let started_at = self.clock.now_iso8601();
let captured = spawn_and_capture(&req)?;
let ended_at = self.clock.now_iso8601();
let exit_code = match captured.status.code() {
Some(c) => c,
None => return Err(killed_by_signal(call.name, &captured.status)),
};
let mut content = captured.stdout.clone();
let is_error = exit_code != 0;
if is_error {
content.extend_from_slice(&captured.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,
ended_at,
};
atomic_write_json(&dir, OUTPUT_FILE, &output_record)?;
Ok(ToolOutcome { content, is_error })
}
}
struct Caller {
workspace: PathBuf,
agent_id: String,
worktree: PathBuf,
}
impl Caller {
fn from_step_dir(step_dir: &Path) -> Option<Self> {
let agent_dir = step_dir.parent()?;
let agent_id = agent_dir.file_name()?.to_str()?.to_string();
let workspace = agent_dir
.parent()
.filter(|p| p.ends_with(STEPS_DIR))?
.parent()?;
let worktree = workspace::agent_worktree(workspace, &agent_id);
worktree.is_dir().then(|| Self {
workspace: workspace.to_path_buf(),
agent_id,
worktree,
})
}
fn env(&self) -> Vec<(&'static str, OsString)> {
vec![
(super::ENV_CONV_BRANCH, OsString::from(&self.agent_id)),
(super::ENV_CONV_REPO, self.workspace.as_os_str().to_owned()),
]
}
}
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,
}
}
pub(super) fn which_in_path(name: &str) -> Option<PathBuf> {
which_in_path_env(name, std::env::var_os("PATH").as_deref())
}
pub(super) fn which_in_path_env(name: &str, path: Option<&OsStr>) -> Option<PathBuf> {
let path = path?;
for dir in std::env::split_paths(path) {
let candidate = dir.join(name);
if candidate.is_file() {
return Some(candidate);
}
}
None
}