"""CDP WS holder + Unix socket relay. One daemon per BU_NAME.""" import asyncio, json, os, socket, sys, urllib.request from collections import deque from pathlib import Path from cdp_use.client import CDPClient def _load_env(): p = Path(__file__).parent / ".env" if not p.exists(): return for line in p.read_text().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 = f"/tmp/bu-{NAME}.sock" LOG = f"/tmp/bu-{NAME}.log" PID = f"/tmp/bu-{NAME}.pid" BUF = 500 PROFILES = [ Path.home() / "Library/Application Support/Google/Chrome", Path.home() / ".config/google-chrome", Path.home() / "AppData/Local/Google/Chrome/User Data", ] 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") API_KEY = os.environ.get("BROWSER_USE_API_KEY") def log(msg): open(LOG, "a").write(f"{msg}\n") def get_ws_url(): if url := os.environ.get("BU_CDP_WS"): return url for base in PROFILES: try: port, path = (base / "DevToolsActivePort").read_text().strip().split("\n", 1) except (FileNotFoundError, NotADirectoryError): continue probe = socket.socket(socket.AF_INET, socket.SOCK_STREAM) probe.settimeout(1) try: probe.connect(("127.0.0.1", int(port.strip()))) except OSError: raise RuntimeError( f"Chrome is not accepting DevTools on 127.0.0.1:{port.strip()} — open Chrome and chrome://inspect/#remote-debugging, then retry" ) finally: probe.close() return f"ws://127.0.0.1:{port.strip()}{path.strip()}" 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 or not API_KEY: return try: req = urllib.request.Request( f"{BU_API}/browsers/{REMOTE_ID}", data=json.dumps({"action": "stop"}).encode(), method="PATCH", headers={"X-Browser-Use-API-Key": API_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) class Daemon: def __init__(self): self.cdp = None self.session = None self.events = deque(maxlen=BUF) 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)] or [t for t in targets if t["type"] == "page"] if not pages: self.session = None return None self.session = (await self.cdp.send_raw( "Target.attachToTarget", {"targetId": pages[0]["targetId"], "flatten": True} ))["sessionId"] log(f"attached {pages[0]['targetId']} ({pages[0].get('url','')[:80]}) session={self.session}") for d in ("Page", "DOM", "Runtime", "Network"): try: await self.cdp.send_raw(f"{d}.enable", session_id=self.session) except Exception as e: log(f"enable {d}: {e}") return pages[0] async def start(self): self.stop = asyncio.Event() url = get_ws_url() log(f"connecting to {url}") self.cdp = CDPClient(url) # Chrome shows a native "Allow debugging" dialog on first connect — user may take a while to click. for attempt in range(12): try: await self.cdp.start(); break except Exception as e: log(f"ws handshake attempt {attempt+1} failed: {e} — retrying") self.cdp = CDPClient(url) await asyncio.sleep(5) else: raise RuntimeError("CDP WS handshake never succeeded — did you accept Chrome's Allow dialog?") await self.attach_first_page() orig = self.cdp._event_registry.handle_event async def tap(method, params, session_id=None): self.events.append({"method": method, "params": params, "session_id": session_id}) return await orig(method, params, session_id) self.cdp._event_registry.handle_event = tap async def handle(self, req): meta = req.get("meta") if meta == "drain_events": out = list(self.events); self.events.clear() return {"events": out} if meta == "session": return {"session_id": self.session} if meta == "set_session": self.session = req.get("session_id"); return {"session_id": self.session} 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): if os.path.exists(SOCK): os.unlink(SOCK) 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() server = await asyncio.start_unix_server(handler, path=SOCK) os.chmod(SOCK, 0o600) log(f"listening on {SOCK} (name={NAME}, remote={REMOTE_ID or 'local'})") async with server: await d.stop.wait() async def main(): d = Daemon() await d.start() await serve(d) def already_running(): try: s = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM); s.settimeout(1) s.connect(SOCK); s.close(); return True except (FileNotFoundError, ConnectionRefusedError, socket.timeout): return False 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