"""
SKYKOI browser agent — a local browser-use service the koi talks to.

WHY THIS EXISTS
---------------
The koi used to browse by having its main model improvise CDP calls inside a
code-mode program. That model is the slowest thing in the system (a device turn
runs 15-40s), it could not see the page as a structure, and on a real call it
opened the same wrong IKEA page three times and then told the user it had done
something it had not. browser-use solves exactly that problem: an indexed view
of the interactive elements, a step loop, and a model that only has to choose
among numbered choices rather than invent selectors.

WHAT THIS ADDS ON TOP OF browser-use
------------------------------------
1. A WARM BROWSER. Launching Chrome costs ~5s of every task. The session is
   kept alive between tasks and reused, so the second task starts instantly.
2. A PERSISTENT PROFILE, so logins survive. The user signs into a site once and
   every later task is already authenticated.
3. VERIFICATION. browser-use will report success while sitting on a homepage —
   measured: with flash_mode the IKEA task "succeeded" in 17s on ikea.com
   instead of the towels page. Every result carries the observable end state
   (url, title) and a `verified` flag, so the caller can tell a claim from a
   fact instead of taking the agent's word for it.
4. STREAMING. Steps are emitted as they happen so the koi can narrate what it
   is doing rather than going silent for half a minute.

MEASURED (IKEA "open the bath towels listing", this machine):
    cold, one action per step ......... 52.9s  correct
    warm, flash_mode .................. 17.3s  WRONG (stopped on the homepage)
    warm, multi-action, no flash ...... 32.0s  correct   <- the default here

Run:
    python server.py            # binds 127.0.0.1:18793
"""

from __future__ import annotations

import asyncio
import json
import os
import re
import shutil
import subprocess
import sys
import time
import traceback
from typing import Any


# ---- self-bootstrap: install our python deps on first run -------------------
# A fresh machine has Python but not our packages (the koi installer provisions
# Node; python deps land here on demand). Missing imports: pip install from the
# requirements.txt beside this file, then re-exec once. The env flag guards
# against an install-fail loop; exit 78 tells the supervisor "configuration".
def _ensure_python_deps(probe_modules) -> None:
    import importlib.util
    import os as _os
    import subprocess as _sp
    import sys as _sys

    if all(importlib.util.find_spec(m) is not None for m in probe_modules):
        return
    if _os.environ.get("SKYKOI_PY_BOOTSTRAPPED") == "1":
        print("[bootstrap] python deps still missing after install; giving up", flush=True)
        _sys.exit(78)
    req = _os.path.join(_os.path.dirname(_os.path.abspath(__file__)), "requirements.txt")
    print(f"[bootstrap] installing python deps from {req} (first run)…", flush=True)
    for cmd in (
        [_sys.executable, "-m", "pip", "install", "-r", req],
        [_sys.executable, "-m", "pip", "install", "--user", "-r", req],
    ):
        try:
            if _sp.run(cmd, timeout=900).returncode == 0:
                break
        except Exception:
            continue
    env = {**_os.environ, "SKYKOI_PY_BOOTSTRAPPED": "1"}
    _os.execve(_sys.executable, [_sys.executable, _os.path.abspath(__file__)], env)

_ensure_python_deps(("aiohttp", "browser_use"))

from aiohttp import web

# ── configuration ────────────────────────────────────────────────────────────


def _load_env_fallback() -> None:
    """
    Fill in missing env vars from the places skykoi keeps them.

    The service inherits its env from whoever spawned it — and "whoever" has
    ranged from the gateway (which has GROQ_API_KEY) to a supervisor, a shell,
    or a person restarting it by hand mid-debug (which usually does NOT). A
    keyless sidecar is the worst kind of broken: /health answers ok, the warm
    browser sits there, and every task dies instantly with "no model to think
    with" — watched live 2026-08-19, where it looked exactly like "the browser
    never opened". The service's ability to think must not depend on how it
    was started, so missing keys are read here from the same files the
    gateway loads. Spawner-provided values always win; only gaps are filled.
    """
    here = os.path.dirname(os.path.abspath(__file__))
    candidates = [
        os.path.join(os.path.expanduser("~"), ".skykoi", ".env"),
        os.path.join(here, "..", "..", ".env.local"),
        os.path.join(here, "..", "..", ".env"),
    ]
    for path in candidates:
        try:
            with open(path, encoding="utf8") as handle:
                for raw in handle:
                    line = raw.strip()
                    if not line or line.startswith("#") or "=" not in line:
                        continue
                    key, _, value = line.partition("=")
                    key = key.strip()
                    # Values in these files are often QUOTED (KEY="gsk_...");
                    # the quotes are file syntax, not part of the secret.
                    value = value.strip().strip('"').strip("'")
                    if key and value and not os.environ.get(key):
                        os.environ[key] = value
        except OSError:
            continue


_load_env_fallback()

HOST = os.environ.get("SKYKOI_BROWSER_AGENT_HOST", "127.0.0.1")
PORT = int(os.environ.get("SKYKOI_BROWSER_AGENT_PORT", "18793"))
# Loopback only by default, and a shared secret when the runtime sets one: this
# process can drive a logged-in browser, so it is not something to leave open.
TOKEN = os.environ.get("SKYKOI_BROWSER_AGENT_TOKEN", "")

PROFILE_DIR = os.environ.get(
    "SKYKOI_BROWSER_AGENT_PROFILE",
    os.path.join(os.path.expanduser("~"), ".skykoi", "browser-agent-profile"),
)
# Headful by default: the user watching their own machine do the thing is the
# point. Set 1 for background jobs.
HEADLESS = os.environ.get("SKYKOI_BROWSER_AGENT_HEADLESS", "0") == "1"

# gpt-oss-120b on Groq: strong enough to choose actions reliably and fast enough
# that a step is not a wait. NOTE: Groq rejects browser-use's array-shaped
# message content, so vision is off on this provider (`use_vision=False`) —
# the DOM index is what the agent actually steers by.
# qwen3.6 on Groq's LPU is the owner's pick for this lane and the default it
# should have had: same model the call's instant layer runs on, and with
# reasoning_effort "none" (see below) it answers the action schema in a few
# hundred ms instead of deliberating on every step.
# A "wafer/..." model routes to Wafer's OpenAI-compatible serverless endpoint
# (pass.wafer.ai) instead of Groq. Only the prefix picks the provider;
# everything after it is the model id as Wafer's /v1/models lists it.
#
# DeepSeek-V4-Flash-Fast is the owner's pick for prod (2026-08-20), CHOSEN
# WITH THE NUMBERS ON THE TABLE: measured ~0.5-1s slower per step than
# qwen3.6 on Groq (RTT 403ms vs 138ms, decode ~280 vs ~530 tok/s — physics,
# not config), in exchange for a stronger model: it lands on the right page
# (IKEA category page, not search results), carries a 1M context, and
# prefix-caches the stable prompt. reasoning_effort "none" is load-bearing;
# see _wafer_client. Groq remains one env var away.
MODEL = os.environ.get("SKYKOI_BROWSER_AGENT_MODEL", "wafer/DeepSeek-V4-Flash-0731-Fast")
WAFER_BASE_URL = os.environ.get("SKYKOI_BROWSER_AGENT_WAFER_URL", "https://pass.wafer.ai/v1")
MAX_STEPS_DEFAULT = int(os.environ.get("SKYKOI_BROWSER_AGENT_MAX_STEPS", "18"))
# How long a browser may sit warm with nobody using it before it is closed.
# Four hours, not fifteen minutes: every idle close costs the NEXT task a
# 3-5s Chrome launch, and "open the towels page" pausing to boot a browser is
# exactly the latency the owner asked to have gone. RAM for a warm idle
# Chrome is cheap; the launch in the middle of a request is not.
IDLE_SHUTDOWN_S = int(os.environ.get("SKYKOI_BROWSER_AGENT_IDLE_S", "14400"))
# How long a browser may take to come up before the profile is treated as bad.
# Generous next to the ~3s a healthy launch takes, tight next to the 30s the
# watchdog spends on a profile it will never open.
LAUNCH_TIMEOUT_S = int(os.environ.get("SKYKOI_BROWSER_AGENT_LAUNCH_TIMEOUT_S", "20"))

# True only while an agent is actually driving. Recovery machinery that would
# heal a crashed tab mid-task must not fight a person closing an idle window.
_task_running = False


def _chromium_path() -> str | None:
    """
    Playwright's plain Chromium, when the machine has it — for the EPHEMERAL
    fallback only. The persistent profile stays on the browser that created it
    (an Edge profile opened by Chromium is a different profile), but a
    throwaway browser has no such tie, and Edge is a bad default for one:
    with no Chrome installed browser-use falls back to Edge, and Edge opens
    every new tab on ntp.msn.com — a feed of headlines, widgets and sign-in
    prompts that floods the element index before the task's first real page
    (watched live 2026-08-18: four steps spent on the MSN page, task over).
    Chromium opens on about:blank: zero elements, zero noise. An explicit
    SKYKOI_BROWSER_AGENT_EXECUTABLE override still wins.
    """
    override = os.environ.get("SKYKOI_BROWSER_AGENT_EXECUTABLE")
    if override:
        return override
    root = os.path.join(os.environ.get("LOCALAPPDATA", ""), "ms-playwright")
    try:
        versions = sorted(
            (name for name in os.listdir(root) if name.startswith("chromium-")),
            reverse=True,
        )
    except OSError:
        return None
    for version in versions:
        for sub in ("chrome-win64", "chrome-win"):
            candidate = os.path.join(root, version, sub, "chrome.exe")
            if os.path.isfile(candidate):
                return candidate
    return None


class Timing:
    """
    Where the wall-clock actually goes.

    A browser task is a loop of: think, act, re-read the page. Without splitting
    those you cannot tell a slow model from a slow site, and every "make it
    faster" is a guess. The model is timed by wrapping the client; the rest of
    each step is everything else — DOM extraction, the CDP calls, page loads.
    """

    def __init__(self) -> None:
        self.llm_calls: list[float] = []
        self.first_llm_ms: float | None = None
        self.started = time.time()

    def record_llm(self, ms: float) -> None:
        if self.first_llm_ms is None:
            self.first_llm_ms = ms
        self.llm_calls.append(ms)

    def summary(self, history: Any) -> dict[str, Any]:
        total_ms = (time.time() - self.started) * 1000
        llm_ms = sum(self.llm_calls)
        steps = []
        for item in getattr(history, "history", []) or []:
            meta = getattr(item, "metadata", None)
            if meta and meta.step_end_time and meta.step_start_time:
                steps.append((meta.step_end_time - meta.step_start_time) * 1000)
        return {
            "totalMs": round(total_ms),
            "llmMs": round(llm_ms),
            "llmCalls": len(self.llm_calls),
            "llmSlowestMs": round(max(self.llm_calls)) if self.llm_calls else 0,
            "llmMedianMs": round(sorted(self.llm_calls)[len(self.llm_calls) // 2]) if self.llm_calls else 0,
            "firstLlmMs": round(self.first_llm_ms) if self.first_llm_ms else 0,
            # Everything that is not the model: page loads, DOM extraction, CDP.
            "browserMs": round(total_ms - llm_ms),
            "stepMs": [round(x) for x in steps],
        }


def _wafer_client():
    """
    Wafer's OpenAI-compatible endpoint, for "wafer/<model-id>" MODEL values.

    Verified against the live catalog (GET /v1/models, 2026-08-20):
    DeepSeek-V4-Flash-0731-Fast exists, tier "serverless_only", "the same
    model served for high TPS", tools + json_schema + streaming, 1,048,576
    ctx, $0.28/M in / $0.56/M out. The ~1,077 tok/s making the rounds is a
    DECODE figure — a browser-use step is a short structured answer dominated
    by time-to-first-token, so the only number that matters here is measured
    llmMs on real steps, which is what Timing exists to expose.
    """
    from browser_use.llm.openai.chat import ChatOpenAI

    key = os.environ.get("WAFER_API_KEY", "")
    if not key:
        raise RuntimeError("WAFER_API_KEY is not set — the browser agent has no model to think with")
    model_id = MODEL.split("/", 1)[1]
    return ChatOpenAI(
        model=model_id,
        api_key=key,
        base_url=WAFER_BASE_URL,
        # Same reasoning discipline as the Groq lane: picking an element index
        # needs an answer, not a deliberation. DeepSeek V4 Flash reasons by
        # DEFAULT (measured: 140 reasoning tokens to answer "hi"), and even
        # "low" still burned 106 — only "none" actually turns it off
        # (verified 0 reasoning tokens, 2026-08-20).
        reasoning_effort="none",
        temperature=0.0,
    )


def _groq_client(model: str | None = None):
    from browser_use.llm.groq.chat import ChatGroq

    key = os.environ.get("GROQ_API_KEY", "")
    if not key:
        raise RuntimeError("GROQ_API_KEY is not set — the browser agent has no model to think with")

    groq_model = model or MODEL
    # temperature 0: the action schema is a FORM, not prose. At 0.2 qwen drifted
    # out of it nine times in one run ("{(" instead of "{", a pseudo-XML
    # <scroll> block, a truncated string) and every drift costs a full round
    # trip to Groq before the agent can retry.
    client = ChatGroq(model=groq_model, api_key=key, temperature=0.0)

    # TURN THE THINKING OFF. browser-use's ChatGroq has no reasoning_effort
    # field, so qwen deliberated on every single step: measured 28 calls at a
    # 3.2s median — 92s of a 130s task — where the same model answers a direct
    # question in 169-305ms. The agent loop does not need chain-of-thought to
    # pick an element index; it needs an answer. Injected at the client, which
    # is the only place the parameter can reach the wire.
    effort = "none" if "qwen" in groq_model.lower() else "low"
    original_get_client = client.get_client

    def get_client():
        groq_client = original_get_client()
        if not getattr(groq_client, "_skykoi_effort_patched", False):
            create = groq_client.chat.completions.create

            async def create_with_effort(*args: Any, **kw: Any):
                kw.setdefault("reasoning_effort", effort)
                return await create(*args, **kw)

            groq_client.chat.completions.create = create_with_effort
            groq_client._skykoi_effort_patched = True
        return groq_client

    client.get_client = get_client  # type: ignore[method-assign]
    return client


FALLBACK_GROQ_MODEL = os.environ.get("SKYKOI_BROWSER_AGENT_FALLBACK_MODEL", "qwen/qwen3.6-27b")


def _llm(timing: Timing | None = None):
    if MODEL.startswith("wafer/"):
        if os.environ.get("WAFER_API_KEY"):
            client = _wafer_client()
        else:
            # A device without a Wafer key must DEGRADE, not die. The default
            # model changing providers cannot be the thing that turns every
            # browser task on someone's machine into an instant error — that
            # is the exact keyless-but-healthy failure this file already got
            # burned by once (2026-08-19).
            client = _groq_client(FALLBACK_GROQ_MODEL)
    else:
        client = _groq_client()

    if timing is None:
        return client

    # Time every model call without reimplementing the client: wrap the one
    # method the agent loop uses.
    original = client.ainvoke

    async def timed(*args: Any, **kw: Any):
        t0 = time.time()
        try:
            return await original(*args, **kw)
        finally:
            timing.record_llm((time.time() - t0) * 1000)

    client.ainvoke = timed  # type: ignore[method-assign]
    return client


# ── the warm browser ─────────────────────────────────────────────────────────


def _let_the_user_close_it() -> None:
    """
    Let the person close the browser. Traced the hard way, twice.

    Closing the last tab put a NEW about:blank in its place within three
    seconds, every time, forever — proven over CDP by closing repeatedly and
    watching the target id change on each round. Two separate pieces of
    browser-use do this, and only the second one was the culprit:

      1. AboutBlankWatchdog.on_TabClosedEvent — "creating new about:blank tab
         to avoid closing entire browser". Neutralised below. NOT the cause:
         the tabs kept coming back with it disabled.
      2. SessionManager._recover_agent_focus — "No tabs remain! Creating new
         tab for agent...". This is the one. It exists so a crashed tab does
         not strand a RUNNING agent, which is correct; the trouble is that a
         warm idle browser still has a live session, so it treats the user
         closing a window as a crash to heal.

    So recovery is kept for the case it was written for — a task in flight —
    and skipped when nothing is running, which is when a closed window means
    the person wants it closed. `_task_running` is set around the agent run.

    Both are patched on the CLASS: the watchdog does not exist yet at launch
    (attach_all_watchdogs runs later), so an instance hook set then binds to
    nothing — that was the first fix, and it silently did nothing.
    """
    try:
        from browser_use.browser.watchdogs.aboutblank_watchdog import AboutBlankWatchdog

        # The handler name matters: browser-use asserts every handler it
        # attaches starts with "on_", so a descriptively-named replacement
        # breaks attachment outright.
        async def on_TabClosedEvent(self, event) -> None:  # noqa: ANN001, N802
            return None

        AboutBlankWatchdog.on_TabClosedEvent = on_TabClosedEvent  # type: ignore[method-assign]
    except Exception as error:  # a future version may move it; never fatal
        print(f"[browser-agent] about:blank watchdog not patched: {error}", flush=True)

    try:
        from browser_use.browser.session_manager import SessionManager

        original = SessionManager._recover_agent_focus

        async def _recover_agent_focus(self, crashed_target_id):  # noqa: ANN001, ANN202
            if not _task_running:
                print("[browser-agent] browser closed while idle; leaving it closed", flush=True)
                return None
            return await original(self, crashed_target_id)

        SessionManager._recover_agent_focus = _recover_agent_focus  # type: ignore[method-assign]
    except Exception as error:
        print(f"[browser-agent] focus recovery not patched: {error}", flush=True)
    print("[browser-agent] closing the window will close it", flush=True)


def _browser_is_alive(browser: Any) -> bool:
    """Whether the session still has a live CDP connection behind it."""
    try:
        return bool(getattr(browser, "is_cdp_connected", False))
    except Exception:
        return False


class BrowserPool:
    """
    One browser, kept alive between tasks.

    Serialized on purpose: two agents driving one browser fight over the same
    tabs. Concurrency here would mean a second Browser, which is a second
    Chrome, which is the kind of thing that quietly eats a laptop.
    """

    def __init__(self) -> None:
        self._browser: Any = None
        self._lock = asyncio.Lock()
        self._last_used = 0.0

    def _quarantine_profile(self) -> str | None:
        """
        Move a profile Chrome will not open out of the way, and start clean.

        A launch killed part-way (a crash, a taskkill, a closed laptop) can
        leave the user-data-dir in a state Chrome HANGS on rather than refuses:
        30 seconds, then a watchdog error that names an event handler and says
        nothing about the profile. Every task dead until someone deletes a
        directory nobody would think to suspect.

        Proven the hard way here — every freshly created profile started in
        ~3s while this one never did, same process, same code, same flags.

        Moved rather than deleted: it holds the user's logins, and a corrupt
        profile is still worth keeping if those are ever worth recovering.
        """
        global PROFILE_DIR
        if not os.path.isdir(PROFILE_DIR):
            return None
        # Close what is still using it FIRST. Skipping this is what left a
        # browser running in every profile ever abandoned, and it is also why
        # the rename below kept failing: the handles belonged to a browser we
        # could have closed all along.
        self._kill_profile_processes(PROFILE_DIR)
        # Renaming is the tidier move and it FAILS on Windows whenever a
        # lingering Chrome still holds a handle into the directory — which is
        # exactly the situation that corrupted it. So the profile is abandoned
        # in place and a fresh one takes over: always possible, no handles
        # required, and the old directory is still there to recover logins from.
        broken = PROFILE_DIR
        try:
            os.rename(PROFILE_DIR, f"{PROFILE_DIR}.broken-{int(time.time())}")
        except OSError:
            PROFILE_DIR = f"{PROFILE_DIR}-{int(time.time())}"
        return broken

    def _marker(self) -> str:
        return os.path.join(PROFILE_DIR, ".skykoi-clean-exit")

    @staticmethod
    async def _find_attachable(profile_dir: str) -> str | None:
        """
        A browser still holding the profile may be perfectly healthy.

        An unclean exit usually means the SERVICE died, not the browser: the
        laptop lid, a taskkill, an update. The browser that survives it has the
        user's logins open and their tabs where they left them — and until now
        the recovery path killed it just to relaunch the same thing colder.
        (Ported from koi-harness, where re-attaching across restarts is how the
        shared window survives.)

        Every browser this service launches carries --remote-debugging-port on
        its command line, so the still-running one can be found by the profile
        name and asked over CDP whether it is alive. If it answers, attach to
        it instead of killing it: the warm window, the logins, and the launch
        cost all come back for free.
        """
        name = os.path.basename(profile_dir)
        if not name or sys.platform != "win32":
            return None
        script = (
            "Get-CimInstance Win32_Process -Filter \"Name='msedge.exe' or Name='chrome.exe'\" "
            f"| Where-Object {{ $_.CommandLine -like '*{name}*' -and $_.CommandLine -like '*--remote-debugging-port=*' }} "
            "| Select-Object -ExpandProperty CommandLine"
        )
        try:
            proc = await asyncio.create_subprocess_exec(
                "powershell", "-NoProfile", "-NonInteractive", "-Command", script,
                stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.DEVNULL,
            )
            out, _ = await asyncio.wait_for(proc.communicate(), timeout=15)
        except (asyncio.TimeoutError, OSError):
            return None
        import aiohttp

        for match in re.finditer(r"--remote-debugging-port=(\d+)", out.decode(errors="replace")):
            url = f"http://127.0.0.1:{match.group(1)}"
            try:
                async with aiohttp.ClientSession() as session:
                    async with session.get(f"{url}/json/version", timeout=aiohttp.ClientTimeout(total=2)) as resp:
                        if resp.status == 200:
                            return url
            except Exception:
                continue
        return None

    @staticmethod
    def _kill_profile_processes(profile_dir: str) -> int:
        """
        Kill whatever browser is still holding `profile_dir`.

        Abandoning a profile does NOT close the browser using it, and nothing
        else ever did: found on the owner's machine, 21 abandoned profile
        directories and 79 live Edge processes across 8 of them, every one of
        them orphaned by a quarantine that walked away. They hold RAM, they hold
        file handles on the very directory that could not be renamed, and one of
        them is the window that "keeps coming back" after being closed.
        """
        name = os.path.basename(profile_dir)
        if not name:
            return 0
        try:
            if sys.platform == "win32":
                # Match on the profile directory in the command line: this is
                # the only thing that distinguishes our browser from the user's
                # own Edge, which must never be touched.
                script = (
                    "Get-CimInstance Win32_Process -Filter \"Name='msedge.exe' or Name='chrome.exe'\" "
                    f"| Where-Object {{ $_.CommandLine -like '*{name}*' }} "
                    "| ForEach-Object { Stop-Process -Id $_.ProcessId -Force -ErrorAction SilentlyContinue }"
                )
                subprocess.run(
                    ["powershell", "-NoProfile", "-NonInteractive", "-Command", script],
                    capture_output=True,
                    timeout=30,
                )
            else:
                subprocess.run(["pkill", "-f", name], capture_output=True, timeout=30)
        except Exception as error:  # never let cleanup break a task
            print(f"[browser-agent] could not clean up {name}: {error}", flush=True)
            return 0
        return 1

    async def acquire(self) -> Any:
        from browser_use import Browser

        # The user is allowed to close it, so the handle we hold may point at a
        # browser that no longer exists. Handing that to a task produces a hang
        # rather than an error, so it is dropped here and relaunched below.
        if self._browser is not None and not _browser_is_alive(self._browser):
            print("[browser-agent] the browser was closed; starting a new one for this task", flush=True)
            self._browser = None

        if self._browser is None:
            # A profile left behind by a killed Chrome does not fail to open —
            # it HANGS, and the watchdog takes 30 seconds to say so. Measured:
            # 39s of "browser time" on a 2-step task, 30 of it this. So the
            # question is asked BEFORE launching: did the last session end
            # cleanly? No marker means no, and a fresh profile costs 3s.
            # An unclean exit USUALLY just means a browser is still holding the
            # directory — the user closed the window, the service was killed,
            # the laptop shut. Throwing the profile away for that costs them
            # every login it existed to keep, so close the holder and try the
            # profile first. Quarantine is the fallback, not the reflex.
            if os.path.isdir(PROFILE_DIR) and not os.path.exists(self._marker()):
                # Before killing whatever holds the profile, ask whether it is a
                # healthy browser we can simply re-attach to — an unclean exit
                # is usually OURS, not the browser's, and killing a live window
                # costs the user its warmth, its tabs, and a relaunch.
                cdp = await self._find_attachable(PROFILE_DIR)
                if cdp:
                    try:
                        attached = Browser(cdp_url=cdp, keep_alive=True)
                        await asyncio.wait_for(attached.start(), timeout=LAUNCH_TIMEOUT_S)
                        print(f"[browser-agent] re-attached to the running browser at {cdp}", flush=True)
                        self._browser = attached
                        self._last_used = time.time()
                        return self._browser
                    except (Exception, asyncio.TimeoutError):
                        print("[browser-agent] running browser would not accept an attach; closing it", flush=True)
                print("[browser-agent] previous session did not exit cleanly; closing anything still holding the profile", flush=True)
                self._kill_profile_processes(PROFILE_DIR)
            os.makedirs(PROFILE_DIR, exist_ok=True)

            async def start_one() -> Any:
                browser = Browser(
                    headless=HEADLESS,
                    user_data_dir=PROFILE_DIR,
                    keep_alive=True,
                    # Watching it work is the point of the headful mode, and on
                    # a machine with no Chrome this drives EDGE — whose window
                    # is indistinguishable from the ones the user already has
                    # open. So it gets its own window, maximised and in front,
                    # rather than quietly joining the pile.
                    # --disable-extensions: the user's Edge account syncs its
                    # extensions into this profile, and their content scripts
                    # inject into every page the agent drives — measured as a
                    # consistent ~15s of "browser time" PER STEP on a page that
                    # loads in under a second. Logins live in the profile and
                    # survive; the extensions have no business in an agent
                    # window anyway.
                    # The occlusion flags keep the compositor producing frames
                    # while the window is covered by others — without them the
                    # embedded view's screencast (browser-view) goes dark the
                    # moment the user puts anything in front of the browser.
                    args=(
                        ["--disable-extensions", "--disable-backgrounding-occluded-windows", "--disable-renderer-backgrounding", "--disable-background-timer-throttling"]
                        if HEADLESS
                        else ["--new-window", "--start-maximized", "--window-position=0,0", "--disable-extensions", "--disable-backgrounding-occluded-windows", "--disable-renderer-backgrounding", "--disable-background-timer-throttling"]
                    ),
                )
                await browser.start()
                return browser

            try:
                # A profile Chrome cannot open HANGS rather than refusing, and
                # the built-in watchdog takes 30s to notice. The timeout turns
                # that into a failure we can actually act on.
                self._browser = await asyncio.wait_for(start_one(), timeout=LAUNCH_TIMEOUT_S)
            except (Exception, asyncio.TimeoutError):
                aside = self._quarantine_profile()
                print(f"[browser-agent] profile would not open; moved aside to {aside}", flush=True)
                os.makedirs(PROFILE_DIR, exist_ok=True)
                try:
                    self._browser = await asyncio.wait_for(start_one(), timeout=LAUNCH_TIMEOUT_S)
                except (Exception, asyncio.TimeoutError):
                    # Even a fresh profile would not open — something machine-
                    # wide (AV lock, disk, a Chrome update mid-flight). A task
                    # without logins still beats no task at all, so the last
                    # resort is an ephemeral profile. It goes through the pool
                    # like any other browser so the idle reaper closes it, and
                    # the acquire after that tries the real profile again.
                    print("[browser-agent] fresh profile would not open either; running on an ephemeral profile", flush=True)
                    self._browser = Browser(headless=HEADLESS, user_data_dir=None, executable_path=_chromium_path(), keep_alive=True)
                    await self._browser.start()
            # The marker is removed while the browser is live and rewritten on a
            # clean close, so it is only ever present for a profile nobody is
            # using and nothing crashed on.
            try:
                os.remove(self._marker())
            except OSError:
                pass
        self._last_used = time.time()
        return self._browser

    async def close(self) -> None:
        browser, self._browser = self._browser, None
        if browser is not None:
            try:
                await browser.kill()
            except Exception:
                pass
            try:
                with open(self._marker(), "w", encoding="utf8") as handle:
                    handle.write(str(int(time.time())))
            except OSError:
                pass

    async def reap_idle(self) -> None:
        while True:
            await asyncio.sleep(60)
            if self._browser is not None and time.time() - self._last_used > IDLE_SHUTDOWN_S:
                if not self._lock.locked():
                    await self.close()

    @property
    def lock(self) -> asyncio.Lock:
        return self._lock


POOL = BrowserPool()


# ===== Embedded browser view — the site's Browser tab =====
#
# The PAGE, not the desktop: a CDP screencast of the active tab's viewport,
# served as multipart JPEG the way the gateway's screen stream is — but scoped
# to the browser page, so the web app embeds the real browser without ever
# showing (or stealing) the user's screen. Input dispatch rides the same CDP
# connection space, so the embed is a browser you can actually use.
#
# Coexists with browser-use: CDP multiplexes clients, so the screencast
# watches the same page the agent drives — when the koi browses, the embed IS
# the koi browsing.

VIEW_MAX_W = 1400
VIEW_MAX_H = 1000


class BrowserView:
    def __init__(self) -> None:
        self._viewers: set[Any] = set()
        self._pump_task: asyncio.Task | None = None
        self._viewport = {"width": 1280, "height": 800}
        # Persistent CDP connection for input. Hover streams ~15 events/sec;
        # a fresh websocket per event would spend more time on handshakes
        # than on the events. Refreshed when the active page changes or the
        # socket errors. Guarded by the lock — CDP multiplexes fine, but our
        # bookkeeping does not.
        self._input_session: Any = None
        self._input_ws: Any = None
        self._input_page_id: str | None = None
        self._input_checked_at = 0.0
        self._input_msg_id = 0
        self._input_lock = asyncio.Lock()
        self._last_viewer_at = 0.0

    async def _cdp_http(self) -> str | None:
        browser = POOL._browser
        if browser is None or not _browser_is_alive(browser):
            async with POOL.lock:
                browser = await POOL.acquire()
        return self._normalize_cdp(getattr(browser, "cdp_url", None))

    @staticmethod
    def _normalize_cdp(url: Any) -> str | None:
        """browser-use may hold the BROWSER WebSocket url (ws://host:port/devtools/...);
        the /json inventory lives at plain http on the same host:port."""
        if not url:
            return None
        match = re.match(r"^(?:wss?|https?)://([^/]+)", str(url))
        return f"http://{match.group(1)}" if match else None

    @staticmethod
    async def _active_page(session: Any, cdp_http: str) -> dict | None:
        import aiohttp

        try:
            async with session.get(f"{cdp_http}/json/list", timeout=aiohttp.ClientTimeout(total=3)) as resp:
                targets = await resp.json()
        except Exception:
            return None
        pages = [
            t
            for t in targets
            if t.get("type") == "page"
            and not str(t.get("url", "")).startswith(("devtools://", "chrome-extension://"))
        ]
        # Chrome lists the most recently active target first.
        return pages[0] if pages else None

    def viewport(self) -> dict:
        return dict(self._viewport)

    def watched(self) -> bool:
        """True while the site's Browser tab is (or was moments ago) streaming
        this browser — a page reload must not open a bring-to-front window."""
        return bool(self._viewers) or time.time() - self._last_viewer_at < 30.0

    async def state(self) -> dict:
        import aiohttp

        cdp_http = None
        if POOL._browser is not None and _browser_is_alive(POOL._browser):
            cdp_http = self._normalize_cdp(getattr(POOL._browser, "cdp_url", None))
        if not cdp_http:
            return {"ok": True, "browser": "asleep", "viewport": self.viewport()}
        async with aiohttp.ClientSession() as session:
            page = await self._active_page(session, cdp_http)
        if not page:
            return {"ok": True, "browser": "starting", "viewport": self.viewport()}
        return {
            "ok": True,
            "browser": "up",
            "url": page.get("url"),
            "title": page.get("title"),
            "viewport": self.viewport(),
        }

    async def _close_input_conn(self) -> None:
        ws, self._input_ws = self._input_ws, None
        session, self._input_session = self._input_session, None
        self._input_page_id = None
        for closable in (ws, session):
            if closable is not None:
                try:
                    await closable.close()
                except Exception:
                    pass

    async def _ensure_input_conn(self) -> Any:
        """Cached CDP socket to the ACTIVE page; re-resolved at most every 2s."""
        import aiohttp

        now = time.time()
        if self._input_ws is not None and not self._input_ws.closed and now - self._input_checked_at < 2.0:
            return self._input_ws
        cdp_http = await self._cdp_http()
        if not cdp_http:
            await self._close_input_conn()
            return None
        if self._input_session is None:
            self._input_session = aiohttp.ClientSession()
        page = await self._active_page(self._input_session, cdp_http)
        if not page or "webSocketDebuggerUrl" not in page:
            await self._close_input_conn()
            return None
        self._input_checked_at = now
        if self._input_ws is not None and not self._input_ws.closed and page.get("id") == self._input_page_id:
            return self._input_ws
        if self._input_ws is not None:
            try:
                await self._input_ws.close()
            except Exception:
                pass
        self._input_ws = await self._input_session.ws_connect(
            page["webSocketDebuggerUrl"], max_msg_size=8 * 1024 * 1024
        )
        self._input_page_id = page.get("id")
        return self._input_ws

    async def input(self, event: dict) -> dict:
        async with self._input_lock:
            try:
                return await self._dispatch_input(event)
            except Exception:
                # One reconnect covers a page swap or a dropped socket.
                await self._close_input_conn()
                try:
                    return await self._dispatch_input(event)
                except Exception as error:  # noqa: BLE001
                    await self._close_input_conn()
                    return {"ok": False, "error": str(error)}

    async def _dispatch_input(self, event: dict) -> dict:
        ws = await self._ensure_input_conn()
        if ws is None:
            return {"ok": False, "error": "no page"}

        async def send(method: str, params: dict | None = None) -> None:
            self._input_msg_id += 1
            await ws.send_json({"id": self._input_msg_id, "method": method, "params": params or {}})

        kind = event.get("kind")
        x = float(event.get("x") or 0)
        y = float(event.get("y") or 0)
        if kind == "click":
            button = event.get("button") or "left"
            count = 2 if event.get("double") else 1
            base = {"x": x, "y": y, "button": button, "clickCount": count}
            await send("Input.dispatchMouseEvent", {"type": "mousePressed", **base})
            await send("Input.dispatchMouseEvent", {"type": "mouseReleased", **base})
        elif kind == "down":
            button = event.get("button") or "left"
            count = 2 if event.get("double") else 1
            await send(
                "Input.dispatchMouseEvent",
                {"type": "mousePressed", "x": x, "y": y, "button": button, "clickCount": count},
            )
        elif kind == "up":
            button = event.get("button") or "left"
            await send(
                "Input.dispatchMouseEvent",
                {"type": "mouseReleased", "x": x, "y": y, "button": button, "clickCount": 1},
            )
        elif kind == "move":
            # buttons bitmask (1=left, 2=right, 4=middle) makes a move a DRAG —
            # text selection, sliders, drag-and-drop all live here. Without it
            # the move is a hover, which is just as load-bearing: menus,
            # tooltips, and :hover states only exist if the pointer moves.
            params: dict = {"type": "mouseMoved", "x": x, "y": y}
            buttons = int(event.get("buttons") or 0)
            if buttons:
                params["buttons"] = buttons
                params["button"] = "left" if buttons & 1 else "right" if buttons & 2 else "middle"
            await send("Input.dispatchMouseEvent", params)
        elif kind == "scroll":
            await send(
                "Input.dispatchMouseEvent",
                {"type": "mouseWheel", "x": x, "y": y, "deltaX": float(event.get("deltaX") or 0), "deltaY": float(event.get("deltaY") or 0)},
            )
        elif kind == "type":
            await send("Input.insertText", {"text": str(event.get("text") or "")[:2000]})
        elif kind == "key":
            key = str(event.get("key") or "")
            # The few keys a page needs beyond typed text.
            codes = {
                "Enter": ("Enter", 13, "\r"),
                "Backspace": ("Backspace", 8, ""),
                "Tab": ("Tab", 9, ""),
                "Escape": ("Escape", 27, ""),
                "Delete": ("Delete", 46, ""),
                "ArrowUp": ("ArrowUp", 38, ""),
                "ArrowDown": ("ArrowDown", 40, ""),
                "ArrowLeft": ("ArrowLeft", 37, ""),
                "ArrowRight": ("ArrowRight", 39, ""),
                "PageUp": ("PageUp", 33, ""),
                "PageDown": ("PageDown", 34, ""),
            }
            modifier_bits = {"alt": 1, "ctrl": 2, "control": 2, "meta": 4, "cmd": 4, "shift": 8}
            modifiers = 0
            for name in event.get("modifiers") or []:
                modifiers |= modifier_bits.get(str(name).lower(), 0)
            if key in codes:
                code, vk, text = codes[key]
                down = {"type": "rawKeyDown", "key": key, "code": code, "windowsVirtualKeyCode": vk, "nativeVirtualKeyCode": vk, "modifiers": modifiers}
                if text and not modifiers:
                    down = {**down, "type": "keyDown", "text": text, "unmodifiedText": text}
                await send("Input.dispatchKeyEvent", down)
                await send("Input.dispatchKeyEvent", {"type": "keyUp", "key": key, "code": code, "windowsVirtualKeyCode": vk, "nativeVirtualKeyCode": vk, "modifiers": modifiers})
            elif len(key) == 1 and modifiers:
                # Shortcuts: ctrl/cmd+A/C/V/X/Z and friends. The page receives
                # the combination exactly as a physical keyboard would send it.
                upper = key.upper()
                code = f"Key{upper}" if upper.isalpha() else f"Digit{upper}" if upper.isdigit() else upper
                vk = ord(upper) if (upper.isalpha() or upper.isdigit()) else 0
                base = {"key": key.lower(), "code": code, "windowsVirtualKeyCode": vk, "nativeVirtualKeyCode": vk, "modifiers": modifiers}
                await send("Input.dispatchKeyEvent", {"type": "rawKeyDown", **base})
                await send("Input.dispatchKeyEvent", {"type": "keyUp", **base})
        elif kind == "navigate":
            url = str(event.get("url") or "").strip()
            if url and url.startswith(("http://", "https://")):
                await send("Page.navigate", {"url": url})
        elif kind == "viewport":
            # Make the page's viewport MATCH the embed's on-page box,
            # so the stream fills the surface with no letterboxing.
            width = int(max(320, min(2560, float(event.get("width") or 0))))
            height = int(max(240, min(1600, float(event.get("height") or 0))))
            await send(
                "Emulation.setDeviceMetricsOverride",
                {"width": width, "height": height, "deviceScaleFactor": 1, "mobile": False},
            )
            self._viewport = {"width": width, "height": height}
        else:
            return {"ok": False, "error": f"unknown kind {kind!r}"}
        POOL._last_used = time.time()
        return {"ok": True}

    async def add_viewer(self, res: Any) -> None:
        self._viewers.add(res)
        self._last_viewer_at = time.time()
        if self._pump_task is None or self._pump_task.done():
            self._pump_task = asyncio.create_task(self._pump())

    def drop_viewer(self, res: Any) -> None:
        self._viewers.discard(res)
        self._last_viewer_at = time.time()

    async def _pump(self) -> None:
        import base64

        import aiohttp

        try:
            print("[browser-view] pump started", flush=True)
            while self._viewers:
                cdp_http = await self._cdp_http()
                if not cdp_http:
                    print("[browser-view] no cdp endpoint yet", flush=True)
                    await asyncio.sleep(1.0)
                    continue
                async with aiohttp.ClientSession() as session:
                    page = await self._active_page(session, cdp_http)
                    if not page or "webSocketDebuggerUrl" not in page:
                        print(f"[browser-view] no page target at {cdp_http}", flush=True)
                        await asyncio.sleep(1.0)
                        continue
                    page_id = page.get("id")
                    print(f"[browser-view] attaching to {str(page.get('url'))[:80]}", flush=True)
                    try:
                        async with session.ws_connect(page["webSocketDebuggerUrl"], max_msg_size=64 * 1024 * 1024) as ws:
                            msg_id = 0
                            send_lock = asyncio.Lock()
                            # ids ≥ this are one-shot screenshot requests — the
                            # fallback for the frames Chrome refuses to composite
                            # (an occluded or minimised window screencasts
                            # NOTHING; a captureScreenshot still renders).
                            SHOT_BASE = 1_000_000
                            shot_id = SHOT_BASE
                            shot_pending = False
                            last_frame = 0.0

                            async def send(method: str, params: dict | None = None, explicit_id: int | None = None) -> None:
                                nonlocal msg_id
                                async with send_lock:
                                    if explicit_id is None:
                                        msg_id += 1
                                    await ws.send_json({"id": explicit_id if explicit_id is not None else msg_id, "method": method, "params": params or {}})

                            async def broadcast(jpeg: bytes) -> None:
                                header = (
                                    b"--skykoiframe\r\nContent-Type: image/jpeg\r\nContent-Length: "
                                    + str(len(jpeg)).encode()
                                    + b"\r\n\r\n"
                                )
                                dead = []
                                for viewer in list(self._viewers):
                                    try:
                                        await viewer.write(header + jpeg + b"\r\n")
                                    except Exception:
                                        dead.append(viewer)
                                for viewer in dead:
                                    self._viewers.discard(viewer)

                            await send("Page.enable")
                            await send(
                                "Page.startScreencast",
                                {"format": "jpeg", "quality": 60, "maxWidth": VIEW_MAX_W, "maxHeight": VIEW_MAX_H, "everyNthFrame": 1},
                            )
                            while self._viewers:
                                try:
                                    msg = await ws.receive(timeout=1.5)
                                except asyncio.TimeoutError:
                                    # No frames flowing. Ask for one directly —
                                    # captureScreenshot renders even when the
                                    # compositor is idle — and check whether
                                    # the agent moved to a different tab.
                                    if not shot_pending and time.time() - last_frame > 1.5:
                                        shot_pending = True
                                        shot_id += 1
                                        await send("Page.captureScreenshot", {"format": "jpeg", "quality": 60}, explicit_id=shot_id)
                                    fresh = await self._active_page(session, cdp_http)
                                    if fresh and fresh.get("id") != page_id:
                                        break
                                    continue
                                if msg.type in (aiohttp.WSMsgType.CLOSED, aiohttp.WSMsgType.CLOSING, aiohttp.WSMsgType.ERROR):
                                    break
                                if msg.type != aiohttp.WSMsgType.TEXT:
                                    continue
                                data = json.loads(msg.data)
                                if data.get("error"):
                                    if data.get("id", 0) >= SHOT_BASE:
                                        shot_pending = False
                                    print(f"[browser-view] CDP error: {data['error']}", flush=True)
                                    continue
                                if data.get("id", 0) >= SHOT_BASE:
                                    shot_pending = False
                                    shot = data.get("result", {}).get("data", "")
                                    if shot:
                                        last_frame = time.time()
                                        await broadcast(base64.b64decode(shot))
                                    continue
                                if data.get("method") != "Page.screencastFrame":
                                    continue
                                params = data.get("params", {})
                                await send("Page.screencastFrameAck", {"sessionId": params.get("sessionId")})
                                meta = params.get("metadata", {})
                                if meta.get("deviceWidth") and meta.get("deviceHeight"):
                                    self._viewport = {"width": meta["deviceWidth"], "height": meta["deviceHeight"]}
                                last_frame = time.time()
                                await broadcast(base64.b64decode(params.get("data", "")))
                    except Exception as error:  # noqa: BLE001
                        print(f"[browser-view] screencast attach failed: {error!r}", flush=True)
                        await asyncio.sleep(1.0)
        except Exception as error:  # noqa: BLE001
            print(f"[browser-view] pump died: {error!r}", flush=True)
        finally:
            self._pump_task = None


VIEW = BrowserView()


async def handle_view_stream(request: web.Request) -> web.StreamResponse:
    if not _authorized(request):
        return web.json_response({"error": "unauthorized"}, status=401)
    res = web.StreamResponse(
        headers={
            "Content-Type": "multipart/x-mixed-replace; boundary=skykoiframe",
            "Cache-Control": "no-store",
        },
    )
    await res.prepare(request)
    await VIEW.add_viewer(res)
    POOL._last_used = time.time()
    try:
        # Hold the response open; aiohttp cancels this handler when the
        # watcher disconnects, which is the only exit we need.
        await asyncio.Event().wait()
    except (asyncio.CancelledError, ConnectionResetError):
        pass
    finally:
        VIEW.drop_viewer(res)
    return res


async def handle_view_input(request: web.Request) -> web.Response:
    if not _authorized(request):
        return web.json_response({"error": "unauthorized"}, status=401)
    try:
        event = await request.json()
    except Exception:
        return web.json_response({"error": "expected JSON"}, status=400)
    result = await VIEW.input(event if isinstance(event, dict) else {})
    return web.json_response(result, status=200 if result.get("ok") else 502)


async def handle_view_state(request: web.Request) -> web.Response:
    if not _authorized(request):
        return web.json_response({"error": "unauthorized"}, status=401)
    return web.json_response(await VIEW.state())

def sweep_abandoned_profiles() -> None:
    """
    Close and delete every profile this service walked away from.

    Each quarantine minted a new directory and left the old one — with its
    browser still running — behind forever. Measured on the owner's machine
    before this existed: 21 directories and 79 live Edge processes, several
    gigabytes of disk, and windows reappearing from sessions nobody remembered
    starting. Nothing here is recoverable state: the live profile is the only
    one with the user's logins in it, and it is the one directory skipped.
    """
    parent = os.path.dirname(PROFILE_DIR) or "."
    keep = os.path.basename(PROFILE_DIR)
    base = "browser-agent-profile"
    try:
        names = os.listdir(parent)
    except OSError:
        return
    closed = 0
    removed = 0
    for name in names:
        if not name.startswith(base) or name == keep:
            continue
        path = os.path.join(parent, name)
        if not os.path.isdir(path):
            continue
        BrowserPool._kill_profile_processes(path)
        closed += 1
        try:
            shutil.rmtree(path, ignore_errors=True)
            removed += 1
        except OSError:
            pass
    if closed:
        print(f"[browser-agent] cleaned up {closed} abandoned profile(s), removed {removed}", flush=True)



# ── running a task ───────────────────────────────────────────────────────────


async def _bring_window_to_front() -> None:
    """
    Raise the agent's browser window when a task starts.

    Headful is the point of this service, and yet on a real call (2026-08-19,
    "find IKEA's cheapest item") the task ran, verified, and the user saw
    NOTHING — indistinguishable from headless. Two things conspire:

    1. The warm-browser design opens the window at BOOT (prewarm), hours before
       anyone asks for anything, so no window appears when a task starts.
    2. browser-use launches with --disable-window-activation and
       --disable-focus-on-load, so the window is never activated — it sits at
       the bottom of the Z-order under whatever the user is actually looking
       at (on a call: the call screen). --start-maximized only mattered in the
       second it launched.

    So the window is raised HERE, at every task start, where "watch it work"
    actually needs it. SetForegroundWindow from a background process is often
    refused by Windows' foreground lock; the minimize/restore fallback forces
    the activation anyway. Best-effort and fire-and-forget: a task must never
    fail, or even wait, because a window could not be raised.
    """
    if HEADLESS or sys.platform != "win32":
        return
    name = os.path.basename(PROFILE_DIR)
    if not name:
        return
    script = (
        "$ErrorActionPreference='SilentlyContinue';"
        "Add-Type -Namespace SkyKoi -Name Win32 -MemberDefinition "
        "'[DllImport(\"user32.dll\")] public static extern bool SetForegroundWindow(IntPtr h);"
        "[DllImport(\"user32.dll\")] public static extern bool ShowWindowAsync(IntPtr h, int c);"
        "[DllImport(\"user32.dll\")] public static extern bool IsIconic(IntPtr h);';"
        # Match on the profile directory in the command line, same as every
        # other place this file finds its own browser among the user's.
        f"$owners = Get-CimInstance Win32_Process -Filter \"Name='msedge.exe' or Name='chrome.exe'\" | Where-Object {{ $_.CommandLine -like '*{name}*' }};"
        "foreach ($o in $owners) {"
        " $h = (Get-Process -Id $o.ProcessId -ErrorAction SilentlyContinue).MainWindowHandle;"
        " if ($h -and $h -ne 0) {"
        "  if ([SkyKoi.Win32]::IsIconic($h)) { [SkyKoi.Win32]::ShowWindowAsync($h, 9) | Out-Null };"
        "  if (-not [SkyKoi.Win32]::SetForegroundWindow($h)) {"
        # 6 = SW_MINIMIZE, 9 = SW_RESTORE: the restore of a just-minimized
        # window is the one activation Windows never refuses, and SW_RESTORE
        # returns a maximized window to maximized, not to a small one.
        "   [SkyKoi.Win32]::ShowWindowAsync($h, 6) | Out-Null;"
        "   Start-Sleep -Milliseconds 150;"
        "   [SkyKoi.Win32]::ShowWindowAsync($h, 9) | Out-Null;"
        "   [SkyKoi.Win32]::SetForegroundWindow($h) | Out-Null"
        "  };"
        "  break"
        " }"
        "}"
    )
    try:
        proc = await asyncio.create_subprocess_exec(
            "powershell", "-NoProfile", "-NonInteractive", "-Command", script,
            stdout=asyncio.subprocess.DEVNULL, stderr=asyncio.subprocess.DEVNULL,
        )
        await asyncio.wait_for(proc.communicate(), timeout=10)
    except (asyncio.TimeoutError, OSError):
        pass


def _observable_end_state(history: Any) -> dict[str, Any]:
    """What is actually true at the end, independent of what the agent claims."""
    urls = [u for u in (history.urls() or []) if u]
    titles = [t for t in (history.titles() or []) if t] if hasattr(history, "titles") else []
    return {
        "url": urls[-1] if urls else None,
        "title": titles[-1] if titles else None,
        "visitedUrls": list(dict.fromkeys(urls))[-6:],
        "steps": len(history.history),
    }


def _verify(task: str, claimed: str | None, end: dict[str, Any]) -> tuple[bool, str | None]:
    """
    Does the end state support the claim?

    Deliberately cheap and mechanical — this is not a second model second-
    guessing the first, it is the difference between "the agent said it opened
    the towels page" and "the browser is on a page that looks like it".

    The HOST IS IGNORED, which the first version of this got wrong: asked for
    IKEA bath towels it ended on ikea.com/us/en/cat/wall-decor and passed,
    because "ikea" appears in every URL on the site. Only the path, the query
    and the page title can testify to what was actually reached.
    """
    url = (end.get("url") or "").lower()
    if not url:
        return False, "the browser did not end on any page"
    after_host = url.split("//", 1)[-1]
    host = after_host.split("/", 1)[0]
    # A search engine is where you go to LOOK for something, never the thing
    # itself — and its URL echoes the task back verbatim in ?q=, so it matches
    # every keyword and "verifies" anything. Caught asking for a dashboard that
    # does not exist: it ended on duckduckgo with the task in the query string
    # and the check passed.
    if any(engine in host for engine in ("google.", "bing.", "duckduckgo.", "ecosia.", "yahoo.", "yandex.")):
        return False, f"ended on a search results page ({end.get('url')}), not on the page itself"

    # THE SITE ASKED FOR IS PART OF THE TASK. Asked for en.wikipedia.org it
    # answered from simple.wikipedia.org and this passed, because the host was
    # deliberately ignored when matching keywords. Ignoring the host for
    # KEYWORDS is right; ignoring which SITE was asked for is not — a price
    # from a different shop is a different answer.
    named = re.findall(r"\b(?:[a-z0-9-]+\.)+[a-z]{2,}\b", task.lower())
    named = [n for n in named if not n.endswith((".jpg", ".png", ".html"))]
    if named and not any(host == n or host.endswith("." + n) or n.endswith("." + host) for n in named):
        return False, f"asked for {named[0]} but ended on {host}"
    evidence = (after_host.split("/", 1)[1] if "/" in after_host else "") + " " + (end.get("title") or "").lower()
    # READING is not NAVIGATING. "report the heading on example.com" is finished
    # when there is an answer, and no URL will ever contain the word "heading" —
    # demanding page evidence for it marks a correct result unverified, which is
    # its own kind of dishonesty. For these, being on the named site with an
    # answer in hand is the evidence.
    read_verbs = ("report", "read", "tell", "what", "summar", "extract", "check", "how many", "price", "list")
    if any(v in task.lower() for v in read_verbs) and (claimed or "").strip():
        # The scheme is not part of the host. "https://news.ycombinator.com"
        # used to split to "https:", which is in no host on earth, so every
        # read task that named its site by URL came back unverified — a
        # correct answer marked as a claim.
        named = [
            w.split("//")[-1].split("/")[0].strip(".,")
            for w in task.lower().split()
            if "." in w and len(w) > 3
        ]
        named = [n for n in named if n]
        if not named or any(n in host or host in n for n in named):
            return True, None
    stop = {
        "the", "a", "an", "to", "on", "go", "open", "and", "for", "of", "page", "please", "then",
        "find", "show", "me", "my", "it", "that", "this", "with", "from", "report", "final", "url",
        "get", "into", "at", "in", "is", "are", "up", "search", "look", "can", "you", "website",
        "com", "www", "listing", "site", "new", "tab", "window", "then", "both", "closest", "equivalent",
    }
    words = [w for w in "".join(c if c.isalnum() else " " for c in task.lower()).split() if len(w) > 2 and w not in stop]
    # Words that name the destination rather than the site it lives on.
    words = [w for w in words if w not in after_host.split("/", 1)[0]]
    if not words:
        return True, None
    hits = [w for w in words if w in evidence]
    if hits:
        return True, None
    return False, f"ended on {end.get('url')}, which shows no sign of: {' '.join(words[:6])}"


async def _attempt(
    task: str, max_steps: int, on_event, attempt: int, timing: Timing, new_tab: bool = False
) -> tuple[dict[str, Any], Any]:
    from browser_use import Agent

    browser = await POOL.acquire()
    # The user watches the browser work — that is the contract of headful mode,
    # and the warm window has been buried since boot. Raised in the background
    # so the first LLM call is not waiting on a PowerShell process.
    # EXCEPT when the embedded Browser tab is watching: they are seeing the
    # page through the site, and yanking the device's Chrome to the front on
    # top of whatever they are doing ruins the very view they are using.
    if not VIEW.watched():
        asyncio.ensure_future(_bring_window_to_front())
    if new_tab:
        # DETERMINISTIC, not a suggestion. Asked in the task text to "open a new
        # tab", the agent navigated the tab it was already in — watched live.
        # Anything the caller states as a requirement should not depend on the
        # model choosing to honour it, so the tab is opened here and the agent
        # simply starts in it.
        before = len(await browser.get_tabs())
        # navigate_to(new_tab=True), NOT new_page(): new_page returned without
        # adding a tab (measured — the count stayed at 1), so the agent carried
        # on in the tab it was already in while its final report claimed
        # otherwise. browser-use's own judge caught the lie in the log; the user
        # caught it watching the screen.
        await browser.navigate_to("about:blank", new_tab=True)
        after = len(await browser.get_tabs())
        opened = after > before
        await on_event({"kind": "tab", "opened": opened, "tabs": after})
        if not opened:
            # Say so rather than letting the agent narrate a tab that is not
            # there. A requirement that silently did not happen is the whole
            # failure mode this service exists to stop.
            await on_event({"kind": "warning", "text": "could not open a new tab; continuing in the current one"})
    step_no = {"n": 0}

    async def on_step(*args: Any, **kwargs: Any) -> None:
        step_no["n"] += 1
        url = None
        try:
            url = await browser.get_current_page_url()
        except Exception:
            pass
        await on_event({"kind": "step", "n": step_no["n"], "attempt": attempt, "url": url})



    stall = {"url": None, "n": 0}

    # A way to click what the element index cannot see. See click_text.py: the
    # colour swatch the agent fought with for 18 steps is a hidden checkbox
    # behind a styled label, and no amount of retrying an index it does not
    # have was going to work.
    from browser_use import Tools

    from click_text import register_click_text, register_read_table

    tools = Tools()
    register_click_text(tools, lambda: browser)
    register_read_table(tools, lambda: browser)

    agent = Agent(
        task=task,
        llm=_llm(timing),
        browser=browser,
        # TOOL SELECTION, not task instructions — the task text belongs to the
        # caller. Traced on a Wikipedia table: 30 steps and 292 seconds spent
        # scrolling a long page up and down hunting for one cell, when `extract`
        # and `find_text` were sitting right there. Nothing was broken; the model
        # simply reached for the wrong tool, and reading a page by dragging a
        # viewport over it cannot work without vision.
        extend_system_message=(
            "READING vs ACTING. To read anything off a page - a table cell, a price, a list, "
            "an article - use extract or find_text or search_page. Do NOT scroll looking for it: "
            "you cannot see the page, only the elements listed for you, and scrolling a long "
            "document to find a value wastes the whole step budget. Scroll only to bring a "
            "CONTROL you intend to click into view."
        ),
        tools=tools,
        # Off for Groq: it rejects array-shaped content. The element index is
        # what the agent steers by in any case.
        use_vision=False,
        # Four actions a step roughly halves wall-clock without the accuracy
        # loss flash_mode brings (flash "finished" the IKEA task on the
        # homepage in 17s).
        max_actions_per_step=4,
        # THE PROMPT IS THE COST. history is unbounded by default, so every
        # step carries every previous step and the model call gets steadily
        # slower: measured 1.6s median on a 2-step task against 2.5s on a
        # 40-step one, same model, same machine. Six is enough for the agent
        # to know what it just tried and not repeat itself.
        max_history_items=6,
        # browser-use's own chain-of-thought field, on top of the model's.
        # Picking an element index does not need it, and it is output tokens
        # on every single step.
        use_thinking=False,
        register_new_step_callback=on_step,
    )
    history = await agent.run(max_steps=max_steps)
    return _observable_end_state(history), history


async def run_task(task: str, max_steps: int, on_event, new_tab: bool = False) -> dict[str, Any]:
    """
    Run the task, and if the end state does not back up the claim, run it once
    more knowing where it ended up.

    The retry is the point. A browser agent that wanders onto the wrong page is
    normal; one that then TELLS the user it succeeded is what made the koi
    untrustworthy. Measured on this machine, the same IKEA task landed on the
    towels listing on one run and on wall decor on another — the difference
    between those two has to be caught here, not by the person listening.
    """
    started = time.time()
    global _task_running
    async with POOL.lock:
        # Marks the window where a closed tab really is a crash to recover from
        # rather than a person shutting a browser they are done with.
        _task_running = True
        try:
            timing = Timing()
            end, history = await _attempt(task, max_steps, on_event, 1, timing, new_tab)
        finally:
            _task_running = False
        claimed = history.final_result()
        verified, why = _verify(task, claimed, end)

        finished = bool(history.is_successful())
        # Retry when the attempt did not FINISH, not only when it cannot be
        # verified. An aborted stall lands on a page that often still looks
        # plausible — sorted results, right site — and the first version of this
        # returned that as the answer to a task that was never completed.
        if not verified or not finished:
            await on_event(
                {"kind": "retry", "reason": why or "did not finish"}
            )
            stuck = f"A previous attempt ended on {end.get('url')} without getting there. "
            retry_task = (
                f"{task}\n\n"
                + stuck
                + "Do not repeat that route — in particular do not fight with filter dialogs or dropdown "
                "menus, which is where the last attempt got stuck. Prefer a direct search, or simply read "
                "the results already on the page, and finish on the page that was asked for."
            )
            end2, history2 = await _attempt(retry_task, max_steps, on_event, 2, timing)
            claimed2 = history2.final_result()
            verified2, why2 = _verify(task, claimed2, end2)
            if verified2 or not verified:
                end, claimed, verified, why = end2, claimed2, verified2, why2
                history = history2

    return {
        "kind": "done",
        # `ok` requires BOTH the agent's own verdict and observable evidence.
        "ok": bool(history.is_successful()) and verified,
        "verified": verified,
        "unverifiedReason": why,
        "claimed": claimed,
        **end,
        "ms": int((time.time() - started) * 1000),
        "timing": timing.summary(history),
    }


# ── http ─────────────────────────────────────────────────────────────────────


def _authorized(request: web.Request) -> bool:
    if not TOKEN:
        return True
    return request.headers.get("authorization", "") == f"Bearer {TOKEN}"


async def handle_run(request: web.Request) -> web.StreamResponse:
    if not _authorized(request):
        return web.json_response({"error": "unauthorized"}, status=401)
    try:
        body = await request.json()
    except Exception:
        return web.json_response({"error": "expected JSON"}, status=400)
    task = str(body.get("task") or "").strip()
    if not task:
        return web.json_response({"error": "expected { task }"}, status=400)
    max_steps = int(body.get("maxSteps") or MAX_STEPS_DEFAULT)
    # Honour an explicit flag, and also the plain-English request — a caller who
    # writes "open a new tab and…" means it.
    new_tab = bool(body.get("newTab")) or bool(
        # NOT r"new tab" written through a shell heredoc: that wrote a
        # literal backspace () into the pattern, which matches nothing,
        # so the tab never opened while the agent cheerfully said it had.
        re.search(r"new\s+tab|another\s+tab|new\s+window", task, re.I)
    )
    print(f"[browser-agent] newTab={new_tab}", flush=True)

    response = web.StreamResponse(
        headers={"Content-Type": "application/x-ndjson", "Cache-Control": "no-store"}
    )
    await response.prepare(request)

    async def emit(event: dict[str, Any]) -> None:
        try:
            await response.write((json.dumps(event) + "\n").encode())
        except Exception:
            pass  # the caller hung up; the task still finishes and the browser stays warm

    try:
        result = await run_task(task, max_steps, emit, new_tab)
    except Exception as error:  # noqa: BLE001 - the caller gets the reason, not a hang
        result = {"kind": "done", "ok": False, "verified": False, "error": str(error)}
        traceback.print_exc()
    await emit(result)
    await response.write_eof()
    return response


async def handle_health(request: web.Request) -> web.Response:
    return web.json_response(
        {
            "ok": True,
            "model": MODEL,
            # False here is the exact state that looks like "the browser never
            # opened" from the outside: /health fine, every task dead on step 0.
            "hasModelKey": bool(
                os.environ.get("WAFER_API_KEY" if MODEL.startswith("wafer/") else "GROQ_API_KEY")
            ),
            "headless": HEADLESS,
            "profile": PROFILE_DIR,
            "browserWarm": POOL._browser is not None,
        }
    )


async def handle_close(request: web.Request) -> web.Response:
    if not _authorized(request):
        return web.json_response({"error": "unauthorized"}, status=401)
    await POOL.close()
    return web.json_response({"ok": True})


async def _serve() -> None:
    app = web.Application(client_max_size=2 * 1024 * 1024)
    app.add_routes(
        [
            web.post("/run", handle_run),
            web.get("/health", handle_health),
            web.post("/close", handle_close),
            web.get("/view/stream", handle_view_stream),
            web.post("/view/input", handle_view_input),
            web.get("/view/state", handle_view_state),
        ]
    )
    runner = web.AppRunner(app)
    await runner.setup()
    site = web.TCPSite(runner, HOST, PORT)
    await site.start()
    reaper = asyncio.create_task(POOL.reap_idle())
    _let_the_user_close_it()
    sweep_abandoned_profiles()

    async def prewarm() -> None:
        """
        Chrome's launch is ~3s that no task should ever pay.

        Headful included, which reverses the old choice on the owner's call:
        "whenever it opens the browser it needs to be instantaneous". The
        window appears once, when the service boots — and closing it is
        respected (the recovery patch above leaves a user-closed idle browser
        closed), so it cannot read as the machine fighting them the way the
        old every-task respawn did. The cost of the alternative was real: the
        first browser task of every session stalled mid-request on a Chrome
        launch, which is exactly where a person is watching the clock.
        """
        try:
            async with POOL.lock:
                await POOL.acquire()
            print("[browser-agent] browser warm", flush=True)
        except Exception as error:  # noqa: BLE001
            print(f"[browser-agent] prewarm failed: {error}", flush=True)

    warm = asyncio.create_task(prewarm())
    try:
        await asyncio.Event().wait()
    finally:
        reaper.cancel()
        warm.cancel()
        await POOL.close()
        await runner.cleanup()


def main() -> None:
    """
    Served from asyncio.run, NOT web.run_app.

    On Windows the loop has to be the Proactor one or Chrome never launches:
    asyncio's Selector loop cannot spawn subprocesses, and browser_use's start
    then hangs until a 30s watchdog fires with an error about "event handler
    timed out" that says nothing about the cause. The same code in a plain
    asyncio.run script worked the whole time, which is what gave it away.
    asyncio.run uses the default policy, which is Proactor on Windows.
    """
    if sys.platform == "win32":
        asyncio.set_event_loop_policy(asyncio.WindowsProactorEventLoopPolicy())
    try:
        asyncio.run(_serve())
    except KeyboardInterrupt:
        pass


if __name__ == "__main__":
    main()
