按授权清空原有内容后重新提交(原 oil-ui-pro 一并移出,可从历史恢复)。 browser-harness 剔除 .venv 等运行环境;根 .gitignore 补记 .venv/ 与 node_modules/。
644 lines
28 KiB
Python
644 lines
28 KiB
Python
"""CDP WS holder + IPC relay (Unix socket on POSIX, TCP loopback on Windows). One daemon per BU_NAME."""
|
|
import asyncio, json, os, platform, socket, sys, time, urllib.error, urllib.request
|
|
from urllib.parse import urlparse
|
|
from collections import deque
|
|
from pathlib import Path
|
|
|
|
from . import _ipc as ipc
|
|
from . import auth
|
|
from . import paths
|
|
from cdp_use.client import CDPClient
|
|
|
|
|
|
def _load_env():
|
|
repo_root = Path(__file__).resolve().parents[2]
|
|
workspace = paths.workspace_dir()
|
|
for p in (repo_root / ".env", workspace / ".env"):
|
|
if not p.exists():
|
|
continue
|
|
_load_env_file(p)
|
|
|
|
|
|
def _load_env_file(p):
|
|
for line in p.read_text(encoding="utf-8-sig", errors="replace").splitlines():
|
|
line = line.strip()
|
|
if not line or line.startswith("#") or "=" not in line:
|
|
continue
|
|
k, v = line.split("=", 1)
|
|
os.environ.setdefault(k.strip(), v.strip().strip('"').strip("'"))
|
|
|
|
|
|
_load_env()
|
|
|
|
NAME = os.environ.get("BU_NAME", "default")
|
|
SOCK = ipc.sock_addr(NAME)
|
|
LOG = str(ipc.log_path(NAME))
|
|
PID = str(ipc.pid_path(NAME))
|
|
BUF = 500
|
|
_MAC_PROFILES = (
|
|
"Library/Application Support/Google/Chrome",
|
|
"Library/Application Support/Google/Chrome Canary",
|
|
"Library/Application Support/Comet",
|
|
"Library/Application Support/Arc/User Data",
|
|
"Library/Application Support/Dia/User Data",
|
|
"Library/Application Support/Microsoft Edge",
|
|
"Library/Application Support/Microsoft Edge Beta",
|
|
"Library/Application Support/Microsoft Edge Dev",
|
|
"Library/Application Support/Microsoft Edge Canary",
|
|
"Library/Application Support/BraveSoftware/Brave-Browser",
|
|
)
|
|
_LINUX_PROFILES = (
|
|
".config/google-chrome",
|
|
".config/chromium",
|
|
".config/chromium-browser",
|
|
".config/microsoft-edge",
|
|
".config/microsoft-edge-beta",
|
|
".config/microsoft-edge-dev",
|
|
".var/app/org.chromium.Chromium/config/chromium",
|
|
".var/app/com.google.Chrome/config/google-chrome",
|
|
".var/app/com.brave.Browser/config/BraveSoftware/Brave-Browser",
|
|
".var/app/com.microsoft.Edge/config/microsoft-edge",
|
|
)
|
|
_WINDOWS_PROFILES = ( # relative to %LOCALAPPDATA%; SxS = Canary channel
|
|
"Google/Chrome/User Data",
|
|
"Google/Chrome SxS/User Data",
|
|
"Google/Chrome Beta/User Data",
|
|
"Google/Chrome Dev/User Data",
|
|
"Chromium/User Data",
|
|
"Microsoft/Edge/User Data",
|
|
"Microsoft/Edge Beta/User Data",
|
|
"Microsoft/Edge Dev/User Data",
|
|
"Microsoft/Edge SxS/User Data",
|
|
"BraveSoftware/Brave-Browser/User Data",
|
|
)
|
|
|
|
|
|
def profile_dirs(system=None):
|
|
system = system or platform.system()
|
|
if system == "Windows":
|
|
local = Path(os.environ.get("LOCALAPPDATA") or Path.home() / "AppData/Local")
|
|
return [local / p for p in _WINDOWS_PROFILES]
|
|
if system == "Darwin":
|
|
return [Path.home() / p for p in _MAC_PROFILES]
|
|
return [Path.home() / p for p in _LINUX_PROFILES]
|
|
|
|
|
|
PROFILES = profile_dirs()
|
|
INTERNAL = ("chrome://", "chrome-untrusted://", "devtools://", "chrome-extension://", "about:")
|
|
BU_API = "https://api.browser-use.com/api/v3"
|
|
REMOTE_ID = os.environ.get("BU_BROWSER_ID")
|
|
BROWSER_KIND = "cloud" if REMOTE_ID else ("cdp" if (os.environ.get("BU_CDP_WS") or os.environ.get("BU_CDP_URL")) else "local")
|
|
# Chrome 144+ shows a per-connection popup. Keep popup open enough to click.
|
|
LOCAL_HANDSHAKE_TIMEOUT = 45
|
|
# How long get_ws_url() keeps waiting for DevToolsActivePort before giving up
|
|
NO_TOGGLE_GRACE = 3
|
|
TOGGLE_BOOT_GRACE = 12
|
|
|
|
|
|
def _devtools_port_live(base):
|
|
"""True when something is listening on the profile's DevToolsActivePort port.
|
|
|
|
A stale file left behind by a closed browser must not count as a running
|
|
instance — it would route recovery to "click Allow" on a popup that can't
|
|
exist."""
|
|
try:
|
|
port = int((base / "DevToolsActivePort").read_text(encoding="utf-8", errors="replace").splitlines()[0].strip())
|
|
except (OSError, ValueError, IndexError):
|
|
return False
|
|
try:
|
|
socket.create_connection(("127.0.0.1", port), timeout=0.5).close()
|
|
return True
|
|
except OSError:
|
|
return False
|
|
|
|
|
|
def remote_debugging_user_enabled():
|
|
"""chrome://inspect's "Allow remote debugging" toggle
|
|
|
|
True only when a toggle-on profile also has a live DevTools port.
|
|
False if a profile records it off, None when no profile records it."""
|
|
seen = None
|
|
for base in PROFILES:
|
|
try:
|
|
state = json.loads((base / "Local State").read_text(encoding="utf-8", errors="replace"))
|
|
enabled = ((state.get("devtools") or {}).get("remote_debugging") or {}).get("user-enabled")
|
|
except (OSError, ValueError, AttributeError):
|
|
continue
|
|
if enabled is True and _devtools_port_live(base):
|
|
return True
|
|
if enabled is False:
|
|
seen = False
|
|
return seen
|
|
|
|
|
|
def remote_debugging_toggle_profiles():
|
|
"""Profile dirs whose chrome://inspect toggle is recorded on in Local State"""
|
|
out = []
|
|
for base in PROFILES:
|
|
try:
|
|
state = json.loads((base / "Local State").read_text(encoding="utf-8", errors="replace"))
|
|
if ((state.get("devtools") or {}).get("remote_debugging") or {}).get("user-enabled") is True:
|
|
out.append(base)
|
|
except (OSError, ValueError, AttributeError):
|
|
continue
|
|
return out
|
|
|
|
|
|
def browser_running_for_profile(base):
|
|
"""True when a running browser instance holds this user-data-dir (POSIX)"""
|
|
try:
|
|
target = os.readlink(str(base / "SingletonLock"))
|
|
except OSError:
|
|
return False
|
|
try:
|
|
pid = int(target.rsplit("-", 1)[-1])
|
|
except ValueError:
|
|
return False
|
|
try:
|
|
os.kill(pid, 0)
|
|
return True
|
|
except ProcessLookupError:
|
|
return False
|
|
except OSError:
|
|
return True # pid exists but belongs to another user
|
|
|
|
|
|
def supported_browser_running():
|
|
"""Is any browser whose profile we scan actually running?"""
|
|
if platform.system() == "Windows":
|
|
# Chromium on Windows uses a named mutex instead of SingletonLock —
|
|
import subprocess
|
|
try:
|
|
out = subprocess.check_output(["tasklist"], text=True, errors="replace", timeout=5).lower()
|
|
except Exception:
|
|
return True # can't tell — assume running so recovery stays on the popup/toggle path
|
|
return any(n in out for n in ("chrome.exe", "msedge.exe", "chromium.exe", "brave.exe", "helium.exe"))
|
|
return any(browser_running_for_profile(base) for base in PROFILES)
|
|
|
|
|
|
def log(msg):
|
|
open(LOG, "a", encoding="utf-8", errors="replace").write(f"{msg}\n")
|
|
|
|
|
|
async def _silent(coro):
|
|
try:
|
|
await coro
|
|
except Exception:
|
|
pass
|
|
|
|
|
|
def _ws_from_devtools_active_port(http_url: str) -> str | None:
|
|
"""When /json/version returns 404 (Chrome 147+ default profile), match DevToolsActivePort by port."""
|
|
p = urlparse(http_url)
|
|
want_port = str(p.port) if p.port else ""
|
|
if not want_port:
|
|
return None
|
|
host = p.hostname or "127.0.0.1"
|
|
if ":" in host: # urlparse strips IPv6 brackets; restore them for the ws:// URL
|
|
host = f"[{host}]"
|
|
for base in PROFILES:
|
|
try:
|
|
active = (base / "DevToolsActivePort").read_text(encoding="utf-8", errors="replace").splitlines()
|
|
except (FileNotFoundError, NotADirectoryError):
|
|
continue
|
|
port = active[0].strip() if active else ""
|
|
ws_path = active[1].strip() if len(active) > 1 else ""
|
|
if port == want_port and ws_path:
|
|
return f"ws://{host}:{port}{ws_path}"
|
|
return None
|
|
|
|
|
|
def get_ws_url():
|
|
if url := os.environ.get("BU_CDP_WS"):
|
|
return url
|
|
if url := os.environ.get("BU_CDP_URL"):
|
|
# HTTP DevTools endpoint (e.g. http://127.0.0.1:9333) — resolve to ws via /json/version.
|
|
# Use this for a dedicated automation Chrome on a non-default profile, which avoids the
|
|
# M144 "Allow remote debugging" dialog and the M136 default-profile lockdown.
|
|
deadline = time.time() + 30
|
|
last_err = None
|
|
base_url = url.rstrip("/")
|
|
while time.time() < deadline:
|
|
try:
|
|
return json.loads(urllib.request.urlopen(f"{base_url}/json/version", timeout=5).read())["webSocketDebuggerUrl"]
|
|
except urllib.error.HTTPError as e:
|
|
last_err = e
|
|
if e.code == 403:
|
|
raise RuntimeError("permission-blocked: Chrome is reachable, but the per-session Allow remote debugging popup has not been accepted")
|
|
if e.code == 404 and (ws := _ws_from_devtools_active_port(url)):
|
|
return ws
|
|
time.sleep(1)
|
|
except Exception as e:
|
|
last_err = e
|
|
time.sleep(1)
|
|
hint = "is the dedicated automation Chrome running? Launch it with --remote-debugging-port=<port> --user-data-dir=<dedicated dir>"
|
|
if platform.system() == "Windows":
|
|
hint += "; on Windows also check that a firewall/antivirus isn't blocking localhost connections"
|
|
raise RuntimeError(f"BU_CDP_URL={url} unreachable after 30s: {last_err} -- {hint}")
|
|
deadline = time.time() + 30
|
|
next_liveness_check = 0.0
|
|
while time.time() < deadline:
|
|
for base in PROFILES:
|
|
try:
|
|
active = (base / "DevToolsActivePort").read_text(encoding="utf-8", errors="replace").splitlines()
|
|
except (FileNotFoundError, NotADirectoryError):
|
|
continue
|
|
port = active[0].strip() if active else ""
|
|
ws_path = active[1].strip() if len(active) > 1 else ""
|
|
if not port:
|
|
continue
|
|
# Resolve the live WS URL via /json/version instead of trusting the path stored
|
|
# alongside the port in DevToolsActivePort: if Chrome was previously launched
|
|
# with a different --user-data-dir on the same port, that file is left behind
|
|
# with a stale browser UUID and the WS upgrade returns 404.
|
|
try:
|
|
return json.loads(urllib.request.urlopen(f"http://127.0.0.1:{port}/json/version", timeout=1).read())["webSocketDebuggerUrl"]
|
|
except urllib.error.HTTPError as e:
|
|
if e.code == 403:
|
|
raise RuntimeError("permission-blocked: Chrome is reachable, but the per-session Allow remote debugging popup has not been accepted")
|
|
# Chrome 147+ disables /json/* HTTP discovery on the default user-data-dir;
|
|
# the ws path Chrome wrote to DevToolsActivePort still works.
|
|
if e.code == 404 and ws_path:
|
|
return f"ws://127.0.0.1:{port}{ws_path}"
|
|
except (OSError, KeyError, ValueError):
|
|
pass
|
|
# Closed browser leaves stale DevToolsActivePort files
|
|
now = time.time()
|
|
if now >= next_liveness_check:
|
|
if not supported_browser_running():
|
|
raise RuntimeError(
|
|
"chrome-not-running: no supported Chromium-family browser is running -- start Chrome, then retry"
|
|
)
|
|
next_liveness_check = now + 2
|
|
# The browser is running but the port isn't up; waiting 30s
|
|
grace = TOGGLE_BOOT_GRACE if remote_debugging_toggle_profiles() else NO_TOGGLE_GRACE
|
|
if now > deadline - 30 + grace:
|
|
break
|
|
time.sleep(0.2)
|
|
for probe_port in (9222, 9223):
|
|
try:
|
|
with urllib.request.urlopen(f"http://127.0.0.1:{probe_port}/json/version", timeout=1) as r:
|
|
return json.loads(r.read())["webSocketDebuggerUrl"]
|
|
except urllib.error.HTTPError as e:
|
|
if e.code == 403:
|
|
raise RuntimeError("permission-blocked: Chrome is reachable, but the per-session Allow remote debugging popup has not been accepted")
|
|
except (OSError, KeyError, ValueError):
|
|
continue
|
|
if remote_debugging_user_enabled() is False:
|
|
raise RuntimeError('remote debugging is turned off for this browser instance — enable chrome://inspect/#remote-debugging (tick "Allow remote debugging for this browser instance")')
|
|
raise RuntimeError(f"DevToolsActivePort not found in {[str(p) for p in PROFILES]} — enable chrome://inspect/#remote-debugging, or set BU_CDP_WS for a remote browser")
|
|
|
|
|
|
def stop_remote():
|
|
if not REMOTE_ID:
|
|
return
|
|
try:
|
|
key = auth.get_browser_use_api_key()
|
|
req = urllib.request.Request(
|
|
f"{BU_API}/browsers/{REMOTE_ID}",
|
|
data=json.dumps({"action": "stop"}).encode(),
|
|
method="PATCH",
|
|
headers={"X-Browser-Use-API-Key": key, "Content-Type": "application/json"},
|
|
)
|
|
urllib.request.urlopen(req, timeout=15).read()
|
|
log(f"stopped remote browser {REMOTE_ID}")
|
|
except Exception as e:
|
|
log(f"stop_remote failed ({REMOTE_ID}): {e}")
|
|
|
|
|
|
def is_real_page(t):
|
|
return t["type"] == "page" and not t.get("url", "").startswith(INTERNAL)
|
|
|
|
|
|
def is_reusable_blank_page(t):
|
|
"""A plain about:blank tab that is safe to attach to and navigate"""
|
|
url = t.get("url", "")
|
|
return (
|
|
t["type"] == "page"
|
|
and (url == "about:blank" or url.startswith("about:blank#"))
|
|
and not t.get("title", "").startswith("Starting agent ")
|
|
)
|
|
|
|
|
|
def is_inspect_tab(t):
|
|
"""A chrome://inspect tab — normally the one the permission flow opened"""
|
|
return t["type"] == "page" and t.get("url", "").startswith("chrome://inspect")
|
|
|
|
|
|
def harness_opened_inspect():
|
|
"""True when admin's recovery flow opened a chrome://inspect tab that is
|
|
still awaiting cleanup (the marker survives until the next connect)."""
|
|
try:
|
|
return paths.inspect_marker().exists()
|
|
except OSError:
|
|
return False
|
|
|
|
|
|
def is_reusable_new_tab_page(t):
|
|
"""The browser's own New Tab Page, ex: from a fresh launch"""
|
|
return t["type"] == "page" and t.get("url", "").startswith(
|
|
("chrome://newtab", "chrome://new-tab-page", "edge://newtab", "about:newtab")
|
|
)
|
|
|
|
|
|
class _PatientCDPClient(CDPClient):
|
|
"""CDPClient with the WS opening handshake stretched to LOCAL_HANDSHAKE_TIMEOUT."""
|
|
|
|
async def start(self):
|
|
import websockets
|
|
if self.ws is not None:
|
|
raise RuntimeError("Client is already started")
|
|
connect_kwargs = {"max_size": self.max_ws_frame_size, "open_timeout": LOCAL_HANDSHAKE_TIMEOUT}
|
|
if self.additional_headers:
|
|
connect_kwargs["additional_headers"] = self.additional_headers
|
|
self.ws = await websockets.connect(self.url, **connect_kwargs)
|
|
self._message_handler_task = asyncio.create_task(self._handle_messages())
|
|
|
|
|
|
class Daemon:
|
|
def __init__(self):
|
|
self.cdp = None
|
|
self.session = None
|
|
self.target_id = None
|
|
self.events = deque(maxlen=BUF)
|
|
self.dialog = None
|
|
self.stop = None # asyncio.Event, set inside start()
|
|
|
|
async def attach_first_page(self):
|
|
"""Attach to a real page (or any page). Sets self.session. Returns attached target or None."""
|
|
targets = (await self.cdp.send_raw("Target.getTargets"))["targetInfos"]
|
|
pages = [t for t in targets if is_real_page(t)]
|
|
if not pages:
|
|
# Fresh browser (ex: BU cloud) starts w about:blank; reuse it
|
|
pages = [t for t in targets if is_reusable_blank_page(t)]
|
|
if not pages:
|
|
# Freshly launched browser (ex: harness relaunching closed Chrome)
|
|
# starts with just the New Tab Page. Reuse it — creating about:blank
|
|
pages = [t for t in targets if is_reusable_new_tab_page(t)]
|
|
take_over = None
|
|
if not pages and harness_opened_inspect():
|
|
# After perms granted, only tab is often chrome://inspect
|
|
# Attach to it instead of creating a new about:blank
|
|
inspect_tabs = [t for t in targets if is_inspect_tab(t)]
|
|
if inspect_tabs:
|
|
pages = [inspect_tabs[0]]
|
|
take_over = inspect_tabs[0]["targetId"]
|
|
if not pages:
|
|
# No usable pages - create one instead of attaching to omnibox popup.
|
|
tid = (await self.cdp.send_raw("Target.createTarget", {"url": "about:blank"}))["targetId"]
|
|
log(f"no real pages found, created about:blank ({tid})")
|
|
pages = [{"targetId": tid, "url": "about:blank", "type": "page"}]
|
|
self.session = (await self.cdp.send_raw(
|
|
"Target.attachToTarget", {"targetId": pages[0]["targetId"], "flatten": True}
|
|
))["sessionId"]
|
|
self.target_id = pages[0]["targetId"]
|
|
log(f"attached {pages[0]['targetId']} ({pages[0].get('url','')[:80]}) session={self.session}")
|
|
if take_over:
|
|
try:
|
|
await self.cdp.send_raw("Page.navigate", {"url": "about:blank"}, session_id=self.session)
|
|
log(f"took over inspect tab {take_over} -> about:blank")
|
|
except Exception as e:
|
|
log(f"take over inspect tab {take_over}: {e}")
|
|
if BROWSER_KIND == "local":
|
|
await self._close_inspect_tabs(targets)
|
|
await self._enable_default_domains(self.session)
|
|
return pages[0]
|
|
|
|
async def _close_inspect_tabs(self, targets):
|
|
"""Close chrome://inspect tabs left open by the permission recovery flow"""
|
|
if not harness_opened_inspect():
|
|
return
|
|
for t in targets:
|
|
if t["targetId"] != self.target_id and is_inspect_tab(t):
|
|
try:
|
|
await self.cdp.send_raw("Target.closeTarget", {"targetId": t["targetId"]})
|
|
log(f"closed leftover chrome://inspect tab {t['targetId']}")
|
|
except Exception as e:
|
|
log(f"close inspect tab {t['targetId']}: {e}")
|
|
try:
|
|
paths.inspect_marker().unlink()
|
|
except OSError:
|
|
pass
|
|
|
|
async def _enable_default_domains(self, session_id):
|
|
"""Enable Page/DOM/Runtime/Network on a CDP session.
|
|
|
|
Used by both initial attach and set_session (called after switch_tab/
|
|
new_tab). Without this, helpers that depend on Network.* events —
|
|
notably wait_for_network_idle() — silently stop receiving events
|
|
after a tab switch, because each fresh CDP session starts with all
|
|
domains disabled.
|
|
|
|
Runs the four enables in parallel via gather so the worst-case time is
|
|
bounded by a single CDP round trip rather than four sequential ones —
|
|
important on the set_session path, where the helper's IPC socket has
|
|
a 5s read timeout.
|
|
"""
|
|
async def enable_one(d):
|
|
try:
|
|
await asyncio.wait_for(
|
|
self.cdp.send_raw(f"{d}.enable", session_id=session_id),
|
|
timeout=4,
|
|
)
|
|
except Exception as e:
|
|
log(f"enable {d} on {session_id}: {e}")
|
|
await asyncio.gather(*(enable_one(d) for d in ("Page", "DOM", "Runtime", "Network")))
|
|
|
|
async def start(self):
|
|
self.stop = asyncio.Event()
|
|
url = get_ws_url()
|
|
log(f"connecting to {url}")
|
|
self.cdp = _PatientCDPClient(url) if BROWSER_KIND == "local" else CDPClient(url)
|
|
if BROWSER_KIND == "local":
|
|
# Allow while this handshake is still parked on the popup
|
|
log("handshake-wait: if Chrome shows an 'Allow remote debugging?' popup, click Allow")
|
|
try:
|
|
await self.cdp.start()
|
|
except Exception as e:
|
|
if os.environ.get("BU_CDP_WS"):
|
|
raise RuntimeError(
|
|
f"CDP WS handshake failed: {e} -- remote browser WebSocket connection failed. "
|
|
"This can happen when network policy blocks the connection, the WS URL is wrong or expired, or the remote endpoint is down. "
|
|
"If you use Browser Use cloud, verify auth and get a fresh URL via start_remote_daemon()."
|
|
)
|
|
if BROWSER_KIND == "local" and ("timed out" in str(e).lower() or "403" in str(e)) and remote_debugging_user_enabled():
|
|
raise RuntimeError(
|
|
f"permission-blocked: Chrome's 'Allow remote debugging?' popup was not accepted within {LOCAL_HANDSHAKE_TIMEOUT}s"
|
|
" -- wait for the user to click Allow, then retry"
|
|
)
|
|
raise RuntimeError(f"CDP WS handshake failed: {e} -- click Allow in Chrome if prompted, then retry")
|
|
await self.attach_first_page()
|
|
orig = self.cdp._event_registry.handle_event
|
|
mark_js = "if(!document.title.startsWith('\U0001F434'))document.title='\U0001F434 '+document.title"
|
|
async def tap(method, params, session_id=None):
|
|
self.events.append({"method": method, "params": params, "session_id": session_id})
|
|
if method == "Page.javascriptDialogOpening":
|
|
self.dialog = params
|
|
elif method == "Page.javascriptDialogClosed":
|
|
self.dialog = None
|
|
elif method in ("Page.loadEventFired", "Page.domContentEventFired"):
|
|
asyncio.create_task(_silent(asyncio.wait_for(self.cdp.send_raw("Runtime.evaluate", {"expression": mark_js}, session_id=self.session), timeout=2)))
|
|
return await orig(method, params, session_id)
|
|
self.cdp._event_registry.handle_event = tap
|
|
|
|
async def handle(self, req):
|
|
# Token guard for Windows TCP loopback: any local process can otherwise
|
|
# connect and issue CDP commands. expected_token() is None on POSIX so
|
|
# this check is a no-op there (AF_UNIX + chmod 600 is the boundary).
|
|
expected = ipc.expected_token()
|
|
if expected is not None and req.get("token") != expected:
|
|
return {"error": "unauthorized"}
|
|
meta = req.get("meta")
|
|
# Liveness probe — lets clients confirm the listener is actually this
|
|
# daemon and not an unrelated process that reused our port post-crash.
|
|
# `pid` lets restart_daemon() verify the live daemon's identity before
|
|
# signaling — protects against SIGTERM-by-stale-pid-file after PID reuse.
|
|
if meta == "ping": return {"pong": True, "pid": os.getpid(), "browser_kind": BROWSER_KIND}
|
|
if meta == "drain_events":
|
|
out = list(self.events); self.events.clear()
|
|
return {"events": out}
|
|
if meta == "session": return {"session_id": self.session}
|
|
if meta == "current_tab":
|
|
# Resolve the attached page's target info server-side. Helpers can't
|
|
# send Target.getTargetInfo themselves: daemon strips session_id for
|
|
# any Target.* method (browser-level call), and without a targetId
|
|
# Chrome silently returns the *browser* target.
|
|
if not self.target_id:
|
|
return {"error": "not_attached"}
|
|
try:
|
|
info = (await self.cdp.send_raw("Target.getTargetInfo", {"targetId": self.target_id}))["targetInfo"]
|
|
except Exception:
|
|
return {"error": "cdp_disconnected"}
|
|
return {"targetId": info.get("targetId"), "url": info.get("url", ""), "title": info.get("title", "")}
|
|
if meta == "connection_status":
|
|
if not self.target_id:
|
|
return {"error": "not_attached"}
|
|
try:
|
|
info = (await self.cdp.send_raw("Target.getTargetInfo", {"targetId": self.target_id}))["targetInfo"]
|
|
except Exception:
|
|
return {"error": "cdp_disconnected"}
|
|
page = None
|
|
if is_real_page(info):
|
|
page = {
|
|
"targetId": info.get("targetId"),
|
|
"title": info.get("title") or "(untitled)",
|
|
"url": info.get("url") or "",
|
|
}
|
|
return {"target_id": self.target_id, "session_id": self.session, "page": page}
|
|
if meta == "set_session":
|
|
old_session = self.session
|
|
self.session = req.get("session_id")
|
|
self.target_id = req.get("target_id") or self.target_id
|
|
# Run the old-session Network.disable (defense in depth — keeps
|
|
# background-tab traffic out of the global event buffer; the
|
|
# consumer-side filter in wait_for_network_idle is the actual
|
|
# correctness gate) in parallel with the four enables on the new
|
|
# session. Different sessions, independent CDP requests. Keeps
|
|
# the synchronous reply under the helper's 5s IPC read timeout
|
|
# even on a remote daemon — sequentially these would have stacked
|
|
# to ~22s worst case.
|
|
tasks = []
|
|
if old_session and old_session != self.session:
|
|
async def disable_old():
|
|
try:
|
|
await asyncio.wait_for(
|
|
self.cdp.send_raw("Network.disable", session_id=old_session),
|
|
timeout=2,
|
|
)
|
|
except Exception: pass
|
|
tasks.append(disable_old())
|
|
tasks.append(self._enable_default_domains(self.session))
|
|
await asyncio.gather(*tasks)
|
|
# 🐴 tab-marker title prefix is purely cosmetic — fire-and-forget so
|
|
# it doesn't add to the synchronous IPC budget.
|
|
asyncio.create_task(_silent(asyncio.wait_for(
|
|
self.cdp.send_raw(
|
|
"Runtime.evaluate",
|
|
{"expression": "if(!document.title.startsWith('\U0001F434'))document.title='\U0001F434 '+document.title"},
|
|
session_id=self.session,
|
|
),
|
|
timeout=2,
|
|
)))
|
|
return {"session_id": self.session}
|
|
if meta == "pending_dialog": return {"dialog": self.dialog}
|
|
if meta == "shutdown": self.stop.set(); return {"ok": True}
|
|
|
|
method = req["method"]
|
|
params = req.get("params") or {}
|
|
# Browser-level Target.* calls must not use a session (stale or otherwise).
|
|
# For everything else, explicit session in req wins; else default.
|
|
sid = None if method.startswith("Target.") else (req.get("session_id") or self.session)
|
|
try:
|
|
return {"result": await self.cdp.send_raw(method, params, session_id=sid)}
|
|
except Exception as e:
|
|
msg = str(e)
|
|
if "Session with given id not found" in msg and sid == self.session and sid:
|
|
log(f"stale session {sid}, re-attaching")
|
|
if await self.attach_first_page():
|
|
return {"result": await self.cdp.send_raw(method, params, session_id=self.session)}
|
|
return {"error": msg}
|
|
|
|
|
|
async def serve(d):
|
|
async def handler(reader, writer):
|
|
try:
|
|
line = await reader.readline()
|
|
if not line: return
|
|
resp = await d.handle(json.loads(line))
|
|
writer.write((json.dumps(resp, default=str) + "\n").encode())
|
|
await writer.drain()
|
|
except Exception as e:
|
|
log(f"conn: {e}")
|
|
try:
|
|
writer.write((json.dumps({"error": str(e)}) + "\n").encode())
|
|
await writer.drain()
|
|
except Exception:
|
|
pass
|
|
finally:
|
|
writer.close()
|
|
|
|
serve_task = asyncio.create_task(ipc.serve(NAME, handler))
|
|
stop_task = asyncio.create_task(d.stop.wait())
|
|
await asyncio.sleep(0.05) # let serve() bind so sock_addr() resolves to the live endpoint
|
|
log(f"listening on {ipc.sock_addr(NAME)} (name={NAME}, remote={REMOTE_ID or 'local'})")
|
|
try:
|
|
await asyncio.wait({serve_task, stop_task}, return_when=asyncio.FIRST_COMPLETED)
|
|
if serve_task.done(): await serve_task # surfaces a serve crash
|
|
finally:
|
|
for t in (serve_task, stop_task):
|
|
t.cancel()
|
|
try: await t
|
|
except (asyncio.CancelledError, Exception): pass
|
|
ipc.cleanup_endpoint(NAME)
|
|
|
|
|
|
async def main():
|
|
d = Daemon()
|
|
await d.start()
|
|
await serve(d)
|
|
|
|
|
|
def already_running():
|
|
# Ping handshake (not a bare connect) so a stale .port file + port reuse
|
|
# after a daemon crash doesn't make us mistake an unrelated listener for ours.
|
|
return ipc.ping(NAME, timeout=1.0)
|
|
|
|
|
|
if __name__ == "__main__":
|
|
if already_running():
|
|
print(f"daemon already running on {SOCK}", file=sys.stderr)
|
|
sys.exit(0)
|
|
open(LOG, "w").close()
|
|
open(PID, "w").write(str(os.getpid()))
|
|
try:
|
|
asyncio.run(main())
|
|
except KeyboardInterrupt:
|
|
pass
|
|
except Exception as e:
|
|
log(f"fatal: {e}")
|
|
sys.exit(1)
|
|
finally:
|
|
stop_remote()
|
|
try: os.unlink(PID)
|
|
except FileNotFoundError: pass
|