use super::chrome_daemon::cli_path;
use super::chrome_tabs::{self, RetryLeg, TabOutcome};
use crate::util::UnwrapPoison;
use serde::{Deserialize, Serialize};
use std::collections::{HashMap, HashSet, VecDeque};
use std::path::{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, warn};
const RELEASE_RETRY_BASE: Duration = Duration::from_secs(5);
const RELEASE_RETRY_CAP: Duration = Duration::from_mins(30);
const MAX_LEFT_OPEN_ATTEMPTS: u32 = 3;
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 = 256;
const EVICTED_RUN_TTL: Duration = Duration::from_hours(1);
const SHUTDOWN_RELEASE_FLUSH_BUDGET: Duration = Duration::from_secs(10);
#[derive(Clone)]
struct PendingRunRelease {
namespace: String,
names: Vec<String>,
attempts: u32,
next_attempt_at: Instant,
held: bool,
run: String,
left_open: u32,
}
#[derive(Serialize, Deserialize)]
struct PersistedRunRelease {
namespace: String,
names: Vec<String>,
#[serde(default)]
attempts: u32,
#[serde(default)]
next_attempt_at: u64,
#[serde(default)]
held: bool,
#[serde(default)]
run: String,
#[serde(default)]
left_open: u32,
}
struct EvictedRun {
namespace: String,
run: String,
at: Instant,
}
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();
static RECLAIM: OnceLock<Mutex<ReclaimSchedule>> = OnceLock::new();
static STUCK_NAMESPACES: OnceLock<Mutex<HashMap<String, Instant>>> = OnceLock::new();
static PENDING_FORGET: OnceLock<Mutex<Vec<(String, u32)>>> = OnceLock::new();
static EVICTED_RUNS: OnceLock<Mutex<VecDeque<EvictedRun>>> = 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()))
}
fn evicted_runs() -> &'static Mutex<VecDeque<EvictedRun>> {
EVICTED_RUNS.get_or_init(|| Mutex::new(VecDeque::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.left_open = entry.left_open.max(incoming.left_open);
if entry.run.is_empty() {
entry.run = incoming.run;
}
let mut seen: HashSet<String> = entry.names.iter().cloned().collect();
for name in incoming.names {
if seen.insert(name.clone()) {
entry.names.push(name);
}
}
}
fn within_hold(record: &PendingRunRelease, now: Instant) -> bool {
record.held && record.next_attempt_at > now
}
fn evict_oldest(queue: &mut VecDeque<PendingRunRelease>) -> Option<PendingRunRelease> {
let now = Instant::now();
let index = queue.iter().position(|record| !within_hold(record, now))?;
queue.remove(index)
}
fn remember_evicted_run(namespace: &str, run: &str) {
let now = Instant::now();
let mut runs = evicted_runs().lock().unwrap_poison();
runs.retain(|entry| now.saturating_duration_since(entry.at) < EVICTED_RUN_TTL);
runs.push_back(EvictedRun {
namespace: namespace.to_string(),
run: run.to_string(),
at: now,
});
while runs.len() > MAX_PENDING_RELEASES {
runs.pop_front();
}
}
fn evicted_run(namespace: &str) -> Option<String> {
let now = Instant::now();
evicted_runs()
.lock()
.unwrap_poison()
.iter()
.rev()
.find(|entry| {
entry.namespace == namespace
&& now.saturating_duration_since(entry.at) < EVICTED_RUN_TTL
})
.map(|entry| entry.run.clone())
}
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 evicted: Vec<PendingRunRelease> = Vec::new();
while queue.len() >= MAX_PENDING_RELEASES
&& let Some(record) = evict_oldest(&mut queue)
{
evicted.push(record);
}
queue.push_back(entry);
drop(queue);
for record in evicted {
remember_evicted_run(&record.namespace, &record.run);
arm_reclaim_at(Instant::now() + STUCK_REVISIT);
info!(
namespace = %record.namespace,
run = %record.run,
sessions = record.names.len(),
"agent-run chrome release queue full — record evicted; the reclaim sweep owns its tabs"
);
}
}
pub(crate) fn queue_run_session_release(
sessions: &super::chrome::ChromeRunSessions,
held: bool,
run: &str,
) {
queue_run_session_release_after(
sessions,
if held {
RELEASE_HOLD_REDRIVEN_RUN
} else {
Duration::ZERO
},
run,
);
}
fn queue_run_session_release_after(
sessions: &super::chrome::ChromeRunSessions,
hold: Duration,
run: &str,
) {
let names = sessions.snapshot();
if names.is_empty() {
return; }
merge_record(
PendingRunRelease {
namespace: sessions.namespace().to_string(),
names,
attempts: 0,
next_attempt_at: Instant::now() + hold,
held: !hold.is_zero(),
run: run.to_string(),
left_open: 0,
},
Eligibility::RunEnd,
);
persist_pending_releases();
schedule_reclaim_now();
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 {
if record.names.is_empty() || record.namespace.is_empty() {
continue;
}
merge_record(
PendingRunRelease {
namespace: record.namespace,
names: record.names,
attempts: record.attempts,
next_attempt_at: restored_eligibility(record.next_attempt_at, record.held),
held: record.held,
run: record.run,
left_open: record.left_open,
},
Eligibility::Retry,
);
restored += 1;
}
if restored > 0 {
debug!(
restored,
"agent-run chrome releases restored from the previous process"
);
schedule_reclaim_now();
release_wake().notify_one();
}
}
async fn prune_tab_records() {
let Some(durable) = durable_resume_namespaces_checked().await else {
return;
};
let protected = ProtectedNamespaces::snapshot(&durable);
crate::tools::chrome_tab_ledger::prune(|namespace| protected.contains(namespace));
}
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();
let mut index: HashMap<String, usize> = HashMap::new();
for record in queue.iter().chain(parked.iter().map(|(_, record)| record)) {
if let Some(at) = index.get(&record.namespace).copied() {
absorb(&mut folded[at], record.clone(), Eligibility::Retry);
} else {
index.insert(record.namespace.clone(), folded.len());
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,
next_attempt_at: eligibility_deadline(entry.next_attempt_at),
held: entry.held,
run: entry.run.clone(),
left_open: entry.left_open,
})
.collect();
let json = match serde_json::to_string(&records) {
Ok(json) => json,
Err(err) => {
warn!(error = %err, "agent-run chrome release record not written: serialization failed");
return;
}
};
if let Err(err) = crate::util::write_json_record(&path, &json) {
warn!(error = %err, path = %path.display(), "agent-run chrome release record not written");
}
}
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();
prune_tab_records().await;
loop {
tokio::select! {
() = wait_for_release_work() => {}
() = release_wake().notified() => {}
() = shutdown.cancelled() => break,
}
release_due().await;
}
}
async fn wait_for_release_work() {
match next_release_deadline() {
Some(deadline) => {
tokio::time::sleep_until(tokio::time::Instant::from_std(deadline)).await;
}
None => std::future::pending().await,
}
}
fn next_release_deadline() -> Option<Instant> {
let sweep = reclaim_schedule().lock().unwrap_poison().next_at;
pending_releases()
.lock()
.unwrap_poison()
.iter()
.filter(|entry| !run_namespace_is_live(&entry.namespace))
.map(|entry| entry.next_attempt_at)
.min()
.into_iter()
.chain(sweep)
.chain(forget_retry_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(mut entry) = queue.pop_front() {
let releasable = !run_namespace_is_live(&entry.namespace) && entry.next_attempt_at <= now;
if releasable {
entry.held = false;
taken.push(entry);
} else {
waiting.push_back(entry);
}
}
*queue = waiting;
taken
}
struct RecordDisposition {
retry: Vec<String>,
left_open: Vec<String>,
unclosable: Vec<String>,
}
fn record_disposition(
entry: &mut PendingRunRelease,
outcomes: &HashMap<String, TabOutcome>,
) -> RecordDisposition {
if entry.names.iter().any(|name| {
matches!(
outcomes.get(name),
Some(TabOutcome::Retry(RetryLeg::LeftOpen))
)
}) {
entry.left_open = entry.left_open.saturating_add(1);
}
let concluded = entry.left_open >= MAX_LEFT_OPEN_ATTEMPTS;
let mut retry = Vec::new();
let mut left_open = Vec::new();
let mut unclosable = Vec::new();
for name in &entry.names {
match outcomes.get(name) {
Some(TabOutcome::Gone) => {}
Some(TabOutcome::Unclosable) => unclosable.push(name.clone()),
Some(TabOutcome::Retry(RetryLeg::LeftOpen)) if concluded => {
unclosable.push(name.clone());
}
Some(TabOutcome::Retry(RetryLeg::LeftOpen)) => {
left_open.push(name.clone());
retry.push(name.clone());
}
Some(TabOutcome::Retry(RetryLeg::Silent)) | None => retry.push(name.clone()),
}
}
RecordDisposition {
retry,
left_open,
unclosable,
}
}
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: &HashMap<String, TabOutcome>,
mut finish: impl FnMut(PendingRunRelease, RecordDisposition),
) {
while let Some(mut entry) = self.records.pop() {
self.in_flight = Some(entry.namespace.clone());
let disposition = record_disposition(&mut entry, outcomes);
finish(entry, disposition);
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 cli = release_cli();
let deadline = Instant::now() + release_settings().attempt_timeout;
let pass = ReleasePass::take();
let taken = pass.is_some();
let sweep_due = reclaim_due();
let durable = if taken || sweep_due {
durable_resume_namespaces().await
} else {
Vec::new()
};
let mut protected = ProtectedNamespaces::snapshot(&durable);
let session = read_session(pass.as_ref());
let ask = reclaim_names(cli.as_deref(), &session, &protected, deadline, sweep_due).await;
let issues = attempt_pass(
cli.as_deref(),
pass,
&session,
&ask,
&mut protected,
deadline,
)
.await;
report_issues(issues).await;
}
fn read_session(pass: Option<&ReleasePass>) -> String {
if let Some(name) = pass
.and_then(|pass| pass.records.first())
.and_then(|record| record.names.first())
{
return name.clone();
}
let parked = parked_releases().lock().unwrap_poison();
let queued = pending_releases().lock().unwrap_poison();
parked
.iter()
.map(|(_, record)| record)
.chain(queued.iter())
.find_map(|record| record.names.first())
.cloned()
.unwrap_or_else(|| crate::tools::chrome_tabs::OWN_SESSION.to_string())
}
async fn close_unprotected(
cli: Option<&Path>,
names: &[String],
session: &str,
protected: &mut ProtectedNamespaces,
deadline: Instant,
) -> HashMap<String, TabOutcome> {
let Some(cli) = cli else {
return HashMap::new();
};
protected.refresh();
let unprotected: Vec<String> = names
.iter()
.filter(|name| !protected.contains(name))
.cloned()
.collect();
let owned_again = ProtectedNamespaces::protects_live_or_held;
chrome_tabs::close_sessions(cli, &unprotected, session, deadline, &owned_again)
.await
.into_iter()
.collect()
}
pub(crate) async fn evict_agent_tab(namespace: &str, session: &str, read_session: &str) -> bool {
let Some(cli) = release_cli() else {
return false;
};
let deadline = Instant::now() + release_settings().attempt_timeout;
let owned_again = |name: &str| {
!names_a_namespace(namespace, name) && ProtectedNamespaces::protects_live_or_held(name)
};
let name = session.to_string();
let outcomes = chrome_tabs::close_sessions(
&cli,
std::slice::from_ref(&name),
read_session,
deadline,
&owned_again,
)
.await;
let gone = outcomes
.first()
.is_some_and(|(_, outcome)| matches!(outcome, chrome_tabs::TabOutcome::Gone));
if gone {
for _ in 0..2 {
let outcome = chrome_tabs::forget_settled_sessions(
&cli,
std::slice::from_ref(&name),
deadline,
&|_: &str| false,
)
.await;
if outcome.tried.is_empty() && outcome.deferred.is_empty() {
break;
}
}
}
gone
}
async fn forget_settled(cli: &Path, just_settled: &[String], deadline: Instant) {
let mut todo: Vec<(String, u32)> = pending_forget().lock().unwrap_poison().clone();
let snapshot: HashSet<String> = todo.iter().map(|(name, _)| name.clone()).collect();
todo.extend(just_settled.iter().map(|name| (name.clone(), 0)));
dedupe_and_bound_forgotten(&mut todo);
if todo.is_empty() {
let queue = pending_forget().lock().unwrap_poison();
if queue.is_empty() {
*forget_retry().lock().unwrap_poison() = None;
}
return;
}
let names: Vec<String> = todo.iter().map(|(name, _)| name.clone()).collect();
let owned_again = ProtectedNamespaces::protects_live_or_held;
let left = chrome_tabs::forget_settled_sessions(cli, &names, deadline, &owned_again).await;
let mut kept: Vec<(String, u32)> = Vec::new();
for (name, attempts) in todo {
if left.owned.contains(&name) {
continue;
}
if left.deferred.contains(&name) || left.busy.contains(&name) {
kept.push((name, attempts));
continue;
}
if !left.tried.contains(&name) {
continue; }
let attempts = attempts.saturating_add(1);
if attempts >= MAX_FORGET_ATTEMPTS && name != chrome_tabs::OWN_SESSION {
info!(
session = %name,
attempts,
"a settled chrome session could not be let go — dropped from the let-go queue \
with chrome-use's own per-session record left behind; a pass that still owes \
this session's let-go queues it again"
);
continue;
}
kept.push((name, attempts));
}
let deferred_behind_call = !left.busy.is_empty();
let mut queue = pending_forget().lock().unwrap_poison();
let mut merged = kept;
merged.extend(
queue
.iter()
.filter(|(name, _)| !snapshot.contains(name))
.cloned(),
);
for (name, attempts) in &mut merged {
if let Some((_, live)) = queue.iter().find(|(queued, _)| queued == name) {
*attempts = (*attempts).max(*live);
}
}
*queue = merged;
dedupe_and_bound_forgotten(&mut queue);
let most_attempts = queue.iter().map(|(_, attempts)| *attempts).max();
let retry_at = most_attempts.map(|attempts| Instant::now() + release_backoff(attempts));
*forget_retry().lock().unwrap_poison() = if deferred_behind_call {
let floor = Instant::now() + BUSY_FORGET_FLOOR;
Some(retry_at.map_or(floor, |at| at.max(floor)))
} else {
retry_at
};
}
fn queue_forget(name: &str) {
let mut queue = pending_forget().lock().unwrap_poison();
if queue.iter().any(|(queued, _)| queued == name) {
return;
}
queue.push((name.to_string(), 0));
dedupe_and_bound_forgotten(&mut queue);
*forget_retry().lock().unwrap_poison() = Some(Instant::now() + release_backoff(0));
}
static FORGET_RETRY: OnceLock<Mutex<Option<Instant>>> = OnceLock::new();
fn forget_retry() -> &'static Mutex<Option<Instant>> {
FORGET_RETRY.get_or_init(|| Mutex::new(None))
}
fn forget_retry_at() -> Option<Instant> {
release_cli()?;
*forget_retry().lock().unwrap_poison()
}
const MAX_PENDING_FORGET: usize = 64;
const BUSY_FORGET_FLOOR: Duration = crate::chrome::SESSION_STOP_TIMEOUT;
const MAX_FORGET_ATTEMPTS: u32 = 5;
fn dedupe_and_bound_forgotten(names: &mut Vec<(String, u32)>) {
let mut seen: HashSet<String> = HashSet::with_capacity(names.len());
names.retain(|(name, _)| seen.insert(name.clone()));
let excess = names.len().saturating_sub(MAX_PENDING_FORGET);
if excess == 0 {
return;
}
let mut dropped = 0usize;
names.retain(|(name, _)| {
if dropped == excess || name == chrome_tabs::OWN_SESSION {
return true;
}
dropped += 1;
false
});
debug!(
dropped,
"settled chrome sessions past the let-go cap were dropped — chrome-use records left \
behind by it, never a run's tab"
);
}
fn pending_forget() -> &'static Mutex<Vec<(String, u32)>> {
PENDING_FORGET.get_or_init(|| Mutex::new(Vec::new()))
}
fn pass_names(ask: &ReclaimAsk, pass: Option<&ReleasePass>) -> Vec<String> {
let mut names = ask.names.clone();
if let Some(pass) = pass {
let mut seen: HashSet<String> = names.iter().cloned().collect();
for name in record_names(pass) {
if seen.insert(name.clone()) {
names.push(name);
}
}
}
names
}
fn record_names(pass: &ReleasePass) -> Vec<String> {
let mut names: Vec<String> = Vec::new();
let mut seen: HashSet<String> = HashSet::new();
for record in &pass.records {
for name in &record.names {
if seen.insert(name.clone()) {
names.push(name.clone());
}
}
}
names
}
fn reclaim_issues(
reclaim: &[String],
answered: bool,
outcomes: &HashMap<String, TabOutcome>,
) -> (usize, Vec<Issue>) {
let mut reclaimed = 0usize;
let mut issues = Vec::new();
let mut silent = false;
for name in reclaim {
let Some(outcome) = outcomes.get(name) else {
continue;
};
let namespace = crate::tools::chrome::session_namespace(name).to_string();
match outcome {
TabOutcome::Gone => {
reclaimed += 1;
stuck_namespaces().lock().unwrap_poison().remove(&namespace);
}
TabOutcome::Unclosable => {
issues.push(reclaimed_issue(name, &namespace, IssueKind::Unclosable));
park_stuck(&namespace);
}
TabOutcome::Retry(RetryLeg::LeftOpen) => {
issues.push(reclaimed_issue(name, &namespace, IssueKind::LeftOpen));
park_stuck(&namespace);
}
TabOutcome::Retry(RetryLeg::Silent) => silent = true,
}
}
if silent {
reclaim_unfinished();
} else if answered {
reset_reclaim_failures();
}
(reclaimed, issues)
}
fn park_stuck(namespace: &str) {
stuck_namespaces()
.lock()
.unwrap_poison()
.insert(namespace.to_string(), Instant::now());
schedule_reclaim_revisit();
}
fn reclaimed_issue(name: &str, namespace: &str, kind: IssueKind) -> Issue {
Issue {
kind,
namespace: namespace.to_string(),
run: evicted_run(namespace).unwrap_or_default(),
name: name.to_string(),
}
}
async fn let_go_own_session(cli: &Path, deadline: Instant) {
if chrome_tabs::release_own_session(cli, deadline).await {
return;
}
queue_forget(chrome_tabs::OWN_SESSION);
if !chrome_tabs::own_scratch_group_possible() {
return;
}
report_once(
OWN_SESSION_MESSAGE,
OWN_SESSION_REASON,
serde_json::json!({
"detail": format!(
"a scratch group titled {} MAY be left in the owner's tab strip — this \
module's own about:blank read group, the session its browser reads run in — \
or the reads that look for it may themselves have failed, so the release is \
not confirmed; the group's own session is queued for a stop and the product's \
session sweeps close that family, and nothing of any run's is left behind by it",
chrome_tabs::OWN_SESSION
),
}),
)
.await;
}
async fn report_unreadable_answers() {
if !chrome_tabs::saw_unreadable_answer() {
return;
}
report_once(
UNREADABLE_MESSAGE,
UNREADABLE_REASON,
serde_json::json!({
"detail": "a read this cleanup decides from settled nothing: one of the extension's \
own answers about the browser — its status, the tab ledger it holds, or \
the browser's tab and group lists — carried a field no parser here will \
guess at, or a helper answer came back with no readable envelope, or with \
neither a success verdict nor a reason for failing. A partial or reasonless \
answer is never guessed at — a dropped tab would read as one that is gone — \
so nothing is concluded from it: another read, of this pass or a later one, \
still settles what it can",
}),
)
.await;
}
fn conclude_released_record(namespace: &str) {
if ProtectedNamespaces::protects_live_or_held(namespace) {
return;
}
crate::tools::chrome_tab_ledger::forget_namespace(namespace);
}
async fn attempt_pass(
cli: Option<&Path>,
mut pass: Option<ReleasePass>,
session: &str,
ask: &ReclaimAsk,
protected: &mut ProtectedNamespaces,
deadline: Instant,
) -> Vec<Issue> {
let names = pass_names(ask, pass.as_ref());
if names.is_empty() {
if ask.answered {
reset_reclaim_failures();
}
if let Some(cli) = cli {
let_go_own_session(cli, deadline).await;
forget_settled(cli, &[], deadline).await;
}
report_unreadable_answers().await;
return Vec::new();
}
let outcomes = close_unprotected(cli, &names, session, protected, deadline).await;
if let Some(cli) = cli {
let_go_own_session(cli, deadline).await;
}
if let Some(cli) = cli {
let live_now = ProtectedNamespaces::live_and_held();
let settled: Vec<String> = names
.iter()
.filter(|name| {
!live_now.contains(name)
&& matches!(outcomes.get(name.as_str()), Some(TabOutcome::Gone))
})
.cloned()
.collect();
forget_settled(cli, &settled, deadline).await;
}
let mut issues: Vec<Issue> = Vec::new();
let (reclaimed, reclaimed_issues) = reclaim_issues(&ask.names, ask.answered, &outcomes);
issues.extend(reclaimed_issues);
if reclaimed > 0 {
info!(
closed = reclaimed,
"agent-run chrome leftovers reclaimed by the sweep"
);
}
if let Some(pass) = pass.as_mut() {
pass.reconcile(&outcomes, |mut entry, disposition| {
let unclosable = disposition.unclosable.len();
let reported = disposition
.unclosable
.into_iter()
.map(|name| (IssueKind::Unclosable, name))
.chain(
disposition
.left_open
.into_iter()
.map(|name| (IssueKind::LeftOpen, name)),
);
for (kind, name) in reported {
issues.push(Issue {
kind,
namespace: entry.namespace.clone(),
run: entry.run.clone(),
name,
});
}
if unclosable > 0 {
arm_reclaim_at(Instant::now() + STUCK_REVISIT);
}
if disposition.retry.is_empty() {
if unclosable == 0 {
conclude_released_record(&entry.namespace);
info!(
namespace = %entry.namespace,
run = %entry.run,
sessions = entry.names.len(),
"agent-run chrome sessions released"
);
} else {
info!(
namespace = %entry.namespace,
run = %entry.run,
unclosable,
left_open = entry.left_open,
"agent-run chrome session release settled — no route in this record \
closes the rest (reported)"
);
}
return;
}
entry.attempts = entry.attempts.saturating_add(1);
entry.next_attempt_at = Instant::now() + release_backoff(entry.attempts);
debug!(
namespace = %entry.namespace,
sessions = disposition.retry.len(),
attempts = entry.attempts,
"agent-run chrome session release still open — retrying"
);
entry.names = disposition.retry;
merge_record(entry, Eligibility::Retry);
});
persist_pending_releases();
}
report_unreadable_answers().await;
issues
}
async fn flush_pending_run_releases(budget: Duration) {
let deadline = Instant::now() + budget;
let cli = release_cli();
let Some(mut pass) = ReleasePass::take() else {
if let Some(cli) = cli.as_deref() {
forget_settled(cli, &[], deadline).await;
}
return;
};
let names = record_names(&pass);
let durable = durable_resume_namespaces().await;
let mut protected = ProtectedNamespaces::snapshot(&durable);
let session = read_session(Some(&pass));
let outcomes =
close_unprotected(cli.as_deref(), &names, &session, &mut protected, deadline).await;
if let Some(cli) = cli.as_deref() {
let_go_own_session(cli, deadline).await;
forget_settled(cli, &[], deadline).await;
}
let mut released = 0usize;
let mut still_open = 0usize;
pass.reconcile(&outcomes, |mut entry, disposition| {
let keep: Vec<String> = disposition
.retry
.into_iter()
.chain(disposition.unclosable)
.collect();
released += entry.names.len() - keep.len();
if keep.is_empty() {
return;
}
still_open += keep.len();
entry.names = keep;
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 reclaim_names(
cli: Option<&Path>,
session: &str,
protected: &ProtectedNamespaces,
deadline: Instant,
due: bool,
) -> ReclaimAsk {
if !due {
return ReclaimAsk::unanswered();
}
prune_stuck();
let Some(cli) = cli else {
reclaim_unfinished();
return ReclaimAsk::unanswered();
};
let claimed = claimed_names();
match chrome_tabs::agent_group_names(cli, session, deadline).await {
Ok(names) => {
reclaim_answered();
ReclaimAsk {
names: names
.into_iter()
.filter(|name| {
!protected.contains(name)
&& !claimed.contains(name)
&& !stuck(crate::tools::chrome::session_namespace(name))
})
.collect(),
answered: true,
}
}
Err(reason) => {
reclaim_unfinished();
debug!(
reason = %reason.text(),
"agent-run chrome reclaim could not read the live tab groups"
);
if chrome_tabs::ownership_door_missing(cli, session, deadline).await {
report_once(
NO_DOOR_MESSAGE,
NO_OWNERSHIP_DOOR_REASON,
serde_json::json!({
"detail": "there is no ownership door on this browser side (no \
extension installed, one installed but disabled, or one \
older than the version that has it), so the sweep cannot \
enumerate the browser's groups and a name a record holds is \
left to chrome-use's own `session stop`",
}),
)
.await;
}
ReclaimAsk::unanswered()
}
}
}
struct ReclaimAsk {
names: Vec<String>,
answered: bool,
}
impl ReclaimAsk {
fn unanswered() -> Self {
Self {
names: Vec::new(),
answered: false,
}
}
}
#[derive(Default)]
struct ReclaimSchedule {
next_at: Option<Instant>,
failures: u32,
}
fn reclaim_schedule() -> &'static Mutex<ReclaimSchedule> {
RECLAIM.get_or_init(|| Mutex::new(ReclaimSchedule::default()))
}
fn arm_reclaim_at(at: Instant) {
let mut schedule = reclaim_schedule().lock().unwrap_poison();
schedule.next_at = Some(schedule.next_at.map_or(at, |old| old.min(at)));
}
fn next_reclaim_revisit() -> Option<Instant> {
stuck_namespaces()
.lock()
.unwrap_poison()
.values()
.filter(|concluded| concluded.elapsed() < STUCK_REVISIT)
.map(|concluded| *concluded + STUCK_REVISIT)
.min()
}
fn schedule_reclaim_now() {
arm_reclaim_at(Instant::now());
}
fn schedule_reclaim_revisit() {
arm_reclaim_at(Instant::now() + STUCK_REVISIT);
}
const STUCK_REVISIT: Duration = Duration::from_mins(30);
fn stuck_namespaces() -> &'static Mutex<HashMap<String, Instant>> {
STUCK_NAMESPACES.get_or_init(|| Mutex::new(HashMap::new()))
}
fn stuck(namespace: &str) -> bool {
let mut stuck = stuck_namespaces().lock().unwrap_poison();
match stuck.get(namespace) {
Some(concluded) if concluded.elapsed() < STUCK_REVISIT => true,
Some(_) => {
stuck.remove(namespace);
false
}
None => false,
}
}
fn prune_stuck() {
let now = Instant::now();
stuck_namespaces()
.lock()
.unwrap_poison()
.retain(|_, concluded| now.saturating_duration_since(*concluded) < STUCK_REVISIT);
}
fn reclaim_due() -> bool {
let mut schedule = reclaim_schedule().lock().unwrap_poison();
match schedule.next_at {
Some(at) if at <= Instant::now() => {
schedule.next_at = None;
true
}
_ => false,
}
}
fn reclaim_answered() {
let revisit = next_reclaim_revisit();
let mut schedule = reclaim_schedule().lock().unwrap_poison();
if let Some(at) = revisit {
schedule.next_at = Some(schedule.next_at.map_or(at, |armed| armed.min(at)));
}
}
fn reset_reclaim_failures() {
reclaim_schedule().lock().unwrap_poison().failures = 0;
}
fn reclaim_unfinished() {
let failures = {
let mut schedule = reclaim_schedule().lock().unwrap_poison();
schedule.failures = schedule.failures.saturating_add(1);
schedule.failures
};
arm_reclaim_at(Instant::now() + release_backoff(failures));
}
async fn durable_resume_namespaces() -> Vec<String> {
durable_resume_namespaces_checked()
.await
.unwrap_or_default()
}
async fn durable_resume_namespaces_checked() -> Option<Vec<String>> {
let store = crate::session::SESSIONS.get()?;
match crate::jobs::resumable_roster_agent_ids(&store.conn).await {
Ok(agent_ids) => Some(
agent_ids
.iter()
.map(|agent_id| crate::tools::chrome::run_session_namespace(agent_id))
.collect(),
),
Err(error) => {
debug!(
%error,
"could not read the resumable runs' agent ids — reclaiming without them"
);
None
}
}
}
fn claimed_names() -> HashSet<String> {
let parked = parked_releases().lock().unwrap_poison();
let queue = pending_releases().lock().unwrap_poison();
queue
.iter()
.chain(parked.iter().map(|(_, record)| record))
.flat_map(|record| record.names.iter().cloned())
.collect()
}
struct ProtectedNamespaces {
namespaces: Vec<String>,
}
impl ProtectedNamespaces {
fn snapshot(durable: &[String]) -> Self {
let now = Instant::now();
let live: Vec<String> = live_run_namespaces()
.lock()
.unwrap_poison()
.keys()
.cloned()
.collect();
let held: Vec<String> = pending_releases()
.lock()
.unwrap_poison()
.iter()
.filter(|record| within_hold(record, now))
.map(|record| record.namespace.clone())
.collect();
let mut protected = Self { namespaces: live };
protected.namespaces.extend(held);
protected.namespaces.extend(durable.iter().cloned());
protected
}
fn live_and_held() -> Self {
Self::snapshot(&[])
}
fn protects_live_or_held(name: &str) -> bool {
let live = live_run_namespaces().lock().unwrap_poison();
if live
.keys()
.any(|namespace| names_a_namespace(namespace, name))
{
return true;
}
drop(live);
let now = Instant::now();
pending_releases()
.lock()
.unwrap_poison()
.iter()
.any(|record| within_hold(record, now) && names_a_namespace(&record.namespace, name))
}
fn refresh(&mut self) {
self.namespaces.extend(Self::live_and_held().namespaces);
}
fn contains(&self, name: &str) -> bool {
self.namespaces
.iter()
.any(|prefix| names_a_namespace(prefix, name))
}
}
fn names_a_namespace(prefix: &str, name: &str) -> bool {
!prefix.is_empty() && name.starts_with(prefix)
}
const ISSUE_TARGET: &str = "chrome-tabs";
const LEFTOVER_MESSAGE: &str = "agent-run browser tabs could not be closed automatically";
const LEFT_OPEN_MESSAGE: &str =
"agent-run browser tabs were left open: the browser answered for them and did not close them";
const NO_DOOR_MESSAGE: &str = "the browser side here cannot list agent-run tab groups, so \
leftovers no record names cannot be closed";
const OWN_SESSION_MESSAGE: &str =
"the product's own browser-read session could not be confirmed released";
const OWN_SESSION_REASON: &str = "chrome-tabs:own-session";
const UNCLOSABLE_DETAIL: &str = "the product could not close these tabs and keeps retrying \
them: the browser extension answered that it does not hold \
them as its own (it lost its ledger — updated, re-installed or \
restarted since they were created — or never created them), or \
the browser answered this run's close attempts and left the \
tabs standing — the removal refused as not its own to make, or \
answered while they still stood";
const LEFT_OPEN_DETAIL: &str = "the browser extension answered this run's close and the tabs \
were left standing: it refused the removal as not its own to \
make, or answered it while the tabs still stood";
const NO_OWNERSHIP_DOOR_REASON: &str = "chrome-tabs:no-ownership-door";
const UNREADABLE_MESSAGE: &str = "the browser's answers gave this cleanup no decision it could act on, so no name was concluded from them";
const UNREADABLE_REASON: &str = "chrome-tabs:unreadable-answer";
const SESSION_NAMES_MAX_CHARS: usize = 300;
static REPORTED_REASONS: OnceLock<Mutex<HashSet<String>>> = OnceLock::new();
fn reported_reasons() -> &'static Mutex<HashSet<String>> {
REPORTED_REASONS.get_or_init(|| Mutex::new(HashSet::new()))
}
#[derive(Clone, Copy, PartialEq, Eq)]
enum IssueKind {
Unclosable,
LeftOpen,
}
impl IssueKind {
fn message(self) -> &'static str {
match self {
Self::Unclosable => LEFTOVER_MESSAGE,
Self::LeftOpen => LEFT_OPEN_MESSAGE,
}
}
fn reason_prefix(self) -> &'static str {
match self {
Self::Unclosable => "chrome-tabs:",
Self::LeftOpen => "chrome-tabs-left-open:",
}
}
fn detail(self) -> &'static str {
match self {
Self::Unclosable => UNCLOSABLE_DETAIL,
Self::LeftOpen => LEFT_OPEN_DETAIL,
}
}
}
struct Issue {
kind: IssueKind,
namespace: String,
run: String,
name: String,
}
struct IssueGroup {
kind: IssueKind,
namespace: String,
run: String,
names: Vec<String>,
}
async fn report_issues(issues: Vec<Issue>) {
let mut groups: Vec<IssueGroup> = Vec::new();
for issue in issues {
match groups
.iter_mut()
.find(|group| group.kind == issue.kind && group.namespace == issue.namespace)
{
Some(group) => {
if group.run.is_empty() {
group.run.clone_from(&issue.run);
}
group.names.push(issue.name);
}
None => groups.push(IssueGroup {
kind: issue.kind,
namespace: issue.namespace,
run: issue.run,
names: vec![issue.name],
}),
}
}
for group in groups {
let reason = format!("{}{}", group.kind.reason_prefix(), group.namespace);
let fields = serde_json::json!({
"detail": group.kind.detail(),
"run": group.run,
"sessions": group.names.len(),
"session": crate::util::truncate(&group.names.join(", "), SESSION_NAMES_MAX_CHARS),
});
report_once(group.kind.message(), &reason, fields).await;
}
}
async fn report_once(message: &str, reason: &str, fields: serde_json::Value) {
if reported_reasons().lock().unwrap_poison().contains(reason) {
return;
}
#[cfg(test)]
if let Some(sink) = release_settings().issues {
sink.lock()
.unwrap_poison()
.push((message.to_string(), reason.to_string(), fields));
return mark_reported(reason);
}
if matches!(
crate::logs::record_issue_once(message, reason, ISSUE_TARGET, fields).await,
crate::logs::IssueWrite::Written | crate::logs::IssueWrite::AlreadyRecorded
) {
mark_reported(reason);
}
}
fn mark_reported(reason: &str) {
reported_reasons()
.lock()
.unwrap_poison()
.insert(reason.to_string());
}
#[cfg(test)]
type CapturedIssues = std::sync::Arc<Mutex<Vec<(String, String, serde_json::Value)>>>;
const ATTEMPT_TIMEOUT: Duration = crate::tools::chrome_tabs::READ_PHASE_RESERVE
.saturating_mul(2)
.saturating_add(crate::chrome::SESSION_STOP_TIMEOUT);
#[derive(Clone, Debug)]
struct ReleaseSettings {
attempt_timeout: Duration,
retry_base: Duration,
boot_grace: Duration,
store: Option<PathBuf>,
#[cfg(test)]
cli: Option<PathBuf>,
#[cfg(test)]
issues: Option<CapturedIssues>,
}
fn release_settings() -> ReleaseSettings {
#[cfg(test)]
if let Some(settings) = test_release_settings() {
return settings;
}
ReleaseSettings {
attempt_timeout: ATTEMPT_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,
#[cfg(test)]
issues: 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();
reported_reasons().lock().unwrap_poison().clear();
*reclaim_schedule().lock().unwrap_poison() = ReclaimSchedule::default();
stuck_namespaces().lock().unwrap_poison().clear();
pending_forget().lock().unwrap_poison().clear();
*forget_retry().lock().unwrap_poison() = None;
evicted_runs().lock().unwrap_poison().clear();
crate::tools::chrome_tabs::reset_unreadable_answer();
}
#[cfg(test)]
fn pending_forget_snapshot() -> Vec<(String, u32)> {
pending_forget().lock().unwrap_poison().clone()
}
#[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::{AGENT_TAB_PREFIX, ChromeRunSessions};
use crate::tools::chrome_tab_ledger;
use crate::tools::chrome_tabs::OWN_SESSION;
use std::fs;
use std::path::Path;
use std::sync::Arc;
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);
}
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
}
fn pending_left_open() -> Vec<u32> {
pending_releases()
.lock()
.unwrap_poison()
.iter()
.map(|entry| entry.left_open)
.collect()
}
fn write_browser_state(path: &Path, names: &[String], owned: bool) {
let each: Vec<(&str, bool)> = names.iter().map(|name| (name.as_str(), owned)).collect();
write_browser_state_each(path, &each);
}
fn write_browser_state_each(path: &Path, owned: &[(&str, bool)]) {
let index = |i: usize| i64::try_from(i).expect("a small test index");
let mut lines: Vec<String> = owned
.iter()
.enumerate()
.map(|(i, (_, held))| {
format!(
"tab {} {} {}",
100 + index(i),
10 + index(i),
u8::from(*held)
)
})
.collect();
lines.extend(
owned
.iter()
.enumerate()
.map(|(i, (name, _))| format!("group {} {name}", 10 + index(i))),
);
fs::write(path, lines.join("\n")).expect("write browser state");
}
#[expect(clippy::too_many_lines)] fn stub_script(log: &Path, mode: &Path, state: &Path, mark: &Path) -> String {
r#"#!/bin/sh
printf '%s\n' "$*" >> __LOG__
mode="$(cat __MODE__)"
case "$mode" in
refuse) printf '%s' '{"success":false,"error":"the browser extension refused"}'; exit 1 ;;
silent) exec sleep 5 ;;
esac
case "$mode" in
no-extension|disabled-extension)
case "$*" in
# The real envelope of a host with no usable extension: the host manifest IS
# installed (this product writes it on every start) while nothing answers an
# `extension call`. For `no-extension` both the live version and chromeExtension are
# null (`installed` is not the extension's state); for `disabled-extension` the
# extension is there with a version that passes the gate and a non-empty
# `disableReasons`, which Chrome reports for one switched off in chrome://extensions.
"extension status"*)
if [ "$mode" = disabled-extension ]; then
printf '%s' '{"success":true,"data":{"installed":true,"liveExtensionVersion":null,"chromeExtension":{"version":"0.5.25","disableReasons":["user"]}}}'
else
printf '%s' '{"success":true,"data":{"installed":true,"liveExtensionVersion":null,"chromeExtension":null}}'
fi
exit 0 ;;
*) printf '%s' '{"success":false,"error":"no browser extension"}'; exit 1 ;;
esac ;;
ledger-unreadable)
# `extension status` reads exactly like `no-extension`, while `tabGroups.query` still
# answers from the state file (below). Both ownership reads fail with the state-less
# envelope, so the close must fall back and map each live group back to its name. The
# fallback's own `session stop` succeeds except for a name carrying `refuse`, so the
# mapping has two different outcomes to place.
case "$*" in
"extension status"*) printf '%s' '{"success":true,"data":{"installed":true,"liveExtensionVersion":null,"chromeExtension":null}}'; exit 0 ;;
"extension state"*|"extension call tabs.query"*) printf '%s' '{"success":false,"error":"no browser extension"}'; exit 1 ;;
"session stop"*refuse*) printf '%s' '{"success":false,"error":"the session is wedged"}'; exit 1 ;;
esac ;;
garbled-answer)
# A live extension that answers one read in a shape no parser can use: the group list
# carries an element without an integer id, which this route refuses to guess a part of.
case "$*" in
"extension call tabGroups.query"*) printf '%s' '{"success":true,"data":{"result":[{"id":"one","title":"agent-tab-0badf00d0000-x"}]}}'; exit 0 ;;
esac ;;
stop-refused)
# The shape a session whose browser the extension no longer reaches answers a stop with:
# every read works, and `session stop` reports that a tab it created could not be closed
# and that ownership was kept — chrome-use's own wording, which nothing else in this stub
# produces, so a test can tell a stop that cannot close a group from one that never ran.
case "$*" in
"session stop"*) printf '%s' '{"success":false,"error":"stopped session daemon, but 1 tab it created could not be closed; ownership was preserved"}'; exit 1 ;;
esac ;;
esac
case "$*" in
"extension call tabGroups.query"*)
printf '{"success":true,"data":{"result":[%s]}}' "$(sed -n 's/^group \([0-9]*\) \(.*\)$/{"id":\1,"title":"\2"}/p' __STATE__ | paste -sd, -)" ;;
"extension call tabs.query"*)
printf '{"success":true,"data":{"result":[%s]}}' "$(sed -n 's/^tab \([0-9]*\) \([0-9-]*\) .*/{"id":\1,"groupId":\2}/p' __STATE__ | paste -sd, -)" ;;
"extension state"*)
# `amnesia-once` answers the first ledger read with an empty one — the shape a
# swallowed storage read takes, which reads exactly like a ledger that lost its tabs.
if [ "$mode" = amnesia-once ] && [ ! -e __MARK__ ]; then
: > __MARK__
printf '%s' '{"success":true,"data":{"ownedTabs":[]}}'
else
printf '{"success":true,"data":{"ownedTabs":[%s]}}' "$(awk '$1=="tab" && $4=="1"{printf "%s%s", (o?",":""), $2; o=1}' __STATE__)"
fi ;;
"extension status"*)
printf '%s' '{"success":true,"data":{"installed":true,"liveExtensionVersion":"0.5.25","chromeExtension":{"version":"0.5.25"}}}' ;;
"extension call tabs.remove"*)
# The ids must arrive as the call's own single argument, a nested array
# ([[1,2]]): the extension spreads the JSON as chrome.tabs.remove's positional
# arguments, so [1,2] would be read as a tab id plus a callback.
case "$4" in
"[["*) ;;
*) printf '%s' '{"success":false,"error":"tabs.remove wants the tab ids as its first argument"}'; exit 1 ;;
esac
# The two cores of the close's promise, told apart purely by the leg's error text.
# `remove-sentence` is the extension's own refusal sentence, relayed verbatim through
# the helper — the one failing leg that counts as the browser answering about the tabs.
# `remove-transient` is a leg that never reached the browser at all, so it must never be
# read as the browser answering for the tabs. Both answer every read above normally, so the
# ledger holds the group's tabs and this leg is actually reached.
if [ "$mode" = remove-sentence ]; then
printf '%s' '{"success":false,"error":"call: tabs.remove refused — tab 7 is not owned by this relay (agent-created or adopted tabs only)"}'
exit 1
fi
if [ "$mode" = remove-transient ]; then
printf '%s' "{\"success\":false,\"error\":\"relay isn't connected\"}"
exit 1
fi
ids="$(printf '%s' "$4" | tr -d '[]' | tr ',' ' ')"
# The real door refuses the WHOLE call unless every id is one it holds, so a
# wrong or foreign id must fail the call here too.
for id in $ids; do
awk -v id="$id" '$1=="tab" && $2==id && $4=="1"{found=1} END{exit !found}' __STATE__ \
|| { printf '%s' '{"success":false,"error":"tab '"$id"' is not owned by this relay"}'; exit 1; }
done
# `keep` is the browser that answers and leaves its own tabs open anyway. Otherwise
# exactly the tabs asked for go, and nothing else: a close that asks for too few — or
# for a held tab of another session's group — must leave the browser changed
# accordingly, so a test can catch it. A group whose last tab went no longer exists.
if [ "$mode" != keep ]; then
for id in $ids; do
sed "/^tab $id /d" __STATE__ > __STATE__.new && mv __STATE__.new __STATE__
done
awk 'NR==FNR{if ($1=="tab") live[$3]=1; next} $1=="group" && !($2 in live){next} {print}' __STATE__ __STATE__ > __STATE__.new && mv __STATE__.new __STATE__
fi
printf '%s' '{"success":true}' ;;
*)
printf '%s' '{"success":true}' ;;
esac
exit 0
"#
.replace("__LOG__", &log.display().to_string())
.replace("__MODE__", &mode.display().to_string())
.replace("__STATE__", &state.display().to_string())
.replace("__MARK__", &mark.display().to_string())
}
struct ReleaseGuard {
state: PathBuf,
mode: PathBuf,
log: PathBuf,
issues: CapturedIssues,
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");
let state = dir.path().join("state.json");
let script = stub_script(&log, &mode, &state, &dir.path().join("mark"));
fs::write(&cli, script).expect("write release stub");
crate::util::test::make_executable(&cli);
Self::write_mode(&mode, "ok");
write_browser_state(&state, &[], true);
let issues = Arc::new(Mutex::new(Vec::new()));
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)),
issues: Some(Arc::clone(&issues)),
});
let guard = Self {
state,
mode,
log,
issues,
previous,
dir,
};
let _ = chrome_tabs::close_sessions(
&cli,
&["warmup".to_string()],
OWN_SESSION,
Instant::now() + attempt_timeout,
&|_: &str| false,
)
.await;
chrome_tabs::release_own_session(&cli, Instant::now() + 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 write_state(&self, names: &[String], owned: bool) {
write_browser_state(&self.state, names, owned);
}
fn write_state_each(&self, owned: &[(&str, bool)]) {
write_browser_state_each(&self.state, owned);
}
fn group_titles(&self) -> Vec<String> {
fs::read_to_string(&self.state)
.unwrap_or_default()
.lines()
.filter_map(|line| line.strip_prefix("group "))
.filter_map(|rest| rest.split_once(' ').map(|(_, title)| title.to_string()))
.collect()
}
#[expect(
clippy::unused_self,
reason = "called as a guard method so every seam installs the same way"
)]
fn set_binary_absent(&self) {
let mut settings = release_settings();
settings.cli = None;
swap_release_settings(settings);
}
#[expect(
clippy::unused_self,
reason = "called as a guard method so every seam installs the same way"
)]
fn arm_sweep(&self) {
schedule_reclaim_now();
}
fn log_lines(&self) -> Vec<String> {
fs::read_to_string(&self.log)
.unwrap_or_default()
.lines()
.map(str::to_string)
.collect()
}
fn invoked(&self, needle: &str) -> bool {
self.log_lines().iter().any(|line| line.contains(needle))
}
fn removed_tab_ids(&self) -> Vec<i64> {
self.log_lines()
.iter()
.filter(|line| line.starts_with("extension call tabs.remove"))
.flat_map(|line| {
let Some(start) = line.find("[[") else {
return Vec::new();
};
let rest = &line[start + 2..];
let Some(end) = rest.find("]]") else {
return Vec::new();
};
rest[..end]
.split(',')
.filter_map(|id| id.parse::<i64>().ok())
.collect()
})
.collect()
}
fn only_swept(&self) -> bool {
self.log_lines().iter().all(|line| {
!line.contains("tabs.remove")
&& (!line.contains("session stop")
|| line.contains(&format!("--session {OWN_SESSION}")))
})
}
fn clear_log(&self) {
fs::write(&self.log, "").expect("reset release log");
}
fn issues(&self) -> Vec<(String, String, serde_json::Value)> {
self.issues.lock().unwrap_poison().clone()
}
fn clear_issues(&self) {
self.issues.lock().unwrap_poison().clear();
}
fn store_path(&self) -> PathBuf {
self.dir.path().join(RELEASE_FILE_NAME)
}
}
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, agent_id);
}
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 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(ATTEMPT_TIMEOUT).await;
let names = run_session_names("run-verified", &["agent-tab-v-default", "agent-tab-v-docs"]);
guard.write_state(&names, true);
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 settled gone"
);
assert!(
guard.invoked("extension call tabs.remove"),
"the close went through the browser's own extension"
);
assert_eq!(
guard.removed_tab_ids(),
vec![100, 101],
"exactly the run's own tabs were asked for: {:?}",
guard.log_lines()
);
assert!(
guard.invoked(&format!(
"session stop --force --json --session {}",
names[0]
)),
"the record's own session the reads ran in is let go once its group is settled, so \
it cannot outlive the pass: {:?}",
guard.log_lines()
);
assert!(
!guard.only_swept(),
"a pass that asked tabs.remove is not a sweep, whatever session it ran in: {:?}",
guard.log_lines()
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_record_the_browser_reports_gone_is_dropped() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
queue_names("run-gone", &["agent-tab-g-default"]);
release_due().await;
assert!(
pending_releases_snapshot().is_empty(),
"the browser holds no group by that name, so its tabs are gone"
);
assert!(!guard.invoked("extension call tabs.remove"));
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_record_the_browser_never_settles_stays_queued_past_five_attempts() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
let names = run_session_names("run-refuse", &["agent-tab-r-stuck"]);
queue_names("run-refuse", &["agent-tab-r-stuck"]);
guard.set_mode("refuse");
for _ in 0..6 {
release_due().await;
assert_eq!(
pending_names_and_attempts().len(),
1,
"the record stays queued across attempts"
);
assert_eq!(
pending_names_and_attempts()[0].0,
names,
"and keeps its names"
);
}
let (queued, attempts) = &pending_names_and_attempts()[0];
assert_eq!(queued, &names, "the names are intact");
assert!(
*attempts > 5,
"past five attempts, still queued: {attempts}"
);
assert_eq!(
record_file_names(&guard.store_path()),
{
let mut sorted = names.clone();
sorted.sort();
sorted
},
"the record file still describes it — nothing drops an unanswered record from the queue"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_record_the_browser_reports_as_not_ours_is_reported_once_and_dropped() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
let namespace = ChromeRunSessions::for_run("run-unclosable")
.namespace()
.to_string();
let names = run_session_names("run-unclosable", &["agent-tab-u-default"]);
let session = names[0].clone();
guard.write_state(&names, false);
queue_names("run-unclosable", &["agent-tab-u-default"]);
release_due().await;
assert!(
pending_releases_snapshot().is_empty(),
"the record's route is done with those tabs, so they are not retried"
);
let issues = guard.issues();
assert_eq!(issues.len(), 1, "the leftover is reported once");
assert_eq!(issues[0].0, LEFTOVER_MESSAGE);
assert_eq!(issues[0].1, format!("chrome-tabs:{namespace}"));
assert_eq!(issues[0].2["run"], "run-unclosable");
assert_eq!(issues[0].2["sessions"], 1);
assert!(
guard.invoked(&format!("session stop --json --session {session}")),
"the session's own route is asked before any tab is called unclosable: {:?}",
guard.log_lines()
);
release_due().await;
assert_eq!(guard.issues().len(), 1, "and never again");
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_transient_empty_ledger_does_not_write_a_run_off() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
let names = run_session_names("run-amnesia", &["agent-tab-a-default"]);
guard.write_state(&names, true);
guard.set_mode("amnesia-once");
queue_names("run-amnesia", &["agent-tab-a-default"]);
release_due().await;
assert_eq!(
pending_names_and_attempts().len(),
1,
"an answer the second read did not repeat is not a verdict"
);
assert!(
guard.issues().is_empty(),
"and nothing is announced about those tabs"
);
release_due().await;
assert!(
pending_releases_snapshot().is_empty(),
"the next pass reads the real ledger and closes them"
);
assert!(guard.invoked("extension call tabs.remove"));
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_close_the_browser_leaves_open_is_reported_once_and_still_retried() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
let namespace = ChromeRunSessions::for_run("run-kept")
.namespace()
.to_string();
let names = run_session_names("run-kept", &["agent-tab-k-default"]);
guard.write_state(&names, true);
guard.set_mode("keep");
queue_names("run-kept", &["agent-tab-k-default"]);
release_due().await;
assert_eq!(
pending_names_and_attempts().len(),
1,
"one left-open answer proves nothing, so the record stays queued"
);
assert_eq!(
record_file_names(&guard.store_path()),
{
let mut sorted = names.clone();
sorted.sort();
sorted
},
"and stays in the durable file"
);
let issues = guard.issues();
assert_eq!(issues.len(), 1, "the left-open answer is reported once");
assert_eq!(issues[0].0, LEFT_OPEN_MESSAGE);
assert_eq!(issues[0].1, format!("chrome-tabs-left-open:{namespace}"));
assert_eq!(issues[0].2["run"], "run-kept");
let before = pending_names_and_attempts()[0].1;
release_due().await;
assert_eq!(
guard.issues().len(),
1,
"and never again for the same run's namespace"
);
assert!(
pending_names_and_attempts()[0].1 > before,
"the retry ladder still moves"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_remove_the_extension_refuses_concludes_after_the_cap() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
let namespace = ChromeRunSessions::for_run("run-refused-sentence")
.namespace()
.to_string();
let names = run_session_names("run-refused-sentence", &["agent-tab-rs-default"]);
guard.write_state(&names, true);
guard.set_mode("remove-sentence");
queue_names("run-refused-sentence", &["agent-tab-rs-default"]);
for pass in 1..MAX_LEFT_OPEN_ATTEMPTS {
release_due().await;
assert_eq!(
pending_names_and_attempts().len(),
1,
"a relayed left-open answer short of the cap keeps the record queued"
);
assert_eq!(
pending_left_open(),
vec![pass],
"each pass is one browser answer, counted"
);
}
release_due().await;
assert!(
pending_names_and_attempts().is_empty(),
"the cap concludes the close route: the record stops retrying"
);
let issues = guard.issues();
let conclusion: Vec<_> = issues
.iter()
.filter(|issue| issue.1 == format!("chrome-tabs:{namespace}"))
.collect();
assert_eq!(
conclusion.len(),
1,
"the leftover is reported once as unclosable: {issues:?}"
);
assert_eq!(conclusion[0].0, LEFTOVER_MESSAGE);
assert_eq!(
conclusion[0].2["run"], "run-refused-sentence",
"and the conclusion names the run it came from"
);
assert!(
issues
.iter()
.any(|issue| issue.1 == format!("chrome-tabs-left-open:{namespace}")),
"the browser's own answer was told early, before the cap: {issues:?}"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_transient_remove_failure_is_never_counted_as_an_answer() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
let namespace = ChromeRunSessions::for_run("run-transient-remove")
.namespace()
.to_string();
let names = run_session_names("run-transient-remove", &["agent-tab-tr-default"]);
guard.write_state(&names, true);
guard.set_mode("remove-transient");
queue_names("run-transient-remove", &["agent-tab-tr-default"]);
for _ in 0..MAX_LEFT_OPEN_ATTEMPTS {
release_due().await;
}
assert_eq!(
pending_names_and_attempts(),
vec![(names.clone(), MAX_LEFT_OPEN_ATTEMPTS)],
"the transient leg leaves the record queued, attempted once per pass"
);
assert_eq!(
pending_left_open(),
vec![0],
"a leg that never reached the browser is never counted as an answer"
);
assert!(
guard.invoked("extension call tabs.remove"),
"the close did reach the removal leg — the transient failure is that leg's own \
answer: {:?}",
guard.log_lines()
);
assert!(
guard.issues().iter().all(|(_, reason, _)| {
reason != &format!("chrome-tabs-left-open:{namespace}")
&& reason != &format!("chrome-tabs:{namespace}")
}),
"a transient failure files neither the left-open nor the leftover row: {:?}",
guard.issues()
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_host_without_the_ownership_door_is_reported_once() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
guard.set_mode("no-extension");
let names = run_session_names("run-no-door", &["agent-tab-n-default"]);
let session = names[0].clone();
guard.write_state(&names, true);
queue_names("run-no-door", &["agent-tab-n-default"]);
release_due().await;
let host_rows = || -> Vec<(String, String, serde_json::Value)> {
guard
.issues()
.into_iter()
.filter(|issue| issue.1 == NO_OWNERSHIP_DOOR_REASON)
.collect()
};
assert_eq!(
host_rows().len(),
1,
"the absent door is reported once: {:?}",
guard.issues()
);
let issues = host_rows();
assert_eq!(issues[0].0, NO_DOOR_MESSAGE);
assert!(
issues[0].2["detail"]
.as_str()
.is_some_and(|d| !d.is_empty())
);
assert_eq!(
pending_names_and_attempts().len(),
1,
"a door-less host never writes the run's tabs off: the record keeps trying"
);
assert!(
guard.invoked(&format!("session stop --json --session {session}")),
"the only route such a host has is the helper's own stop"
);
assert!(
guard
.issues()
.iter()
.all(|issue| issue.1 != OWN_SESSION_REASON),
"no scratch-group row: the reads no longer run in the product's own session: {:?}",
guard.issues()
);
guard.arm_sweep();
release_due().await;
assert_eq!(host_rows().len(), 1, "and never again");
assert!(
guard
.issues()
.iter()
.all(|issue| issue.1 != OWN_SESSION_REASON),
"and no scratch-group row either: the reads run in the record's session"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_scratch_group_the_let_go_could_not_confirm_is_surfaced_once() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
guard.set_mode("stop-refused");
guard.write_state(&[OWN_SESSION.to_string()], false);
guard.arm_sweep();
release_due().await;
release_due().await;
let rows: Vec<(String, String, serde_json::Value)> = guard
.issues()
.into_iter()
.filter(|issue| issue.1 == OWN_SESSION_REASON)
.collect();
assert_eq!(
rows.len(),
1,
"one row for the product's own session: {:?}",
guard.issues()
);
assert_eq!(rows[0].0, OWN_SESSION_MESSAGE);
let detail = rows[0].2["detail"].as_str().unwrap_or_default();
assert!(
detail.contains("the release is not confirmed"),
"and it says only that the release was not confirmed: {detail}"
);
assert_eq!(
pending_forget_snapshot(),
vec![(OWN_SESSION.to_string(), 2)],
"and the session stays queued with one charged attempt per pass that asked its \
stop — the one route left that closes its group was answered with a refusal"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn the_fallback_maps_each_live_name_back_to_its_own_slot() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
guard.set_mode("ledger-unreadable");
let names = run_session_names(
"run-ledger-unreadable",
&["agent-tab-missing", "agent-tab-refuse", "agent-tab-live"],
);
guard.write_state(&names[1..], true);
queue_names(
"run-ledger-unreadable",
&["agent-tab-missing", "agent-tab-refuse", "agent-tab-live"],
);
release_due().await;
assert!(
guard.invoked(&format!("session stop --json --session {}", names[1])),
"the fallback's stop is asked for every name with a live group: {:?}",
guard.log_lines()
);
assert!(guard.invoked(&format!("session stop --json --session {}", names[2])));
assert!(
!guard.invoked(&format!("session stop --json --session {}", names[0])),
"and never for the name with no live group, which the one group read settled"
);
assert_eq!(
pending_names_and_attempts(),
vec![(vec![names[1].clone()], 1)],
"each name keeps its own outcome: the refused stop stays queued with its attempt, \
the answered one settles"
);
assert!(
guard.issues().is_empty(),
"and no leftover row is filed — the fallback's retry is silent: {:?}",
guard.issues()
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_host_whose_answers_are_unreadable_is_reported_once() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
guard.set_mode("garbled-answer");
let names = run_session_names("run-garbled", &["agent-tab-garbled"]);
guard.write_state(&names, true);
queue_names("run-garbled", &["agent-tab-garbled"]);
release_due().await;
release_due().await;
let rows: Vec<(String, String, serde_json::Value)> = guard
.issues()
.into_iter()
.filter(|issue| issue.1 == UNREADABLE_REASON)
.collect();
assert_eq!(
rows.len(),
1,
"one row for the host, whatever the passes: {:?}",
guard.issues()
);
assert_eq!(rows[0].0, UNREADABLE_MESSAGE);
let detail = rows[0].2["detail"].as_str().unwrap_or_default();
assert!(
detail.contains("settled nothing") && detail.contains("no parser here will guess at"),
"and it says what settled nothing, and why it is not guessed at: {detail}"
);
assert_eq!(
pending_names_and_attempts().len(),
1,
"while the record keeps the name queued: nothing was settled"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_host_with_a_disabled_extension_is_reported_once() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
guard.set_mode("disabled-extension");
let names = run_session_names("run-disabled-ext", &["agent-tab-d-disabled"]);
let session = names[0].clone();
guard.write_state(&names, true);
queue_names("run-disabled-ext", &["agent-tab-d-disabled"]);
release_due().await;
let issues: Vec<(String, String, serde_json::Value)> = guard
.issues()
.into_iter()
.filter(|issue| issue.1 == NO_OWNERSHIP_DOOR_REASON)
.collect();
assert_eq!(
issues.len(),
1,
"a disabled extension's door is missing too, and reported once: {:?}",
guard.issues()
);
assert_eq!(issues[0].0, NO_DOOR_MESSAGE);
assert!(
issues[0].2["detail"]
.as_str()
.is_some_and(|detail| detail.contains("disabled")),
"the row names the disabled case rather than only a missing extension: {issues:?}"
);
assert_eq!(
pending_names_and_attempts().len(),
1,
"and the record still keeps trying rather than being written off"
);
assert!(
guard.invoked(&format!("session stop --json --session {session}")),
"the fallback close is attempted on this host too: {:?}",
guard.log_lines()
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn the_reclaim_spares_every_group_a_run_may_come_back_to() {
crate::util::test::init_test_stores().await;
let store = crate::session::store();
let resumable_id = "ticket_reclaim_0_resume_analyst";
crate::util::test::JobRowBuilder::new(
&store.conn,
"job-reclaim-cu",
"analyze",
"analyst",
"ws",
)
.timestamps("2026-01-01T00:00:00Z")
.insert()
.await
.expect("insert a job that has not been terminalized");
store
.conn
.execute(
crate::jobs::AGENT_INSERT_SQL,
crate::jobs::agent_params(
"job-reclaim-cu",
resumable_id,
crate::jobs::AgentKind::Analyst,
Some(0),
"task",
),
)
.await
.expect("insert a roster row");
let finished_id = "ticket_reclaim_1_resume_analyst";
store
.conn
.execute(
crate::jobs::AGENT_INSERT_SQL,
crate::jobs::agent_params(
"job-reclaim-cu",
finished_id,
crate::jobs::AgentKind::Analyst,
Some(1),
"task",
),
)
.await
.expect("insert a finished roster row");
crate::jobs::write_agent_outcome(
&store.conn,
"job-reclaim-cu",
finished_id,
crate::jobs::RowStatus::Done,
Some("{}"),
)
.await
.expect("mark the roster slot done");
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
let resumable = format!(
"{}default",
crate::tools::chrome::run_session_namespace(resumable_id)
);
let finished = format!(
"{}default",
crate::tools::chrome::run_session_namespace(finished_id)
);
let held = run_session_names("run-held-rec", &["agent-tab-h-default"]).remove(0);
queue_names_with_hold(
"run-held-rec",
&["agent-tab-h-default"],
Duration::from_hours(1),
);
let claimed = run_session_names("run-claimed-rec", &["agent-tab-c-default"]).remove(0);
queue_names("run-claimed-rec", &["agent-tab-c-default"]);
let live = ChromeRunSessions::for_run("run-live-rec");
let live_name = format!("{}default", live.namespace());
let orphan = format!("{AGENT_TAB_PREFIX}0123456789ab-orphan");
guard.write_state(
&[
resumable.clone(),
finished.clone(),
held.clone(),
claimed.clone(),
live_name.clone(),
orphan.clone(),
],
true,
);
let cli = release_cli();
let protected = ProtectedNamespaces::snapshot(&durable_resume_namespaces().await);
let reclaim = reclaim_names(
cli.as_deref(),
OWN_SESSION,
&protected,
Instant::now() + Duration::from_secs(10),
true,
)
.await;
drop(live);
store
.conn
.execute(
"DELETE FROM jobs WHERE id = ?1",
crate::db::params!["job-reclaim-cu"],
)
.await
.expect("delete the test job");
assert!(
reclaim.names.contains(&orphan),
"the group nothing accounts for is the one to reclaim: {:?}",
reclaim.names
);
for spared in [&resumable, &held, &claimed, &live_name] {
assert!(
!reclaim.names.contains(spared),
"{spared} belongs to a run that may come back: {:?}",
reclaim.names
);
}
assert!(
reclaim.names.contains(&finished),
"a slot that already finished is never re-dispatched, so its group is reclaimed like \
any other one nothing will work under again: {:?}",
reclaim.names
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_reclaimed_group_the_browser_leaves_open_is_reported() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
let orphan = format!("{AGENT_TAB_PREFIX}0123456789ab-orphan");
guard.write_state(&[orphan], true);
guard.set_mode("keep");
guard.arm_sweep();
release_due().await;
let issues = guard.issues();
assert_eq!(
issues.len(),
1,
"the left-open answer is reported: {issues:?}"
);
assert_eq!(issues[0].0, LEFT_OPEN_MESSAGE);
assert_eq!(
issues[0].1,
format!("chrome-tabs-left-open:{AGENT_TAB_PREFIX}0123456789ab-")
);
assert_eq!(
issues[0].2["run"], "",
"no record was evicted for a group nothing recorded, so no run can be named"
);
guard.arm_sweep();
release_due().await;
assert_eq!(guard.issues().len(), 1, "and never again");
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_stuck_leftover_is_not_interrogated_again_every_sweep() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
let orphan = format!("{AGENT_TAB_PREFIX}0123456789ab-orphan");
guard.write_state(std::slice::from_ref(&orphan), false);
guard.arm_sweep();
release_due().await;
let issues = guard.issues();
assert_eq!(issues.len(), 1, "the leftover is reported: {issues:?}");
assert_eq!(issues[0].0, LEFTOVER_MESSAGE);
assert_eq!(
issues[0].1,
format!("chrome-tabs:{AGENT_TAB_PREFIX}0123456789ab-")
);
guard.clear_log();
guard.arm_sweep();
release_due().await;
assert!(
!guard.invoked(&format!("session stop --json --session {orphan}")),
"a remembered verdict is not re-confirmed every sweep: {:?}",
guard.log_lines()
);
assert_eq!(guard.issues().len(), 1, "and never reported again either");
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn an_expired_stuck_verdict_is_dropped_and_its_group_examined_again() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
let orphan = format!("{AGENT_TAB_PREFIX}0123456789ab-orphan");
guard.write_state(std::slice::from_ref(&orphan), false);
guard.arm_sweep();
release_due().await;
assert_eq!(guard.issues().len(), 1, "the leftover is reported once");
let expired = Instant::now()
.checked_sub(STUCK_REVISIT)
.expect("the process clock reaches back the revisit window");
stuck_namespaces().lock().unwrap_poison().insert(
crate::tools::chrome::session_namespace(&orphan).to_string(),
expired,
);
guard.clear_log();
guard.arm_sweep();
release_due().await;
assert!(
guard.invoked(&format!("session stop --json --session {orphan}")),
"the group is examined again, not skipped by an expired verdict: {:?}",
guard.log_lines()
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn the_pass_lets_go_of_the_session_its_reads_ran_in() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
guard.write_state(&[OWN_SESSION.to_string()], true);
let cli = release_cli().expect("the stub binary");
let deadline = Instant::now() + Duration::from_secs(10);
chrome_tabs::agent_group_names(&cli, OWN_SESSION, deadline)
.await
.expect("the live groups");
chrome_tabs::release_own_session(&cli, deadline).await;
assert_eq!(
guard.removed_tab_ids(),
vec![100],
"the scratch group is closed through the extension's own door: {:?}",
guard.log_lines()
);
assert!(
guard.invoked(&format!(
"session stop --force --json --session {OWN_SESSION}"
)),
"the group the ledger door confirmed gone is followed by the record-dropping stop: \
{:?}",
guard.log_lines()
);
guard.clear_log();
chrome_tabs::release_own_session(&cli, deadline).await;
assert!(
guard.log_lines().is_empty(),
"an idle pass spawns nothing: {:?}",
guard.log_lines()
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_scratch_group_the_ledger_cannot_close_goes_to_its_session_stop() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
guard.set_mode("stop-refused");
guard.write_state(&[OWN_SESSION.to_string()], false);
let cli = release_cli().expect("the stub binary");
let deadline = Instant::now() + ATTEMPT_TIMEOUT;
chrome_tabs::agent_group_names(&cli, OWN_SESSION, deadline)
.await
.expect("the live groups");
let_go_own_session(&cli, deadline).await;
assert!(
guard.removed_tab_ids().is_empty(),
"a tab the extension does not hold is never asked for: {:?}",
guard.log_lines()
);
assert!(
!guard.invoked("session stop"),
"and the unconfirmed let-go does not stop the session itself — the pass's let-go \
step is the one stop site: {:?}",
guard.log_lines()
);
assert!(
pending_forget_snapshot() == vec![(OWN_SESSION.to_string(), 0)],
"it hands the session to the driver's retry path instead: {:?}",
guard.log_lines()
);
guard.clear_log();
forget_settled(&cli, &[], deadline).await;
let stops = guard
.log_lines()
.iter()
.filter(|line| line.starts_with("session stop"))
.count();
assert_eq!(
stops,
1,
"one stop for the pass, with the verb that closes the group: {:?}",
guard.log_lines()
);
assert!(
guard.invoked(&format!("session stop --json --session {OWN_SESSION}")),
"and it is the graceful form: {:?}",
guard.log_lines()
);
assert!(
!guard.invoked(&format!(
"session stop --force --json --session {OWN_SESSION}"
)),
"never the form that would drop the record and leave the group standing: {:?}",
guard.log_lines()
);
assert!(
pending_forget_snapshot() == vec![(OWN_SESSION.to_string(), 1)],
"a whole-bound stop that could not close it charges one attempt and keeps the name \
queued: {:?}",
guard.log_lines()
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_settled_scratch_session_is_not_started_again_by_the_next_pass() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
guard.write_state(&[OWN_SESSION.to_string()], false);
let cli = release_cli().expect("the stub binary");
let deadline = Instant::now() + ATTEMPT_TIMEOUT;
chrome_tabs::agent_group_names(&cli, OWN_SESSION, deadline)
.await
.expect("the live groups");
let_go_own_session(&cli, deadline).await;
assert_eq!(
pending_forget_snapshot(),
vec![(OWN_SESSION.to_string(), 0)],
"the unconfirmed let-go is queued"
);
guard.clear_log();
forget_settled(&cli, &[], deadline).await;
assert!(
guard.invoked(&format!("session stop --json --session {OWN_SESSION}")),
"the let-go step asks the graceful stop: {:?}",
guard.log_lines()
);
assert!(
pending_forget_snapshot().is_empty(),
"which settles the session and drops the name"
);
guard.clear_log();
chrome_tabs::release_own_session(&cli, deadline).await;
assert!(
guard.log_lines().is_empty(),
"and the next pass owes nothing for those reads: {:?}",
guard.log_lines()
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn overlapping_let_go_passes_keep_each_others_names() {
let budget = Duration::from_millis(300);
let guard = ReleaseGuard::install(budget).await;
guard.set_mode("silent");
let cli = release_cli().expect("the stub binary");
let deadline = Instant::now() + budget;
let first = run_session_names("run-forget-overlap-a", &["agent-tab-fo-a"]).remove(0);
let second = run_session_names("run-forget-overlap-b", &["agent-tab-fo-b"]).remove(0);
tokio::join!(
forget_settled(&cli, std::slice::from_ref(&first), deadline),
forget_settled(&cli, std::slice::from_ref(&second), deadline),
);
let queued: Vec<String> = pending_forget_snapshot()
.into_iter()
.map(|(name, _)| name)
.collect();
assert!(
queued.contains(&first) && queued.contains(&second),
"both passes' names survive: {queued:?}"
);
}
#[test]
fn the_let_go_queue_is_deduped_and_bounded() {
let mut names = vec![
("s0".to_string(), 3),
("s1".to_string(), 0),
("s0".to_string(), 0),
("s2".to_string(), 1),
];
dedupe_and_bound_forgotten(&mut names);
assert_eq!(
names,
vec![
("s0".to_string(), 3),
("s1".to_string(), 0),
("s2".to_string(), 1)
],
"a repeat keeps its first place, attempts and all"
);
let mut full: Vec<(String, u32)> = (0..MAX_PENDING_FORGET)
.map(|i| (format!("s{i}"), 0))
.collect();
full.push(("newest".to_string(), 0));
dedupe_and_bound_forgotten(&mut full);
assert_eq!(full.len(), MAX_PENDING_FORGET, "the queue is bounded");
assert_eq!(
full[0],
("s1".to_string(), 0),
"the oldest name is what the bound drops"
);
assert_eq!(full.last(), Some(&("newest".to_string(), 0)));
let mut with_scratch: Vec<(String, u32)> = (0..MAX_PENDING_FORGET)
.map(|i| (format!("s{i}"), 0))
.collect();
with_scratch.insert(0, (OWN_SESSION.to_string(), 0));
dedupe_and_bound_forgotten(&mut with_scratch);
assert_eq!(with_scratch[0], (OWN_SESSION.to_string(), 0));
assert!(
!with_scratch.iter().any(|(name, _)| name == "s0"),
"so the oldest bookkeeping name is what the cap drops instead"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_settled_session_left_over_is_let_go_by_a_later_pass() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
let cli = release_cli().expect("the stub binary");
let name = run_session_names("run-forget", &["agent-tab-f-default"]).remove(0);
forget_settled(&cli, std::slice::from_ref(&name), Instant::now()).await;
assert_eq!(pending_forget_snapshot(), vec![(name.clone(), 0)]);
assert!(
!guard.invoked("session stop"),
"and spawned nothing: {:?}",
guard.log_lines()
);
assert!(
forget_retry_at().is_some(),
"while the queue arms its own retry, so a budget-starved pass is not the last word"
);
for _ in 0..MAX_FORGET_ATTEMPTS {
forget_settled(&cli, &[], Instant::now()).await;
}
assert_eq!(pending_forget_snapshot(), vec![(name.clone(), 0)]);
forget_settled(&cli, &[], Instant::now() + ATTEMPT_TIMEOUT).await;
assert!(pending_forget_snapshot().is_empty());
assert!(
guard.invoked("session stop"),
"the helper's own recipe is what clears the stale record: {:?}",
guard.log_lines()
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_stop_cut_to_the_pass_is_not_charged_an_attempt() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
guard.set_mode("refuse");
let cli = release_cli().expect("the stub binary");
let name = run_session_names("run-cut-stop", &["agent-tab-cs-default"]).remove(0);
*pending_forget().lock().unwrap_poison() = vec![(name.clone(), 2)];
forget_settled(&cli, &[], Instant::now() + Duration::from_secs(10)).await;
assert!(
guard.invoked("session stop"),
"the stop is still attempted: {:?}",
guard.log_lines()
);
assert_eq!(
pending_forget_snapshot(),
vec![(name, 2)],
"and its refused answer charges nothing: the rung is where it was"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn the_shutdown_flush_reaches_the_let_go_queue() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
let name = run_session_names("run-flush-forget", &["agent-tab-ff-default"]).remove(0);
queue_forget(&name);
assert!(pending_releases_snapshot().is_empty(), "no record is due");
flush_pending_run_releases(SHUTDOWN_RELEASE_FLUSH_BUDGET).await;
assert!(
guard.invoked(&format!("session stop --force --json --session {name}")),
"the queued session is stopped before the process goes down: {:?}",
guard.log_lines()
);
assert!(
pending_forget_snapshot().is_empty(),
"and a helper that answers inside the flush's budget lets the name go"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn the_let_go_queue_gives_up_after_its_attempt_bound() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
guard.set_mode("refuse");
let cli = release_cli().expect("the stub binary");
let name = run_session_names("run-forget-bound", &["agent-tab-fb-default"]).remove(0);
let deadline = Instant::now() + ATTEMPT_TIMEOUT;
for attempts in 1..MAX_FORGET_ATTEMPTS {
forget_settled(&cli, std::slice::from_ref(&name), deadline).await;
assert_eq!(
pending_forget_snapshot(),
vec![(name.clone(), attempts)],
"the name waits with one more failed attempt on it"
);
}
forget_settled(&cli, &[], deadline).await;
assert!(
pending_forget_snapshot().is_empty(),
"the bound drops it instead of retrying it for the whole process"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn the_scratch_session_is_retried_past_the_attempt_bound() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
guard.set_mode("stop-refused");
let cli = release_cli().expect("the stub binary");
let deadline = Instant::now() + ATTEMPT_TIMEOUT;
for _ in 0..MAX_FORGET_ATTEMPTS {
forget_settled(&cli, &[OWN_SESSION.to_string()], deadline).await;
}
assert_eq!(
pending_forget_snapshot(),
vec![(OWN_SESSION.to_string(), MAX_FORGET_ATTEMPTS)],
"the scratch session stays queued past the bound"
);
assert!(
forget_retry_at().is_some(),
"and the queue's own deadline is what retries it, not a restart"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_let_go_behind_a_call_of_ours_waits_the_stop_bound() {
let _guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
let cli = release_cli().expect("the stub binary");
let deadline = Instant::now() + ATTEMPT_TIMEOUT;
let call = chrome_tabs::hold_door_call_in_flight();
let outcome = chrome_tabs::forget_settled_sessions(
&cli,
std::slice::from_ref(&OWN_SESSION.to_string()),
deadline,
&|_: &str| false,
)
.await;
assert_eq!(outcome.busy, vec![OWN_SESSION.to_string()]);
assert!(
outcome.deferred.is_empty() && outcome.tried.is_empty(),
"and it is neither failed nor a deferral of a bound that did not fit"
);
forget_settled(&cli, &[OWN_SESSION.to_string()], deadline).await;
assert_eq!(
pending_forget_snapshot(),
vec![(OWN_SESSION.to_string(), 0)],
"the name is kept, at the rung it had"
);
let wait = forget_retry_at()
.expect("while the call holds the session the queue arms its own retry")
.saturating_duration_since(Instant::now());
assert!(
BUSY_FORGET_FLOOR.saturating_sub(wait) < Duration::from_secs(1),
"and the wake waits out one whole stop, not the base gap: {wait:?}"
);
drop(call);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_session_a_run_owns_again_leaves_the_let_go_queue() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
guard.set_mode("refuse");
let cli = release_cli().expect("the stub binary");
let name = run_session_names("run-owned-again", &["agent-tab-ro-default"]).remove(0);
let deadline = Instant::now() + Duration::from_secs(10);
let tracker = ChromeRunSessions::for_run("run-owned-again");
forget_settled(&cli, std::slice::from_ref(&name), deadline).await;
assert!(
!guard.invoked("session stop"),
"the stop would end the run's own session: {:?}",
guard.log_lines()
);
assert!(
pending_forget_snapshot().is_empty(),
"and nothing waits for a name the run owns"
);
drop(tracker);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_let_go_retry_cannot_wake_a_driver_with_no_binary() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
*forget_retry().lock().unwrap_poison() = Some(Instant::now());
assert!(
next_release_deadline().is_some(),
"with a binary to run it, the retry is a real deadline"
);
guard.set_binary_absent();
assert_eq!(
next_release_deadline(),
None,
"with none, nothing reads the queue's past deadline as work"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn the_reclaim_asks_the_browser_once_per_armed_ask() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
let orphan = format!("{AGENT_TAB_PREFIX}0123456789ab-orphan");
guard.write_state(&[orphan], true);
guard.arm_sweep();
release_due().await;
assert!(
guard.invoked("extension call tabs.remove"),
"the first pass reclaims"
);
assert!(
next_release_deadline().is_none(),
"a clean ask leaves nothing scheduled — no ask, no queued retry"
);
guard.clear_log();
release_due().await;
assert!(
guard.log_lines().is_empty(),
"a successful ask schedules nothing, so no interval re-asks: {:?}",
guard.log_lines()
);
schedule_reclaim_now();
release_due().await;
assert!(
guard.invoked("extension call tabGroups.query"),
"an event — a run end queuing names — is what arms the next ask"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn repeated_silent_sweeps_climb_the_ladder() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
widen_retry_base(Duration::from_secs(5));
let orphan = format!("{AGENT_TAB_PREFIX}0123456789ab-orphan");
guard.write_state(&[orphan], true);
guard.set_mode("remove-transient");
guard.arm_sweep();
release_due().await;
let after_first = reclaim_schedule().lock().unwrap_poison().failures;
let first = next_release_deadline()
.expect("an unfinished ask keeps the next one armed")
.saturating_duration_since(Instant::now());
reclaim_schedule().lock().unwrap_poison().next_at = Some(Instant::now());
release_due().await;
let second = next_release_deadline()
.expect("the second unfinished ask keeps one armed too")
.saturating_duration_since(Instant::now());
assert_eq!(
after_first, 1,
"the Silent answer charged the ladder's first rung"
);
assert!(
first > Duration::ZERO,
"the first ask is spaced by the base, not at the floor: {first:?}"
);
assert!(
second > first + Duration::from_secs(2),
"the second Silent pass took a further rung rather than re-arming at the floor: \
{first:?} then {second:?}"
);
guard.set_mode("");
reclaim_schedule().lock().unwrap_poison().next_at = Some(Instant::now());
release_due().await;
assert_eq!(
reclaim_schedule().lock().unwrap_poison().failures,
0,
"a clean answered ask resets the ladder"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn one_pass_with_several_silent_leftovers_takes_one_rung() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
widen_retry_base(Duration::from_secs(5));
let first = format!("{AGENT_TAB_PREFIX}0123456789ab-orphan-one");
let second = format!("{AGENT_TAB_PREFIX}0123456789ab-orphan-two");
guard.write_state(&[first, second], true);
guard.set_mode("remove-transient");
guard.arm_sweep();
release_due().await;
assert_eq!(
reclaim_schedule().lock().unwrap_poison().failures,
1,
"one pass, one rung, however many names it left unanswered"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_leftover_the_sweep_could_not_settle_keeps_an_ask_armed() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
let orphan = format!("{AGENT_TAB_PREFIX}0123456789ab-orphan");
guard.write_state(std::slice::from_ref(&orphan), false);
guard.arm_sweep();
release_due().await;
let parked = reclaim_schedule()
.lock()
.unwrap_poison()
.next_at
.expect("an ask is armed");
let revisit = next_reclaim_revisit().expect("the parked leftover owes a revisit");
let drift = parked
.saturating_duration_since(revisit)
.max(revisit.saturating_duration_since(parked));
assert!(
drift < Duration::from_secs(1),
"the parked leftover's revisit is what the ask carries: {parked:?} vs {revisit:?}"
);
assert!(
parked > Instant::now(),
"and it is in the future — the sweep does not spin on it"
);
assert!(
next_release_deadline().is_some(),
"the driver wakes for that ask by itself, with no event"
);
guard.clear_log();
release_due().await;
assert!(
guard.log_lines().is_empty(),
"and nothing asks before it is due: {:?}",
guard.log_lines()
);
reclaim_schedule().lock().unwrap_poison().next_at = Some(Instant::now());
guard.clear_log();
release_due().await;
assert!(
guard.invoked("extension call tabGroups.query"),
"the armed ask is what asks the browser again, with no event: {:?}",
guard.log_lines()
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_record_settled_unclosable_arms_the_sweep() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
let names = run_session_names("run-unclosable", &["agent-tab-u-default"]);
guard.write_state(&names, false);
queue_names("run-unclosable", &["agent-tab-u-default"]);
release_due().await;
assert!(
pending_names_and_attempts().is_empty(),
"the record stops retrying the names no route can close"
);
let issues = guard.issues();
assert_eq!(issues.len(), 1, "{issues:?}");
assert_eq!(issues[0].0, LEFTOVER_MESSAGE);
assert_eq!(issues[0].2["run"], "run-unclosable");
let remaining = next_release_deadline()
.expect("the driver wakes for the group the record left behind")
.saturating_duration_since(Instant::now());
assert!(
STUCK_REVISIT.saturating_sub(remaining) < Duration::from_secs(1),
"the sweep is armed for the revisit, not at once: {remaining:?}"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_name_a_record_dropped_is_still_swept() {
let guard = ReleaseGuard::install(ATTEMPT_TIMEOUT).await;
let names = run_session_names("run-mixed", &["agent-tab-m-stuck", "agent-tab-m-keep"]);
guard.write_state_each(&[(&names[0], false), (&names[1], true)]);
guard.set_mode("remove-transient");
queue_names("run-mixed", &["agent-tab-m-stuck", "agent-tab-m-keep"]);
release_due().await;
assert_eq!(
pending_names_and_attempts(),
vec![(vec![names[1].clone()], 1)],
"the record drops the unclosable name and keeps retrying the other"
);
guard.write_state(std::slice::from_ref(&names[0]), true);
guard.set_mode("ok");
guard.clear_log();
reclaim_schedule().lock().unwrap_poison().next_at = Some(Instant::now());
release_due().await;
assert!(
guard.removed_tab_ids().contains(&100),
"the sweep examines the group the record dropped instead of skipping it as \
claimed: {:?}",
guard.log_lines()
);
assert!(
pending_names_and_attempts().is_empty(),
"the retried name's group is gone as well, so the record settles"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_protected_name_is_never_reclaimed() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
guard.arm_sweep();
let live = ChromeRunSessions::for_run("run-live-prot");
let live_name = format!("{}default", live.namespace());
let held = run_session_names("run-held-prot", &["agent-tab-h-prot"]).remove(0);
queue_names_with_hold(
"run-held-prot",
&["agent-tab-h-prot"],
Duration::from_hours(1),
);
guard.write_state(&[live_name, held], true);
release_due().await;
assert!(
!guard.invoked("extension call tabs.remove"),
"a protected group is never closed: {:?}",
guard.log_lines()
);
drop(live);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn an_overflow_eviction_drops_the_oldest_unheld_record() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
clear_pending_releases();
for i in 0..MAX_PENDING_RELEASES {
queue_names(&format!("run-evict-{i}"), &["agent-tab-e-default"]);
}
guard.clear_issues();
queue_names("run-evict-overflow", &["agent-tab-e-overflow"]);
guard.set_binary_absent();
release_due().await;
let queued = pending_names_and_attempts();
assert_eq!(queued.len(), MAX_PENDING_RELEASES, "the cap is kept");
let evicted = ChromeRunSessions::for_run("run-evict-0")
.namespace()
.to_string();
assert!(
!queued
.iter()
.any(|(names, _)| names.iter().any(|name| name.starts_with(&evicted))),
"the oldest unheld record was the one evicted"
);
assert!(
queued.iter().any(|(names, _)| names
.iter()
.any(|name| name.ends_with("agent-tab-e-overflow"))),
"the record the cap made room for is queued"
);
assert_eq!(
evicted_run(&evicted).as_deref(),
Some("run-evict-0"),
"the evicted record's run is remembered for the sweep's report"
);
assert!(
guard.issues().is_empty(),
"an eviction closes nothing and reports nothing: the sweep owns those tabs"
);
clear_pending_releases();
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn an_evicted_records_run_is_named_by_the_sweeps_report() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
clear_pending_releases();
for i in 0..MAX_PENDING_RELEASES {
queue_names(&format!("run-evict-{i}"), &["agent-tab-e-default"]);
}
queue_names("run-evict-overflow", &["agent-tab-e-overflow"]);
let evicted = run_session_names("run-evict-0", &["agent-tab-e-default"]);
guard.write_state(&evicted, true);
guard.set_mode("keep");
release_due().await;
let namespace = ChromeRunSessions::for_run("run-evict-0")
.namespace()
.to_string();
let issues = guard.issues();
assert_eq!(issues.len(), 1, "{issues:?}");
assert_eq!(issues[0].0, LEFT_OPEN_MESSAGE);
assert_eq!(issues[0].1, format!("chrome-tabs-left-open:{namespace}"));
assert_eq!(
issues[0].2["run"], "run-evict-0",
"the eviction's remembered run is what the sweep's report carries"
);
clear_pending_releases();
}
#[test]
#[serial_test::serial(chrome_release)]
fn an_evicted_run_is_remembered_for_the_sweeps_report() {
clear_pending_releases();
assert_eq!(evicted_run("agent-tab-x-"), None, "nothing remembered yet");
remember_evicted_run("agent-tab-x-", "run-x");
remember_evicted_run("agent-tab-y-", "run-y");
remember_evicted_run("agent-tab-x-", "run-x-again");
assert_eq!(evicted_run("agent-tab-x-").as_deref(), Some("run-x-again"));
assert_eq!(evicted_run("agent-tab-y-").as_deref(), Some("run-y"));
assert_eq!(evicted_run("agent-tab-z-"), None);
let stale = Instant::now()
.checked_sub(EVICTED_RUN_TTL)
.expect("the process clock reaches back an hour");
evicted_runs().lock().unwrap_poison()[1].at = stale;
remember_evicted_run("agent-tab-z-", "run-z");
assert_eq!(evicted_run("agent-tab-y-"), None, "the stale entry is gone");
assert_eq!(evicted_run("agent-tab-z-").as_deref(), Some("run-z"));
clear_pending_releases();
}
#[test]
fn eviction_prefers_an_unheld_record_over_a_held_one() {
let record = |held: bool, until: Instant| PendingRunRelease {
namespace: "agent-tab-h-".to_string(),
names: vec!["agent-tab-h-default".to_string()],
attempts: 0,
next_attempt_at: until,
held,
run: String::new(),
left_open: 0,
};
let held = Instant::now() + Duration::from_hours(1);
let mut mixed = VecDeque::from([record(true, held), record(false, Instant::now())]);
let evicted = evict_oldest(&mut mixed).expect("the unheld record is evicted first");
assert!(!evicted.held, "the unheld record is evicted first");
let mut held_only = VecDeque::from([record(true, held), record(true, held)]);
assert!(
evict_oldest(&mut held_only).is_none(),
"with only held records, nothing is evicted"
);
let mut elapsed = VecDeque::from([record(true, Instant::now())]);
assert!(
evict_oldest(&mut elapsed).is_some(),
"a hold that has elapsed is evictable again"
);
}
#[test]
#[serial_test::serial(chrome_release)]
fn the_queue_cap_never_evicts_a_held_record() {
clear_pending_releases();
let held = |namespace: String| PendingRunRelease {
namespace: namespace.clone(),
names: vec![format!("{namespace}default")],
attempts: 0,
next_attempt_at: Instant::now() + Duration::from_hours(1),
held: true,
run: String::new(),
left_open: 0,
};
for i in 0..MAX_PENDING_RELEASES {
merge_record(held(format!("held-{i}")), Eligibility::RunEnd);
}
merge_record(held("held-overflow".to_string()), Eligibility::RunEnd);
assert_eq!(
pending_releases_snapshot().len(),
MAX_PENDING_RELEASES + 1,
"the cap yields to a queue in which every record is held"
);
clear_pending_releases();
}
#[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("silent");
let started = Instant::now();
release_due().await;
assert!(
started.elapsed() < Duration::from_secs(3),
"the 1 s attempt bound must cut the silent 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"
);
guard.set_mode("ok");
widen_attempt_bound(ATTEMPT_TIMEOUT);
release_due().await;
assert!(
pending_releases_snapshot().is_empty(),
"dropped on the next pass, when the browser answers"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_release_with_no_chrome_use_binary_keeps_the_record_queued() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
queue_names("run-absent", &["agent-tab-a-default"]);
guard.set_binary_absent();
release_due().await;
assert_eq!(
pending_names_and_attempts(),
vec![(run_session_names("run-absent", &["agent-tab-a-default"]), 1)],
"nothing the browser could answer, so the name stays queued"
);
assert!(guard.log_lines().is_empty(), "no child was spawned");
}
#[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("refuse");
release_due().await;
assert_eq!(
pending_names_and_attempts(),
vec![(
run_session_names("run-durable", &["agent-tab-d-default"]),
1
)],
"one unanswered pass, 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();
fs::write(
guard.store_path(),
r#"[{"namespace":"","names":["agent-tab-ns-less-default"]},
{"namespace":"agent-tab-nameless-","names":[]}]"#,
)
.expect("write refused record file");
restore_pending_releases();
assert!(
pending_releases_snapshot().is_empty(),
"neither entry is a record the queue would hold"
);
assert!(
reclaim_schedule().lock().unwrap_poison().next_at.is_none(),
"and none of them arms a read: a start with no work opens no page"
);
}
#[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, 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 browser holds no group by those names, so the record is settled"
);
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.only_swept(),
"nothing was asked for the record: {:?}",
guard.log_lines()
);
tokio::time::sleep(grace * 2).await;
release_due().await;
assert!(
pending_releases_snapshot().is_empty(),
"settled once the boot grace elapsed"
);
assert!(guard.invoked("extension call tabGroups.query"));
}
#[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.only_swept(),
"no close while the run it belongs to may still come back: {:?}",
guard.log_lines()
);
tokio::time::sleep(hold * 2).await;
release_due().await;
assert!(
pending_releases_snapshot().is_empty(),
"settled once the hold elapsed"
);
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"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn a_spent_hold_stops_protecting_its_record() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
widen_retry_base(Duration::from_secs(600));
let names = run_session_names("run-spent-hold", &["agent-tab-sh-default"]);
guard.write_state(&names, true);
queue_names_with_hold(
"run-spent-hold",
&["agent-tab-sh-default"],
Duration::from_millis(1),
);
tokio::time::sleep(Duration::from_millis(20)).await;
guard.set_mode("keep");
release_due().await;
assert_eq!(
pending_names_and_attempts(),
vec![(names.clone(), 1)],
"the browser answered and left the tabs, so the record is retried"
);
assert!(
!ProtectedNamespaces::live_and_held().contains(&names[0]),
"a spent hold no longer protects its namespace"
);
assert!(
queue_has_evictable_record(),
"and the record is an ordinary, evictable one again"
);
}
fn queue_has_evictable_record() -> bool {
let mut queue = pending_releases().lock().unwrap_poison().clone();
evict_oldest(&mut queue).is_some()
}
#[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, "run-live");
release_due().await;
assert_eq!(
pending_releases_snapshot().len(),
1,
"a live run's record stays queued and is not an attempt"
);
assert!(
guard.only_swept(),
"no close pass is ever run for a live run's sessions: {:?}",
guard.log_lines()
);
drop(live);
release_due().await;
assert!(
pending_releases_snapshot().is_empty(),
"settled once the run that owns them is gone"
);
assert!(guard.invoked("extension call tabGroups.query"));
}
#[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 names = run_session_names("run-midpass", &["agent-tab-m-default"]);
guard.write_state(&names, true);
queue_names("run-midpass", &["agent-tab-m-default"]);
let pass = ReleasePass::take().expect("the record is due before the run comes back");
let live = ChromeRunSessions::for_run("run-midpass");
let cli = release_cli();
let mut protected = ProtectedNamespaces::snapshot(&durable_resume_namespaces().await);
let issues = attempt_pass(
cli.as_deref(),
Some(pass),
OWN_SESSION,
&ReclaimAsk::unanswered(),
&mut protected,
Instant::now() + Duration::from_secs(10),
)
.await;
assert_eq!(
pending_names_and_attempts(),
vec![(names, 1)],
"a run that came back mid-pass keeps its sessions, retried on the next pass"
);
assert!(
guard.log_lines().is_empty(),
"no child was spawned for a run that is live again"
);
assert!(issues.is_empty());
drop(live);
}
#[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, "run-driver");
let shutdown = CancellationToken::new();
let driver = tokio::spawn(release_queue(shutdown.clone()));
tokio::time::sleep(Duration::from_millis(300)).await;
assert!(
guard.only_swept(),
"the driver leaves a live run's session alone: {:?}",
guard.log_lines()
);
assert!(
guard.invoked(&format!("--session {}", names[0])),
"and runs its own read in the run's session, not the product's own: {:?}",
guard.log_lines()
);
assert!(
!guard.invoked(&format!("--session {OWN_SESSION}")),
"so a run end opens nothing of ours at all: {:?}",
guard.log_lines()
);
assert_eq!(pending_releases_snapshot().len(), 1);
drop(live);
tokio::time::timeout(Duration::from_secs(5), async {
while !record_file_names(&guard.store_path()).is_empty() {
tokio::time::sleep(Duration::from_millis(20)).await;
}
})
.await
.expect("the driver settles the record once the run that owned it is gone");
assert!(
pending_releases_snapshot().is_empty(),
"and the queue holds nothing either"
);
assert!(
guard.invoked("extension call tabGroups.query"),
"the driver asked the browser for the record's groups"
);
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, "run-merge");
let second = ChromeRunSessions::for_run("run-merge");
second.track(&names[0]);
second.track(&names[1]);
queue_run_session_release(&second, false, "run-merge");
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("silent");
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 unsettled 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 settle and leaves the held record for its resumed run"
);
assert!(
guard.invoked("extension call tabGroups.query"),
"one browser-driven close pass, and never one for the held record"
);
}
#[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(&HashMap::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,
next_attempt_at: 1,
held: true,
run: String::new(),
left_open: 0,
}],
);
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]
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn an_over_cap_eviction_closes_the_named_tab_and_leaves_others_alone() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
let ours = ChromeRunSessions::for_run("evict-ours");
let ours_session = format!("{}victim", ours.namespace());
let other = ChromeRunSessions::for_run("evict-other");
let other_session = format!("{}bystander", other.namespace());
guard.write_state(&[ours_session.clone(), other_session.clone()], true);
let gone = evict_agent_tab(ours.namespace(), &ours_session, &other_session).await;
assert!(gone, "the browser confirmed the victim's group gone");
assert!(
!guard.group_titles().contains(&ours_session),
"the victim's group is gone: {:?}",
guard.group_titles()
);
assert!(
guard.group_titles().contains(&other_session),
"another live run's group — the close's own read session — is untouched: {:?}",
guard.group_titles()
);
assert!(
guard.invoked("session stop --force"),
"the settled session's chrome-use record was retired: {:?}",
guard.log_lines()
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn an_eviction_the_browser_leaves_open_is_not_confirmed() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
guard.set_mode("keep");
let ours = ChromeRunSessions::for_run("evict-kept");
let ours_session = format!("{}victim", ours.namespace());
guard.write_state(std::slice::from_ref(&ours_session), true);
let gone = evict_agent_tab(ours.namespace(), &ours_session, &ours_session).await;
assert!(
!gone,
"the browser answered and left the group standing: {:?}",
guard.log_lines()
);
assert!(
guard.group_titles().contains(&ours_session),
"the group is still in the browser"
);
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn an_eviction_whose_session_stop_fails_is_retried() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
guard.set_mode("stop-refused");
let ours = ChromeRunSessions::for_run("evict-stop-refused");
let ours_session = format!("{}victim", ours.namespace());
guard.write_state(std::slice::from_ref(&ours_session), true);
let gone = evict_agent_tab(ours.namespace(), &ours_session, &ours_session).await;
assert!(gone, "the group itself was closed: {:?}", guard.log_lines());
let stops = guard
.log_lines()
.iter()
.filter(|line| line.starts_with("session stop"))
.count();
assert_eq!(stops, 2, "the stop that did not settle was asked again");
}
#[tokio::test]
#[serial_test::serial(chrome_release, chrome_tabs)]
async fn the_tools_sixth_appearance_closes_the_least_recently_addressed_tab() {
let guard = ReleaseGuard::install(Duration::from_secs(10)).await;
let _ledger = chrome_tab_ledger::TestLedgerGuard::install();
let sessions = ChromeRunSessions::for_run("cap-tool-wiring");
let tool = crate::tools::chrome::ChromeTool::new(Arc::clone(&sessions));
let namespace = sessions.namespace().to_string();
let victim_session = format!("{namespace}a");
guard.write_state(std::slice::from_ref(&victim_session), true);
for tab in ["a", "b", "c", "d", "e"] {
tool.record_addressed(tab, None);
}
assert!(
!chrome_tab_ledger::is_closed(&namespace, "a"),
"the fifth tab is inside the cap"
);
tool.record_addressed("f", None);
let deadline = Instant::now() + Duration::from_secs(5);
while !chrome_tab_ledger::is_closed(&namespace, "a") && Instant::now() < deadline {
tokio::time::sleep(Duration::from_millis(10)).await;
}
assert!(
chrome_tab_ledger::is_closed(&namespace, "a"),
"the run's least recently addressed tab is the one the cap closed"
);
assert!(
!guard.group_titles().contains(&victim_session),
"the confirmed close took the tab out of the browser: {:?}",
guard.group_titles()
);
}
#[tokio::test]
#[serial_test::serial(chrome_release, chrome_tabs)]
async fn a_namespace_no_run_claims_is_dropped_at_boot() {
crate::util::test::init_test_stores().await;
let ledger = chrome_tab_ledger::TestLedgerGuard::install();
let namespace = ChromeRunSessions::for_run("prune-orphan")
.namespace()
.to_string();
assert!(
chrome_tab_ledger::note_addressed(&namespace, "tab", true).is_none(),
"one tab is inside the cap"
);
prune_tab_records().await;
assert!(
!ledger.holds(&namespace),
"no run claims the namespace: its record is dropped"
);
let live = ChromeRunSessions::for_run("prune-orphan");
assert!(
chrome_tab_ledger::note_addressed(&namespace, "tab", true).is_none(),
"one tab is inside the cap"
);
prune_tab_records().await;
assert!(
ledger.holds(&namespace),
"a live run's record survives the boot prune"
);
drop(live);
}
}