"""One isolated Botasaurus @browser/driver.get per frozen URL; private captures.
Use exact requirements.txt, documentation.json and an existing matching Chrome binary.
"""
import argparse
import asyncio
from collections import Counter
from datetime import datetime, timezone
import hashlib
from html.parser import HTMLParser
from importlib.metadata import version
import json
import os
from pathlib import Path
import platform
import re
import shlex
import signal
import subprocess
import sys
import time
import tempfile
from urllib.parse import urlsplit
from botasaurus.browser import browser, Driver, cdp
from botasaurus_driver.core.env import is_docker
from contextlib import contextmanager
import base64

CONFIG = {
 "proxy_mode": "none", "login": False, "captcha_solving": False,
 "concurrency": 1, "gap_seconds": 2, "application_retries": 0,
 "mode": "@browser -> driver.get", "headless": True, "fresh_browser_per_target": True,
 "profile": None, "tiny_profile": False, "reuse_driver": False,
 "cache": False, "output": None, "bypass_cloudflare": False, "human_mode": False,
 "navigation_timeout_seconds": 20, "attempt_timeout_seconds": 40,
 "slice_timeout_seconds": 1260, "wait_for_complete_page_load": False, "settle_seconds": 2,
 "max_capture_bytes": 5242880, "network_response_byte_cap": None,
 "redirects": "follow", "max_redirects": 20,
 "redirect_cap_basis": "Same pinned Chrome 153 binary as Patchright; Chromium native URLRequest cap 20. HTTP redirects and other URL changes recorded separately.",
 "resources": "Normal Chrome JavaScript/images/styles/iframes/workers; block_images=False, block_images_and_css=False; downloads denied via CDP",
 "cookies": "Fresh temporary Chrome profile per worker; no imported cookies/login; ordinary session cookies possible and discarded by driver teardown",
 "fingerprint": "Default driver headless UA removes Headless via Network.setUserAgentOverride. user_agent=None, window_size=None, lang=None; no custom fingerprint, locale/timezone overrides or JS injection. Driver default browser flags retained.",
 "launch_args": ["--no-proxy-server"], "pss_sample_interval_seconds": 0.1,
 "pss_scope": "Isolated Python worker and all descendants, including Botasaurus driver/Chrome startup, DOM capture, viewport screenshot and teardown; excludes batch parent",
 "min_host_available_mb": 3072, "max_attempt_tree_pss_mb": 4096,
 "startup_retry_policy": "Driver 4.0.101 retries Chrome startup up to 3 total launches on startup errors; not target navigation retries. @browser max_retry=0; one driver.get only.",
 "cleanup_policy": "Short private mkdtemp directory under /tmp per worker avoids Chrome Unix socket path limits and scopes driver profile garbage collection to owned folders; refuse Docker global zombie cleanup mode. Only recorded PID/start-time child identities are signalled.",
}

BLOCK_PHRASES = (
    "verify you are human", "robot or human", "robot check", "access denied",
    "pardon our interruption", "please complete the captcha", "enter the characters",
    "additional verification required", "unusual traffic", "checking your browser",
    "enable javascript and cookies to continue", "just a moment",
)
GATE_PHRASES = ("log in to continue", "sign in to continue", "sign in to linkedin", "consent required", "before you continue")


def utc():
    return datetime.now(timezone.utc).isoformat()


def sha(data):
    return hashlib.sha256(data).hexdigest()


class PageText(HTMLParser):
    """Source-text approximation, not rendered visibility or structured extraction."""
    SKIP = {"head", "script", "style", "template", "noscript"}

    def __init__(self):
        super().__init__(convert_charrefs=True)
        self.hidden = []
        self.parts = []

    def handle_starttag(self, tag, attrs):
        if tag in self.SKIP:
            self.hidden.append(tag)

    def handle_endtag(self, tag):
        if tag in self.hidden:
            self.hidden = self.hidden[:self.hidden.index(tag)]

    def handle_data(self, data):
        if not self.hidden:
            self.parts.append(data)


def evaluate(body, target):
    parser = PageText()
    parser.feed(body.decode("utf-8", errors="replace"))
    text = " ".join(" ".join(parser.parts).split())
    lower = text.lower()
    checks = [
        {"alternatives": group, "matched": [m for m in group if m.lower() in lower]}
        for group in target["required_text_groups"]
    ]
    regex_checks = [
        {**check, "matches": len(re.findall(check["pattern"], text, flags=re.I))}
        for check in target["text_regex_checks"]
    ]
    expected = (
        len(text) >= target["min_text_chars"]
        and all(check["matched"] for check in checks)
        and all(check["matches"] >= check["min_matches"] for check in regex_checks)
    )
    return {
        "text_chars": len(text), "required_text_checks": checks, "text_regex_checks": regex_checks,
        "min_text_chars": target["min_text_chars"], "has_required_content": bool(expected),
        "soft_block_checks": {"phrases_checked": list(BLOCK_PHRASES), "matched": [p for p in BLOCK_PHRASES if p in lower]},
        "gate_checks": {"phrases_checked": list(GATE_PHRASES), "matched": [p for p in GATE_PHRASES if p in lower]},
    }



def result_class(row):
    if row["error"]:
        return row["error"]["classification"]
    if row["navigation_error"]:
        return "navigation_error"
    if row["capture_error"]:
        return "capture_error"
    if row["status"] is None:
        return "status_unobserved"
    if row["status"] in (401, 403, 999):
        return "http_access_denied"
    if row["status"] == 429:
        return "http_rate_limited"
    if row["status"] >= 500:
        return "http_server_error"
    if row["status"] >= 400:
        return "http_client_error"
    if 300 <= row["status"] < 400:
        return "unresolved_redirect"
    if row["checks"]["soft_block_checks"]["matched"]:
        return "soft_block"
    if not row["checks"]["has_required_content"]:
        return "login_or_consent_gate" if row["checks"]["gate_checks"]["matched"] else "missing_required_content"
    return "usable" if row["status"] == 200 and row["capture_complete"] else "incomplete_capture"


def env_direct():
    env = os.environ.copy()
    for key in list(env):
        if key.lower() in ("http_proxy", "https_proxy", "all_proxy"):
            env.pop(key)
    env.update(PYTHONDONTWRITEBYTECODE="1")
    return env


def memory():
    readings = {line.split(":")[0]: int(line.split()[1]) / 1024
                for line in Path("/proc/meminfo").read_text().splitlines() if line.split(":")[0] in ("MemTotal", "MemAvailable", "SwapTotal", "SwapFree")}
    return {**readings, "unit": "MiB", "measured_at": utc()}


def base_row(target):
    return {"target_id": target["id"], "site": target["name"], "category": target["category"],
            "url": target["url"], "original_url": target["url"], "final_url": None,
            "expected_required_content": target["expected_required_content"], "started_at": utc(),
            "navigation_started_at": None, "status": None, "navigation_outcome": "not_started",
            "redirect_history": [], "redirect_outcome": "unobserved", "document_responses": [],
            "browser_version": None, "browser_cdp_version": None, "error": None,
            "navigation_error": None, "capture_error": None, "screenshot_error": None,
            "screenshot_private_path": None, "screenshot_sha256": None,
            "capture_complete": False, "raw_html_sha256": None, "raw_html_private_path": None,
            "html_bytes": 0, "navigation_seconds": None, "attempt_seconds": None}


@contextmanager
def deadline(seconds):
    def expired(_sig, _frame): raise TimeoutError(f"Stage deadline exceeded {seconds} seconds")
    previous = signal.signal(signal.SIGALRM, expired)
    signal.setitimer(signal.ITIMER_REAL, seconds)
    try: yield
    finally:
        signal.setitimer(signal.ITIMER_REAL, 0)
        signal.signal(signal.SIGALRM, previous)


async def single(args, target):
    row = base_row(target); capture = b""; start = time.perf_counter()
    row_path = args.output / f"{target['id']}.json"
    def save(): row_path.write_text(json.dumps(row, indent=2) + "\n")
    save()
    @browser(headless=True, chrome_executable_path=str(args.binary),
             add_arguments=CONFIG["launch_args"], proxy=None, profile=None, tiny_profile=False,
             reuse_driver=False, cache=False, output=None, max_retry=0, close_on_crash=True,
             raise_exception=True, create_error_logs=False, beep=False, parallel=1,
             wait_for_complete_page_load=False, block_images=False, block_images_and_css=False)
    def attempt(driver: Driver, data):
        nonlocal capture
        versions = driver.run_cdp_command(cdp.browser.get_version())
        row["browser_cdp_version"] = dict(zip(["protocolVersion", "product", "revision", "userAgent", "jsVersion"], versions))
        row["browser_version"] = row["browser_cdp_version"]["product"]
        row["effective_browser_arguments"] = driver._browser._process.args[1:]
        driver.run_cdp_command(cdp.browser.set_download_behavior(behavior="deny"))
        main_frame = driver.run_cdp_command(cdp.page.get_frame_tree()).frame.id_
        def request_sent(_id, request, event):
            if event.type_ == cdp.network.ResourceType.DOCUMENT and event.frame_id == main_frame and event.redirect_response:
                old = event.redirect_response
                row["redirect_history"].append({"previous":old.url,"url":request.url,"status":int(old.status)})
        def received(_id, response, event):
            if event.type_ == cdp.network.ResourceType.DOCUMENT and event.frame_id == main_frame:
                row["document_responses"].append({"url":response.url,"status":int(response.status)})
        driver.before_request_sent(request_sent); driver.after_response_received(received)
        driver.run_cdp_command(cdp.network.enable())
        row["navigation_started_at"] = utc(); save(); nav_start = time.perf_counter()
        try:
            with deadline(20): driver.get(target["url"], bypass_cloudflare=False, timeout=20)
            row["navigation_outcome"] = "returned"
        except Exception as exc:
            row["navigation_outcome"] = "error"
            row["navigation_error"] = {"type":type(exc).__name__,"message":str(exc).splitlines()[0][:400]}
        row["navigation_seconds"] = round(time.perf_counter()-nav_start,3)
        time.sleep(CONFIG["settle_seconds"])
        try:
            with deadline(5): row["final_url"] = driver.current_url
        except Exception as exc:
            row["final_url_error"] = {"type":type(exc).__name__,"message":str(exc)[:400]}
        if row["document_responses"]: row["status"] = row["document_responses"][-1]["status"]
        row["redirect_outcome"] = "followed" if row["redirect_history"] else ("none" if row["final_url"]==target["url"] else "final_url_changed_history_unavailable")
        try:
            with deadline(5): html = driver.page_html
            capture = html.encode("utf-8")
            if len(capture)>CONFIG["max_capture_bytes"]:
                capture=capture[:CONFIG["max_capture_bytes"]]
                row["capture_error"]={"type":"DOMSizeLimit","message":"DOM exceeded 5 MiB; saved prefix only"}
            else: row["capture_complete"]=True
            dest=args.output/f"{target['id']}.html";dest.write_bytes(capture)
            row.update(raw_html_private_path=str(dest),raw_html_sha256=sha(capture),html_bytes=len(capture))
        except Exception as exc: row["capture_error"]={"type":type(exc).__name__,"message":str(exc).splitlines()[0][:400]}
        row["latency_seconds"]=round(time.perf_counter()-nav_start,3);save()
        try:
            with deadline(5):
                image=driver.run_cdp_command(cdp.page.capture_screenshot(format_="png",capture_beyond_viewport=False))
                dest=args.output/f"{target['id']}.png";dest.write_bytes(base64.b64decode(image))
                row.update(screenshot_private_path=str(dest),screenshot_sha256=sha(dest.read_bytes()))
            with deadline(3): row["observed_fingerprint"]=driver.run_js("return {userAgent:navigator.userAgent,webdriver:navigator.webdriver,language:navigator.language,width:innerWidth,height:innerHeight,timezone:Intl.DateTimeFormat().resolvedOptions().timeZone}",timeout=3)
        except Exception as exc: row["screenshot_error"]={"type":type(exc).__name__,"message":str(exc).splitlines()[0][:400]}
        save()
    try: attempt(target["url"])
    except Exception as exc: row["error"]={"type":type(exc).__name__,"message":str(exc).splitlines()[0][:400],"classification":"browser_or_driver_error"}
    row.update(finished_at=utc(),attempt_seconds=round(time.perf_counter()-start,3),checks=evaluate(capture,target))
    row.setdefault("latency_seconds",row["attempt_seconds"])
    row["empty_content"]=row["checks"]["text_chars"]==0
    row["classification"]=result_class(row);row["usable_content"]=row["classification"]=="usable";save()


def tree_members(root):
    pending, seen, members = [root], set(), {}
    while pending:
        pid = pending.pop()
        if pid in seen:
            continue
        seen.add(pid)
        try:
            for task in Path(f"/proc/{pid}/task").iterdir():
                pending.extend(int(p) for p in (task / "children").read_text().split())
            stat = Path(f"/proc/{pid}/stat").read_text().rsplit(")", 1)[1].split()
            members[pid] = stat[19]  # start time; protects cleanup from PID reuse
        except (OSError, ValueError, IndexError):
            continue
    return members


def pss_members(members):
    total, read = 0, 0
    for pid in members:
        try:
            match = re.search(r"^Pss:\s+(\d+)", Path(f"/proc/{pid}/smaps_rollup").read_text(), re.M)
            if match:
                total += int(match.group(1))
                read += 1
        except OSError:
            continue
    return total / 1024 if read else None


def signal_owned(members, sig):
    # Chromium may use its own process group. Signal only recorded descendants.
    for pid, birth in reversed(list(members.items())):
        try:
            current = Path(f"/proc/{pid}/stat").read_text().rsplit(")", 1)[1].split()[19]
            if current == birth:
                os.kill(pid, sig)
        except (OSError, IndexError):
            pass


async def batch(args):
    manifest_bytes = args.manifest.read_bytes()
    manifest = json.loads(manifest_bytes)
    targets = manifest["targets"]
    assert len(targets) == len({t["id"] for t in targets}) == len({urlsplit(t["url"]).hostname for t in targets}) == 30
    assert version("botasaurus") == "4.0.97" and version("botasaurus-driver") == "4.0.101"
    assert not is_docker, "Setup blocker: installed driver Docker teardown can run global zombie cleanup"
    assert sha(args.binary.read_bytes()) == "8c599d43aec53f2460a31ae2f4af6bd863f8258b34ff519564bc5d4726bfaa1e"
    slice_start = time.perf_counter()
    initial_ram = memory()
    if initial_ram["MemAvailable"] < CONFIG["min_host_available_mb"]:
        raise RuntimeError("Unsafe RAM blocker before batch: less than 3 GiB available")
    os.umask(0o077)
    args.output.mkdir(parents=True, exist_ok=False, mode=0o700)
    run = {"run_id": args.output.name, "tool": "botasaurus", "catalog_slug": "botosaurus", "tool_version": version("botasaurus"), "driver_version": version("botasaurus-driver"),
           "requirements_sha256": sha((Path(__file__).parent/"requirements.txt").read_bytes()),
           "installed_packages": dict(line.split("==",1) for line in (Path(__file__).parent/"requirements.txt").read_text().splitlines()),
           "documentation": json.loads((Path(__file__).parent/"documentation.json").read_text()),
           "browser_executable": str(args.binary), "binary_sha256": sha(args.binary.read_bytes()),
           "executable_version": subprocess.check_output([str(args.binary), "--version"], text=True).strip(),
           "python_version": platform.python_version(), "platform": platform.platform(),
           "network": "Same local host direct outbound; proxy environment removed and --no-proxy-server; different UTC window from prior tools",
           "manifest_id": manifest["id"], "manifest_sha256": sha(manifest_bytes),
           "script_sha256": sha(Path(__file__).read_bytes()), "config": CONFIG,
           "executed_command": shlex.join([sys.executable, *sys.argv]), "started_at": utc(),
           "host_ram_preflight": initial_ram, "attempts": [],
           "body_hash_basis": "UTF-8 serialized DOM, not HTTP response bytes. Full HTML/screenshots/logs stay private in ignored runner/data.",
           "pss_caveat": "100 ms samples may miss brief peaks or changing processes. Includes isolated Python, Botasaurus driver, Chrome startup, screenshots and teardown; excludes batch parent. Sampled totals, not exact peaks or pure browser memory.",
           "baseline_sha256": {tool:sha((Path(__file__).parent.parent/folder/"results.json").read_bytes()) for tool,folder in {"patchright":"patchright-public-30-20260929","camoufox":"camoufox-public-30-20260929","lightpanda":"lightpanda-public-30-20260929","scrapling":"scrapling-public-30-20260929","wreq":"wreq-public-30-redirects-20260929"}.items()},
           "policy_differences": ["Document-ready-state wait (interactive or complete), not an exact DOMContentLoaded-event wait; two-second settle follows. Timing includes navigation, settle, URL read and DOM capture; startup/screenshot/teardown excluded.", "Botasaurus-driver uses its own CDP connection, default flags and headless UA override; same binary as earlier Patchright but not same launch/control policy.", "Temporary disk Chrome profile per worker rather than a new Playwright context; ordinary session cookies allowed, no imported account state.", "Full resources and JS, no Google referrer, custom fingerprint, human cursor or Cloudflare bypass. Request mode and other documented modes are not tested.", "30 later sequential attempts, not repeated success-rate measurements or a causal feature test."]}

    def save_run():
        run["summary"] = {"attempted_distinct_sites": len(run["attempts"]), "usable_results": sum(a["usable_content"] for a in run["attempts"]),
                          "classifications": dict(Counter(a["classification"] for a in run["attempts"]))}
        (args.output / "results.json").write_text(json.dumps(run, indent=2) + "\n")
    for index, target in enumerate(targets):
        if time.perf_counter()-slice_start>CONFIG["slice_timeout_seconds"]:
            run["blocker"]="Whole slice time guard";save_run();raise RuntimeError(run["blocker"])
        if index:
            await asyncio.sleep(CONFIG["gap_seconds"])
        before_ram = memory()
        if before_ram["MemAvailable"] < CONFIG["min_host_available_mb"]:
            run["blocker"] = "Unsafe host RAM before next navigation"
            save_run()
            raise RuntimeError(run["blocker"])
        command = [sys.executable, str(Path(__file__).resolve()), "--binary", str(args.binary),
                   "--manifest", str(args.manifest), "--output", str(args.output), "--single", target["id"]]
        with open(args.output / f"{target['id']}-driver.log", "wb") as log:
            temporary=Path(tempfile.mkdtemp(prefix="sev-bota-",dir="/tmp"))
            worker_env=env_direct();worker_env["TMPDIR"]=str(temporary)
            proc = await asyncio.create_subprocess_exec(*command, stdout=log, stderr=log, env=worker_env, cwd=args.output, start_new_session=True)
            started = time.perf_counter()
            samples, guard, owned = [], None, {}
            while proc.returncode is None:
                members = tree_members(proc.pid)
                owned.update(members)
                reading = pss_members(members)
                ram = memory()
                if reading is not None:
                    samples.append({"elapsed_seconds": round(time.perf_counter() - started, 3), "pss_mb": round(reading, 3), "host_available_mb": round(ram["MemAvailable"], 3)})
                if ram["MemAvailable"] < CONFIG["min_host_available_mb"] or (reading or 0) > CONFIG["max_attempt_tree_pss_mb"]:
                    guard = "unsafe_memory_guard"
                elif time.perf_counter() - started > CONFIG["attempt_timeout_seconds"]:
                    guard = "attempt_timeout"
                if guard:
                    signal_owned(owned, signal.SIGTERM)
                    await asyncio.sleep(0.2)
                    signal_owned(owned, signal.SIGKILL)
                    break
                await asyncio.sleep(CONFIG["pss_sample_interval_seconds"])
            await proc.wait()
            signal_owned(owned, signal.SIGKILL)
        await asyncio.sleep(0.1)
        survivors=[]
        for pid,birth in owned.items():
            try:
                stat=Path(f"/proc/{pid}/stat").read_text().rsplit(")",1)[1].split()
                if stat[19]==birth and stat[0]!="Z":survivors.append(pid)
            except (OSError,IndexError):pass
        row_path = args.output / f"{target['id']}.json"
        row = json.loads(row_path.read_text()) if row_path.exists() else base_row(target)
        if guard or proc.returncode != 0:
            row["error"] = {"type": "AttemptGuard" if guard else "AttemptProcessExit", "message": guard or f"Attempt process exit {proc.returncode}", "classification": guard or "attempt_process_error"}
        saved_body = args.output / f"{target['id']}.html"
        row.setdefault("checks", evaluate(saved_body.read_bytes() if saved_body.exists() else b"", target))
        row["empty_content"] = row["checks"]["text_chars"] == 0
        row.setdefault("latency_seconds", round(time.perf_counter() - started, 3))
        row.setdefault("finished_at", utc())
        row.update(peak_sampled_tree_pss_mb=max((s["pss_mb"] for s in samples), default=None), pss_samples=len(samples),
                   attempt_process_exit=proc.returncode, cleanup_surviving_owned_pids=survivors, host_ram_before=before_ram)
        row["classification"] = result_class(row)
        row["usable_content"] = row["classification"] == "usable"
        (args.output / f"{target['id']}-pss.json").write_text(json.dumps(samples) + "\n")
        row_path.write_text(json.dumps(row, indent=2) + "\n")
        run["attempts"].append(row)
        save_run()
        print(f"{target['id']}: HTTP {row['status']} | {row['classification']} | {row['html_bytes']} DOM bytes | PSS {row['peak_sampled_tree_pss_mb']} MiB", flush=True)
        if row["navigation_started_at"] is None or guard == "unsafe_memory_guard":
            run["blocker"] = "Setup or unsafe memory blocker: " + str(row["error"])
            save_run()
            raise RuntimeError(run["blocker"])
    run["finished_at"] = utc()
    save_run()
    print(json.dumps(run["summary"]), flush=True)


if __name__ == "__main__":
    parser = argparse.ArgumentParser(description=__doc__)
    parser.add_argument("--binary", type=Path, required=True)
    parser.add_argument("--manifest", type=Path, required=True)
    parser.add_argument("--output", type=Path, required=True)
    parser.add_argument("--single")
    args = parser.parse_args()
    args.binary = args.binary.resolve()
    args.manifest=args.manifest.resolve();args.output=args.output.resolve()
    if args.single:
        target = next(t for t in json.loads(args.manifest.read_text())["targets"] if t["id"] == args.single)
        asyncio.run(single(args, target))
    else:
        asyncio.run(batch(args))
