use super::chrome_daemon::{
NoVerdict, SESSION_RECOVERY_TIMEOUT, cli_path, recover_unresponsive_session_via,
run_cli_json_at,
};
use crate::chrome::contract::is_session_unresponsive_error;
use crate::util::UnwrapPoison;
use futures_util::StreamExt;
use serde::{Deserialize, Serialize};
use std::collections::{HashMap, VecDeque};
use std::path::PathBuf;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Mutex, OnceLock};
use std::time::{Duration, Instant};
use tokio_util::sync::CancellationToken;
use tracing::{debug, info};
const RELEASE_CONCURRENCY: usize = 8;
const RELEASE_MAX_ATTEMPTS: u32 = 5;
const RELEASE_MAX_SKIPS: u32 = 8;
const RELEASE_RETRY_BASE: Duration = Duration::from_secs(5);
const RELEASE_RETRY_CAP: Duration = Duration::from_mins(1);
pub(crate) const RELEASE_HOLD_REDRIVEN_RUN: Duration = Duration::from_mins(30);
pub(crate) fn holds_run_end(classification: &str) -> bool {
matches!(classification, "drain" | "shutdown" | "pause")
}
const RELEASE_BOOT_GRACE: Duration = Duration::from_mins(1);
const RELEASE_FILE_NAME: &str = "chrome-run-releases.json";
const MAX_PENDING_RELEASES: usize = 64;
const SHUTDOWN_RELEASE_FLUSH_BUDGET: Duration = Duration::from_secs(10);
#[derive(Clone)]
struct PendingRunRelease {
namespace: String,
names: Vec<String>,
attempts: u32,
skips: u32,
next_attempt_at: Instant,
held: bool,
}
#[derive(Serialize, Deserialize)]
struct PersistedRunRelease {
namespace: String,
names: Vec<String>,
#[serde(default)]
attempts: u32,
#[serde(default)]
skips: u32,
#[serde(default)]
next_attempt_at: u64,
#[serde(default)]
held: bool,
}
static PENDING_RELEASES: OnceLock<Mutex<VecDeque<PendingRunRelease>>> = OnceLock::new();
static PARKED_RELEASES: OnceLock<Mutex<Vec<(u64, PendingRunRelease)>>> = OnceLock::new();
static NEXT_PASS_ID: AtomicU64 = AtomicU64::new(0);
static RELEASE_WAKE: OnceLock<tokio::sync::Notify> = OnceLock::new();
static LIVE_RUN_NAMESPACES: OnceLock<Mutex<HashMap<String, usize>>> = OnceLock::new();
fn pending_releases() -> &'static Mutex<VecDeque<PendingRunRelease>> {
PENDING_RELEASES.get_or_init(|| Mutex::new(VecDeque::new()))
}
fn parked_releases() -> &'static Mutex<Vec<(u64, PendingRunRelease)>> {
PARKED_RELEASES.get_or_init(|| Mutex::new(Vec::new()))
}
fn release_wake() -> &'static tokio::sync::Notify {
RELEASE_WAKE.get_or_init(tokio::sync::Notify::new)
}
fn live_run_namespaces() -> &'static Mutex<HashMap<String, usize>> {
LIVE_RUN_NAMESPACES.get_or_init(|| Mutex::new(HashMap::new()))
}
pub(crate) fn register_run_namespace(namespace: &str) {
if namespace.is_empty() {
return;
}
*live_run_namespaces()
.lock()
.unwrap_poison()
.entry(namespace.to_string())
.or_insert(0) += 1;
}
pub(crate) fn unregister_run_namespace(namespace: &str) {
if namespace.is_empty() {
return;
}
{
let mut live = live_run_namespaces().lock().unwrap_poison();
match live.get_mut(namespace) {
Some(count) if *count > 1 => *count -= 1,
Some(_) => {
live.remove(namespace);
}
None => return, }
}
release_wake().notify_one();
}
fn run_namespace_is_live(namespace: &str) -> bool {
!namespace.is_empty()
&& live_run_namespaces()
.lock()
.unwrap_poison()
.contains_key(namespace)
}
#[derive(Clone, Copy)]
enum Eligibility {
RunEnd,
Retry,
}
impl Eligibility {
fn keep(self, queued: (Instant, bool), incoming: (Instant, bool)) -> (Instant, bool) {
match self {
Self::RunEnd => incoming,
Self::Retry if incoming.0 > queued.0 => incoming,
Self::Retry => queued,
}
}
}
fn absorb(entry: &mut PendingRunRelease, incoming: PendingRunRelease, keep: Eligibility) {
(entry.next_attempt_at, entry.held) = keep.keep(
(entry.next_attempt_at, entry.held),
(incoming.next_attempt_at, incoming.held),
);
entry.attempts = entry.attempts.max(incoming.attempts);
entry.skips = entry.skips.max(incoming.skips);
for name in incoming.names {
if !entry.names.contains(&name) {
entry.names.push(name);
}
}
}
fn merge_record(entry: PendingRunRelease, eligibility: Eligibility) {
if entry.namespace.is_empty() || entry.names.is_empty() {
return;
}
let mut queue = pending_releases().lock().unwrap_poison();
if let Some(existing) = queue.iter_mut().find(|e| e.namespace == entry.namespace) {
absorb(existing, entry, eligibility);
return;
}
let mut dropped = 0usize;
while queue.len() >= MAX_PENDING_RELEASES {
queue.pop_front();
dropped += 1;
}
if dropped > 0 {
info!(
dropped,
"agent-run chrome release queue full — oldest records dropped"
);
}
queue.push_back(entry);
}
pub(crate) fn queue_run_session_release(sessions: &super::chrome::ChromeRunSessions, held: bool) {
queue_run_session_release_after(
sessions,
if held {
RELEASE_HOLD_REDRIVEN_RUN
} else {
Duration::ZERO
},
);
}
fn queue_run_session_release_after(sessions: &super::chrome::ChromeRunSessions, hold: Duration) {
let names = sessions.snapshot();
if names.is_empty() {
return; }
merge_record(
PendingRunRelease {
namespace: sessions.namespace().to_string(),
names,
attempts: 0,
skips: 0,
next_attempt_at: Instant::now() + hold,
held: !hold.is_zero(),
},
Eligibility::RunEnd,
);
persist_pending_releases();
release_wake().notify_one();
}
fn restore_pending_releases() {
let Some(path) = release_settings().store else {
return;
};
let Ok(json) = std::fs::read_to_string(&path) else {
return;
};
let records: Vec<PersistedRunRelease> = match serde_json::from_str(&json) {
Ok(records) => records,
Err(error) => {
info!(
path = %path.display(),
%error,
"agent-run chrome release file unreadable — ignoring it"
);
return;
}
};
let mut restored = 0usize;
for record in records.into_iter().take(MAX_PENDING_RELEASES) {
if record.names.is_empty() {
continue;
}
merge_record(
PendingRunRelease {
namespace: record.namespace,
names: record.names,
attempts: record.attempts,
skips: record.skips,
next_attempt_at: restored_eligibility(record.next_attempt_at, record.held),
held: record.held,
},
Eligibility::Retry,
);
restored += 1;
}
if restored > 0 {
debug!(
restored,
"agent-run chrome releases restored from the previous process"
);
}
}
fn persist_pending_releases() {
let Some(path) = release_settings().store else {
return;
};
let parked = parked_releases().lock().unwrap_poison();
let queue = pending_releases().lock().unwrap_poison();
let mut folded: Vec<PendingRunRelease> = Vec::new();
for record in queue.iter().chain(parked.iter().map(|(_, record)| record)) {
match folded
.iter_mut()
.find(|entry| entry.namespace == record.namespace)
{
Some(entry) => absorb(entry, record.clone(), Eligibility::Retry),
None => folded.push(record.clone()),
}
}
let records: Vec<PersistedRunRelease> = folded
.iter()
.filter(|entry| !entry.names.is_empty())
.map(|entry| PersistedRunRelease {
namespace: entry.namespace.clone(),
names: entry.names.clone(),
attempts: entry.attempts,
skips: entry.skips,
next_attempt_at: eligibility_deadline(entry.next_attempt_at),
held: entry.held,
})
.collect();
let json = serde_json::to_string(&records).unwrap_or_default();
let tmp = path.with_extension("json.tmp");
if let Some(parent) = path.parent() {
let _ = std::fs::create_dir_all(parent);
}
let _ = std::fs::write(&tmp, json);
let _ = std::fs::rename(&tmp, &path);
}
fn epoch_secs() -> u64 {
crate::util::unix_millis() / 1000
}
fn eligibility_deadline(next_attempt_at: Instant) -> u64 {
epoch_secs().saturating_add(
next_attempt_at
.saturating_duration_since(Instant::now())
.as_secs(),
)
}
fn restored_eligibility(deadline: u64, held: bool) -> Instant {
if held {
return Instant::now() + RELEASE_HOLD_REDRIVEN_RUN;
}
let remaining =
Duration::from_secs(deadline.saturating_sub(epoch_secs())).min(RELEASE_RETRY_CAP);
Instant::now() + release_settings().boot_grace.max(remaining)
}
pub async fn run_session_release_queue() {
release_queue(crate::shutdown::shutdown_token()).await;
}
async fn release_queue(shutdown: CancellationToken) {
restore_pending_releases();
loop {
tokio::select! {
() = sleep_until_due(next_release_deadline()) => {}
() = release_wake().notified() => {}
() = shutdown.cancelled() => break,
}
release_due().await;
}
}
async fn sleep_until_due(deadline: Option<Instant>) {
match deadline {
Some(at) => tokio::time::sleep_until(tokio::time::Instant::from_std(at)).await,
None => std::future::pending::<()>().await,
}
}
fn next_release_deadline() -> Option<Instant> {
pending_releases()
.lock()
.unwrap_poison()
.iter()
.filter(|entry| !run_namespace_is_live(&entry.namespace))
.map(|entry| entry.next_attempt_at)
.min()
}
fn take_releasable() -> Vec<PendingRunRelease> {
let now = Instant::now();
let mut queue = pending_releases().lock().unwrap_poison();
let mut taken = Vec::new();
let mut waiting = VecDeque::with_capacity(queue.len());
while let Some(entry) = queue.pop_front() {
let releasable = !run_namespace_is_live(&entry.namespace) && entry.next_attempt_at <= now;
if releasable {
taken.push(entry);
} else {
waiting.push_back(entry);
}
}
*queue = waiting;
taken
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum ReleaseOutcome {
Released,
Skipped(SkipReason),
Failed,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum SkipReason {
NoBinary,
Untried,
}
async fn attempt_releases(
records: &[PendingRunRelease],
per_name: impl Fn() -> Duration,
) -> Vec<(usize, String, ReleaseOutcome)> {
let max_names = records
.iter()
.map(|record| record.names.len())
.max()
.unwrap_or(0);
let mut jobs: Vec<(usize, String)> = Vec::new();
for name_index in 0..max_names {
for (index, record) in records.iter().enumerate() {
if let Some(name) = record.names.get(name_index) {
jobs.push((index, name.clone()));
}
}
}
futures_util::stream::iter(jobs)
.map(|(index, name)| {
let timeout = per_name();
async move {
if run_namespace_is_live(&records[index].namespace) {
return (index, name, ReleaseOutcome::Skipped(SkipReason::Untried));
}
if timeout.is_zero() {
return (index, name, ReleaseOutcome::Skipped(SkipReason::Untried));
}
let outcome = release_one(&name, timeout).await;
(index, name, outcome)
}
})
.buffer_unordered(RELEASE_CONCURRENCY)
.collect()
.await
}
struct PassResult {
open: Vec<String>,
attempted: bool,
unavailable: bool,
}
fn results_by_record(
records: &[PendingRunRelease],
outcomes: Vec<(usize, String, ReleaseOutcome)>,
) -> Vec<PassResult> {
let mut results: Vec<PassResult> = records
.iter()
.map(|_| PassResult {
open: Vec::new(),
attempted: false,
unavailable: false,
})
.collect();
for (index, name, outcome) in outcomes {
match outcome {
ReleaseOutcome::Released => {}
ReleaseOutcome::Failed => {
results[index].open.push(name);
results[index].attempted = true;
}
ReleaseOutcome::Skipped(SkipReason::NoBinary) => {
results[index].open.push(name);
results[index].unavailable = true;
}
ReleaseOutcome::Skipped(_) => results[index].open.push(name),
}
}
results
}
struct ReleasePass {
id: u64,
records: Vec<PendingRunRelease>,
in_flight: Option<String>,
reconciled: bool,
}
impl ReleasePass {
fn take() -> Option<Self> {
let mut parked = parked_releases().lock().unwrap_poison();
let records = take_releasable();
if records.is_empty() {
return None;
}
let id = NEXT_PASS_ID.fetch_add(1, Ordering::Relaxed);
parked.extend(records.iter().cloned().map(|record| (id, record)));
Some(Self {
id,
records,
in_flight: None,
reconciled: false,
})
}
fn reconcile(
&mut self,
outcomes: Vec<(usize, String, ReleaseOutcome)>,
mut finish: impl FnMut(PendingRunRelease, PassResult),
) {
let mut results = results_by_record(&self.records, outcomes);
while let Some((entry, result)) = self.records.pop().zip(results.pop()) {
self.in_flight = Some(entry.namespace.clone());
finish(entry, result);
self.in_flight = None;
}
self.reconciled = true;
self.unpark(None);
}
fn unpark(&self, keep: Option<&str>) {
parked_releases()
.lock()
.unwrap_poison()
.retain(|(pass, record)| *pass != self.id || keep == Some(record.namespace.as_str()));
}
}
impl Drop for ReleasePass {
fn drop(&mut self) {
if self.reconciled {
return;
}
for entry in self.records.drain(..) {
merge_record(entry, Eligibility::Retry);
}
let in_flight = self.in_flight.take();
self.unpark(in_flight.as_deref());
persist_pending_releases();
}
}
async fn release_due() {
let Some(mut pass) = ReleasePass::take() else {
return;
};
let timeout = release_settings().attempt_timeout;
let outcomes = attempt_releases(&pass.records, || timeout).await;
pass.reconcile(outcomes, |mut entry, result| {
if result.open.is_empty() {
debug!(
namespace = %entry.namespace,
sessions = entry.names.len(),
"agent-run chrome sessions released"
);
return;
}
if !result.attempted {
entry.names = result.open;
if result.unavailable {
entry.skips = entry.skips.saturating_add(1);
entry.next_attempt_at = Instant::now() + release_backoff(entry.skips);
if entry.skips >= RELEASE_MAX_SKIPS {
info!(
namespace = %entry.namespace,
sessions = entry.names.len(),
skips = entry.skips,
"giving up on releasing agent-run chrome sessions — chrome-use was never available to run"
);
return;
}
}
merge_record(entry, Eligibility::Retry);
return;
}
entry.attempts = entry.attempts.saturating_add(1);
entry.next_attempt_at = Instant::now() + release_backoff(entry.attempts);
if entry.attempts >= RELEASE_MAX_ATTEMPTS {
info!(
namespace = %entry.namespace,
sessions = result.open.len(),
attempts = entry.attempts,
"giving up on releasing agent-run chrome sessions — their tabs stay open in the browser"
);
return;
}
debug!(
namespace = %entry.namespace,
sessions = result.open.len(),
attempts = entry.attempts,
"agent-run chrome session release failed — retrying"
);
entry.names = result.open;
merge_record(entry, Eligibility::Retry);
});
persist_pending_releases();
}
async fn flush_pending_run_releases(budget: Duration) {
let deadline = Instant::now() + budget;
let Some(mut pass) = ReleasePass::take() else {
return;
};
let outcomes = attempt_releases(&pass.records, || {
deadline.saturating_duration_since(Instant::now())
})
.await;
let mut released = 0usize;
let mut still_open = 0usize;
pass.reconcile(outcomes, |mut entry, result| {
released += entry.names.len() - result.open.len();
if result.open.is_empty() {
return;
}
still_open += result.open.len();
entry.names = result.open;
merge_record(entry, Eligibility::Retry);
});
persist_pending_releases();
debug!(
released,
still_open, "shutdown: last-chance agent-run chrome session release"
);
}
pub async fn flush_and_close_all_chrome_sessions() {
flush_pending_run_releases(SHUTDOWN_RELEASE_FLUSH_BUDGET).await;
crate::tools::chrome::close_all_chrome_sessions().await;
}
async fn release_one(name: &str, timeout: Duration) -> ReleaseOutcome {
let Some(cli) = release_cli() else {
debug!(
session = name,
"agent-run chrome session release skipped — chrome-use is not available"
);
return ReleaseOutcome::Skipped(SkipReason::NoBinary);
};
let full_budget = timeout >= crate::chrome::SESSION_STOP_TIMEOUT;
match run_cli_json_at(&cli, &["session", "stop"], Some(name), timeout).await {
Ok(_) => {
debug!(session = name, "agent-run chrome session released");
ReleaseOutcome::Released
}
Err(verdict) => {
if is_wedged_release_failure(&verdict, full_budget) {
let recovery =
recover_unresponsive_session_via(&cli, name, SESSION_RECOVERY_TIMEOUT).await;
debug!(
session = name,
recovery = recovery.summary(),
"agent-run chrome session release hit a wedged session — recovered before the attempt is charged"
);
}
debug!(
session = name,
error = verdict.text(),
"agent-run chrome session release failed — the session still owns its tabs, or chrome-use could not answer"
);
ReleaseOutcome::Failed
}
}
}
fn is_wedged_release_failure(verdict: &NoVerdict, full_budget: bool) -> bool {
full_budget
&& match verdict {
NoVerdict::Reported(text) => is_session_unresponsive_error(text),
NoVerdict::TimedOut | NoVerdict::Unreadable => true,
NoVerdict::SpawnFailure => false,
}
}
#[derive(Clone, Debug)]
struct ReleaseSettings {
attempt_timeout: Duration,
retry_base: Duration,
boot_grace: Duration,
store: Option<PathBuf>,
#[cfg(test)]
cli: Option<PathBuf>,
}
fn release_settings() -> ReleaseSettings {
#[cfg(test)]
if let Some(settings) = test_release_settings() {
return settings;
}
ReleaseSettings {
attempt_timeout: crate::chrome::SESSION_STOP_TIMEOUT,
retry_base: RELEASE_RETRY_BASE,
boot_grace: RELEASE_BOOT_GRACE,
store: crate::config::CONFIG
.try_storage_root()
.map(|root| root.join(RELEASE_FILE_NAME)),
#[cfg(test)]
cli: None,
}
}
fn release_cli() -> Option<PathBuf> {
#[cfg(test)]
if let Some(settings) = test_release_settings() {
return settings.cli;
}
cli_path()
}
#[cfg(test)]
static RELEASE_SETTINGS: Mutex<Option<ReleaseSettings>> = Mutex::new(None);
#[cfg(test)]
fn test_release_settings() -> Option<ReleaseSettings> {
RELEASE_SETTINGS.lock().unwrap_poison().clone()
}
#[cfg(test)]
fn swap_release_settings(settings: ReleaseSettings) -> Option<ReleaseSettings> {
RELEASE_SETTINGS.lock().unwrap_poison().replace(settings)
}
#[cfg(test)]
fn restore_release_settings(previous: Option<ReleaseSettings>) {
*RELEASE_SETTINGS.lock().unwrap_poison() = previous;
}
fn release_backoff(attempts: u32) -> Duration {
release_settings()
.retry_base
.saturating_mul(2u32.saturating_pow(attempts.saturating_sub(1)))
.min(RELEASE_RETRY_CAP)
}
#[cfg(test)]
fn pending_releases_snapshot() -> Vec<(Vec<String>, u32, Duration)> {
let now = Instant::now();
pending_releases()
.lock()
.unwrap_poison()
.iter()
.map(|entry| {
(
entry.names.clone(),
entry.attempts,
entry.next_attempt_at.saturating_duration_since(now),
)
})
.collect()
}
#[cfg(test)]
pub(crate) fn pending_names_and_attempts() -> Vec<(Vec<String>, u32)> {
pending_releases_snapshot()
.into_iter()
.map(|(names, attempts, _)| (names, attempts))
.collect()
}
#[cfg(test)]
pub(crate) fn pending_release_delays() -> Vec<Duration> {
pending_releases_snapshot()
.into_iter()
.map(|(_, _, delay)| delay)
.collect()
}
#[cfg(test)]
pub(crate) fn clear_pending_releases() {
parked_releases().lock().unwrap_poison().clear();
pending_releases().lock().unwrap_poison().clear();
}
#[cfg(test)]
pub(crate) struct NoReleaseStore(Option<ReleaseSettings>);
#[cfg(test)]
#[must_use = "the store is turned back on when the guard drops"]
pub(crate) fn no_release_store() -> NoReleaseStore {
let mut settings = release_settings();
settings.store = None;
NoReleaseStore(swap_release_settings(settings))
}
#[cfg(test)]
impl Drop for NoReleaseStore {
fn drop(&mut self) {
restore_release_settings(self.0.take());
}
}
#[cfg(all(test, unix))]
mod tests {
use super::*;
use crate::tools::chrome::ChromeRunSessions;
use std::fs;
use std::path::Path;
fn widen_attempt_bound(attempt_timeout: Duration) {
let mut settings = release_settings();
settings.attempt_timeout = attempt_timeout;
swap_release_settings(settings);
}
fn widen_retry_base(retry_base: Duration) {
let mut settings = release_settings();
settings.retry_base = retry_base;
swap_release_settings(settings);
}
async fn wait_for_invocations(guard: &ReleaseGuard, lines: usize) {
tokio::time::timeout(Duration::from_secs(5), async {
while guard.log_lines().len() < lines {
tokio::time::sleep(Duration::from_millis(20)).await;
}
})
.await
.expect("the stub was invoked");
}
fn record_file_names(path: &Path) -> Vec<String> {
let json = fs::read_to_string(path).expect("record file");
let mut names: Vec<String> = serde_json::from_str::<Vec<PersistedRunRelease>>(&json)
.expect("the file stays a record file")
.into_iter()
.flat_map(|record| record.names)
.collect();
names.sort();
names
}
struct ReleaseGuard {
mode: PathBuf,
log: PathBuf,
previous: Option<ReleaseSettings>,
dir: tempfile::TempDir,
}
impl ReleaseGuard {
async fn install(attempt_timeout: Duration) -> Self {
Self::install_with(attempt_timeout, Duration::ZERO).await
}
async fn install_with(attempt_timeout: Duration, boot_grace: Duration) -> Self {
let dir = tempfile::tempdir().expect("release stub dir");
let cli = dir.path().join("chrome-use");
let mode = dir.path().join("mode");
let log = dir.path().join("log");
fs::write(
&cli,
format!(
"#!/bin/sh\nprintf '%s\\n' \"$*\" >> {log}\ncase \"$(cat {mode})\" in\n \
hang) exec sleep 5 ;;\n \
fail) printf '%s' '{{\"success\":false,\"error\":\"boom\"}}'; exit 1 ;;\n \
wedged) printf '%s' '{{\"success\":false,\"error\":\"session unresponsive\"}}'; exit 1 ;;\n \
ok) printf '%s' '{{\"success\":true}}' ;;\n\
esac\nexit 0\n",
log = log.display(),
mode = mode.display(),
),
)
.expect("write release stub");
crate::util::test::make_executable(&cli);
Self::write_mode(&mode, "ok");
let previous = swap_release_settings(ReleaseSettings {
cli: Some(cli.clone()),
attempt_timeout,
retry_base: Duration::ZERO,
boot_grace,
store: Some(dir.path().join(RELEASE_FILE_NAME)),
});
let guard = Self {
mode,
log,
previous,
dir,
};
let _ = release_one("warmup", attempt_timeout).await;
guard.clear_log();
guard
}
fn set_mode(&self, mode: &str) {
Self::write_mode(&self.mode, mode);
}
fn write_mode(path: &Path, mode: &str) {
fs::write(path, mode).expect("write release mode");
}
fn log_lines(&self) -> Vec<String> {
fs::read_to_string(&self.log)
.unwrap_or_default()
.lines()
.map(str::to_string)
.collect()
}
fn clear_log(&self) {
fs::write(&self.log, "").expect("reset release log");
}
fn store_path(&self) -> PathBuf {
self.dir.path().join(RELEASE_FILE_NAME)
}
}
fn sorted(mut lines: Vec<String>) -> Vec<String> {
lines.sort();
lines
}
impl Drop for ReleaseGuard {
fn drop(&mut self) {
restore_release_settings(self.previous.take());
clear_pending_releases();
}
}
fn queue_names(agent_id: &str, names: &[&str]) {
queue_names_with_hold(agent_id, names, Duration::ZERO);
}
fn queue_names_with_hold(agent_id: &str, names: &[&str], hold: Duration) {
let sessions = ChromeRunSessions::for_run(agent_id);
for name in run_session_names(agent_id, names) {
sessions.track(&name);
}
queue_run_session_release_after(&sessions, hold);
}
fn run_session_names(agent_id: &str, logical: &[&str]) -> Vec<String> {
let ns = ChromeRunSessions::for_run(agent_id).namespace().to_string();
logical.iter().map(|name| format!("{ns}{name}")).collect()
}
fn stop_invocations(names: &[String]) -> Vec<String> {
sorted(
names
.iter()
.map(|name| format!("session stop --json --session {name}"))
.collect(),
)
}
fn write_record_file(path: &Path, records: &[PersistedRunRelease]) {
fs::write(
path,
serde_json::to_string(records).expect("serialize release records"),
)
.expect("write release record file");
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_verified_release_empties_the_record() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
queue_names("run-verified", &["agent-tab-v-default", "agent-tab-v-docs"]);
assert_eq!(pending_releases_snapshot().len(), 1);
release_due().await;
assert!(
pending_releases_snapshot().is_empty(),
"both names released"
);
assert_eq!(
sorted(guard.log_lines()),
stop_invocations(&run_session_names(
"run-verified",
&["agent-tab-v-default", "agent-tab-v-docs"]
))
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_refused_release_is_retried_then_given_up_on_after_bounded_attempts() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
queue_names("run-refused", &["agent-tab-r-stuck"]);
guard.set_mode("fail");
release_due().await;
assert_eq!(
pending_names_and_attempts(),
vec![(run_session_names("run-refused", &["agent-tab-r-stuck"]), 1)],
"a refused release stays queued as one failed attempt"
);
let mut passes = 1;
while !pending_releases_snapshot().is_empty() {
release_due().await;
passes += 1;
assert!(
passes <= RELEASE_MAX_ATTEMPTS,
"the record must be given up on within the bounded attempts"
);
}
assert_eq!(passes, RELEASE_MAX_ATTEMPTS, "exactly the bounded attempts");
assert_eq!(
guard.log_lines().len(),
RELEASE_MAX_ATTEMPTS as usize,
"one `session stop` invocation per attempt"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_timed_out_release_counts_as_one_failed_attempt() {
let guard = ReleaseGuard::install(Duration::from_secs(1)).await;
queue_names("run-slow", &["agent-tab-s-slow"]);
guard.set_mode("hang");
let started = Instant::now();
release_due().await;
assert!(
started.elapsed() < Duration::from_secs(3),
"the 1 s attempt bound must cut the hanging child off"
);
assert_eq!(
pending_names_and_attempts(),
vec![(run_session_names("run-slow", &["agent-tab-s-slow"]), 1)],
"a timed-out attempt is one failed attempt, and the name stays queued"
);
assert_eq!(
guard.log_lines().len(),
1,
"a truncated attempt spawns no recovery stop"
);
guard.set_mode("ok");
widen_attempt_bound(Duration::from_secs(10));
release_due().await;
assert!(
pending_releases_snapshot().is_empty(),
"released on the next pass"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_session_unresponsive_release_runs_the_shared_recovery() {
let _health = crate::tools::chrome_daemon::with_health_test_lock().await;
let guard = ReleaseGuard::install(crate::chrome::SESSION_STOP_TIMEOUT).await;
queue_names("run-wedged", &["agent-tab-w-wedged"]);
guard.set_mode("wedged");
release_due().await;
assert_eq!(
pending_names_and_attempts(),
vec![(run_session_names("run-wedged", &["agent-tab-w-wedged"]), 1)],
"the failed attempt is charged once"
);
assert_eq!(
guard.log_lines().len(),
2,
"the `session stop` that reported the wedge plus the recovery's own"
);
crate::tools::chrome_daemon::reset_health();
}
#[test]
fn only_a_silent_or_unresponsive_stop_is_the_wedge_signature() {
assert!(is_wedged_release_failure(
&NoVerdict::Reported("session unresponsive".into()),
true
));
assert!(is_wedged_release_failure(&NoVerdict::TimedOut, true));
assert!(is_wedged_release_failure(&NoVerdict::Unreadable, true));
assert!(!is_wedged_release_failure(
&NoVerdict::Reported("session stop failed".into()),
true
));
assert!(!is_wedged_release_failure(&NoVerdict::SpawnFailure, true));
for verdict in [
NoVerdict::Reported("session unresponsive".into()),
NoVerdict::TimedOut,
NoVerdict::Unreadable,
] {
assert!(!is_wedged_release_failure(&verdict, false));
}
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_truncated_wedge_attempt_budget_spawns_no_recovery_stop() {
let guard = ReleaseGuard::install(Duration::from_secs(1)).await;
queue_names("run-wedged-flush", &["agent-tab-w-flush"]);
guard.set_mode("wedged");
release_due().await;
assert_eq!(
pending_names_and_attempts(),
vec![(
run_session_names("run-wedged-flush", &["agent-tab-w-flush"]),
1
)],
"the truncated wedge is still one failed attempt"
);
assert_eq!(
guard.log_lines().len(),
1,
"a truncated wedge spawns no recovery stop"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_release_with_no_chrome_use_binary_is_not_charged_an_attempt() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
queue_names("run-absent", &["agent-tab-a-default"]);
let mut settings = release_settings();
settings.cli = None;
settings.retry_base = Duration::from_secs(30);
swap_release_settings(settings);
release_due().await;
assert_eq!(
pending_names_and_attempts(),
vec![(run_session_names("run-absent", &["agent-tab-a-default"]), 0)],
"nothing was attempted, so nothing is charged"
);
assert!(guard.log_lines().is_empty(), "no child was spawned");
assert!(
pending_release_delays()[0] >= Duration::from_secs(25),
"re-armed by a skip gap rather than left instantly due"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_record_nothing_can_release_is_dropped_after_bounded_skips() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
queue_names("run-absent", &["agent-tab-a-default"]);
let mut settings = release_settings();
settings.cli = None;
swap_release_settings(settings);
for _ in 1..RELEASE_MAX_SKIPS {
release_due().await;
assert_eq!(
pending_names_and_attempts(),
vec![(run_session_names("run-absent", &["agent-tab-a-default"]), 0)],
"still queued, and no attempt charged for an unavailable binary"
);
}
release_due().await;
assert!(
pending_releases_snapshot().is_empty(),
"the record is dropped after its bounded skips"
);
assert!(
guard.log_lines().is_empty(),
"and nothing was ever spawned for it"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn overlapping_passes_keep_each_others_records_in_the_file() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
queue_names("run-overlap-a", &["agent-tab-oa-default"]);
let first = ReleasePass::take().expect("the first record is due");
queue_names("run-overlap-b", &["agent-tab-ob-default"]);
let second = ReleasePass::take().expect("the second record is due");
let mut both = run_session_names("run-overlap-a", &["agent-tab-oa-default"]);
both.extend(run_session_names(
"run-overlap-b",
&["agent-tab-ob-default"],
));
both.sort();
persist_pending_releases();
assert_eq!(
record_file_names(&guard.store_path()),
both,
"both passes' records are in the file"
);
drop(first);
assert_eq!(
record_file_names(&guard.store_path()),
both,
"the pass that ended left the other pass's record in the file"
);
assert_eq!(
pending_names_and_attempts(),
vec![(
run_session_names("run-overlap-a", &["agent-tab-oa-default"]),
0
)]
);
drop(second);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_queued_release_survives_the_process_that_queued_it() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
queue_names("run-durable", &["agent-tab-d-default"]);
guard.set_mode("fail");
release_due().await;
assert_eq!(
pending_names_and_attempts(),
vec![(
run_session_names("run-durable", &["agent-tab-d-default"]),
1
)],
"one refused attempt, re-queued"
);
clear_pending_releases();
assert!(pending_releases_snapshot().is_empty());
restore_pending_releases();
assert_eq!(
pending_names_and_attempts(),
vec![(
run_session_names("run-durable", &["agent-tab-d-default"]),
1
)],
"names and attempt count restored from the record file"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_corrupt_record_file_is_ignored() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
fs::write(guard.store_path(), "not a release record{").expect("write corrupt record file");
restore_pending_releases();
assert!(
pending_releases_snapshot().is_empty(),
"a corrupt file restores nothing and fails nothing"
);
fs::write(
guard.store_path(),
r#"[{"namespace":"agent-tab-absurd-","names":["agent-tab-absurd-default"],"next_attempt_at":18446744073709551615}]"#,
)
.expect("write absurd record file");
restore_pending_releases();
let delays = pending_release_delays();
assert_eq!(delays.len(), 1, "the absurd record is restored");
assert!(
delays[0] <= RELEASE_RETRY_CAP,
"its deadline is clamped to the ladder's cap: {delays:?}"
);
clear_pending_releases();
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_restored_record_releases_its_names_exactly_as_written() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
clear_pending_releases();
fs::write(
guard.store_path(),
r#"[{"namespace":"agent-tab-written-","names":["default","link-enricher-x","","agent-tab-other-default","agent-tab-written-kept"]},{"namespace":"agent-tab-nameless-","names":[]}]"#,
)
.expect("write record file naming foreign sessions");
restore_pending_releases();
let restored: Vec<String> = [
"default",
"link-enricher-x",
"",
"agent-tab-other-default",
"agent-tab-written-kept",
]
.iter()
.map(|name| (*name).to_string())
.collect();
assert_eq!(
pending_names_and_attempts(),
vec![(restored.clone(), 0)],
"every name the entry carries is restored as written; the entry left with no names is not queued"
);
release_due().await;
assert!(
pending_releases_snapshot().is_empty(),
"the record is released like any other"
);
assert_eq!(
sorted(guard.log_lines()),
stop_invocations(&restored),
"one `session stop` per name, each exactly as the record wrote it"
);
clear_pending_releases();
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_restored_release_waits_out_the_boot_grace() {
let grace = Duration::from_millis(800);
let guard = ReleaseGuard::install_with(Duration::from_secs(10), grace).await;
queue_names("run-grace", &["agent-tab-g-default"]);
clear_pending_releases();
restore_pending_releases();
release_due().await;
assert_eq!(
pending_releases_snapshot().len(),
1,
"restored at boot, the record is not attempted before the runs it may belong to exist"
);
assert!(guard.log_lines().is_empty());
tokio::time::sleep(grace * 2).await;
release_due().await;
assert!(
pending_releases_snapshot().is_empty(),
"released once the boot grace elapsed"
);
assert_eq!(
sorted(guard.log_lines()),
stop_invocations(&run_session_names("run-grace", &["agent-tab-g-default"]))
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_held_redriven_run_release_waits_out_its_hold() {
let hold = Duration::from_millis(800);
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
queue_names_with_hold("run-held", &["agent-tab-h-default"], hold);
release_due().await;
assert_eq!(
pending_releases_snapshot().len(),
1,
"a held record is left alone, and the skip is not an attempt"
);
assert!(
guard.log_lines().is_empty(),
"no `session stop` while the run it belongs to may still come back"
);
tokio::time::sleep(hold * 2).await;
release_due().await;
assert!(
pending_releases_snapshot().is_empty(),
"released once the hold elapsed"
);
assert_eq!(
sorted(guard.log_lines()),
stop_invocations(&run_session_names("run-held", &["agent-tab-h-default"]))
);
guard.clear_log();
queue_names_with_hold("run-shrunk", &["agent-tab-s-default"], hold * 4);
queue_names("run-shrunk", &["agent-tab-s-default"]);
release_due().await;
assert!(
pending_releases_snapshot().is_empty(),
"a re-queue without a hold makes the record eligible at once"
);
assert_eq!(
sorted(guard.log_lines()),
stop_invocations(&run_session_names("run-shrunk", &["agent-tab-s-default"]))
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_record_whose_run_is_live_again_is_never_attempted() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
let live = ChromeRunSessions::for_run("run-live");
live.track(&format!("{}default", live.namespace()));
queue_run_session_release(&live, false);
release_due().await;
assert_eq!(
pending_releases_snapshot().len(),
1,
"a live run's record stays queued and is not an attempt"
);
assert!(
guard.log_lines().is_empty(),
"no `session stop` is ever spawned for a live run's sessions"
);
drop(live);
release_due().await;
assert!(
pending_releases_snapshot().is_empty(),
"released once the run that owns them is gone"
);
assert_eq!(guard.log_lines().len(), 1);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn the_driver_releases_a_queued_record_and_waits_for_a_live_run_to_go() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
widen_retry_base(Duration::from_secs(600));
let live = ChromeRunSessions::for_run("run-driver");
let names = run_session_names("run-driver", &["agent-tab-driver-default"]);
live.track(&names[0]);
queue_run_session_release(&live, false);
let shutdown = CancellationToken::new();
let driver = tokio::spawn(release_queue(shutdown.clone()));
tokio::time::sleep(Duration::from_millis(300)).await;
assert!(
guard.log_lines().is_empty(),
"the driver leaves a live run's session alone"
);
assert_eq!(pending_releases_snapshot().len(), 1);
drop(live);
wait_for_invocations(&guard, 1).await;
assert_eq!(
guard.log_lines(),
stop_invocations(&names),
"the driver releases it through the run's own `session stop`"
);
assert!(
pending_releases_snapshot().is_empty(),
"and drops the record once it is released"
);
shutdown.cancel();
tokio::time::timeout(Duration::from_secs(5), driver)
.await
.expect("the driver stops on the token")
.expect("the driver task does not panic");
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_requeue_of_the_same_run_does_not_double_count() {
let _guard = ReleaseGuard::install(Duration::from_secs(10)).await;
clear_pending_releases();
let first = ChromeRunSessions::for_run("run-merge");
let names = run_session_names("run-merge", &["agent-tab-m-default", "agent-tab-m-docs"]);
first.track(&names[0]);
queue_run_session_release(&first, false);
let second = ChromeRunSessions::for_run("run-merge");
second.track(&names[0]);
second.track(&names[1]);
queue_run_session_release(&second, false);
drop((first, second));
assert_eq!(
pending_names_and_attempts(),
vec![(names, 0)],
"one record carrying both ends' names"
);
clear_pending_releases();
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn the_shutdown_flush_releases_the_queue_within_its_budget() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
queue_names("run-flush", &["agent-tab-f-default"]);
queue_names_with_hold(
"run-flush-held",
&["agent-tab-f-held"],
Duration::from_hours(1),
);
guard.set_mode("hang");
let started = Instant::now();
flush_pending_run_releases(Duration::from_millis(500)).await;
assert!(
started.elapsed() < Duration::from_secs(3),
"the flush stops at its budget"
);
assert_eq!(
pending_releases_snapshot().len(),
2,
"an unreleased name stays queued"
);
guard.set_mode("ok");
guard.clear_log();
flush_pending_run_releases(Duration::from_secs(10)).await;
assert_eq!(
pending_names_and_attempts(),
vec![(
run_session_names("run-flush-held", &["agent-tab-f-held"]),
0
)],
"the flush drains what it can release and leaves the held record for its resumed run"
);
assert_eq!(
sorted(guard.log_lines()),
stop_invocations(&run_session_names("run-flush", &["agent-tab-f-default"])),
"one verified `session stop` per name, and never one for the held record"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_run_that_comes_back_mid_pass_keeps_its_sessions() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
let live = ChromeRunSessions::for_run("run-midpass");
let name = format!("{}default", live.namespace());
live.track(&name);
let record = PendingRunRelease {
namespace: live.namespace().to_string(),
names: vec![name.clone()],
attempts: 0,
skips: 0,
next_attempt_at: Instant::now(),
held: false,
};
let outcomes =
attempt_releases(std::slice::from_ref(&record), || Duration::from_secs(10)).await;
assert_eq!(
outcomes,
vec![(0, name, ReleaseOutcome::Skipped(SkipReason::Untried))]
);
assert!(
guard.log_lines().is_empty(),
"no `session stop` is spawned for a run that is live again"
);
drop(live);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_hold_queued_while_a_pass_holds_the_record_survives_in_the_file() {
let guard = ReleaseGuard::install_with(Duration::from_secs(10), Duration::ZERO).await;
queue_names("run-fold", &["agent-tab-f-default"]);
let pass = ReleasePass::take().expect("a due record is taken");
queue_names_with_hold(
"run-fold",
&["agent-tab-f-default"],
Duration::from_mins(30),
);
assert_eq!(
record_file_names(&guard.store_path()),
run_session_names("run-fold", &["agent-tab-f-default"]),
"one entry per namespace, so a restore cannot pick the stale pair"
);
clear_pending_releases();
restore_pending_releases();
let delays = pending_release_delays();
assert_eq!(delays.len(), 1, "the record is restored once");
assert!(
delays[0] >= RELEASE_HOLD_REDRIVEN_RUN.saturating_sub(Duration::from_secs(5)),
"the file kept the hold rather than the stale unheld snapshot: {:?}",
delays[0]
);
drop(pass);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_panicking_finish_leaves_its_record_for_the_next_boot() {
let _guard = ReleaseGuard::install(Duration::from_secs(10)).await;
queue_names("run-panic", &["agent-tab-p-default"]);
let mut pass = ReleasePass::take().expect("a due record is taken");
let panicked = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
pass.reconcile(Vec::new(), |_, _| panic!("finish panicked"));
}))
.is_err();
drop(pass);
assert!(panicked, "the panic reached the pass");
assert!(
pending_releases_snapshot().is_empty(),
"the record is out of the live queue"
);
restore_pending_releases();
assert_eq!(
pending_names_and_attempts(),
vec![(run_session_names("run-panic", &["agent-tab-p-default"]), 0)],
"but the file kept it, so the next boot restores it"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn an_aborted_pass_puts_its_records_back() {
let _guard = ReleaseGuard::install(Duration::from_secs(10)).await;
queue_names("run-aborted", &["agent-tab-a-default"]);
let pass = ReleasePass::take().expect("a due record is taken");
assert!(
pending_releases_snapshot().is_empty(),
"the pass holds the record"
);
drop(pass);
assert_eq!(
pending_names_and_attempts(),
vec![(
run_session_names("run-aborted", &["agent-tab-a-default"]),
0
)],
"a dropped pass puts the record back untried"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_restored_held_record_holds_from_boot() {
let guard = ReleaseGuard::install_with(Duration::from_secs(10), Duration::ZERO).await;
let namespace = ChromeRunSessions::for_run("run-resume")
.namespace()
.to_string();
write_record_file(
&guard.store_path(),
&[PersistedRunRelease {
namespace,
names: run_session_names("run-resume", &["agent-tab-r-default"]),
attempts: 0,
skips: 0,
next_attempt_at: 1,
held: true,
}],
);
restore_pending_releases();
let delays = pending_release_delays();
assert_eq!(delays.len(), 1, "the held record is restored");
assert!(
delays[0] >= RELEASE_HOLD_REDRIVEN_RUN.saturating_sub(Duration::from_secs(5)),
"the hold is measured from this boot, not from the expired deadline: {:?}",
delays[0]
);
}
}