use crate::loop_runtime::{DispatchCtx, Dispatcher, LoopError, Session, Unit};
use std::os::unix::process::ExitStatusExt;
use std::path::{Path, PathBuf};
use std::process::{Child, Command};
pub fn preflight(
driver: &str,
lib_dir: &Path,
cli_alias: Option<&str>,
) -> Result<PathBuf, LoopError> {
const ALLOWED: &[&str] = &["claude-code", "hermes", "openclaw", "opencode"];
if !ALLOWED.contains(&driver) {
return Err(LoopError::Config(format!(
"invalid dispatcher '{driver}': must be one of {:?} (whitelist)",
ALLOWED
)));
}
let lib_path = lib_dir.join(format!("driver-{driver}.sh"));
if !lib_path.exists() {
return Err(LoopError::Config(format!(
"driver lib not found: {}",
lib_path.display()
)));
}
let binary = resolve_driver_binary(driver, cli_alias);
if which_binary(&binary).is_none() {
return Err(LoopError::Dispatch(format!(
"missing binary '{binary}': required by dispatcher '{driver}' but not found on PATH"
)));
}
{
let lib_str = lib_path.to_str().ok_or_else(|| {
LoopError::Config(format!(
"driver lib path is not valid UTF-8: {}",
lib_path.display()
))
})?;
let probe_script = r#"source "$1" && type driver_invoke >/dev/null 2>&1"#;
let probe = std::process::Command::new("bash")
.arg("-c")
.arg(probe_script)
.arg("_")
.arg(lib_str)
.output()
.map_err(|e| LoopError::Config(format!("driver_invoke probe bash failed: {e}")))?;
if !probe.status.success() {
return Err(LoopError::Config(format!(
"driver lib '{}' does not define driver_invoke (required function missing)",
lib_path.display()
)));
}
}
Ok(lib_path)
}
pub fn driver_default_max(lib: &Path) -> Result<u64, LoopError> {
let lib_str = lib.to_str().ok_or_else(|| {
LoopError::Config(format!(
"driver lib path is not valid UTF-8: {}",
lib.display()
))
})?;
let script = format!("source {:?} && driver_default_max", lib_str);
let out = Command::new("bash")
.arg("-c")
.arg(&script)
.output()
.map_err(|e| LoopError::Dispatch(format!("bash shellout for driver_default_max: {e}")))?;
let raw = String::from_utf8_lossy(&out.stdout).trim().to_string();
raw.parse::<u64>().map_err(|_| {
LoopError::Dispatch(format!(
"driver_default_max returned non-integer stdout: {:?}",
raw
))
})
}
pub(crate) fn fno_cmd(fno_bin: &str) -> Command {
let binary = std::env::var("FNO_BIN")
.ok()
.filter(|s| !s.is_empty())
.unwrap_or_else(|| fno_bin.to_string());
Command::new(binary)
}
pub(crate) fn retry_etxtbsy<T>(
mut spawn: impl FnMut() -> std::io::Result<T>,
) -> std::io::Result<T> {
const MAX_RETRIES: u32 = 5;
let mut attempt: u32 = 0;
loop {
match spawn() {
Err(e) if e.raw_os_error() == Some(libc::ETXTBSY) && attempt < MAX_RETRIES => {
attempt += 1;
std::thread::sleep(std::time::Duration::from_millis(2 * u64::from(attempt)));
}
other => return other,
}
}
}
pub fn resolve_driver_binary(driver: &str, cli_alias: Option<&str>) -> String {
match driver {
"claude-code" => {
if let Ok(v) = std::env::var("CLAUDE_CLI") {
if !v.is_empty() {
return v;
}
}
if let Some(a) = cli_alias {
if !a.is_empty() {
return a.to_string();
}
}
if let Ok(v) = std::env::var("CLI") {
if !v.is_empty() {
return v;
}
}
"claude".to_string()
}
"hermes" => "hermes-agent".to_string(),
"openclaw" => "openclaw".to_string(),
"opencode" => "opencode".to_string(),
_ => "claude".to_string(), }
}
pub fn which_binary(name: &str) -> Option<PathBuf> {
if name.contains('/') {
let p = PathBuf::from(name);
if p.is_file() {
return Some(p);
}
return None;
}
let path_var = std::env::var("PATH").unwrap_or_default();
for dir in path_var.split(':') {
if dir.is_empty() {
continue;
}
let candidate = PathBuf::from(dir).join(name);
if candidate.is_file() {
use std::os::unix::fs::PermissionsExt;
if let Ok(meta) = std::fs::metadata(&candidate) {
if meta.permissions().mode() & 0o111 != 0 {
return Some(candidate);
}
}
}
}
None
}
const PICKED_ENV_KEY: &str = "CLAUDE_CONFIG_DIR";
pub const PICK_ARGV: [&str; 5] = ["config", "accounts", "pick", "--if-armed", "--print-env"];
type PickedEnv = Vec<(String, String)>;
fn interpret_pick(ok: bool, stdout: &str, stderr: &str) -> Result<PickedEnv, String> {
if !ok {
let reason = stderr
.lines()
.map(str::trim)
.filter(|l| !l.is_empty())
.next_back()
.unwrap_or("no reason given");
return Err(reason.to_string());
}
let mut env: PickedEnv = Vec::new();
let mut pinned = false;
for line in stdout.lines().map(str::trim).filter(|l| !l.is_empty()) {
match line.split_once('=') {
Some((k, v)) if !k.is_empty() => {
if !v.is_empty() {
pinned = true;
}
env.push((k.to_string(), v.to_string()));
}
_ => return Err(format!("unparseable pick output: {line:?}")),
}
}
if !pinned {
return Err("pick output carried no account pin".to_string());
}
Ok(env)
}
fn drives_claude(driver_lib: &Path) -> bool {
driver_lib
.file_name()
.and_then(|n| n.to_str())
.is_some_and(|n| n == "driver-claude-code.sh")
}
fn pick_would_undo_a_route(picked: &[(String, String)], static_env: &[(String, String)]) -> bool {
picked.iter().filter(|(_, v)| v.is_empty()).any(|(k, _)| {
static_env.iter().any(|(ek, _)| ek == k)
|| std::env::var_os(k).is_some_and(|v| !v.is_empty())
})
}
fn pick_account_env() -> Result<PickedEnv, String> {
let out = Command::new("fno")
.args(PICK_ARGV)
.output()
.map_err(|e| format!("could not run `fno config accounts pick`: {e}"))?;
interpret_pick(
out.status.success(),
&String::from_utf8_lossy(&out.stdout),
&String::from_utf8_lossy(&out.stderr),
)
}
pub struct ShelloutSession {
child: Child,
output_file: Option<PathBuf>,
}
impl Session for ShelloutSession {
fn wait(&mut self) -> Result<i32, LoopError> {
let status = self.child.wait().map_err(LoopError::Io)?;
Ok(status
.code()
.unwrap_or_else(|| 128 + status.signal().unwrap_or(0)))
}
fn output_tail(&self) -> Option<String> {
use std::io::{Read, Seek, SeekFrom};
const MAX_TAIL: u64 = 8 * 1024;
let path = self.output_file.as_ref()?;
let mut file = std::fs::File::open(path).ok()?;
if let Ok(len) = file.metadata().map(|m| m.len()) {
let _ = file.seek(SeekFrom::Start(len.saturating_sub(MAX_TAIL)));
}
let mut buf = Vec::new();
file.take(MAX_TAIL).read_to_end(&mut buf).ok()?;
Some(String::from_utf8_lossy(&buf).into_owned())
}
}
pub struct ShelloutDispatcher {
driver_lib: PathBuf,
env: Vec<(String, String)>,
cwd: PathBuf,
}
impl ShelloutDispatcher {
pub fn new(driver_lib: PathBuf, env: Vec<(String, String)>, cwd: PathBuf) -> Self {
Self {
driver_lib,
env,
cwd,
}
}
}
impl Dispatcher for ShelloutDispatcher {
fn run(&self, _unit: &Unit, ctx: &DispatchCtx) -> Result<Box<dyn Session>, LoopError> {
let lib_str = self
.driver_lib
.to_str()
.ok_or_else(|| LoopError::Dispatch("driver lib path is not valid UTF-8".to_string()))?;
let script = r#"source "$FNO_DRIVER_LIB" && (driver_invoke); rc=$?; driver_persist_history >/dev/null 2>&1 || true; exit $rc"#;
let mut cmd = Command::new("bash");
cmd.arg("-c").arg(script);
cmd.env("FNO_DRIVER_LIB", lib_str);
cmd.env("CURRENT_ITER", ctx.iteration.to_string());
cmd.current_dir(&self.cwd);
for (k, v) in &self.env {
cmd.env(k, v);
}
if drives_claude(&self.driver_lib) && !self.env.iter().any(|(k, _)| k == PICKED_ENV_KEY) {
let iter = ctx.iteration;
match pick_account_env() {
Ok(picked) if pick_would_undo_a_route(&picked, &self.env) => {
eprintln!(
"loop: iteration {iter} account not picked \
(this run pins its own provider route)"
);
}
Ok(picked) => {
if let Some((key, value)) = picked.iter().find(|(_, v)| !v.is_empty()) {
if key == PICKED_ENV_KEY {
eprintln!("loop: iteration {iter} account picked -> {value}");
} else {
eprintln!("loop: iteration {iter} account picked -> pinned via {key}");
}
}
for (key, value) in &picked {
if value.is_empty() {
cmd.env_remove(key);
} else {
cmd.env(key, value);
}
}
}
Err(reason) => {
eprintln!("loop: iteration {iter} account not picked ({reason})");
}
}
}
let child = cmd
.spawn()
.map_err(|e| LoopError::Dispatch(format!("spawn bash driver_invoke: {e}")))?;
let output_file = self
.env
.iter()
.find(|(k, _)| k == "OUTPUT_FILE")
.map(|(_, v)| PathBuf::from(v));
Ok(Box::new(ShelloutSession { child, output_file }))
}
}
#[cfg(test)]
mod tests {
use super::{
interpret_pick, pick_would_undo_a_route, resolve_driver_binary, retry_etxtbsy,
PICKED_ENV_KEY,
};
fn pair(k: &str, v: &str) -> (String, String) {
(k.to_string(), v.to_string())
}
#[test]
fn retry_etxtbsy_passes_success_through_without_retry() {
let mut calls = 0u32;
let r: std::io::Result<u8> = retry_etxtbsy(|| {
calls += 1;
Ok(7)
});
assert_eq!(r.unwrap(), 7);
assert_eq!(calls, 1, "a successful spawn must not retry");
}
#[test]
fn retry_etxtbsy_retries_then_succeeds() {
let mut calls = 0u32;
let r: std::io::Result<u8> = retry_etxtbsy(|| {
calls += 1;
if calls < 3 {
Err(std::io::Error::from_raw_os_error(libc::ETXTBSY))
} else {
Ok(42)
}
});
assert_eq!(r.unwrap(), 42);
assert_eq!(calls, 3, "must retry past transient ETXTBSY");
}
#[test]
fn retry_etxtbsy_does_not_swallow_other_errors() {
let mut calls = 0u32;
let r: std::io::Result<u8> = retry_etxtbsy(|| {
calls += 1;
Err(std::io::Error::from_raw_os_error(libc::ENOENT))
});
assert_eq!(r.unwrap_err().raw_os_error(), Some(libc::ENOENT));
assert_eq!(calls, 1, "a non-ETXTBSY error must not retry");
}
#[test]
fn retry_etxtbsy_gives_up_after_max_retries() {
let mut calls = 0u32;
let r: std::io::Result<u8> = retry_etxtbsy(|| {
calls += 1;
Err(std::io::Error::from_raw_os_error(libc::ETXTBSY))
});
assert_eq!(r.unwrap_err().raw_os_error(), Some(libc::ETXTBSY));
assert_eq!(calls, 6, "1 initial attempt + MAX_RETRIES(5)");
}
#[test]
fn a_pick_that_would_scrub_a_pinned_route_is_declined() {
let picked = vec![
pair("ANTHROPIC_BASE_URL", ""),
pair("ANTHROPIC_AUTH_TOKEN", ""),
pair("CLAUDE_CONFIG_DIR", "/alt"),
];
let routed = vec![
pair(
"ANTHROPIC_BASE_URL",
"https://open.bigmodel.cn/api/anthropic",
),
pair("OUTPUT_FILE", "/tmp/out"),
];
assert!(pick_would_undo_a_route(&picked, &routed));
}
#[test]
fn an_unrouted_loop_still_gets_its_pick() {
let picked = vec![
pair("ANTHROPIC_BASE_URL", ""),
pair("CLAUDE_CONFIG_DIR", "/alt"),
];
let plain = vec![pair("OUTPUT_FILE", "/tmp/out"), pair("CLI", "claude")];
assert!(!pick_would_undo_a_route(&picked, &plain));
}
#[test]
fn only_the_clear_list_blocks_a_pick_not_the_pin_itself() {
let picked = vec![pair("CLAUDE_CONFIG_DIR", "/alt")];
let same_key = vec![pair("CLAUDE_CONFIG_DIR", "/other")];
assert!(!pick_would_undo_a_route(&picked, &same_key));
}
#[test]
fn a_picked_account_yields_its_config_dir() {
let env = interpret_pick(true, "CLAUDE_CONFIG_DIR=/Users/x/.claude-alt\n", "")
.expect("a pinned config dir is a successful pick");
assert_eq!(
env,
vec![(
"CLAUDE_CONFIG_DIR".to_string(),
"/Users/x/.claude-alt".to_string()
)]
);
}
#[test]
fn auth_vars_to_clear_are_carried_as_empty_values() {
let stdout = "ANTHROPIC_API_KEY=\nANTHROPIC_BASE_URL=\nCLAUDE_CONFIG_DIR=/alt\n";
let env = interpret_pick(true, stdout, "").expect("overlay parses");
assert_eq!(env.len(), 3);
assert_eq!(env[0], ("ANTHROPIC_API_KEY".to_string(), String::new()));
assert_eq!(env[1], ("ANTHROPIC_BASE_URL".to_string(), String::new()));
assert_eq!(
env[2],
("CLAUDE_CONFIG_DIR".to_string(), "/alt".to_string())
);
}
#[test]
fn a_disarmed_picker_declines_with_its_reason() {
let stderr = "pick: launch picking is not armed (providers.quota.pick_on_launch = false)\n";
assert_eq!(
interpret_pick(false, "", stderr),
Err(
"pick: launch picking is not armed (providers.quota.pick_on_launch = false)"
.to_string()
)
);
}
#[test]
fn a_refusal_surfaces_its_reason_instead_of_erroring_out() {
let stderr = " readyrule: exhausted\npick: every launchable candidate is exhausted\n";
assert_eq!(
interpret_pick(false, "", stderr),
Err("pick: every launchable candidate is exhausted".to_string())
);
assert!(interpret_pick(false, "", "").is_err());
}
#[test]
fn only_the_config_dir_pin_is_safe_to_echo() {
assert_eq!(PICKED_ENV_KEY, "CLAUDE_CONFIG_DIR");
for secret_key in ["ANTHROPIC_API_KEY", "CLAUDE_CODE_OAUTH_TOKEN"] {
assert_ne!(
secret_key, PICKED_ENV_KEY,
"a secret-bearing pin must not take the echo-the-value branch"
);
}
}
#[test]
fn unparseable_output_is_declined_rather_than_guessed() {
assert!(interpret_pick(true, "readyrule\n", "").is_err());
assert!(interpret_pick(true, "CLAUDE_CONFIG_DIR=\n", "").is_err());
assert!(interpret_pick(true, "ANTHROPIC_API_KEY=\n", "").is_err());
assert!(interpret_pick(true, "=/tmp\n", "").is_err());
assert!(interpret_pick(true, "", "").is_err());
}
#[test]
fn no_driver_lib_calls_the_picker_itself() {
let lib_dir = std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("../../scripts/lib");
let mut checked = 0;
for entry in std::fs::read_dir(&lib_dir).expect("scripts/lib is readable") {
let path = entry.expect("dir entry").path();
let name = path
.file_name()
.and_then(|n| n.to_str())
.unwrap_or("")
.to_string();
if !name.starts_with("driver-") || !name.ends_with(".sh") {
continue;
}
let body = std::fs::read_to_string(&path).expect("driver lib is readable");
assert!(
!body.contains("providers pick"),
"{name} calls the picker itself; the loop dispatcher is the one call site"
);
checked += 1;
}
assert!(checked >= 4, "expected the driver libs, scanned {checked}");
}
#[test]
fn loop_wrapper_drivers_resolve_to_fixed_binaries() {
assert_eq!(resolve_driver_binary("opencode", None), "opencode");
assert_eq!(resolve_driver_binary("openclaw", None), "openclaw");
assert_eq!(resolve_driver_binary("hermes", None), "hermes-agent");
}
}