use crate::chrome::contract::{
ChromeResponse, is_daemon_unavailable_error, is_relay_unavailable_error,
is_unreachable_tab_error,
};
use crate::chrome::spawn::{CliRun, CliSpawn, CliTimeout, ensure_chrome_env, spawn_cli};
use crate::util::UnwrapPoison;
use serde_json::Value;
use std::fs;
use std::path::{Path, PathBuf};
use std::process::Stdio;
use std::sync::{Mutex, OnceLock};
use std::time::{Duration, Instant};
use tokio::process::Command;
use tracing::{debug, error, info, warn};
const CLI_TIMEOUT: Duration = Duration::from_secs(8);
const HEALTH_TTL: Duration = Duration::from_secs(10);
const UNHEALTHY_TTL: Duration = Duration::from_mins(1);
const WATCHDOG_INTERVAL: Duration = Duration::from_secs(30);
const CLI_RECHECK: Duration = Duration::from_mins(5);
const CLI_MISSING_THRESHOLD: u32 = 2;
const MAX_RESTART_ATTEMPTS: u32 = 3;
const MAX_LAUNCH_ATTEMPTS: u32 = 3;
const SUSTAINED_HEALTHY_WINDOW: Duration = Duration::from_mins(1);
const RESTART_BACKOFF: [Duration; 3] = [
Duration::from_secs(30),
Duration::from_mins(2),
Duration::from_mins(10),
];
const LAUNCH_BACKOFF: [Duration; 3] = [
Duration::from_secs(30),
Duration::from_mins(2),
Duration::from_mins(10),
];
const HALT_COOLDOWN: Duration = Duration::from_mins(30);
const RELAY_REVIVE_WAIT: Duration = Duration::from_secs(40);
const SWEEP_TOTAL_BUDGET: Duration = Duration::from_secs(15);
const SWEEP_MAX_ROUNDS: u32 = 5;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ProbeFailure {
NotInstalled,
HostBroken,
ExtensionDisabled,
ExtensionAbsent,
RelayDown,
ChromeNotRunning,
UnreachableTab,
DaemonWedge,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ProbeOutcome {
Healthy,
Down(ProbeFailure),
}
impl ProbeOutcome {
fn is_healthy(self) -> bool {
matches!(self, ProbeOutcome::Healthy)
}
fn failure(self) -> Option<ProbeFailure> {
match self {
ProbeOutcome::Healthy => None,
ProbeOutcome::Down(f) => Some(f),
}
}
}
impl ProbeFailure {
fn is_unfixable(self) -> bool {
matches!(
self,
ProbeFailure::NotInstalled
| ProbeFailure::HostBroken
| ProbeFailure::ExtensionDisabled
| ProbeFailure::ExtensionAbsent
| ProbeFailure::UnreachableTab
)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ChromeLaunchOutcome {
Launched,
Failed,
NoDisplay,
}
#[derive(Default)]
struct AttemptBudget {
attempts: u32,
next_at: Option<Instant>,
halted: bool,
halted_until: Option<Instant>,
}
impl AttemptBudget {
fn gate(
&mut self,
max: u32,
backoff: &[Duration],
cooldown: Duration,
now: Instant,
) -> RecoveryGate {
if self.halted {
if self.halted_until.is_some_and(|until| now < until) {
return RecoveryGate::Cooldown;
}
self.halted = false;
self.attempts = 0;
self.halted_until = None;
self.next_at = None;
}
if self.next_at.is_some_and(|next| now < next) {
return RecoveryGate::Backoff;
}
if self.attempts >= max {
self.halted = true;
self.halted_until = Some(now + cooldown);
return RecoveryGate::Halted;
}
let attempt = self.attempts + 1;
self.attempts = attempt;
self.next_at = Some(now + backoff[(attempt as usize - 1).min(backoff.len() - 1)]);
RecoveryGate::Allowed(attempt)
}
fn is_waiting(&self, now: Instant) -> bool {
self.halted_until.is_some_and(|until| now < until)
|| self.next_at.is_some_and(|next| now < next)
}
fn reset(&mut self) {
self.attempts = 0;
self.next_at = None;
self.halted = false;
self.halted_until = None;
}
}
#[derive(Default)]
struct DaemonHealth {
healthy: Option<bool>,
last_probe: Option<Instant>,
restart_budget: AttemptBudget,
launch_budget: AttemptBudget,
launch_outcome: Option<ChromeLaunchOutcome>,
last_failure: Option<ProbeFailure>,
last_cause_warned: Option<ProbeFailure>,
healthy_since: Option<Instant>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum RecoveryGate {
Allowed(u32),
Backoff,
Cooldown,
Halted,
}
impl DaemonHealth {
fn gate_restart(&mut self, now: Instant) -> RecoveryGate {
self.restart_budget
.gate(MAX_RESTART_ATTEMPTS, &RESTART_BACKOFF, HALT_COOLDOWN, now)
}
fn gate_launch(&mut self, now: Instant) -> RecoveryGate {
self.launch_budget
.gate(MAX_LAUNCH_ATTEMPTS, &LAUNCH_BACKOFF, HALT_COOLDOWN, now)
}
fn apply_outcome(&mut self, outcome: ProbeOutcome, now: Instant, seed_window: bool) {
let healthy = outcome.is_healthy();
if healthy {
if self
.healthy_since
.is_some_and(|since| now.duration_since(since) >= SUSTAINED_HEALTHY_WINDOW)
{
self.restart_budget.reset();
self.launch_budget.reset();
self.last_cause_warned = None;
}
if seed_window && self.healthy_since.is_none() {
self.healthy_since = Some(now);
}
self.launch_outcome = None;
} else {
self.healthy_since = None;
}
self.last_failure = outcome.failure();
self.healthy = Some(healthy);
self.last_probe = Some(now);
}
}
static HEALTH: OnceLock<Mutex<DaemonHealth>> = OnceLock::new();
static WAKE: OnceLock<tokio::sync::Notify> = OnceLock::new();
fn health() -> &'static Mutex<DaemonHealth> {
HEALTH.get_or_init(|| Mutex::new(DaemonHealth::default()))
}
fn wake() -> &'static tokio::sync::Notify {
WAKE.get_or_init(tokio::sync::Notify::new)
}
fn record_launch_outcome(outcome: ChromeLaunchOutcome) {
health().lock().unwrap_poison().launch_outcome = Some(outcome);
}
fn classify_failure_text(msg: &str) -> Option<ProbeFailure> {
if is_unreachable_tab_error(msg) {
Some(ProbeFailure::UnreachableTab)
} else if is_relay_unavailable_error(msg) {
Some(ProbeFailure::RelayDown)
} else if is_daemon_unavailable_error(msg) {
Some(ProbeFailure::DaemonWedge)
} else {
None
}
}
pub(crate) const fn chrome_bin() -> &'static str {
if cfg!(target_os = "windows") {
"chrome-use.exe"
} else {
"chrome-use"
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum CliStatus {
Available,
Missing,
Transient(CliProbeFailure),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum CliProbeFailure {
Spawn(String),
BadVersion(String),
Timeout,
}
impl std::fmt::Display for CliProbeFailure {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
CliProbeFailure::Spawn(reason) => write!(f, "spawn failed ({reason})"),
CliProbeFailure::BadVersion(status) => write!(f, "--version check failed ({status})"),
CliProbeFailure::Timeout => write!(f, "probe timed out"),
}
}
}
const CHROME_USE_RELEASE_REPO: &str = "leeguooooo/chrome-use";
const CHROME_USE_DOWNLOAD_TIMEOUT: Duration = Duration::from_secs(300);
const CHROME_USE_RELEASE_TIMEOUT: Duration = Duration::from_secs(30);
const CHROME_USE_INSTALL_TIMEOUT: Duration = Duration::from_secs(120);
pub(crate) const CHROME_USE_INSTALL_HINT: &str = "It installs automatically in the background at \
startup — if it is still missing the quiet install failed; check the logs \
and it will retry on the next boot.";
static CLI_PATH: OnceLock<Mutex<Option<PathBuf>>> = OnceLock::new();
pub(crate) fn cli_path() -> Option<PathBuf> {
let mut cache = CLI_PATH
.get_or_init(|| Mutex::new(None))
.lock()
.unwrap_poison();
if let Some(path) = cache.as_ref().filter(|p| crate::util::is_executable(p)) {
return Some(path.clone());
}
let found = find_cli_binary();
cache.clone_from(&found);
found
}
fn invalidate_cli_path() {
*CLI_PATH
.get_or_init(|| Mutex::new(None))
.lock()
.unwrap_poison() = None;
}
fn find_cli_binary() -> Option<PathBuf> {
let name = chrome_bin();
if let Some(dir) = crate::util::managed_bin::storage_bin_dir() {
let candidate = dir.join(name);
if crate::util::is_executable(&candidate) {
return Some(candidate);
}
}
if let Some(paths) = std::env::var_os("PATH") {
for dir in std::env::split_paths(&paths) {
let candidate = dir.join(name);
if crate::util::is_executable(&candidate) {
return Some(candidate);
}
}
}
if !cfg!(target_os = "windows") {
let home = directories::UserDirs::new().map(|d| d.home_dir().to_path_buf());
let literal_cargo_bin = match (std::env::var_os("CARGO_HOME"), home.as_deref()) {
(Some(cargo_home), Some(h)) if !cargo_home.is_empty() => Some(h.join(".cargo/bin")),
_ => None,
};
for base in [
home.as_deref().map(|h| h.join(".local/bin")),
crate::util::cargo_bin_dir(),
literal_cargo_bin,
Some(PathBuf::from("/usr/local/bin")),
Some(PathBuf::from("/opt/homebrew/bin")),
]
.into_iter()
.flatten()
{
let candidate = base.join(name);
if crate::util::is_executable(&candidate) {
return Some(candidate);
}
}
}
None
}
fn classify_spawn_error(e: &std::io::Error) -> CliStatus {
if e.kind() == std::io::ErrorKind::NotFound {
CliStatus::Missing
} else {
debug!("chrome-use CLI probe spawn failed: {e}");
CliStatus::Transient(CliProbeFailure::Spawn(e.to_string()))
}
}
pub(crate) async fn cli_probe() -> CliStatus {
let Some(path) = cli_path() else {
return CliStatus::Missing;
};
let mut cmd = Command::new(&path);
ensure_chrome_env(&mut cmd);
cmd.arg("--version")
.stdout(Stdio::null())
.stderr(Stdio::null())
.kill_on_drop(true);
let status = match tokio::time::timeout(CLI_TIMEOUT, cmd.status()).await {
Ok(Ok(status)) => status,
Ok(Err(e)) => return classify_spawn_error(&e),
Err(_) => {
debug!("chrome-use CLI probe timed out");
return CliStatus::Transient(CliProbeFailure::Timeout);
}
};
if status.success() {
CliStatus::Available
} else {
debug!("chrome-use CLI probe failed: --version exited with {status}");
CliStatus::Transient(CliProbeFailure::BadVersion(status.to_string()))
}
}
fn parse_cli_version(stdout: &str) -> Option<semver::Version> {
stdout
.split_whitespace()
.find_map(|t| crate::util::managed_bin::parse_tag_version(t, "v"))
}
#[must_use]
fn release_asset_platform(os: &str, arch: &str, musl: bool) -> Option<String> {
let asset = match (os, arch, musl) {
("macos", "x86_64", _) => "darwin-x64",
("macos", "aarch64", _) => "darwin-arm64",
("linux", "x86_64", false) => "linux-x64",
("linux", "aarch64", false) => "linux-arm64",
("linux", "x86_64", true) => "linux-musl-x64",
("linux", "aarch64", true) => "linux-musl-arm64",
("windows", "x86_64", _) => "win32-x64",
_ => return None,
};
Some(asset.to_string())
}
fn release_asset_name() -> Result<String, String> {
let (os, arch) = crate::util::managed_bin::host_os_arch()?;
let musl = os == "linux" && crate::util::managed_bin::linux_host_is_musl();
release_asset_platform(os, arch, musl)
.ok_or_else(|| format!("chrome-use has no release asset for {os}-{arch}"))
}
pub(crate) async fn cli_version() -> Option<semver::Version> {
let path = cli_path()?;
let mut cmd = Command::new(&path);
ensure_chrome_env(&mut cmd);
cmd.arg("--version")
.stdout(Stdio::piped())
.stderr(Stdio::null())
.kill_on_drop(true);
let out = tokio::time::timeout(CLI_TIMEOUT, cmd.output())
.await
.ok()?
.ok()?;
if !out.status.success() {
return None;
}
parse_cli_version(&String::from_utf8_lossy(&out.stdout))
}
async fn run_cli_bounded(args: &[&str], session: Option<&str>) -> Option<std::process::Output> {
let path = cli_path()?;
run_cli_bounded_at(&path, args, session, CLI_TIMEOUT).await
}
async fn run_cli_bounded_at(
path: &Path,
args: &[&str],
session: Option<&str>,
timeout: Duration,
) -> Option<std::process::Output> {
match spawn_cli(CliSpawn {
path,
args,
session,
json: true,
capture_stderr: false,
timeout: CliTimeout::Bounded(timeout),
cancel_kills: true,
input: None,
chrome_deadline: None,
})
.await
{
CliRun::Output(out) => Some(out),
CliRun::SpawnFailure | CliRun::TimedOut => None,
}
}
async fn run_cli_json_opt(args: &[&str], session: Option<&str>) -> Result<Value, Option<String>> {
run_cli_bounded(args, session)
.await
.as_ref()
.ok_or(None)
.and_then(json_outcome)
}
pub(crate) async fn run_cli_json_at(
path: &Path,
args: &[&str],
session: Option<&str>,
timeout: Duration,
) -> Result<Value, Option<String>> {
run_cli_bounded_at(path, args, session, timeout)
.await
.as_ref()
.ok_or(None)
.and_then(json_outcome)
}
fn json_outcome(out: &std::process::Output) -> Result<Value, Option<String>> {
let v: Value = serde_json::from_slice(&out.stdout).map_err(|_| None)?;
let env = ChromeResponse::from_value(&v);
if !out.status.success() || env.verdict() != Some(true) {
return Err(env.error.filter(|e| !e.is_empty()));
}
Ok(v)
}
async fn run_cli_json(args: &[&str]) -> Option<Value> {
run_cli_json_opt(args, None).await.ok()
}
async fn service_state() -> Option<ProbeFailure> {
let status = run_cli_json(&["status"]).await?;
classify_service_state(&status).await
}
async fn classify_service_state(status: &Value) -> Option<ProbeFailure> {
let ext = status.get("data")?.get("extension")?;
if ext.get("hostInstalled").and_then(Value::as_bool) == Some(false) {
return Some(ProbeFailure::NotInstalled);
}
if ext.get("hostHealthy").and_then(Value::as_bool) == Some(false) {
return Some(ProbeFailure::HostBroken);
}
if ext.get("relayUp").and_then(Value::as_bool) == Some(false) {
let ext_state = extension_state().await;
let chrome = if matches!(ext_state, ExtensionState::Present | ExtensionState::Unknown) {
chrome_running().await
} else {
None
};
return Some(classify_relay_down(ext_state, chrome));
}
None
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ExtensionState {
Present,
Disabled,
Absent,
Unknown,
}
fn extension_state_from(status: &Value) -> ExtensionState {
match status.get("data").and_then(|d| d.get("chromeExtension")) {
None => ExtensionState::Unknown,
Some(c) if c.is_null() => ExtensionState::Absent,
Some(c) => {
if c.get("disableReasons")
.and_then(Value::as_array)
.is_some_and(|r| !r.is_empty())
{
ExtensionState::Disabled
} else {
ExtensionState::Present
}
}
}
}
async fn extension_state() -> ExtensionState {
run_cli_json_opt(&["extension", "status"], None)
.await
.map_or(ExtensionState::Unknown, |v| extension_state_from(&v))
}
fn classify_relay_down(ext: ExtensionState, chrome_running: Option<bool>) -> ProbeFailure {
match ext {
ExtensionState::Disabled => ProbeFailure::ExtensionDisabled,
ExtensionState::Absent => ProbeFailure::ExtensionAbsent,
ExtensionState::Present | ExtensionState::Unknown => {
if chrome_running == Some(false) {
ProbeFailure::ChromeNotRunning
} else {
ProbeFailure::RelayDown
}
}
}
}
#[cfg(not(target_os = "windows"))]
const CHROME_PROCESS_NAMES: [&str; 6] = [
"Google Chrome",
"Chromium",
"Brave Browser",
"chrome",
"chromium",
"brave",
];
#[cfg(target_os = "windows")]
const WINDOWS_CHROME_PROCESS_NAMES: [&str; 3] = ["chrome.exe", "chromium.exe", "brave.exe"];
pub(crate) async fn chrome_running() -> Option<bool> {
#[cfg(not(target_os = "windows"))]
{
pgrep_running(&CHROME_PROCESS_NAMES.join("|")).await
}
#[cfg(target_os = "windows")]
{
for name in WINDOWS_CHROME_PROCESS_NAMES {
if tasklist_has(name).await? {
return Some(true);
}
}
Some(false)
}
}
#[cfg(not(target_os = "windows"))]
async fn pgrep_running(pattern: &str) -> Option<bool> {
let mut cmd = Command::new("pgrep");
cmd.args(["-x", pattern])
.stdout(Stdio::null())
.stderr(Stdio::null())
.kill_on_drop(true);
let status = tokio::time::timeout(CLI_TIMEOUT, cmd.status())
.await
.ok()?
.ok()?;
match status.code() {
Some(0) => Some(true),
Some(1) => Some(false), _ => None,
}
}
#[cfg(target_os = "windows")]
async fn tasklist_has(name: &str) -> Option<bool> {
let mut cmd = Command::new("tasklist");
cmd.args(["/FI", &format!("IMAGENAME eq {name}"), "/NH", "/FO", "CSV"])
.stdout(Stdio::piped())
.stderr(Stdio::null())
.kill_on_drop(true);
let out = tokio::time::timeout(CLI_TIMEOUT, cmd.output())
.await
.ok()?
.ok()?;
if !out.status.success() {
return None;
}
Some(String::from_utf8_lossy(&out.stdout).contains(name))
}
async fn evaluate_health() -> ProbeOutcome {
let Some(status) = run_cli_json(&["status"]).await else {
return ProbeOutcome::Healthy;
};
if let Some(failure) = classify_service_state(&status).await {
return ProbeOutcome::Down(failure);
}
advise_extension_skew(&status);
ProbeOutcome::Healthy
}
struct SweepTab {
tab_id: String,
target_id: String,
}
pub(crate) async fn sweep_session(name: &str) {
if !is_mahbot_session_name(name) {
warn!(
session = name,
"tab sweep refused: not a mahbot-owned session (user/default/other-agent sessions are never touched)"
);
return;
}
let deadline = Instant::now() + SWEEP_TOTAL_BUDGET;
if let Some(failure) = service_state().await {
debug!(
session = name,
?failure,
"tab sweep skipped — chrome service unavailable"
);
return;
}
let mut scratch: Option<String> = None;
let mut stopped = false;
for _round in 1..=SWEEP_MAX_ROUNDS {
if Instant::now() >= deadline {
break;
}
let Some(tabs) = session_tab_list(name, deadline).await else {
return; };
if tabs.is_empty() {
let _ = stop_session_daemon(name, deadline).await;
stopped = true;
scratch = None;
continue;
}
if stopped {
if tabs.len() == 1 && scratch.as_deref() != Some(tabs[0].target_id.as_str()) {
clear_sweep_warn();
let _ = stop_session_daemon(name, deadline).await;
return;
}
stopped = false;
scratch = None; }
if tabs.len() == 1 && scratch.as_deref() == Some(tabs[0].target_id.as_str()) {
let _ = stop_session_daemon(name, deadline).await;
stopped = true;
continue;
}
if scratch.is_none() {
let Some(target_id) = session_tab_new_scratch(name, deadline).await else {
return; };
scratch = Some(target_id);
}
for tab in &tabs {
if tab.target_id == *scratch.as_deref().unwrap_or_default() {
continue; }
if Instant::now() >= deadline {
break;
}
let _ = session_close_tab(name, &tab.tab_id, deadline).await;
}
}
sweep_warn_transition(SweepWarn::Deferred);
}
pub(crate) fn is_mahbot_session_name(name: &str) -> bool {
name.starts_with("link-enricher-")
|| name.starts_with(crate::chrome::CLI_EPHEMERAL_PREFIX)
|| name.starts_with("mahbot-browser-ephemeral-")
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum SweepWarn {
UnreachableTab,
CannotEnumerate,
Deferred,
}
static LAST_SWEEP_WARN: OnceLock<Mutex<Option<SweepWarn>>> = OnceLock::new();
fn sweep_warn_transition(cause: SweepWarn) {
let mut last = LAST_SWEEP_WARN
.get_or_init(|| Mutex::new(None))
.lock()
.unwrap_poison();
if *last == Some(cause) {
return;
}
*last = Some(cause);
match cause {
SweepWarn::UnreachableTab => warn!(
"tab sweep: a leftover tab is unreachable (the extension lost its debugger attach; \
about:blank tabs are never re-attached) — close the leftover tab in Chrome to \
unblock this session; the sweep keeps retrying"
),
SweepWarn::CannotEnumerate => warn!(
"tab sweep: cannot enumerate session tabs (relay/daemon unreachable or malformed \
response) — deferring to the next sweep"
),
SweepWarn::Deferred => {
warn!("tab sweep: group not clean within budget — deferring to the next sweep");
}
}
}
fn clear_sweep_warn() {
*LAST_SWEEP_WARN
.get_or_init(|| Mutex::new(None))
.lock()
.unwrap_poison() = None;
}
fn sweep_none_on_cli_error<T>(name: &str, err: Option<&str>) -> Option<T> {
let msg = err.unwrap_or_default();
if is_unreachable_tab_error(msg) {
tracing::debug!(
session = name,
error = msg,
"tab sweep: unreachable-tab detail"
);
sweep_warn_transition(SweepWarn::UnreachableTab);
} else {
sweep_warn_transition(SweepWarn::CannotEnumerate);
}
None
}
async fn session_tab_list(name: &str, deadline: Instant) -> Option<Vec<SweepTab>> {
if Instant::now() >= deadline {
sweep_warn_transition(SweepWarn::Deferred);
return None;
}
let v = match run_session_cli_json(&["tab", "list"], name).await {
Ok(v) => v,
Err(err) => return sweep_none_on_cli_error(name, err.as_deref()),
};
let Some(tabs) = v
.get("data")
.and_then(|d| d.get("tabs"))
.and_then(Value::as_array)
else {
sweep_warn_transition(SweepWarn::CannotEnumerate);
return None;
};
let parsed: Option<Vec<SweepTab>> = tabs
.iter()
.map(|t| {
Some(SweepTab {
tab_id: t.get("tabId")?.as_str()?.to_string(),
target_id: t.get("targetId")?.as_str()?.to_string(),
})
})
.collect();
parsed.or_else(|| {
sweep_warn_transition(SweepWarn::CannotEnumerate);
None
})
}
async fn session_tab_new_scratch(name: &str, deadline: Instant) -> Option<String> {
if Instant::now() >= deadline {
sweep_warn_transition(SweepWarn::Deferred);
return None;
}
let resp = match run_session_cli_json(&["tab", "new"], name).await {
Ok(v) => v,
Err(err) => return sweep_none_on_cli_error(name, err.as_deref()),
};
let Some(tab_id) = resp
.get("data")
.and_then(|d| d.get("tabId"))
.and_then(Value::as_str)
.map(String::from)
else {
sweep_warn_transition(SweepWarn::CannotEnumerate);
return None;
};
let after = session_tab_list(name, deadline).await?; after
.iter()
.find(|t| t.tab_id == tab_id)
.map(|t| t.target_id.clone())
.or_else(|| {
sweep_warn_transition(SweepWarn::CannotEnumerate);
None
})
}
async fn session_close_tab(name: &str, tab_id: &str, deadline: Instant) -> Option<()> {
if Instant::now() >= deadline {
return None;
}
run_session_cli_json(&["close", tab_id], name)
.await
.ok()
.map(|_| ())
}
async fn stop_session_daemon(name: &str, deadline: Instant) -> Option<()> {
if Instant::now() >= deadline {
return None;
}
run_session_cli_json(&["session", "stop"], name)
.await
.ok()
.map(|_| ())
}
async fn run_session_cli_json(args: &[&str], session: &str) -> Result<Value, Option<String>> {
run_cli_json_opt(args, Some(session)).await
}
fn set_health(outcome: ProbeOutcome) {
let mut h = health().lock().unwrap_poison();
h.apply_outcome(outcome, Instant::now(), true);
}
fn set_health_after_recovery(outcome: ProbeOutcome) {
let mut h = health().lock().unwrap_poison();
h.apply_outcome(outcome, Instant::now(), false);
}
pub(crate) async fn is_available() -> bool {
let cached = {
let h = health().lock().unwrap_poison();
let ttl = if h.healthy == Some(false) {
UNHEALTHY_TTL
} else {
HEALTH_TTL
};
h.last_probe
.filter(|t| t.elapsed() < ttl)
.map(|_| h.healthy)
};
if let Some(Some(healthy)) = cached {
return healthy;
}
let outcome = evaluate_health().await;
let healthy = outcome.is_healthy();
set_health(outcome);
if !healthy {
wake().notify_one();
}
healthy
}
pub(crate) fn is_advertised() -> bool {
health().lock().unwrap_poison().healthy != Some(false)
}
pub(crate) fn note_unhealthy(error: &str) {
set_health(ProbeOutcome::Down(
classify_failure_text(error).unwrap_or(ProbeFailure::DaemonWedge),
));
wake().notify_one();
}
pub(crate) fn daemon_down_message() -> String {
let h = health().lock().unwrap_poison();
let cause = match h.last_failure {
Some(ProbeFailure::NotInstalled) => {
"The chrome-use extension or native host is not installed — the chrome daemon \
cannot run. Enable the chrome-use extension at chrome://extensions (the CLI \
re-installs itself in the background at startup); health recovers automatically \
once it is installed."
}
Some(ProbeFailure::HostBroken) => {
"The chrome-use native host launcher is broken — run `chrome-use doctor`; health \
recovers automatically once it is fixed."
}
Some(ProbeFailure::ExtensionDisabled) => {
"The chrome-use extension is disabled — enable it at chrome://extensions. Daemon \
restarts cannot fix a Chrome-side disable; health recovers automatically once \
it is enabled."
}
Some(ProbeFailure::ExtensionAbsent) => {
"The chrome-use extension is not installed in Chrome — install/enable the ab-connect \
extension from the Chrome Web Store. Daemon restarts and launches cannot fix an \
absent extension; health recovers automatically once it is installed."
}
Some(ProbeFailure::RelayDown)
if h.launch_outcome == Some(ChromeLaunchOutcome::Launched) =>
{
"The extension relay is down even after Chrome was auto-launched — Chrome is \
running, but the ab-connect extension is not republishing. If the extension is \
only installed in a non-default profile, install it in Chrome's default profile."
}
Some(ProbeFailure::RelayDown) => {
"The chrome-use extension relay is down (the extension itself is enabled). \
Auto-recovery restarts the session daemons and waits for the extension to \
reconnect."
}
Some(ProbeFailure::ChromeNotRunning) => match h.launch_outcome {
None | Some(ChromeLaunchOutcome::Launched) => "Chrome is not running.",
Some(ChromeLaunchOutcome::Failed) => {
"Chrome is not running and the auto-launch attempt failed (binary not found or \
could not start) — start Chrome manually."
}
Some(ChromeLaunchOutcome::NoDisplay) => {
"Chrome is not running and this host has no display — start Chrome manually."
}
},
Some(ProbeFailure::UnreachableTab) => {
"A browser tab the session was driving is unreachable (the extension lost its \
debugger attach; about:blank tabs are never re-attached) — close the leftover \
tab in Chrome to unblock the session."
}
Some(ProbeFailure::DaemonWedge) | None => "The chrome daemon is down or unresponsive.",
};
let recovery = if matches!(h.last_failure, Some(ProbeFailure::ChromeNotRunning)) {
match h.launch_outcome {
Some(ChromeLaunchOutcome::NoDisplay) => {
" Auto-recovery is paused for this cause — no launch will be attempted; it \
resumes automatically once a display is available."
}
_ if h.launch_budget.halted => {
" Auto-recovery exhausted its launch attempts and is in a 30-minute cooldown \
(thrash protection); it will retry after the cooldown — start Chrome manually \
if Chrome stays down."
}
Some(ChromeLaunchOutcome::Failed) => {
" Auto-recovery will retry the launch with backoff."
}
None | Some(ChromeLaunchOutcome::Launched) => {
" Auto-recovery will launch Chrome, backing off between attempts — no manual \
action is needed."
}
}
} else if h.restart_budget.halted {
" Auto-recovery exhausted its restart attempts and is in a 30-minute cooldown (thrash \
protection); it will retry after the cooldown."
} else if h.last_failure.is_some_and(ProbeFailure::is_unfixable) {
" Auto-recovery is paused for this cause — no restart will be attempted; it resumes \
automatically once the underlying issue is resolved."
} else {
" Auto-recovery was triggered and will restart it automatically — no manual action is \
needed (note: the restart resets chrome sessions)."
};
format!(
"{cause}{recovery} While it's down, use web_search, or shell `curl` for page fetches, \
instead of the chrome tool."
)
}
pub(crate) async fn health_after_call_timeout(session: &str) -> Option<String> {
if !is_available().await {
return Some(daemon_down_message());
}
if run_cli_bounded(&["get", "url"], Some(session))
.await
.is_none()
{
note_unhealthy("chrome-use CLI call timed out twice — daemon unresponsive");
return Some(daemon_down_message());
}
None
}
pub async fn run_watchdog() {
let mut cli_present: Option<bool> = None;
let mut last_cli_check = Instant::now();
let mut cli_missing: u32 = 0;
let mut last_transient: Option<CliProbeFailure> = None;
let mut cleaned = false;
let mut woken = false;
loop {
let mut sleep = WATCHDOG_INTERVAL;
let mut skip_health = false;
let cli_due = last_cli_check.elapsed() >= CLI_RECHECK;
if cli_present != Some(true) || cli_due {
last_cli_check = Instant::now();
match cli_probe().await {
CliStatus::Available => {
cli_present = Some(true);
cli_missing = 0;
last_transient = None;
}
CliStatus::Transient(failure) => {
if last_transient.as_ref() != Some(&failure) {
warn!("chrome-use CLI probe transient: {failure}");
last_transient = Some(failure);
}
cli_missing = 0;
}
CliStatus::Missing => {
cli_missing += 1;
last_transient = None;
if cli_missing < CLI_MISSING_THRESHOLD {
cli_present = None;
skip_health = true;
} else {
if cli_present != Some(false) {
cli_present = Some(false);
warn!(
"chrome-use CLI not found — chrome daemon watchdog standing down"
);
}
sleep = CLI_RECHECK;
skip_health = true;
}
}
}
}
if !skip_health {
if !cleaned {
cleaned = true;
cleanup_stale_sessions().await;
}
let failure = if woken {
health().lock().unwrap_poison().last_failure
} else {
None
};
if let Some(failure) = failure {
attempt_recovery(failure).await;
} else {
let outcome = evaluate_health().await;
set_health(outcome);
if let ProbeOutcome::Down(failure) = outcome {
attempt_recovery(failure).await;
}
}
}
let shutdown = crate::shutdown::shutdown_token();
woken = tokio::select! {
() = tokio::time::sleep(sleep) => false,
() = wake().notified() => true,
() = shutdown.cancelled() => break,
};
}
}
static LAST_EXTENSION_SKEW: OnceLock<Mutex<Option<(String, String)>>> = OnceLock::new();
fn advise_extension_skew(status: &Value) {
let ext = status.get("data").and_then(|d| d.get("extension"));
let (Some(expected), Some(live)) = (
ext.and_then(|e| e.get("expectedVersion"))
.and_then(Value::as_str),
ext.and_then(|e| e.get("liveVersion"))
.and_then(Value::as_str),
) else {
return;
};
if expected.is_empty() || live.is_empty() {
return;
}
let mut last = LAST_EXTENSION_SKEW
.get_or_init(|| Mutex::new(None))
.lock()
.unwrap_poison();
if expected == live {
*last = None;
return;
}
if *last == Some((expected.to_string(), live.to_string())) {
return;
}
*last = Some((expected.to_string(), live.to_string()));
info!(
"chrome-use browser extension version mismatch: installed {live}, CLI expects {expected} — \
Chrome updates Web Store extensions automatically in the background; reloading the \
extension at chrome://extensions (or waiting for the Store auto-update) clears this"
);
}
async fn download_chrome_use_binary(tag: &str) -> Result<(tempfile::TempDir, PathBuf), String> {
use crate::util::http::{DownloadSizeCheck, build_download_client, download_verified};
let platform = release_asset_name()?;
let asset = format!("chrome-use-{platform}.tar.gz");
let base = format!("https://github.com/{CHROME_USE_RELEASE_REPO}/releases/download/{tag}");
let tgz_url = format!("{base}/{asset}");
let sha_url = format!("{tgz_url}.sha256");
let client = build_download_client(CHROME_USE_DOWNLOAD_TIMEOUT)
.map_err(|e| format!("failed to build download client: {e}"))?;
let sidecar = client
.get(&sha_url)
.send()
.await
.map_err(|e| format!("failed to fetch sha256 sidecar {sha_url}: {e}"))?;
if !sidecar.status().is_success() {
return Err(format!(
"failed to fetch sha256 sidecar {sha_url}: HTTP {}",
sidecar.status()
));
}
let body = sidecar
.text()
.await
.map_err(|e| format!("failed to read sha256 sidecar {sha_url}: {e}"))?;
let (hash, sidecar_name) =
crate::util::managed_bin::parse_sha256_sidecar(&body).ok_or_else(|| {
format!("sha256 sidecar {sha_url} is malformed (no `64-hex-hash filename` pair)")
})?;
if sidecar_name != asset {
return Err(format!(
"sha256 sidecar {sha_url} names '{sidecar_name}', expected '{asset}'"
));
}
let dir = tempfile::tempdir().map_err(|e| format!("failed to create temp dir: {e}"))?;
let archive_path = dir.path().join("archive.tar.gz");
download_verified(
&client,
&tgz_url,
&archive_path,
&hash,
None,
DownloadSizeCheck::None,
|_, _| {},
)
.await
.map_err(|e| format!("failed to download {tgz_url}: {e}"))?;
let out_path = crate::util::managed_bin::extract_single_file_tar_gz(
&archive_path,
dir.path(),
chrome_bin(),
)?;
Ok((dir, out_path))
}
pub(crate) async fn install_chrome_use() -> Result<(), String> {
let tag = crate::util::managed_bin::fetch_latest_tag(
CHROME_USE_RELEASE_REPO,
CHROME_USE_RELEASE_TIMEOUT,
)
.await?;
let (_temp, fresh) = download_chrome_use_binary(&tag).await?;
let dest = crate::util::managed_bin::storage_bin_dir()
.ok_or_else(|| "managed chrome-use bin dir unavailable (storage root not set)".to_string())?
.join(chrome_bin());
let parent = dest
.parent()
.ok_or_else(|| format!("invalid chrome-use install path {}", dest.display()))?;
fs::create_dir_all(parent)
.map_err(|e| format!("failed to create {}: {e}", parent.display()))?;
let fresh_install = !dest.exists();
crate::util::managed_bin::swap_binary_in_place(&fresh, &dest)?;
crate::util::managed_bin::set_executable(&dest)?;
invalidate_cli_path();
crate::util::managed_bin::ensure_rc_path_block();
let mut host = Command::new(&dest);
host.args(["extension", "install", "--no-profile"]);
match run_install_step("`chrome-use extension install --no-profile`", host).await {
Ok(()) => Ok(()),
Err(e) => {
if fresh_install {
let _ = fs::remove_file(&dest);
invalidate_cli_path();
Err(format!(
"{e}\nThe chrome-use binary was placed at {} but the native-host \
registration failed, so it was removed to leave no half-installed state.",
dest.display()
))
} else {
Err(format!(
"{e}\nThe chrome-use binary at {} was updated, but the native-host \
re-registration failed — the previous registration still references this \
path and keeps working.",
dest.display()
))
}
}
}
}
async fn run_install_step(label: &str, mut cmd: Command) -> Result<(), String> {
let out = tokio::time::timeout(CHROME_USE_INSTALL_TIMEOUT, cmd.kill_on_drop(true).output())
.await
.map_err(|_| format!("{label} timed out"))?
.map_err(|e| format!("{label} failed to spawn: {e}"))?;
if out.status.success() {
return Ok(());
}
Err(format!(
"{label} failed ({}).\nstdout: {}\nstderr: {}",
out.status,
crate::util::truncate(&String::from_utf8_lossy(&out.stdout), 2048),
crate::util::truncate(&String::from_utf8_lossy(&out.stderr), 2048),
))
}
pub async fn run_chrome_use_management() {
if cli_path().is_none() {
match install_chrome_use().await {
Ok(()) => {
info!("chrome-use installed automatically at startup");
return;
}
Err(e) => {
warn!(
"chrome-use auto-install failed (non-fatal, retried on next boot): {}",
crate::util::truncate(&e, 1024)
);
return;
}
}
}
tokio::time::sleep(Duration::from_mins(5)).await;
run_update_check().await;
}
async fn run_update_check() {
if cli_path().is_none() {
debug!("chrome-use not installed — update check skipped (management only installs)");
return;
}
crate::util::managed_bin::ensure_rc_path_block();
let local = cli_version().await;
let Ok(tag) = crate::util::managed_bin::fetch_latest_tag(
CHROME_USE_RELEASE_REPO,
CHROME_USE_RELEASE_TIMEOUT,
)
.await
else {
debug!("chrome-use auto-update skipped: release check failed (offline?)");
return;
};
let Some(latest) = crate::util::managed_bin::parse_tag_version(&tag, "v") else {
info!(
"chrome-use auto-update: latest release tag '{tag}' is not a semver version; giving up"
);
return;
};
if let Some(local) = local.as_ref()
&& local >= &latest
{
debug!("chrome-use is up to date ({local})");
return;
}
let local = local.map_or_else(|| "unknown version".to_string(), |v| v.to_string());
debug!("chrome-use auto-update: updating {local} → {latest}");
let Some(dest) = cli_path() else {
debug!("chrome-use auto-update skipped: CLI path vanished");
return;
};
let (_temp, fresh) = match download_chrome_use_binary(&tag).await {
Ok(pair) => pair,
Err(e) => {
info!(
"chrome-use auto-update failed: {}",
crate::util::truncate(&e, 1024)
);
return;
}
};
if let Err(e) = crate::util::managed_bin::swap_binary_in_place(&fresh, &dest) {
info!(
"chrome-use auto-update failed: {}",
crate::util::truncate(&e, 1024)
);
return;
}
if let Err(e) = crate::util::managed_bin::set_executable(&dest) {
info!(
"chrome-use auto-update failed: {}",
crate::util::truncate(&e, 1024)
);
return;
}
info!("chrome-use auto-updated to {latest} (binary in place)");
}
async fn cleanup_stale_sessions() {
let Some(sessions) = registered_sessions().await else {
return;
};
for name in sessions {
if is_mahbot_session_name(&name) {
sweep_session(&name).await;
}
}
}
async fn registered_sessions() -> Option<Vec<String>> {
let status = run_cli_json(&["status"]).await?;
Some(
status
.get("data")?
.get("sessions")?
.as_array()?
.iter()
.filter_map(|s| s.get("name").and_then(Value::as_str).map(String::from))
.collect(),
)
}
async fn wait_for_relay(budget: Duration) {
let deadline = Instant::now() + budget;
while Instant::now() < deadline {
if relay_up().await == Some(true) {
return;
}
tokio::time::sleep(Duration::from_secs(2)).await;
}
}
pub(crate) async fn relay_up() -> Option<bool> {
let status = run_cli_json(&["status"]).await?;
status
.get("data")?
.get("extension")?
.get("relayUp")
.and_then(Value::as_bool)
}
const CHROME_LAUNCH_FLAGS: [&str; 3] = [
"--no-first-run",
"--no-default-browser-check",
"--silent-debugger-extension-api",
];
async fn attempt_chrome_launch() {
if !display_available() {
record_launch_outcome(ChromeLaunchOutcome::NoDisplay);
debug!("chrome daemon: headless host — Chrome launch skipped; start Chrome manually");
return;
}
let gate = { health().lock().unwrap_poison().gate_launch(Instant::now()) };
let RecoveryGate::Allowed(attempt) = gate else {
log_gate_denied(gate, "chrome launch", MAX_LAUNCH_ATTEMPTS);
return;
};
let Some(binary) = chrome_binary() else {
record_launch_outcome(ChromeLaunchOutcome::Failed);
warn!("no Chrome/Chromium binary found — start Chrome manually");
return;
};
if let Err(e) = spawn_chrome_detached(&binary) {
record_launch_outcome(ChromeLaunchOutcome::Failed);
warn!(error = %e, "failed to launch Chrome — start Chrome manually");
return;
}
info!(
attempt,
max = MAX_LAUNCH_ATTEMPTS,
"launched the user's Chrome; waiting for the extension relay to come up"
);
wait_for_relay(RELAY_REVIVE_WAIT).await;
let outcome = evaluate_health().await;
set_health_after_recovery(outcome);
if outcome.is_healthy() {
info!("chrome daemon: relay recovered after Chrome launch");
} else {
record_launch_outcome(ChromeLaunchOutcome::Launched);
warn!(
"Chrome launched but the extension relay is still down — if the ab-connect \
extension is only installed in a non-default profile, install it in Chrome's \
default profile"
);
}
}
pub(crate) fn display_available() -> bool {
#[cfg(target_os = "linux")]
{
std::env::var_os("DISPLAY").is_some() || std::env::var_os("WAYLAND_DISPLAY").is_some()
}
#[cfg(not(target_os = "linux"))]
{
true
}
}
fn chrome_binary() -> Option<PathBuf> {
#[cfg(target_os = "macos")]
{
let home = directories::UserDirs::new().map(|d| d.home_dir().to_path_buf());
let mut candidates = Vec::new();
for app in ["Google Chrome", "Chromium", "Brave Browser"] {
candidates.push(PathBuf::from(format!(
"/Applications/{app}.app/Contents/MacOS/{app}"
)));
if let Some(home) = home.as_deref() {
candidates.push(home.join(format!("Applications/{app}.app/Contents/MacOS/{app}")));
}
}
candidates
.into_iter()
.find(|p| crate::util::is_executable(p))
}
#[cfg(target_os = "linux")]
{
let names = [
"google-chrome",
"google-chrome-stable",
"chromium",
"chromium-browser",
"brave-browser",
"brave",
];
let mut dirs: Vec<PathBuf> = Vec::new();
if let Some(paths) = std::env::var_os("PATH") {
dirs.extend(std::env::split_paths(&paths));
}
dirs.extend([
PathBuf::from("/usr/bin"),
PathBuf::from("/usr/local/bin"),
PathBuf::from("/snap/bin"),
]);
for dir in dirs {
for name in names {
let candidate = dir.join(name);
if crate::util::is_executable(&candidate) {
return Some(candidate);
}
}
}
None
}
#[cfg(target_os = "windows")]
{
let mut candidates = Vec::new();
for base in ["ProgramFiles", "ProgramFiles(x86)", "LOCALAPPDATA"] {
let Some(base) = std::env::var_os(base) else {
continue;
};
let base = PathBuf::from(base);
for rel in [
Path::new("Google/Chrome/Application/chrome.exe"),
Path::new("Chromium/Application/chrome.exe"),
Path::new("BraveSoftware/Brave-Browser/Application/brave.exe"),
] {
candidates.push(base.join(rel));
}
}
candidates
.into_iter()
.find(|p| crate::util::is_executable(p))
}
#[cfg(not(any(target_os = "macos", target_os = "linux", target_os = "windows")))]
{
None
}
}
fn spawn_chrome_detached(binary: &Path) -> std::io::Result<()> {
let mut cmd = std::process::Command::new(binary);
cmd.args(CHROME_LAUNCH_FLAGS)
.stdin(Stdio::null())
.stdout(Stdio::null())
.stderr(Stdio::null());
#[cfg(unix)]
{
use std::os::unix::process::CommandExt;
cmd.process_group(0);
}
#[cfg(windows)]
{
use std::os::windows::process::CommandExt;
const DETACHED_PROCESS: u32 = 0x0000_0008;
const CREATE_NO_WINDOW: u32 = 0x0800_0000;
cmd.creation_flags(DETACHED_PROCESS | CREATE_NO_WINDOW);
}
cmd.spawn().map(|_| ())
}
fn warn_transition(failure: ProbeFailure) {
let mut h = health().lock().unwrap_poison();
if h.last_cause_warned == Some(failure) {
return;
}
h.last_cause_warned = Some(failure);
match failure {
ProbeFailure::NotInstalled => warn!(
"chrome-use extension or native host is not installed — the browser \
daemon cannot run. Enable the chrome-use extension at \
chrome://extensions (the CLI re-installs itself in the background at \
startup). Auto-recovery paused until it is installed."
),
ProbeFailure::HostBroken => warn!(
"chrome-use native host launcher is broken — run `chrome-use doctor`. \
Auto-recovery paused until it is fixed."
),
ProbeFailure::ExtensionDisabled => warn!(
"chrome-use extension is disabled — enable it at chrome://extensions. \
Daemon restarts cannot fix a Chrome-side disable; auto-recovery paused \
until it is enabled."
),
ProbeFailure::ExtensionAbsent => warn!(
"chrome-use extension is not installed in Chrome — install/enable the \
ab-connect extension from the Chrome Web Store. Auto-recovery paused \
until it is installed."
),
ProbeFailure::RelayDown => warn!(
"chrome-use extension relay is down (the extension is enabled) — waiting \
for the extension to reconnect and restarting session daemons to clear \
stale relay bindings."
),
ProbeFailure::ChromeNotRunning => {
if !display_available() {
warn!(
"Chrome is not running and no display is available (headless host) — \
start Chrome manually; auto-recovery paused."
);
} else if h.launch_budget.halted {
warn!(
"Chrome is not running — launch attempts are paused in a thrash-protection \
cooldown; start Chrome manually."
);
} else {
warn!(
"Chrome is not running — auto-recovery will launch it (bounded launch budget)."
);
}
}
ProbeFailure::UnreachableTab => warn!(
"a browser tab the session was driving is unreachable (the extension lost its \
debugger attach; about:blank tabs are never re-attached) — close the leftover \
tab in Chrome to unblock the session"
),
ProbeFailure::DaemonWedge => {
warn!("chrome daemon is unresponsive — restarting it (bounded backoff).");
}
}
}
fn log_gate_denied(gate: RecoveryGate, noun: &str, max: u32) {
match gate {
RecoveryGate::Halted => error!(
attempts = max,
"chrome daemon: {max} consecutive failed {noun} attempts; \
auto-recovery halted for 30 min (thrash protection)"
),
RecoveryGate::Backoff => {
debug!("chrome daemon: still down; waiting out {noun} backoff");
}
RecoveryGate::Cooldown => {
debug!("chrome daemon: still down; {noun} cooldown in progress (thrash protection)");
}
RecoveryGate::Allowed(_) => unreachable!(),
}
}
async fn attempt_recovery(mut failure: ProbeFailure) {
warn_transition(failure);
if failure.is_unfixable() {
return;
}
if failure == ProbeFailure::ChromeNotRunning {
attempt_chrome_launch().await;
return;
}
let now = Instant::now();
let throttled = health()
.lock()
.unwrap_poison()
.restart_budget
.is_waiting(now);
if throttled {
return;
}
if failure == ProbeFailure::RelayDown {
wait_for_relay(RELAY_REVIVE_WAIT).await;
let outcome = evaluate_health().await;
set_health(outcome);
match outcome {
ProbeOutcome::Healthy => {
info!("chrome daemon: relay recovered without a restart");
return;
}
ProbeOutcome::Down(f) => {
warn_transition(f);
if f.is_unfixable() {
return;
}
if f == ProbeFailure::ChromeNotRunning {
attempt_chrome_launch().await;
return;
}
failure = f;
}
}
}
let gate = {
let mut h = health().lock().unwrap_poison();
h.gate_restart(Instant::now())
};
let RecoveryGate::Allowed(attempt) = gate else {
log_gate_denied(gate, "restart", MAX_RESTART_ATTEMPTS);
return;
};
info!(
attempt,
max = MAX_RESTART_ATTEMPTS,
"chrome daemon: attempting auto-recovery"
);
let _ = run_cli(&["daemon", "restart"]).await;
if failure == ProbeFailure::RelayDown {
wait_for_relay(RELAY_REVIVE_WAIT).await;
}
let outcome = evaluate_health().await;
set_health_after_recovery(outcome);
if outcome.is_healthy() {
info!("chrome daemon: recovered after restart");
} else {
warn!(
attempt,
"chrome daemon: restart attempt did not restore health"
);
}
}
async fn run_cli(args: &[&str]) -> bool {
let Some(path) = cli_path() else {
return false;
};
let mut cmd = Command::new(path);
ensure_chrome_env(&mut cmd);
cmd.args(args)
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null());
cmd.kill_on_drop(true);
tokio::time::timeout(Duration::from_mins(1), cmd.status())
.await
.is_ok_and(|r| r.is_ok_and(|st| st.success()))
}
#[cfg(test)]
pub(crate) async fn with_health_test_lock() -> tokio::sync::MutexGuard<'static, ()> {
static LOCK: OnceLock<tokio::sync::Mutex<()>> = OnceLock::new();
LOCK.get_or_init(|| tokio::sync::Mutex::new(()))
.lock()
.await
}
#[cfg(test)]
pub(crate) fn reset_health() {
*health().lock().unwrap_poison() = DaemonHealth::default();
}
#[cfg(test)]
mod tests {
use super::*;
use crate::chrome::contract::is_daemon_unavailable_code;
#[tokio::test]
async fn advertisement_and_availability_reflect_daemon_state() {
let _guard = with_health_test_lock().await;
set_health(ProbeOutcome::Down(ProbeFailure::DaemonWedge));
assert!(!is_advertised());
assert!(!is_available().await);
set_health(ProbeOutcome::Healthy);
assert!(is_advertised());
assert!(is_available().await);
reset_health();
assert!(is_advertised());
}
#[tokio::test]
async fn daemon_down_message_reflects_launch_state() {
let _guard = with_health_test_lock().await;
reset_health();
set_health(ProbeOutcome::Down(ProbeFailure::RelayDown));
let msg = daemon_down_message();
assert!(msg.contains("relay is down"));
assert!(!msg.contains("auto-launched"));
reset_health();
set_health(ProbeOutcome::Down(ProbeFailure::RelayDown));
record_launch_outcome(ChromeLaunchOutcome::Launched);
let msg = daemon_down_message();
assert!(msg.contains("even after Chrome was auto-launched"));
assert!(msg.contains("non-default profile"));
reset_health();
set_health(ProbeOutcome::Down(ProbeFailure::ChromeNotRunning));
let msg = daemon_down_message();
assert!(msg.contains("Chrome is not running."));
assert!(msg.contains("will launch Chrome"));
reset_health();
set_health(ProbeOutcome::Down(ProbeFailure::ChromeNotRunning));
record_launch_outcome(ChromeLaunchOutcome::Failed);
let msg = daemon_down_message();
assert!(msg.contains("auto-launch attempt failed"));
assert!(msg.contains("start Chrome manually"));
assert!(msg.contains("retry the launch with backoff"));
reset_health();
set_health(ProbeOutcome::Down(ProbeFailure::ChromeNotRunning));
record_launch_outcome(ChromeLaunchOutcome::Failed);
health().lock().unwrap_poison().launch_budget.halted = true;
let msg = daemon_down_message();
assert!(msg.contains("exhausted its launch attempts"));
assert!(!msg.contains("retry the launch with backoff"));
reset_health();
set_health(ProbeOutcome::Down(ProbeFailure::ChromeNotRunning));
record_launch_outcome(ChromeLaunchOutcome::NoDisplay);
let msg = daemon_down_message();
assert!(msg.contains("no display"));
assert!(msg.contains("paused for this cause"));
reset_health();
set_health(ProbeOutcome::Down(ProbeFailure::ExtensionAbsent));
let msg = daemon_down_message();
assert!(msg.contains("Chrome Web Store"));
assert!(msg.contains("no restart will be attempted"));
reset_health();
set_health(ProbeOutcome::Down(ProbeFailure::DaemonWedge));
let msg = daemon_down_message();
assert!(msg.contains("down or unresponsive"));
assert!(msg.contains("restart it automatically"));
reset_health();
}
#[test]
fn sustained_health_resets_restart_attempts() {
let now = Instant::now();
let mut h = DaemonHealth {
restart_budget: AttemptBudget {
attempts: 2,
next_at: Some(now),
halted: true,
halted_until: Some(now),
},
..DaemonHealth::default()
};
h.apply_outcome(ProbeOutcome::Healthy, now, false);
assert_eq!(h.restart_budget.attempts, 2);
assert!(h.restart_budget.halted);
h.apply_outcome(ProbeOutcome::Healthy, now + WATCHDOG_INTERVAL, true);
assert_eq!(h.restart_budget.attempts, 2);
assert!(h.healthy_since.is_some());
h.apply_outcome(ProbeOutcome::Healthy, now + WATCHDOG_INTERVAL * 2, true);
assert_eq!(h.restart_budget.attempts, 2);
assert!(h.restart_budget.halted);
h.apply_outcome(ProbeOutcome::Healthy, now + WATCHDOG_INTERVAL * 3, true);
assert_eq!(h.restart_budget.attempts, 0);
assert_eq!(h.restart_budget.next_at, None);
assert!(!h.restart_budget.halted);
assert!(h.restart_budget.halted_until.is_none());
assert_eq!(h.last_failure, None);
}
#[test]
fn cause_flapping_and_transient_health_do_not_reset_restart_budget() {
let now = Instant::now();
let mut h = DaemonHealth {
restart_budget: AttemptBudget {
attempts: 2,
next_at: Some(now),
..AttemptBudget::default()
},
last_failure: Some(ProbeFailure::DaemonWedge),
..DaemonHealth::default()
};
h.apply_outcome(ProbeOutcome::Down(ProbeFailure::RelayDown), now, true);
assert_eq!(h.last_failure, Some(ProbeFailure::RelayDown));
assert_eq!(h.restart_budget.attempts, 2);
assert!(h.restart_budget.next_at.is_some());
h.apply_outcome(ProbeOutcome::Down(ProbeFailure::DaemonWedge), now, true);
h.apply_outcome(ProbeOutcome::Down(ProbeFailure::RelayDown), now, true);
assert_eq!(h.last_failure, Some(ProbeFailure::RelayDown));
assert_eq!(h.restart_budget.attempts, 2);
h.apply_outcome(ProbeOutcome::Healthy, now, true);
assert_eq!(h.restart_budget.attempts, 2);
assert!(h.restart_budget.next_at.is_some());
h.apply_outcome(ProbeOutcome::Healthy, now + SUSTAINED_HEALTHY_WINDOW, true);
assert_eq!(h.last_failure, None);
assert_eq!(h.restart_budget.attempts, 0);
assert_eq!(h.restart_budget.next_at, None);
assert!(!h.restart_budget.halted);
}
#[test]
fn gate_honors_backoff_before_halt() {
let now = Instant::now();
let mut h = DaemonHealth::default();
assert_eq!(h.gate_restart(now), RecoveryGate::Allowed(1));
assert_eq!(h.gate_restart(now), RecoveryGate::Backoff);
assert_eq!(
h.gate_restart(now + RESTART_BACKOFF[0]),
RecoveryGate::Allowed(2)
);
let t2 = now + RESTART_BACKOFF[0] + RESTART_BACKOFF[1];
assert_eq!(h.gate_restart(t2), RecoveryGate::Allowed(3));
assert_eq!(h.gate_restart(t2), RecoveryGate::Backoff);
let t3 = t2 + RESTART_BACKOFF[2];
assert_eq!(h.gate_restart(t3), RecoveryGate::Halted);
assert!(h.restart_budget.halted);
assert_eq!(h.gate_restart(t3), RecoveryGate::Cooldown);
assert_eq!(h.gate_restart(t3 + HALT_COOLDOWN), RecoveryGate::Allowed(1));
assert_eq!(h.restart_budget.attempts, 1);
assert!(!h.restart_budget.halted);
}
#[test]
fn relay_down_classification_precedence() {
assert_eq!(
classify_relay_down(ExtensionState::Disabled, Some(true)),
ProbeFailure::ExtensionDisabled
);
assert_eq!(
classify_relay_down(ExtensionState::Disabled, Some(false)),
ProbeFailure::ExtensionDisabled
);
assert_eq!(
classify_relay_down(ExtensionState::Disabled, None),
ProbeFailure::ExtensionDisabled
);
assert_eq!(
classify_relay_down(ExtensionState::Absent, Some(false)),
ProbeFailure::ExtensionAbsent
);
assert_eq!(
classify_relay_down(ExtensionState::Absent, None),
ProbeFailure::ExtensionAbsent
);
assert_eq!(
classify_relay_down(ExtensionState::Present, Some(false)),
ProbeFailure::ChromeNotRunning
);
assert_eq!(
classify_relay_down(ExtensionState::Present, Some(true)),
ProbeFailure::RelayDown
);
assert_eq!(
classify_relay_down(ExtensionState::Present, None),
ProbeFailure::RelayDown
);
assert_eq!(
classify_relay_down(ExtensionState::Unknown, Some(false)),
ProbeFailure::ChromeNotRunning
);
assert_eq!(
classify_relay_down(ExtensionState::Unknown, None),
ProbeFailure::RelayDown
);
assert!(!ProbeFailure::ChromeNotRunning.is_unfixable());
assert!(ProbeFailure::ExtensionAbsent.is_unfixable());
}
#[test]
fn extension_state_parsing() {
assert_eq!(
extension_state_from(&serde_json::json!({})),
ExtensionState::Unknown
);
assert_eq!(
extension_state_from(&serde_json::json!({ "data": {} })),
ExtensionState::Unknown
);
assert_eq!(
extension_state_from(&serde_json::json!({ "data": { "chromeExtension": null } })),
ExtensionState::Absent
);
assert_eq!(
extension_state_from(&serde_json::json!({
"data": { "chromeExtension": { "disableReasons": [] } }
})),
ExtensionState::Present
);
assert_eq!(
extension_state_from(&serde_json::json!({
"data": { "chromeExtension": { "enabled": true } }
})),
ExtensionState::Present
);
assert_eq!(
extension_state_from(&serde_json::json!({
"data": { "chromeExtension": { "disableReasons": ["user"] } }
})),
ExtensionState::Disabled
);
}
#[test]
fn launch_budget_is_independent_of_restart_budget() {
let now = Instant::now();
let mut h = DaemonHealth {
restart_budget: AttemptBudget {
attempts: MAX_RESTART_ATTEMPTS,
halted: true,
halted_until: Some(now),
..AttemptBudget::default()
},
..DaemonHealth::default()
};
assert!(matches!(h.gate_launch(now), RecoveryGate::Allowed(_)));
assert_eq!(h.launch_budget.attempts, 1);
let mut h2 = DaemonHealth {
launch_budget: AttemptBudget {
attempts: MAX_LAUNCH_ATTEMPTS,
..AttemptBudget::default()
},
..DaemonHealth::default()
};
assert!(!matches!(h2.gate_launch(now), RecoveryGate::Allowed(_)));
assert!(matches!(h2.gate_restart(now), RecoveryGate::Allowed(_)));
let mut h3 = DaemonHealth {
restart_budget: AttemptBudget {
attempts: 2,
next_at: Some(now),
halted: true,
halted_until: Some(now),
},
launch_budget: AttemptBudget {
attempts: 2,
next_at: Some(now),
halted: true,
halted_until: Some(now),
},
launch_outcome: Some(ChromeLaunchOutcome::Failed),
..DaemonHealth::default()
};
h3.apply_outcome(ProbeOutcome::Healthy, now + WATCHDOG_INTERVAL, true);
assert_eq!(h3.restart_budget.attempts, 2);
assert_eq!(h3.launch_budget.attempts, 2);
assert!(h3.restart_budget.halted);
assert!(h3.launch_budget.halted);
assert_eq!(h3.launch_outcome, None);
h3.apply_outcome(ProbeOutcome::Healthy, now + WATCHDOG_INTERVAL * 3, true);
assert_eq!(h3.restart_budget.attempts, 0);
assert_eq!(h3.launch_budget.attempts, 0);
assert!(!h3.restart_budget.halted);
assert!(!h3.launch_budget.halted);
}
#[test]
fn daemon_unavailable_error_signature_detected() {
for msg in [
"Failed to read: Resource temporarily unavailable (os error 35) (after 5 retries - daemon may be busy or unresponsive)",
"Failed to connect: No such file or directory (os error 2) (after 5 retries - daemon may be busy or unresponsive)",
"session unresponsive: no response within 45s",
"Daemon failed to start (socket: /tmp/x.sock)",
"session unresponsive: the stuck '__mahbot_probe' daemon was stopped automatically",
"Failed to connect: the daemon endpoint for session '__mahbot_probe' disappeared (/tmp/x.sock).",
"CDP session is unresponsive after attaching (Connection reset).",
"Auto-launch failed: Could not drive your Chrome through the ab-connect extension.",
] {
assert!(is_daemon_unavailable_error(msg), "should detect: {msg}");
}
for msg in [
"chrome-use error: Element not found",
"chrome-use error: Evaluation error: ReferenceError",
"chrome-use error: Navigation failed",
"Failed to connect to example.com: Connection timed out",
] {
assert!(
!is_daemon_unavailable_error(msg),
"should NOT detect: {msg}"
);
}
}
#[test]
fn relay_unavailable_signature_detected() {
for msg in [
"The chrome-use extension is installed, but its relay isn't connected.",
"Could not drive your Chrome through the ab-connect extension.",
"Chrome relay dropped — reconnecting…",
] {
assert!(is_relay_unavailable_error(msg), "should detect: {msg}");
}
for msg in [
"chrome-use error: Element not found",
"Failed to read: Resource temporarily unavailable (os error 35)",
] {
assert!(!is_relay_unavailable_error(msg), "should NOT detect: {msg}");
}
}
#[test]
fn daemon_unavailable_code_detected() {
assert!(is_daemon_unavailable_code(Some("browser_not_launched")));
assert!(!is_daemon_unavailable_code(Some("connection_failed")));
assert!(!is_daemon_unavailable_code(Some("timeout")));
assert!(!is_daemon_unavailable_code(Some("element_not_found")));
assert!(!is_daemon_unavailable_code(None));
}
#[test]
fn failure_text_classification_is_shared_between_detection_paths() {
assert_eq!(
classify_failure_text(
"Auto-launch failed: Could not drive your Chrome through the ab-connect \
extension. The tab this session was driving can no longer be resolved (it \
was closed, or a flaky relay dropped it)"
),
Some(ProbeFailure::UnreachableTab)
);
assert_eq!(
classify_failure_text(
"Auto-launch failed: Could not drive your Chrome through the ab-connect extension."
),
Some(ProbeFailure::RelayDown)
);
assert_eq!(
classify_failure_text(
"Failed to connect: the daemon endpoint for session '__mahbot_probe' \
disappeared (/tmp/x.sock)."
),
Some(ProbeFailure::DaemonWedge)
);
assert_eq!(
classify_failure_text("chrome-use error: Element not found"),
None
);
}
#[test]
fn spawn_error_classification_distinguishes_missing_from_transient() {
let not_found = std::io::Error::from(std::io::ErrorKind::NotFound);
assert_eq!(classify_spawn_error(¬_found), CliStatus::Missing);
for kind in [
std::io::ErrorKind::WouldBlock, std::io::ErrorKind::OutOfMemory, std::io::ErrorKind::PermissionDenied, std::io::ErrorKind::StorageFull, std::io::ErrorKind::TimedOut,
] {
let err = std::io::Error::from(kind);
assert!(
matches!(
classify_spawn_error(&err),
CliStatus::Transient(CliProbeFailure::Spawn(_))
),
"kind {kind:?} must classify as transient, not missing"
);
}
}
#[test]
fn cli_version_parsing_scans_the_real_banner() {
let banner = "chrome-use 1.5.100\n\
report bugs / rough edges: https://github.com/leeguooooo/chrome-use/issues";
assert_eq!(
parse_cli_version(banner),
Some(semver::Version::new(1, 5, 100))
);
assert_eq!(
parse_cli_version("chrome-use v1.5.99"),
Some(semver::Version::new(1, 5, 99))
);
assert_eq!(parse_cli_version("chrome-use\nno version here"), None);
}
#[test]
fn release_asset_platform_maps_supported_combos() {
assert_eq!(
release_asset_platform("macos", "x86_64", false).as_deref(),
Some("darwin-x64")
);
assert_eq!(
release_asset_platform("macos", "aarch64", false).as_deref(),
Some("darwin-arm64")
);
assert_eq!(
release_asset_platform("linux", "x86_64", false).as_deref(),
Some("linux-x64")
);
assert_eq!(
release_asset_platform("linux", "aarch64", false).as_deref(),
Some("linux-arm64")
);
assert_eq!(
release_asset_platform("linux", "x86_64", true).as_deref(),
Some("linux-musl-x64")
);
assert_eq!(
release_asset_platform("linux", "aarch64", true).as_deref(),
Some("linux-musl-arm64")
);
assert_eq!(
release_asset_platform("windows", "x86_64", true).as_deref(),
Some("win32-x64")
);
assert_eq!(
release_asset_platform("macos", "aarch64", true).as_deref(),
Some("darwin-arm64")
);
assert_eq!(release_asset_platform("windows", "aarch64", false), None);
assert_eq!(release_asset_platform("freebsd", "x86_64", false), None);
}
}