use std::path::Path;
use std::sync::{Arc, Mutex};
use std::time::Duration;
use tokio::process::Command;
use tracing::debug;
use crate::chrome::{CliRecovery, RELAY_RECOVERY_BUDGET};
use crate::util::UnwrapPoison;
#[derive(Debug, Clone, Copy)]
pub(crate) enum CliTimeout {
Unbounded,
Bounded(Duration),
}
pub(crate) struct CliSpawn<'a> {
pub(crate) path: &'a Path,
pub(crate) args: &'a [&'a str],
pub(crate) session: Option<&'a str>,
pub(crate) json: bool,
pub(crate) capture_stderr: bool,
pub(crate) cancel_kills: bool,
pub(crate) input: Option<Vec<u8>>,
pub(crate) timeout: CliTimeout,
pub(crate) chrome_side: Duration,
pub(crate) recovery: CliRecovery,
}
pub(crate) struct CliOutput {
pub(crate) status: std::process::ExitStatus,
pub(crate) stdout: Vec<u8>,
pub(crate) stderr: Vec<u8>,
pub(crate) leftover_pipes: bool,
pub(crate) truncated: bool,
}
pub(crate) enum CliRun {
Output(CliOutput),
SpawnFailure,
TimedOut,
}
pub(crate) fn ensure_chrome_env(cmd: &mut Command) {
clear_browser_env(cmd);
if std::env::var_os("HOME").is_none() {
cmd.env("HOME", "/tmp");
}
if std::env::var_os("CHROMIUM_FLAGS").is_none() {
cmd.env(
"CHROMIUM_FLAGS",
"--no-first-run --no-default-browser-check --disable-gpu",
);
}
cmd.env("AGENT_BROWSER_IDLE_TIMEOUT_MS", "300000");
cmd.env("AGENT_BROWSER_HUMANIZE", "human");
cmd.env("CHROME_USE_NO_UPDATE_CHECK", "1");
cmd.env("AGENT_BROWSER_NO_UPDATE_CHECK", "1");
pin_real_browser_env(cmd);
}
fn clear_browser_env(cmd: &mut Command) {
for (name, _) in std::env::vars_os() {
if let Some(name) = name.to_str() {
let upper = name.to_ascii_uppercase();
if (upper.starts_with("AGENT_BROWSER_") || upper.starts_with("CHROME_USE_"))
&& !matches!(
upper.as_str(),
"AGENT_BROWSER_RELAY_DIR" | "CHROME_USE_RELAY_DIR"
)
{
cmd.env_remove(name);
}
}
}
}
fn pin_real_browser_env(cmd: &mut Command) {
cmd.env_remove("CI");
cmd.env("AGENT_BROWSER_AUTO_CONNECT", "1");
}
fn apply_recovery_env(cmd: &mut Command, recovery: CliRecovery) {
match recovery {
CliRecovery::Allowed => {
cmd.env(
"AGENT_BROWSER_RELAY_REVIVE_SECS",
RELAY_RECOVERY_BUDGET.as_secs().to_string(),
);
}
CliRecovery::Suppressed => {
cmd.env("AGENT_BROWSER_NO_AUTO_RECONNECT", "1");
cmd.env("AGENT_BROWSER_RELAY_REVIVE_SECS", "0");
}
}
}
pub(crate) fn apply_chrome_side(cmd: &mut Command, chrome_side: Duration) {
cmd.env(
"AGENT_BROWSER_DEFAULT_TIMEOUT",
chrome_side.as_millis().to_string(),
);
}
fn build_argv(args: &[&str], json: bool, session: Option<&str>) -> Vec<String> {
let mut global: Vec<String> = Vec::new();
if json {
global.push("--json".to_string());
}
if let Some(s) = session {
global.extend(["--session".to_string(), s.to_string()]);
}
if global.is_empty() {
return args.iter().map(|a| (*a).to_string()).collect();
}
let insert_at = args.iter().position(|a| *a == "--").unwrap_or(args.len());
let mut folded: Vec<String> = args[..insert_at].iter().map(|a| (*a).to_string()).collect();
folded.extend(global);
folded.extend(args[insert_at..].iter().map(|a| (*a).to_string()));
folded
}
pub(crate) async fn spawn_cli(spec: CliSpawn<'_>) -> CliRun {
let mut cmd = Command::new(spec.path);
#[cfg(windows)]
cmd.creation_flags(windows_sys::Win32::System::Threading::CREATE_NO_WINDOW);
ensure_chrome_env(&mut cmd);
apply_chrome_side(&mut cmd, spec.chrome_side);
apply_recovery_env(&mut cmd, spec.recovery);
cmd.args(build_argv(spec.args, spec.json, spec.session));
if spec.json {
cmd.stdout(std::process::Stdio::piped());
} else {
cmd.stdout(std::process::Stdio::null());
}
if spec.capture_stderr {
cmd.stderr(std::process::Stdio::piped());
} else {
cmd.stderr(std::process::Stdio::null());
}
if spec.input.is_some() {
cmd.stdin(std::process::Stdio::piped());
} else {
cmd.stdin(std::process::Stdio::null());
}
cmd.kill_on_drop(spec.cancel_kills);
run_child(&mut cmd, spec.input, spec.timeout).await
}
async fn run_child(cmd: &mut Command, input: Option<Vec<u8>>, timeout: CliTimeout) -> CliRun {
let Ok(mut child) = cmd.spawn() else {
return CliRun::SpawnFailure;
};
if let (Some(data), Some(mut stdin)) = (input, child.stdin.take()) {
tokio::task::spawn(async move {
use tokio::io::AsyncWriteExt;
let _ = stdin.write_all(&data).await;
let _ = stdin.shutdown().await;
});
}
let stdout = child.stdout.take().map(PipeDrain::start);
let stderr = child.stderr.take().map(PipeDrain::start);
if let Some(drain) = stdout.as_ref() {
drain.started().await;
}
if let Some(drain) = stderr.as_ref() {
drain.started().await;
}
let status = match timeout {
CliTimeout::Unbounded => match child.wait().await {
Ok(s) => s,
Err(_) => return CliRun::SpawnFailure,
},
CliTimeout::Bounded(deadline) => match tokio::time::timeout(deadline, child.wait()).await {
Ok(Ok(status)) => status,
Ok(Err(_)) => return CliRun::SpawnFailure,
Err(_) => return CliRun::TimedOut,
},
};
let (stdout_held, stderr_held) = tokio::join!(
held_after_exit(stdout.as_ref()),
held_after_exit(stderr.as_ref())
);
let leftover_pipes = stdout_held || stderr_held;
let truncated = [stdout.as_ref(), stderr.as_ref()]
.into_iter()
.flatten()
.any(PipeDrain::truncated);
let output = CliOutput {
status,
stdout: stdout.map_or_else(Vec::new, |drain| drain.buf.take()),
stderr: stderr.map_or_else(Vec::new, |drain| drain.buf.take()),
leftover_pipes,
truncated,
};
if output.leftover_pipes {
debug!(
"chrome-use exited but a process it left behind still holds its output channel — \
the exit status and the bytes drained so far are this call's result"
);
}
CliRun::Output(output)
}
async fn held_after_exit(drain: Option<&PipeDrain>) -> bool {
match drain {
Some(drain) => drain.held_after_exit().await,
None => false,
}
}
const PIPE_DRAIN_GRACE: Duration = Duration::from_secs(2);
const PIPE_DRAIN_MAX: Duration = Duration::from_secs(10);
const PIPE_DRAIN_CHUNK_BYTES: usize = 8 * 1024;
pub(crate) const PIPE_READ_CAP: usize = 64 * 1024 * 1024;
struct PipeDrain {
buf: Arc<PipeBuf>,
started: Arc<tokio::sync::Notify>,
finished: Arc<tokio::sync::Notify>,
progress: Arc<std::sync::atomic::AtomicU64>,
truncated: Arc<std::sync::atomic::AtomicBool>,
task: Option<tokio::task::JoinHandle<()>>,
}
impl PipeDrain {
fn start<R>(mut pipe: R) -> Self
where
R: tokio::io::AsyncRead + Unpin + Send + 'static,
{
let buf = Arc::new(PipeBuf(Mutex::new(Some(Vec::new()))));
let started = Arc::new(tokio::sync::Notify::new());
let finished = Arc::new(tokio::sync::Notify::new());
let progress = Arc::new(std::sync::atomic::AtomicU64::new(0));
let truncated = Arc::new(std::sync::atomic::AtomicBool::new(false));
let task = tokio::task::spawn({
let buf = Arc::clone(&buf);
let started = Arc::clone(&started);
let finished = Arc::clone(&finished);
let progress = Arc::clone(&progress);
let truncated = Arc::clone(&truncated);
async move {
use tokio::io::AsyncReadExt;
started.notify_one();
let mut chunk = [0u8; PIPE_DRAIN_CHUNK_BYTES];
while let Ok(read) = pipe.read(&mut chunk).await {
if read == 0 {
break;
}
progress.fetch_add(read as u64, std::sync::atomic::Ordering::Relaxed);
if let Some(bytes) = buf.0.lock().unwrap_poison().as_mut() {
let room = PIPE_READ_CAP.saturating_sub(bytes.len()).min(read);
if room < read {
truncated.store(true, std::sync::atomic::Ordering::Relaxed);
}
bytes.extend_from_slice(&chunk[..room]);
}
}
finished.notify_one();
}
});
Self {
buf,
started,
finished,
progress,
truncated,
task: Some(task),
}
}
async fn started(&self) {
let _ = tokio::time::timeout(PIPE_DRAIN_GRACE, self.started.notified()).await;
}
async fn held_after_exit(&self) -> bool {
let mut last = self.bytes();
let mut waited = Duration::ZERO;
loop {
if tokio::time::timeout(PIPE_DRAIN_GRACE, self.finished.notified())
.await
.is_ok()
{
return false;
}
waited = waited.saturating_add(PIPE_DRAIN_GRACE);
let now = self.bytes();
if now == last {
return true;
}
if waited >= PIPE_DRAIN_MAX {
self.truncated
.store(true, std::sync::atomic::Ordering::Relaxed);
return true;
}
last = now;
}
}
fn bytes(&self) -> u64 {
self.progress.load(std::sync::atomic::Ordering::Relaxed)
}
fn truncated(&self) -> bool {
self.truncated.load(std::sync::atomic::Ordering::Relaxed)
}
}
impl Drop for PipeDrain {
fn drop(&mut self) {
self.buf.take();
if let Some(task) = self.task.take() {
crate::util::leftover_channels::retain_channel(task);
}
}
}
struct PipeBuf(Mutex<Option<Vec<u8>>>);
impl PipeBuf {
fn take(&self) -> Vec<u8> {
self.0.lock().unwrap_poison().take().unwrap_or_default()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::chrome::{
CHROME_USE_DECLARED_BUDGET, KILL_SLACK, RELAY_RECOVERY_BUDGET, kill_bound,
};
use crate::util::test::set_env_var;
use tokio::process::Command;
fn cmd_env(cmd: &Command, key: &str) -> Option<String> {
cmd.as_std()
.get_envs()
.find(|(k, _)| *k == std::ffi::OsStr::new(key))
.and_then(|(_, v)| v.map(|v| v.to_string_lossy().into_owned()))
}
fn cmd_removes(cmd: &Command, key: &str) -> bool {
cmd.as_std()
.get_envs()
.any(|(k, v)| v.is_none() && k == std::ffi::OsStr::new(key))
}
fn cmd_inherits(cmd: &Command, key: &str) -> bool {
cmd.as_std()
.get_envs()
.all(|(k, _)| k != std::ffi::OsStr::new(key))
}
#[test]
fn build_argv_folds_global_flags_before_dash_shield() {
assert_eq!(
build_argv(
&["fill", "#q", "--", "-tail"],
true,
Some("mahbot-chrome-abc")
),
[
"fill",
"#q",
"--json",
"--session",
"mahbot-chrome-abc",
"--",
"-tail"
]
);
assert_eq!(
build_argv(&["open", "https://x"], true, Some("s")),
["open", "https://x", "--json", "--session", "s"]
);
assert_eq!(
build_argv(&["session", "list"], false, None),
["session", "list"]
);
assert_eq!(
build_argv(&["session", "list"], false, Some("s")),
["session", "list", "--session", "s"]
);
}
#[test]
fn ensure_chrome_env_defaults_home_only_when_missing() {
{
let _guard = set_env_var("HOME", None);
let mut cmd = Command::new("true");
ensure_chrome_env(&mut cmd);
assert_eq!(cmd_env(&cmd, "HOME").as_deref(), Some("/tmp"));
}
{
let _guard = set_env_var("HOME", Some("/home/user"));
let mut cmd = Command::new("true");
ensure_chrome_env(&mut cmd);
assert_eq!(cmd_env(&cmd, "HOME"), None);
}
}
#[test]
fn ensure_chrome_env_defaults_chromium_flags_only_when_missing() {
{
let _guard = set_env_var("CHROMIUM_FLAGS", None);
let mut cmd = Command::new("true");
ensure_chrome_env(&mut cmd);
assert_eq!(
cmd_env(&cmd, "CHROMIUM_FLAGS").as_deref(),
Some("--no-first-run --no-default-browser-check --disable-gpu")
);
}
{
let _guard = set_env_var("CHROMIUM_FLAGS", Some("--headless"));
let mut cmd = Command::new("true");
ensure_chrome_env(&mut cmd);
assert_eq!(cmd_env(&cmd, "CHROMIUM_FLAGS"), None);
}
}
#[test]
fn ensure_chrome_env_sets_fixed_env_vars() {
let mut cmd = Command::new("true");
ensure_chrome_env(&mut cmd);
assert_eq!(
cmd_env(&cmd, "AGENT_BROWSER_IDLE_TIMEOUT_MS").as_deref(),
Some("300000")
);
assert_eq!(
cmd_env(&cmd, "AGENT_BROWSER_HUMANIZE").as_deref(),
Some("human")
);
assert_eq!(
cmd_env(&cmd, "AGENT_BROWSER_NO_UPDATE_CHECK").as_deref(),
Some("1")
);
assert!(cmd_inherits(&cmd, "AGENT_BROWSER_DEFAULT_TIMEOUT"));
assert!(cmd_inherits(&cmd, "AGENT_BROWSER_RELAY_REVIVE_SECS"));
assert!(cmd_inherits(&cmd, "AGENT_BROWSER_NO_AUTO_RECONNECT"));
}
#[test]
fn browser_targeting_env_cannot_divert_a_call_off_the_owners_browser() {
for key in [
"CI",
"AGENT_BROWSER_CONFIG",
"AGENT_BROWSER_BROWSER",
"AGENT_BROWSER_NO_AUTO_CONNECT",
"AGENT_BROWSER_PROVIDER",
"AGENT_BROWSER_EXECUTABLE_PATH",
"AGENT_BROWSER_PROFILE",
"AGENT_BROWSER_ENGINE",
"AGENT_BROWSER_FORCE_LAUNCH",
"CHROME_USE_PROFILE",
"CHROME_USE_EXECUTABLE_PATH",
"chrome_use_profile",
] {
let _guard = set_env_var(key, Some("hostile"));
let mut cmd = Command::new("true");
ensure_chrome_env(&mut cmd);
assert!(
cmd_removes(&cmd, key),
"{key} must be removed from the child's environment"
);
}
for (key, expected) in [
("AGENT_BROWSER_IDLE_TIMEOUT_MS", "300000"),
("AGENT_BROWSER_HUMANIZE", "human"),
("AGENT_BROWSER_NO_UPDATE_CHECK", "1"),
("AGENT_BROWSER_AUTO_CONNECT", "1"),
("CHROME_USE_NO_UPDATE_CHECK", "1"),
] {
let _guard = set_env_var(key, Some("hostile"));
let mut cmd = Command::new("true");
ensure_chrome_env(&mut cmd);
assert_eq!(
cmd_env(&cmd, key).as_deref(),
Some(expected),
"{key} must carry the product's own value, not the inherited one"
);
}
for key in ["AGENT_BROWSER_RELAY_DIR", "CHROME_USE_RELAY_DIR"] {
let mut cmd = Command::new("true");
{
let _guard = set_env_var(key, Some("/tmp/relay"));
ensure_chrome_env(&mut cmd);
}
assert!(
cmd_inherits(&cmd, key),
"{key} must ride the service's environment through"
);
}
}
#[test]
fn recovery_policy_gates_the_relay_self_heal() {
let _guard = set_env_var("AGENT_BROWSER_NO_AUTO_RECONNECT", Some("1"));
let mut allowed = Command::new("true");
ensure_chrome_env(&mut allowed);
apply_recovery_env(&mut allowed, CliRecovery::Allowed);
assert!(cmd_removes(&allowed, "AGENT_BROWSER_NO_AUTO_RECONNECT"));
assert_eq!(
cmd_env(&allowed, "AGENT_BROWSER_RELAY_REVIVE_SECS").as_deref(),
Some("45")
);
let mut suppressed = Command::new("true");
ensure_chrome_env(&mut suppressed);
apply_recovery_env(&mut suppressed, CliRecovery::Suppressed);
assert_eq!(
cmd_env(&suppressed, "AGENT_BROWSER_NO_AUTO_RECONNECT").as_deref(),
Some("1")
);
assert_eq!(
cmd_env(&suppressed, "AGENT_BROWSER_RELAY_REVIVE_SECS").as_deref(),
Some("0")
);
}
#[test]
fn apply_chrome_side_declares_the_deadline_on_every_call() {
let mut cmd = Command::new("true");
ensure_chrome_env(&mut cmd);
apply_chrome_side(&mut cmd, Duration::from_secs(18));
assert_eq!(
cmd_env(&cmd, "AGENT_BROWSER_DEFAULT_TIMEOUT").as_deref(),
Some("18000")
);
let mut cmd = Command::new("true");
ensure_chrome_env(&mut cmd);
apply_chrome_side(&mut cmd, Duration::from_millis(250));
assert_eq!(
cmd_env(&cmd, "AGENT_BROWSER_DEFAULT_TIMEOUT").as_deref(),
Some("250")
);
}
#[test]
fn kill_bound_rides_above_the_declared_deadline_and_the_recovery() {
let allowed = kill_bound(CHROME_USE_DECLARED_BUDGET, CliRecovery::Allowed);
let suppressed = kill_bound(CHROME_USE_DECLARED_BUDGET, CliRecovery::Suppressed);
assert_eq!(allowed, suppressed + RELAY_RECOVERY_BUDGET);
assert_eq!(suppressed, CHROME_USE_DECLARED_BUDGET + KILL_SLACK);
assert_eq!(
allowed,
CHROME_USE_DECLARED_BUDGET + KILL_SLACK + RELAY_RECOVERY_BUDGET
);
assert!(kill_bound(Duration::from_secs(6), CliRecovery::Allowed) > Duration::from_secs(6));
assert!(
kill_bound(Duration::from_secs(6), CliRecovery::Suppressed) > Duration::from_secs(6)
);
}
#[cfg(unix)]
#[tokio::test]
async fn child_exit_decides_the_call_when_a_leftover_holds_the_pipe() {
let run = spawn_cli(CliSpawn {
path: Path::new("/bin/sh"),
args: &["-c", "sleep 5 & echo answered"],
session: None,
json: true,
capture_stderr: true,
cancel_kills: true,
input: None,
timeout: CliTimeout::Bounded(Duration::from_secs(10)),
chrome_side: Duration::from_secs(8),
recovery: CliRecovery::Suppressed,
})
.await;
let CliRun::Output(out) = run else {
panic!("expected output");
};
assert!(out.status.success());
assert_eq!(String::from_utf8_lossy(&out.stdout).trim(), "answered");
assert!(out.leftover_pipes, "the background sleep holds the channel");
let run = spawn_cli(CliSpawn {
path: Path::new("/bin/sh"),
args: &["-c", "echo answered"],
session: None,
json: true,
capture_stderr: true,
cancel_kills: true,
input: None,
timeout: CliTimeout::Bounded(Duration::from_secs(10)),
chrome_side: Duration::from_secs(8),
recovery: CliRecovery::Suppressed,
})
.await;
let CliRun::Output(out) = run else {
panic!("expected output");
};
assert!(out.status.success());
assert_eq!(String::from_utf8_lossy(&out.stdout).trim(), "answered");
assert!(!out.leftover_pipes);
}
}