feat(browser): deterministic private_browser lifecycle (Wave 5A)

Own each agent-browser session as a browser tree: the daemon's POSIX
session, its runtime files and its Chrome profile. Timeouts, launch
failures, bootstrap recovery, cancellation and shutdown clean that tree
and verify nothing survives, instead of killing only the daemon and
orphaning Chrome. Per-call cleanup no longer sweeps every Chrome under
the runtime TMPDIR.

Sessionless calls get an ephemeral browser closed before returning.
Actions on one session are serialized. Recovery is bounded by one
deadline with at most one retry for local HTML open, and the retry flag
is no longer model-visible. Observations after a failed navigation are
marked stale. read URL navigates and extracts in one batch because
agent-browser has no read command. Results carry a browser_lifecycle
receipt with stages, timings, ownership and cleanup evidence.

Playwright MCP tool calls are bounded by
ODYSSEUS_BROWSER_MCP_CALL_TIMEOUT_S and are not retried. research_navigator
now passes timeout_ms.
This commit is contained in:
Alexandre Teixeira
2026-10-01 20:59:27 +01:00
parent cb5b81022b
commit 576abb012d
9 changed files with 1552 additions and 78 deletions
+378 -70
View File
@@ -1,6 +1,7 @@
import asyncio
import base64
import contextlib
import contextvars
import inspect
import io
import json
@@ -10,15 +11,18 @@ import signal
import shutil
import sys
import tempfile
import time
import html
import hashlib
import pwd
import urllib.parse
import urllib.request
import uuid
from pathlib import Path
from typing import Dict, Any
from core import platform_compat
from src import browser_lifecycle
from src.constants import MAX_OUTPUT_CHARS
PDF_EXTRACT_MAX_BYTES = 80_000_000
@@ -90,6 +94,28 @@ def _scoped_browser_session(namespace: str, session_id: str) -> str:
return f"ody-{_bounded_browser_identity(scope)}"
def _browser_namespace(env: dict[str, str] | None) -> str:
return str(
(env or {}).get("ODYSSEUS_BROWSER_NAMESPACE")
or os.getenv("ODYSSEUS_BROWSER_NAMESPACE", "odysseus-ui")
).strip() or "odysseus-ui"
# Browser CLI processes started by the current private_browser call, so a
# cancellation can stop every client it spawned, not only the main command.
_BROWSER_CALL_PROCS: contextvars.ContextVar[list | None] = contextvars.ContextVar(
"_BROWSER_CALL_PROCS", default=None
)
async def _spawn_browser_cli(*command, **kwargs):
proc = await asyncio.create_subprocess_exec(*command, **kwargs)
tracked = _BROWSER_CALL_PROCS.get()
if tracked is not None:
tracked.append(proc)
return proc
def _browser_pid_file_candidates(
runtime_dir: Path, namespace: str, session_id: str | None
) -> list[Path]:
@@ -2433,15 +2459,30 @@ class PrivateBrowserTool:
@staticmethod
def _terminate_owned_daemon(
env: dict[str, str], session_id: str | None = None
) -> None:
"""Terminate detached agent-browser daemon(s) for this runtime."""
) -> dict[str, Any] | None:
"""Terminate the detached agent-browser tree owned by one session.
namespace = str(
env.get("ODYSSEUS_BROWSER_NAMESPACE")
or os.getenv("ODYSSEUS_BROWSER_NAMESPACE", "odysseus-ui")
).strip() or "odysseus-ui"
Killing only the daemon reparents its Chrome children, so the whole
POSIX session the daemon leads is cleaned together with the session's
runtime files and browser profile. Returns the cleanup receipt.
"""
namespace = _browser_namespace(env)
receipt = None
pid_files = []
if session_id:
key = _scoped_browser_session(namespace, session_id)
root = browser_lifecycle.runtime_root(env)
receipt = browser_lifecycle.force_cleanup(
root, key, pid_alive=lambda pid: _process_is_alive(pid)
).as_dict()
pid_files.append(root / f"{key}.pid")
runtime_dir = Path(os.getenv("XDG_RUNTIME_DIR") or f"/run/user/{os.getuid()}")
pid_files = _browser_pid_file_candidates(runtime_dir, namespace, session_id)
pid_files = [
path
for path in _browser_pid_file_candidates(runtime_dir, namespace, session_id)
if path not in pid_files
]
for pid_file in pid_files:
try:
pid = int(pid_file.read_text().strip())
@@ -2463,6 +2504,7 @@ class PrivateBrowserTool:
os.kill(pid, signal.SIGKILL)
with contextlib.suppress(FileNotFoundError, PermissionError, OSError):
pid_file.unlink()
return receipt
@staticmethod
def _owned_daemon_exists(env: dict[str, str], session_id: str | None) -> bool:
@@ -2475,12 +2517,17 @@ class PrivateBrowserTool:
if not session_id:
return False
namespace = str(
env.get("ODYSSEUS_BROWSER_NAMESPACE")
or os.getenv("ODYSSEUS_BROWSER_NAMESPACE", "odysseus-ui")
).strip() or "odysseus-ui"
namespace = _browser_namespace(env)
key = _scoped_browser_session(namespace, session_id)
root = browser_lifecycle.runtime_root(env)
if browser_lifecycle.has_live_daemon(
root, key, pid_alive=lambda pid: _process_is_alive(pid)
):
return True
runtime_dir = Path(os.getenv("XDG_RUNTIME_DIR") or f"/run/user/{os.getuid()}")
for pid_file in _browser_pid_file_candidates(runtime_dir, namespace, session_id):
if pid_file == root / f"{key}.pid":
continue
try:
pid = int(pid_file.read_text().strip())
except (OSError, ValueError):
@@ -2553,10 +2600,184 @@ class PrivateBrowserTool:
resolved = cls._resolve_workspace_path(raw_path)
return resolved.as_uri()
# Time allowed beyond the action timeout for one bounded recovery attempt.
_RECOVERY_BUDGET_S = 75
_CLOSE_TIMEOUT_S = 10
async def execute(self, content: str, ctx: dict) -> dict:
"""Run one browser action inside its session's lifecycle.
Actions on one session are serialized. A call without an owning
Odysseus session gets a browser of its own that is closed before the
call returns; it never falls back to agent-browser's shared default
session. Cancellation stops every CLI client the call started and
cleans the session's browser tree, because its state is unknown.
"""
ctx = dict(ctx) if isinstance(ctx, dict) else {}
session_id = str(ctx.get("session_id") or "").strip()
ephemeral = not session_id
if ephemeral:
session_id = f"ephemeral-{uuid.uuid4().hex}"
ctx["session_id"] = session_id
runtime_env = ctx.get("subproc_env") if isinstance(ctx.get("subproc_env"), dict) else {}
key = _scoped_browser_session(_browser_namespace(runtime_env), session_id)
browser = browser_lifecycle.session_for(key, ephemeral)
clock = browser_lifecycle.StageClock()
procs: list = []
token = _BROWSER_CALL_PROCS.set(procs)
lock = browser.lock()
acquired = False
try:
await lock.acquire()
acquired = True
result = await self._execute_unlocked(content, ctx, browser=browser, clock=clock)
if ephemeral:
await self._release_session(browser, session_id, clock)
if isinstance(result, dict) and browser.env is not None:
result["browser_lifecycle"] = browser.receipt(clock)
return result
except asyncio.CancelledError:
if acquired:
for proc in procs:
if getattr(proc, "returncode", None) is None:
self._terminate_subprocess(proc)
if browser.env is not None:
self._terminate_owned_daemon(browser.env, session_id)
browser.discarded("cancelled")
raise
except Exception:
if acquired and ephemeral and browser.env is not None:
self._terminate_owned_daemon(browser.env, session_id)
raise
finally:
if acquired:
lock.release()
_BROWSER_CALL_PROCS.reset(token)
if ephemeral:
browser_lifecycle.forget(key)
_ACTIVE_BROWSER_SESSIONS.discard(key)
async def _release_session(
self,
browser: browser_lifecycle.BrowserSession,
session_id: str,
clock: browser_lifecycle.StageClock,
) -> None:
"""Close a session gracefully, then verify nothing it owned survives."""
if browser.env is None:
return
started = time.monotonic()
graceful = False
if self._owned_daemon_exists(browser.env, session_id):
proc = None
try:
proc = await _spawn_browser_cli(
*browser.command_prefix,
"close",
stdout=asyncio.subprocess.DEVNULL,
stderr=asyncio.subprocess.DEVNULL,
env=browser.env,
start_new_session=True,
)
await asyncio.wait_for(proc.wait(), timeout=self._CLOSE_TIMEOUT_S)
graceful = (proc.returncode or 0) == 0
except Exception:
if proc is not None:
self._terminate_subprocess(proc)
receipt = self._terminate_owned_daemon(browser.env, session_id)
verified = bool(receipt.get("verified")) if isinstance(receipt, dict) else False
clock.record("close", started, verified or graceful)
clock.extra["cleanup"] = {"graceful_close": graceful, **(receipt or {})}
if browser.page_url:
clock.extra["closed_page_url"] = browser.page_url
browser.discarded("closed")
def _discard_session(
self,
browser: browser_lifecycle.BrowserSession,
env: dict[str, str],
session_id: str,
clock: browser_lifecycle.StageClock,
state: str,
) -> None:
"""Force-clean a session whose browser state can no longer be trusted."""
started = time.monotonic()
receipt = self._terminate_owned_daemon(env, session_id or None)
verified = bool(receipt.get("verified")) if isinstance(receipt, dict) else False
clock.record("forced_cleanup", started, verified, reason=state)
clock.extra["cleanup"] = receipt
browser.discarded(state)
@staticmethod
def _navigation_target(action: str, args: dict) -> str:
"""URL this action navigates the session to, or ``""``."""
if action == "read" and any(
str(args.get(key) or "").strip() for key in ("selector", "target", "ref")
):
return ""
if action in {"open", "read"}:
return str(args.get("url") or "").strip()
if action == "batch" and isinstance(args.get("commands"), list):
target = ""
for command in args["commands"]:
if (
isinstance(command, list)
and len(command) > 1
and str(command[0]).lower() in {"open", "goto", "navigate"}
):
target = str(command[1]).strip()
return target
return ""
@staticmethod
def _read_page_from_rows(output: str) -> dict[str, Any]:
"""Page text from an open + ``get text`` batch, only if both succeeded."""
try:
rows = json.loads(output)
except (ValueError, TypeError):
rows = None
if not isinstance(rows, list) or len(rows) != 2 or not all(isinstance(r, dict) for r in rows):
return {"ok": False, "error": "private_browser read returned no structured page result"}
opened, extracted = rows
for row in rows:
if row.get("success") is not True:
return {"ok": False, "error": f"private_browser read failed: {row.get('error') or 'unknown error'}"}
opened_result = opened.get("result") if isinstance(opened.get("result"), dict) else {}
extracted_result = extracted.get("result") if isinstance(extracted.get("result"), dict) else {}
text = extracted_result.get("text")
if not isinstance(text, str):
return {"ok": False, "error": "private_browser read observed no page text"}
url = str(opened_result.get("url") or extracted_result.get("origin") or "")
title = str(opened_result.get("title") or "")
header = "\n".join(part for part in (title, url) if part)
return {"ok": True, "url": url, "text": f"{header}\n\n{text}".strip()}
@staticmethod
def _navigated_url(output: str) -> str:
"""Final URL reported by ``open`` (after redirects), when present."""
match = re.search(r"^\s+([a-z][a-z0-9+.-]*:\S+)\s*$", str(output or ""), re.MULTILINE)
return match.group(1) if match else ""
async def _execute_unlocked(
self,
content: str,
ctx: dict,
*,
browser: browser_lifecycle.BrowserSession,
clock: browser_lifecycle.StageClock,
retry: bool = False,
deadline: float | None = None,
) -> dict:
args, err = self._parse_args(content)
if err:
return {"error": err, "exit_code": 1}
args.pop("_odysseus_browser_retry", None)
action = str(args.get("action") or "").strip().lower()
if action not in self._ACTIONS:
@@ -2697,8 +2918,23 @@ class PrivateBrowserTool:
# only errors emitted by this page. Opening a URL replaces prior page
# state anyway; cookies are irrelevant for confined file:// artifacts.
session_id = str((ctx or {}).get("session_id") or "").strip()
if verifies_local_html and self._owned_daemon_exists(env, session_id):
browser.bind(env, cmd_prefix)
loop = asyncio.get_running_loop()
if deadline is None:
deadline = loop.time() + timeout_s + self._RECOVERY_BUDGET_S
warm = self._owned_daemon_exists(env, session_id)
if verifies_local_html and warm:
reset_started = time.monotonic()
await self._reset_browser_session(cmd_prefix, env, timeout_s)
clock.record("reset", reset_started, True)
browser.discarded("reset")
navigation_url = self._navigation_target(action, command_args)
stale_note = (
browser.stale_observation_note()
if action in browser_lifecycle.OBSERVATION_ACTIONS and not navigation_url
else ""
)
command_started = time.monotonic()
# agent-browser starts a persistent daemon which can inherit the
# client's stdout/stderr descriptors. Pipes therefore never reach
@@ -2708,8 +2944,10 @@ class PrivateBrowserTool:
# depend on the client process, not its detached daemon.
stdout_file = tempfile.TemporaryFile()
stderr_file = tempfile.TemporaryFile()
attempt_timeout = max(1.0, min(float(timeout_s), deadline - loop.time()))
proc = None
try:
proc = await asyncio.create_subprocess_exec(
proc = await _spawn_browser_cli(
*command,
stdin=asyncio.subprocess.PIPE if stdin_data is not None else None,
stdout=stdout_file,
@@ -2719,40 +2957,47 @@ class PrivateBrowserTool:
)
await asyncio.wait_for(
proc.communicate(stdin_data.encode("utf-8") if stdin_data is not None else None),
timeout=timeout_s,
timeout=attempt_timeout,
)
stdout_file.seek(0)
stderr_file.seek(0)
stdout = stdout_file.read()
stderr = stderr_file.read()
except asyncio.TimeoutError:
with contextlib.suppress(Exception):
self._terminate_subprocess(proc)
self._terminate_owned_chrome(env)
self._terminate_owned_daemon(
env, str((ctx or {}).get("session_id") or "").strip() or None
)
if proc is not None:
with contextlib.suppress(Exception):
self._terminate_subprocess(proc)
clock.record(action, command_started, False, cold_start=not warm, failure="timeout")
self._discard_session(browser, env, session_id, clock, "timed_out")
# A failed local-page verification can leave agent-browser's
# persistent session between a page-error response and the next
# repair attempt. Reopen exactly once after resetting that session;
# never retry mutating browser actions or arbitrary URLs.
if (
verifies_local_html
and action == "open"
and not args.get("_odysseus_browser_retry")
):
# repair attempt. The session was cleaned above, so reopen exactly
# once within the call's deadline; never retry mutating browser
# actions or arbitrary URLs.
remaining = deadline - loop.time()
if verifies_local_html and action == "open" and not retry and remaining >= 10:
retry_args = dict(args)
retry_args["_odysseus_browser_retry"] = True
retry_args["timeout_ms"] = max(60_000, timeout_s * 1000)
if self._owned_daemon_exists(env, session_id):
await self._reset_browser_session(cmd_prefix, env, timeout_s)
return await self.execute(json.dumps(retry_args), ctx)
return {"error": f"private_browser timed out after {timeout_s}s", "exit_code": 1}
retry_args["timeout_ms"] = int(
min(max(60.0, float(timeout_s)), remaining - 5) * 1000
)
clock.extra["recovery_attempts"] = 1
return await self._execute_unlocked(
json.dumps(retry_args), ctx,
browser=browser, clock=clock, retry=True, deadline=deadline,
)
return {
"error": f"private_browser timed out after {int(attempt_timeout)}s",
"exit_code": 1,
}
except Exception as e:
self._terminate_owned_chrome(env)
self._terminate_owned_daemon(
env, str((ctx or {}).get("session_id") or "").strip() or None
)
if proc is not None:
with contextlib.suppress(Exception):
self._terminate_subprocess(proc)
clock.record(action, command_started, False, cold_start=not warm, failure=type(e).__name__)
if proc is not None:
# The client reached the daemon, so the session's state is
# unknown. A client that never started left it untouched.
self._discard_session(browser, env, session_id, clock, "failed")
return {"error": f"private_browser failed: {type(e).__name__}: {e}", "exit_code": 1}
finally:
stdout_file.close()
@@ -2763,6 +3008,47 @@ class PrivateBrowserTool:
combined = out
if err_text:
combined = f"{combined}\n\n[stderr]\n{err_text}".strip()
command_ok = (proc.returncode or 0) == 0
read_page = None
if action == "read" and navigation_url:
read_page = self._read_page_from_rows(out)
command_ok = command_ok and read_page.get("ok", False)
clock.record(action, command_started, command_ok, cold_start=not warm)
if not command_ok and browser_lifecycle.LAUNCH_FAILURE_RE.search(combined):
# The browser never became ready. The daemon outlives this failure
# and a later close cannot reach a browser, so clean it here.
self._discard_session(browser, env, session_id, clock, "launch_failed")
return {
"output": combined[:4000],
"error": (
"private_browser could not launch the browser; no page was "
"opened or observed. The browser session was cleaned up."
),
"exit_code": 1,
"untrusted_content": True,
}
if navigation_url:
if command_ok:
browser.navigated(
(read_page or {}).get("url") or self._navigated_url(out) or navigation_url
)
else:
browser.navigation_failed(navigation_url)
if read_page is not None:
if not command_ok:
return {
"output": combined[:4000],
"error": read_page.get("error") or "private_browser read failed; no page text was observed",
"exit_code": 1,
"untrusted_content": True,
}
out = read_page["text"]
combined = out if not err_text else f"{out}\n\n[stderr]\n{err_text}"
elif command_ok and browser.state in {"idle", "closed", "reset"}:
browser.state = "ready"
if stale_note and command_ok:
combined = f"[{stale_note}]\n\n{combined}".strip()
clock.extra["stale_observation"] = True
from src.turn_contract import active_turn_contract
contract = active_turn_contract()
model_choice = getattr(contract, 'routing_experiment', '') == 'recent_model_choice'
@@ -2806,13 +3092,15 @@ class PrivateBrowserTool:
and action == "open"
and (proc.returncode or 0) != 0
and self._retryable_local_open_failure(combined)
and not args.get("_odysseus_browser_retry")
and not retry
and deadline - loop.time() >= 10
):
self._terminate_owned_chrome(env)
self._terminate_owned_daemon(env, session_id or None)
retry_args = dict(args)
retry_args["_odysseus_browser_retry"] = True
return await self.execute(json.dumps(retry_args), ctx)
self._discard_session(browser, env, session_id, clock, "bootstrap_failed")
clock.extra["recovery_attempts"] = 1
return await self._execute_unlocked(
json.dumps(args), ctx,
browser=browser, clock=clock, retry=True, deadline=deadline,
)
page_errors = ""
if (
verifies_local_html
@@ -2893,7 +3181,7 @@ class PrivateBrowserTool:
proc = None
try:
async with asyncio.timeout(max(0, deadline - loop.time())):
proc = await asyncio.create_subprocess_exec(
proc = await _spawn_browser_cli(
*cmd_prefix, "batch", "--json",
stdin=asyncio.subprocess.PIPE,
stdout=asyncio.subprocess.PIPE,
@@ -3044,7 +3332,7 @@ class PrivateBrowserTool:
"""Best-effort reset of state retained by a persistent browser session."""
proc = None
try:
proc = await asyncio.create_subprocess_exec(
proc = await _spawn_browser_cli(
*cmd_prefix,
"close",
stdout=asyncio.subprocess.PIPE,
@@ -3070,7 +3358,7 @@ class PrivateBrowserTool:
"""Return bounded JavaScript errors from the current browser page."""
proc = None
try:
proc = await asyncio.create_subprocess_exec(
proc = await _spawn_browser_cli(
*cmd_prefix,
"errors",
stdout=asyncio.subprocess.PIPE,
@@ -3114,7 +3402,7 @@ class PrivateBrowserTool:
return None
proc = None
try:
proc = await asyncio.create_subprocess_exec(
proc = await _spawn_browser_cli(
*command,
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.PIPE,
@@ -3331,10 +3619,13 @@ class PrivateBrowserTool:
if target:
return [*prefix, "get", "text", target], None, None
url = str(args.get("url") or "").strip()
command = [*prefix, "read"]
# agent-browser has no `read` command. Navigate and extract in one
# client call so the text is observed after this navigation.
if url:
command.append(url)
return command, None, None
return [*prefix, "batch", "--json"], json.dumps(
[["open", url], ["get", "text", "body"]]
), None
return [*prefix, "get", "text", "body"], None, None
if action == "snapshot":
return [*prefix, "snapshot"], None, None
if action == "find":
@@ -3465,7 +3756,13 @@ class PrivateBrowserTool:
async def shutdown_private_browser_sessions() -> None:
"""Close browser sessions owned by this Odysseus runtime namespace."""
"""Close and verify every browser session this runtime started.
Each session is closed with the environment it was launched with, so its
runtime directory resolves to the daemon's own. ``close`` is sent only to
a verified live daemon, because against a missing one it bootstraps a new
browser. Forced cleanup of the session's own browser tree always follows.
"""
binary = shutil.which("agent-browser") or PrivateBrowserTool._local_agent_browser_binary()
command_prefix = [binary] if binary else (
@@ -3474,26 +3771,37 @@ async def shutdown_private_browser_sessions() -> None:
sessions = sorted(_ACTIVE_BROWSER_SESSIONS)
if not command_prefix or not sessions:
return
env = dict(os.environ)
env["HOME"] = str(_service_home())
env.setdefault("AGENT_BROWSER_IDLE_TIMEOUT_MS", "300000")
base_env = dict(os.environ)
base_env["HOME"] = str(_service_home())
base_env.setdefault("AGENT_BROWSER_IDLE_TIMEOUT_MS", "300000")
try:
for session in sessions:
proc = None
try:
proc = await asyncio.create_subprocess_exec(
*command_prefix, "--session", session, "close",
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.PIPE,
env=env,
start_new_session=True,
)
await asyncio.wait_for(proc.communicate(), timeout=20)
except Exception:
if proc is not None:
with contextlib.suppress(Exception):
PrivateBrowserTool._terminate_subprocess(proc)
record = browser_lifecycle.registered(session)
env = record.env if record is not None and record.env is not None else base_env
root = browser_lifecycle.runtime_root(env)
if browser_lifecycle.has_live_daemon(
root, session, pid_alive=lambda pid: _process_is_alive(pid)
):
proc = None
try:
proc = await asyncio.create_subprocess_exec(
*command_prefix, "--session", session, "close",
stdout=asyncio.subprocess.DEVNULL,
stderr=asyncio.subprocess.DEVNULL,
env=env,
start_new_session=True,
)
await asyncio.wait_for(proc.communicate(), timeout=20)
except Exception:
if proc is not None:
with contextlib.suppress(Exception):
PrivateBrowserTool._terminate_subprocess(proc)
browser_lifecycle.force_cleanup(
root, session, method="shutdown",
pid_alive=lambda pid: _process_is_alive(pid),
)
browser_lifecycle.forget(session)
finally:
_ACTIVE_BROWSER_SESSIONS.difference_update(sessions)
PrivateBrowserTool._terminate_owned_chrome(env)
PrivateBrowserTool._terminate_owned_daemon(env)
PrivateBrowserTool._terminate_owned_chrome(base_env)
PrivateBrowserTool._terminate_owned_daemon(base_env)
+397
View File
@@ -0,0 +1,397 @@
"""Lifecycle ownership for private_browser's agent-browser sessions.
agent-browser runs a short-lived CLI client against a detached daemon. The
daemon calls ``setsid`` and every Chrome process it launches stays in that
POSIX session, so the daemon pid recorded in the session's own pid file
identifies the complete browser tree. Cleanup here is limited to that tree,
the session's runtime files and its ``agent-browser-chrome-*`` profile.
This is browser-specific ownership only. Generic process containment and
lifecycle primitives belong to the shared process layer; when those exist,
``kill_browser_tree`` is the single seam to replace.
"""
from __future__ import annotations
import asyncio
import os
import re
import shutil
import signal
import time
from dataclasses import dataclass, field
from pathlib import Path
from typing import Any, Callable
from core import platform_compat
PROFILE_PREFIX = "agent-browser-chrome-"
RUNTIME_SUFFIXES = (".pid", ".sock", ".stream", ".version", ".engine")
# Output that proves the browser never became ready. The daemon survives such
# a failure and a later ``close`` cannot reach a browser to shut down.
LAUNCH_FAILURE_RE = re.compile(
r"Chrome exited early|DevToolsActivePort|No usable sandbox|"
r"Failed to launch (?:the )?browser|Browser (?:process )?exited before",
re.IGNORECASE,
)
OBSERVATION_ACTIONS = frozenset(
{"snapshot", "read", "find", "evaluate", "screenshot", "scroll", "wait"}
)
def runtime_root(env: dict[str, str] | None) -> Path:
"""Directory where agent-browser keeps ``<session>.pid`` and its socket.
Mirrors agent-browser's own resolution for the environment the daemon is
launched with: an explicit socket directory, then the XDG runtime
directory, then ``$HOME/.agent-browser``.
"""
source = env or {}
def _get(name: str) -> str:
return str(source.get(name) or os.environ.get(name) or "").strip()
socket_dir = _get("AGENT_BROWSER_SOCKET_DIR")
if socket_dir:
return Path(socket_dir)
xdg = _get("XDG_RUNTIME_DIR")
if xdg:
return Path(xdg) / "agent-browser"
home = _get("HOME") or str(Path.home())
return Path(home) / ".agent-browser"
def _read_cmdline(pid: int) -> str | None:
try:
return (platform_compat.PROC_ROOT / str(pid) / "cmdline").read_bytes().replace(
b"\0", b" "
).decode("utf-8", errors="replace")
except (OSError, UnicodeError):
return None
def _read_stat(pid: int) -> tuple[str, int, int] | None:
"""Return ``(state, pgid, sid)`` for a pid, or ``None`` when unreadable."""
try:
raw = (platform_compat.PROC_ROOT / str(pid) / "stat").read_text()
except (OSError, UnicodeError):
return None
_, _, rest = raw.rpartition(")")
fields = rest.split()
if len(fields) < 4:
return None
try:
return fields[0], int(fields[2]), int(fields[3])
except ValueError:
return None
def _live_pids() -> list[int]:
if not platform_compat.has_procfs():
return []
pids = []
for entry in platform_compat.PROC_ROOT.iterdir():
if entry.name.isdigit():
pids.append(int(entry.name))
return pids
def daemon_pid(root: Path, key: str) -> int | None:
"""Pid recorded in this session's pid file, whether or not it is alive."""
try:
return int((root / f"{key}.pid").read_text().strip())
except (OSError, ValueError):
return None
def is_verified_daemon(pid: int | None) -> bool:
"""Whether ``pid`` is a live agent-browser process (requires procfs)."""
if not pid:
return False
stat = _read_stat(pid)
if stat is not None and stat[0] == "Z":
return False
command_line = _read_cmdline(pid)
return bool(command_line) and "agent-browser" in command_line
def browser_tree(leader: int) -> list[int]:
"""Processes owned by the browser session whose daemon pid is ``leader``.
While the daemon is verified alive, every member of its POSIX session is
owned. Once the daemon is gone the pid may be reused, so only Chrome
process groups whose root carries an agent-browser profile are claimed.
"""
members: list[tuple[int, int]] = []
for pid in _live_pids():
stat = _read_stat(pid)
if stat is None or stat[0] == "Z" or stat[2] != leader:
continue
members.append((pid, stat[1]))
if not members:
return []
if is_verified_daemon(leader):
return sorted(pid for pid, _ in members)
owned_groups = {
pgid for pid, pgid in members if _profile_dirs([pid])
}
return sorted(pid for pid, pgid in members if pgid in owned_groups)
def _profile_dirs(pids: list[int]) -> set[Path]:
profiles: set[Path] = set()
for pid in pids:
for token in (_read_cmdline(pid) or "").split():
if not token.startswith("--user-data-dir="):
continue
path = Path(token.split("=", 1)[1])
if path.name.startswith(PROFILE_PREFIX):
profiles.add(path)
return profiles
def kill_browser_tree(leader: int, *, settle_s: float = 1.0) -> tuple[list[int], list[int], set[Path]]:
"""SIGKILL one browser session tree and wait briefly for it to exit.
Returns ``(killed, survivors, profile_dirs)``. Synchronous so it can run
from cancellation and shutdown paths without awaiting.
"""
members = browser_tree(leader)
profiles = _profile_dirs(members)
ordered = [pid for pid in members if pid != leader]
if leader in members:
ordered.append(leader)
killed = []
for pid in ordered:
try:
os.kill(pid, signal.SIGKILL)
killed.append(pid)
except (ProcessLookupError, PermissionError, OSError):
continue
deadline = time.monotonic() + settle_s
survivors = list(killed)
while survivors and time.monotonic() < deadline:
survivors = [pid for pid in survivors if _alive(pid)]
if survivors:
time.sleep(0.02)
return killed, survivors, profiles
def _alive(pid: int) -> bool:
stat = _read_stat(pid)
if stat is None:
return False
return stat[0] != "Z"
@dataclass
class CleanupReceipt:
method: str
daemon_pid: int | None = None
killed: int = 0
survivors: list[int] = field(default_factory=list)
removed_files: list[str] = field(default_factory=list)
removed_profiles: int = 0
verified: bool = False
note: str = ""
def as_dict(self) -> dict[str, Any]:
return {
"method": self.method,
"daemon_pid": self.daemon_pid,
"killed": self.killed,
"survivors": list(self.survivors),
"removed_files": list(self.removed_files),
"removed_profiles": self.removed_profiles,
"verified": self.verified,
**({"note": self.note} if self.note else {}),
}
def force_cleanup(
root: Path,
key: str,
*,
method: str = "forced",
pid_alive: Callable[[int], bool] = platform_compat.pid_alive,
) -> CleanupReceipt:
"""Kill this session's browser tree and remove its owned resources.
Without procfs nothing can be attributed safely, so live processes are
left alone and only the pid file of a dead daemon is forgotten.
"""
pid = daemon_pid(root, key)
receipt = CleanupReceipt(method=method, daemon_pid=pid)
if not platform_compat.has_procfs():
if pid and not pid_alive(pid):
_remove_runtime_files(root, key, receipt)
receipt.verified = True
else:
receipt.note = "procfs unavailable; browser ownership could not be verified"
return receipt
profiles: set[Path] = set()
if pid:
killed, survivors, profiles = kill_browser_tree(pid)
receipt.killed = len(killed)
receipt.survivors = survivors
if not receipt.survivors:
_remove_runtime_files(root, key, receipt)
for profile in profiles:
if profile.name.startswith(PROFILE_PREFIX) and profile.is_dir():
shutil.rmtree(profile, ignore_errors=True)
if not profile.exists():
receipt.removed_profiles += 1
receipt.verified = not receipt.survivors and not (pid and browser_tree(pid))
return receipt
def _remove_runtime_files(root: Path, key: str, receipt: CleanupReceipt) -> None:
for suffix in RUNTIME_SUFFIXES:
path = root / f"{key}{suffix}"
try:
path.unlink()
receipt.removed_files.append(path.name)
except FileNotFoundError:
continue
except OSError:
continue
class StageClock:
"""Ordered stage timings for one browser call."""
def __init__(self) -> None:
self.stages: list[dict[str, Any]] = []
self.extra: dict[str, Any] = {}
self._start = time.monotonic()
def record(self, stage: str, started: float, ok: bool, **detail: Any) -> None:
entry = {
"stage": stage,
"ms": int((time.monotonic() - started) * 1000),
"ok": bool(ok),
}
entry.update({k: v for k, v in detail.items() if v not in (None, "")})
self.stages.append(entry)
def total_ms(self) -> int:
return int((time.monotonic() - self._start) * 1000)
@dataclass
class BrowserSession:
"""In-process lifecycle record for one owned agent-browser session."""
key: str
ephemeral: bool
root: Path | None = None
env: dict[str, str] | None = field(default=None, repr=False)
command_prefix: list[str] = field(default_factory=list, repr=False)
state: str = "idle"
navigation_generation: int = 0
page_url: str = ""
failed_navigation_url: str = ""
_lock: asyncio.Lock | None = field(default=None, repr=False)
_lock_loop: Any = field(default=None, repr=False)
def bind(self, env: dict[str, str], command_prefix: list[str]) -> None:
"""Record the environment and CLI prefix the daemon is launched with."""
self.env = dict(env)
self.root = runtime_root(env)
self.command_prefix = list(command_prefix)
def lock(self) -> asyncio.Lock:
loop = asyncio.get_running_loop()
if self._lock is None or self._lock_loop is not loop:
self._lock = asyncio.Lock()
self._lock_loop = loop
return self._lock
def navigated(self, url: str) -> None:
self.navigation_generation += 1
self.page_url = url
self.failed_navigation_url = ""
self.state = "ready"
def navigation_failed(self, url: str) -> None:
self.failed_navigation_url = url
self.state = "navigation_failed"
def discarded(self, state: str) -> None:
"""The browser and its page are gone; nothing earlier is observable."""
self.page_url = ""
self.failed_navigation_url = ""
self.state = state
def stale_observation_note(self) -> str:
if not self.failed_navigation_url:
return ""
shown = self.page_url or "an earlier page"
return (
f"Browser lifecycle: the most recent navigation to {self.failed_navigation_url} "
f"failed. This observation shows {shown} (navigation "
f"#{self.navigation_generation}), not {self.failed_navigation_url}."
)
def receipt(self, clock: StageClock) -> dict[str, Any]:
payload = {
"session": self.key,
"ownership": "ephemeral" if self.ephemeral else "retained",
"state": self.state,
"navigation_generation": self.navigation_generation,
"page_url": self.page_url,
"stages": clock.stages,
"elapsed_ms": clock.total_ms(),
}
payload.update({k: v for k, v in clock.extra.items() if v is not None})
return payload
_SESSIONS: dict[str, BrowserSession] = {}
def session_for(key: str, ephemeral: bool) -> BrowserSession:
record = _SESSIONS.get(key)
if record is None:
record = BrowserSession(key=key, ephemeral=ephemeral)
_SESSIONS[key] = record
return record
def forget(key: str) -> None:
_SESSIONS.pop(key, None)
def registered(key: str) -> BrowserSession | None:
return _SESSIONS.get(key)
def has_live_daemon(
root: Path,
key: str,
*,
pid_alive: Callable[[int], bool] = platform_compat.pid_alive,
) -> bool:
"""Whether this session has a daemon a ``close`` command could reach.
Without procfs a live pid from our own pid file is treated as a match,
because answering "no daemon" lets ``close`` bootstrap a fresh browser.
"""
pid = daemon_pid(root, key)
if not pid:
return False
if not platform_compat.has_procfs():
return pid_alive(pid)
return is_verified_daemon(pid)
+30
View File
@@ -17,6 +17,18 @@ from src.runtime_paths import get_app_root
logger = logging.getLogger(__name__)
BROWSER_MCP_SERVER_ID = "builtin_browser"
def browser_mcp_call_timeout() -> float:
"""Upper bound for one Playwright MCP tool call, in seconds."""
try:
value = float(os.environ.get("ODYSSEUS_BROWSER_MCP_CALL_TIMEOUT_S", "90"))
except ValueError:
return 90.0
return value if value > 0 else 90.0
def _format_mcp_connection_error(name: str, command: str = "", args: Optional[List[str]] = None, error: Exception = None) -> str:
"""Return a user-actionable MCP connection error message."""
args = args or []
@@ -521,6 +533,24 @@ class McpManager:
return {"error": f"MCP server not connected: {server_id}", "exit_code": 1}
try:
if server_id == BROWSER_MCP_SERVER_ID:
# The shared Playwright browser must not hold a turn forever.
# The call is abandoned, not retried: page state is unknown.
limit = browser_mcp_call_timeout()
try:
return await asyncio.wait_for(
self._do_call(session, tool_name, arguments), timeout=limit
)
except asyncio.TimeoutError:
logger.warning("Browser MCP call %s timed out after %ss", tool_name, limit)
return {
"error": (
f"Browser call {tool_name} timed out after {limit:g}s and was "
"not retried. The current page state is unknown; navigate "
"again before relying on any observation."
),
"exit_code": 1,
}
result = await self._do_call(session, tool_name, arguments)
except Exception as e:
# Auto-reconnect for builtin servers whose subprocess may have died
+1 -1
View File
@@ -211,7 +211,7 @@ class ResearchNavigator:
self._progress({"phase": "navigating", "url": url, "title": url})
tool = PrivateBrowserTool()
result = await tool.execute(
json.dumps({"action": "read", "url": url, "timeout": timeout}),
json.dumps({"action": "read", "url": url, "timeout_ms": int(timeout * 1000)}),
{"session_id": self.session_id or "research"},
)
output = str(result.get("output") or "").strip()