use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicUsize, Ordering};
use car_fleet::{DeclineReason, DispatchOutcome, SubtaskDispatch, WorkGuard};
static IN_FLIGHT: AtomicUsize = AtomicUsize::new(0);
fn work_guard() -> &'static std::sync::Mutex<WorkGuard> {
static GUARD: std::sync::OnceLock<std::sync::Mutex<WorkGuard>> = std::sync::OnceLock::new();
GUARD.get_or_init(|| std::sync::Mutex::new(WorkGuard::default()))
}
const SAME_ACCOUNT_FLEET_BUDGET: &str = "same-account-fleet";
fn budget_key(_caller: &str) -> &'static str {
SAME_ACCOUNT_FLEET_BUDGET
}
struct InFlightGuard;
impl Drop for InFlightGuard {
fn drop(&mut self) {
IN_FLIGHT.fetch_sub(1, Ordering::SeqCst);
}
}
struct WorkerWorktree {
repo: PathBuf,
path: PathBuf,
}
impl WorkerWorktree {
fn provision(repo: &Path, base: &Path, name: &str, commit: &str) -> Result<Self, String> {
let path = base.join(name);
let _ = git(
repo,
&["worktree", "remove", "--force", &path.to_string_lossy()],
);
let _ = git(repo, &["worktree", "prune"]);
if path.exists() {
let _ = std::fs::remove_dir_all(&path);
}
std::fs::create_dir_all(base).map_err(|e| format!("create {}: {e}", base.display()))?;
git(
repo,
&[
"worktree",
"add",
"--detach",
&path.to_string_lossy(),
commit,
],
)?;
Ok(Self {
repo: repo.to_path_buf(),
path,
})
}
}
impl Drop for WorkerWorktree {
fn drop(&mut self) {
let _ = git(
&self.repo,
&[
"worktree",
"remove",
"--force",
&self.path.to_string_lossy(),
],
);
if self.path.exists() {
let _ = std::fs::remove_dir_all(&self.path);
}
}
}
fn git(cwd: &Path, args: &[&str]) -> Result<String, String> {
let out = std::process::Command::new("git")
.arg("-C")
.arg(cwd)
.args(args)
.output()
.map_err(|e| format!("run git {}: {e}", args.join(" ")))?;
if !out.status.success() {
return Err(format!(
"git {} failed: {}",
args.join(" "),
String::from_utf8_lossy(&out.stderr).trim()
));
}
Ok(String::from_utf8_lossy(&out.stdout).trim().to_string())
}
pub async fn run_dispatch(
dispatch: SubtaskDispatch,
caller: &str,
) -> Result<DispatchOutcome, String> {
let started = std::time::Instant::now();
let config = super::FleetWorkerConfig::load();
if !config.accepts_work {
return Ok(declined(
caller,
&dispatch,
DeclineReason::NotAcceptingWork,
"this instance is not enrolled as a fleet worker (`fleet.worker.set`)",
));
}
let admitted = IN_FLIGHT
.fetch_update(Ordering::SeqCst, Ordering::SeqCst, |n| {
(n < config.max_parallel as usize).then_some(n + 1)
})
.is_ok();
if !admitted {
return Ok(declined(
caller,
&dispatch,
DeclineReason::Busy,
&format!("already running {} subtasks", config.max_parallel),
));
}
let _guard = InFlightGuard;
let repo = match car_fleet::resolve_repo(&dispatch.repo, &config.repos) {
Ok(path) => path.to_path_buf(),
Err(DeclineReason::CommitUnavailable) if config.fetch_missing_base => {
let clones: Vec<std::path::PathBuf> =
car_fleet::clones_of(&dispatch.repo, &config.repos)
.into_iter()
.map(|p| p.to_path_buf())
.collect();
let remote = config.fetch_remote.clone();
let wanted = dispatch.repo.head_commit.clone();
let fetched = tokio::task::spawn_blocking(move || {
clones
.into_iter()
.find(|path| car_fleet::fetch_base(path, &remote, &wanted))
})
.await
.map_err(|e| format!("fetch task panicked: {e}"))?;
match fetched {
Some(path) => path,
None => {
return Ok(declined(
caller,
&dispatch,
DeclineReason::CommitUnavailable,
"the base commit is not on this worker's remote either — push the \
branch before distributing",
))
}
}
}
Err(reason) => {
return Ok(declined(
caller,
&dispatch,
reason,
match reason {
DeclineReason::CommitUnavailable => {
"the repository is here but the base commit is not — push it somewhere \
this worker can fetch from, or run the subtask locally"
}
_ => "no configured checkout of that repository",
},
))
}
};
let adapters = super::detected_adapters().await;
let adapter = match dispatch.adapter.as_deref() {
Some(wanted) => adapters
.iter()
.find(|a| a.id == wanted)
.map(|a| a.id.clone()),
None => adapters.first().map(|a| a.id.clone()),
};
let Some(adapter) = adapter else {
return Ok(declined(
caller,
&dispatch,
DeclineReason::AdapterUnavailable,
"no installed coding CLI could run this subtask",
));
};
let verdict = {
let now = car_fleet::now_ms();
let mut guard = work_guard().lock().unwrap_or_else(|e| e.into_inner());
guard.set_limit(config.dispatches_per_hour);
guard.evict_idle(now);
guard.admit(budget_key(caller), now)
};
if !verdict.is_accept() {
return Ok(declined(
caller,
&dispatch,
DeclineReason::RateLimited,
&verdict.reason(),
));
}
let base = super::worktree_base();
let name = sanitize(&format!("{}-{}", dispatch.run_id, dispatch.subtask_id));
let worktree = match WorkerWorktree::provision(&repo, &base, &name, &dispatch.repo.head_commit)
{
Ok(w) => w,
Err(e) => return Err(format!("provision worktree: {e}")),
};
let subtask = car_multi::Subtask::files_only(
dispatch.subtask_id.clone(),
dispatch.prompt.clone(),
dispatch.files.clone(),
);
let request = car_multi::WorktreeAgentRequest {
subtask: &subtask,
cwd: &worktree.path,
allowed_tools: narrow_tools(
dispatch.allowed_tools.clone(),
config.allowed_tools.as_deref(),
),
mcp_endpoint: None,
mcp_config_dir: None,
};
let mut agent = car_external_agents::ForemanExternalAgent::new(adapter.clone());
agent.timeout_secs = Some(
dispatch
.timeout_secs
.unwrap_or(config.max_subtask_secs)
.min(config.max_subtask_secs),
);
let answer = match car_multi::WorktreeAgent::run_in(&agent, &request).await {
Ok(summary) => summary.answer,
Err(e) => {
audit(caller, &dispatch, "error", &e.to_string());
return Err(format!("subtask failed: {e}"));
}
};
let patch = {
let cwd = worktree.path.clone();
match tokio::task::spawn_blocking(move || car_multi::capture_patch(&cwd)).await {
Ok(Ok(p)) => p,
Ok(Err(e)) => {
audit(caller, &dispatch, "error", &e.to_string());
return Err(format!("capture patch: {e}"));
}
Err(e) => return Err(format!("capture task panicked: {e}")),
}
};
audit(
caller,
&dispatch,
"completed",
&format!("{} bytes via {adapter}", patch.len()),
);
Ok(DispatchOutcome::Completed {
patch,
answer,
adapter,
duration_ms: started.elapsed().as_millis() as u64,
})
}
fn narrow_tools(requested: Option<Vec<String>>, worker: Option<&[String]>) -> Option<Vec<String>> {
match (requested, worker) {
(_, None) => None,
(None, Some(worker)) => Some(worker.to_vec()),
(Some(requested), Some(worker)) => Some(
requested
.into_iter()
.filter(|t| worker.iter().any(|w| w == t))
.collect(),
),
}
}
fn declined(
caller: &str,
dispatch: &SubtaskDispatch,
reason: DeclineReason,
detail: &str,
) -> DispatchOutcome {
audit(caller, dispatch, reason.as_str(), detail);
DispatchOutcome::Declined {
reason,
detail: detail.to_string(),
}
}
fn audit(caller: &str, dispatch: &SubtaskDispatch, outcome: &str, detail: &str) {
use std::io::Write;
let Some(dir) = car_home::root() else {
return;
};
if std::fs::create_dir_all(&dir).is_err() {
return;
}
let record = serde_json::json!({
"ts": chrono::Utc::now().to_rfc3339(),
"peer": caller,
"run_id": dispatch.run_id,
"subtask_id": dispatch.subtask_id,
"repo_root_commit": dispatch.repo.root_commit,
"base_commit": dispatch.repo.head_commit,
"outcome": outcome,
"detail": detail,
});
let Ok(line) = serde_json::to_string(&record) else {
return;
};
if let Ok(mut f) = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(dir.join("fleet-work.jsonl"))
{
let _ = writeln!(f, "{line}");
}
}
fn sanitize(name: &str) -> String {
let s: String = name
.chars()
.map(|c| {
if c.is_ascii_alphanumeric() || c == '-' || c == '_' {
c
} else {
'-'
}
})
.collect();
let trimmed = s.trim_matches('-').to_string();
if trimmed.is_empty() {
"subtask".to_string()
} else {
trimmed.chars().take(96).collect()
}
}
#[cfg(test)]
mod tests {
use super::*;
use car_fleet::RepoFingerprint;
fn dispatch() -> SubtaskDispatch {
SubtaskDispatch::new(
"run",
"sub",
"do the thing",
RepoFingerprint {
root_commit: "r".into(),
head_commit: "h".into(),
name: None,
},
)
}
#[test]
fn a_peers_ids_cannot_escape_the_worktree_base() {
assert_eq!(sanitize("../../etc/passwd"), "etc-passwd");
assert_eq!(sanitize("run/../..\\x"), "run-------x");
for hostile in ["../../etc/passwd", "run/../..\\x", "a/b/c", "..\\..\\win"] {
let safe = sanitize(hostile);
assert!(!safe.contains('/') && !safe.contains('\\'), "{safe}");
assert!(!safe.contains(".."), "{safe}");
}
assert_eq!(sanitize(""), "subtask");
assert_eq!(sanitize("---"), "subtask");
assert!(!sanitize(&"x".repeat(500)).contains('/'));
assert_eq!(sanitize(&"x".repeat(500)).len(), 96);
}
#[test]
fn every_trusted_caller_spends_one_budget_not_one_per_machine() {
assert_eq!(budget_key("aa:bb:cc:dd"), budget_key("11:22:33:44"));
}
#[test]
fn the_stricter_tool_list_wins_in_both_directions() {
assert_eq!(
narrow_tools(
Some(vec!["Read".into(), "Bash".into()]),
Some(&["Read".to_string()])
),
Some(vec!["Read".into()])
);
assert_eq!(
narrow_tools(None, Some(&["Read".to_string()])),
Some(vec!["Read".into()])
);
assert_eq!(narrow_tools(Some(vec!["Bash".into()]), None), None);
assert_eq!(
narrow_tools(Some(vec!["Bash".into()]), Some(&["Read".to_string()])),
Some(Vec::new())
);
}
#[test]
fn a_senders_timeout_cannot_exceed_this_machines_ceiling() {
let ceiling = 1800u64;
let clamp = |requested: Option<u64>| requested.unwrap_or(ceiling).min(ceiling);
assert_eq!(clamp(Some(60)), 60, "a shorter request is honoured");
assert_eq!(clamp(Some(99_999)), ceiling, "a longer one is clamped");
assert_eq!(
clamp(None),
ceiling,
"omitted is the ceiling, not unbounded"
);
}
#[tokio::test]
async fn a_daemon_that_did_not_opt_in_declines_before_touching_anything() {
if super::super::FleetWorkerConfig::load().accepts_work {
return;
}
let outcome = run_dispatch(dispatch(), "aa:bb:cc:dd")
.await
.expect("no internal error");
assert!(matches!(
outcome,
DispatchOutcome::Declined {
reason: DeclineReason::NotAcceptingWork,
..
}
));
}
}