"""One local Lightpanda navigation per frozen target; private captures only.
Requires Python >=3.11 and playwright==1.63.0. Uses an existing local 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 socket
import subprocess
import sys
import time
from urllib.parse import urlsplit
from playwright.async_api import async_playwright

CONFIG = {
    "proxy_mode": "none", "login": False, "captcha_solving": False,
    "concurrency": 1, "gap_seconds": 2, "application_retries": 0,
    "navigation_timeout_seconds": 20, "attempt_timeout_seconds": 30,
    "wait_until": "domcontentloaded", "settle_seconds": 2,
    "max_response_bytes": 5242880, "max_capture_bytes": 5242880,
    "redirects": "follow", "max_redirects": 10,
    "redirect_cap_basis": "Hard-coded httpMaxRedirects() in exact binary source src/Config.zig",
    "robots": "default off, matching wreq; differs from earlier Lightpanda two-page slice",
    "resources": "Lightpanda defaults: scripts/fetch/XHR; images, stylesheets, iframes and workers not opted in",
    "cookies": "fresh in-memory browser context; no imported account state",
    "telemetry": "disabled", "pss_sample_interval_seconds": 0.1,
    "pss_scope": "Isolated attempt Python process and all descendants, including Playwright driver and Lightpanda; excludes batch parent",
}

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 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,
            "redirect_history": [], "redirect_outcome": "unobserved", "document_responses": [],
            "browser_version": None, "browser_cdp_version": None, "error": None,
            "navigation_error": None, "capture_error": None, "capture_complete": False,
            "raw_html_sha256": None, "raw_html_private_path": None, "html_bytes": 0}


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(LIGHTPANDA_DISABLE_TELEMETRY="true", PYTHONDONTWRITEBYTECODE="1")
    return env


async def single(args, target):
    row = base_row(target)
    capture = b""
    (args.output / f"{target['id']}.json").write_text(json.dumps(row, indent=2) + "\n")
    with socket.socket() as sock:
        sock.bind(("127.0.0.1", 0))
        port = sock.getsockname()[1]
    command = [str(args.binary), "serve", "--host", "127.0.0.1", "--port", str(port),
               "--http-max-response-size", str(CONFIG["max_response_bytes"]),
               "--block-private-networks", "--log-format", "logfmt", "--log-level", "warn"]
    row["browser_command"] = command
    start = time.perf_counter()
    log = open(args.output / f"{target['id']}-browser.log", "wb")
    proc = subprocess.Popen(command, env=env_direct(), stdout=log, stderr=log)
    page = None
    try:
        for _ in range(100):
            if proc.poll() is not None:
                raise RuntimeError(f"Lightpanda exited during startup with code {proc.returncode}")
            try:
                reader, writer = await asyncio.open_connection("127.0.0.1", port)
                writer.close()
                await writer.wait_closed()
                break
            except OSError:
                await asyncio.sleep(0.05)
        else:
            raise TimeoutError("Local CDP listener not ready after five seconds")
        async with async_playwright() as p:
            browser = await p.chromium.connect_over_cdp(f"ws://127.0.0.1:{port}", timeout=5000)
            row["browser_version"] = browser.version
            session = await browser.new_browser_cdp_session()
            row["browser_cdp_version"] = await session.send("Browser.getVersion")
            context = await browser.new_context()
            page = await context.new_page()
            def received(response):
                try:
                    req = response.request
                    if req.is_navigation_request() and req.frame == page.main_frame:
                        row["document_responses"].append({"url": response.url, "status": response.status})
                except Exception:
                    pass
            page.on("response", received)
            row["navigation_started_at"] = utc()
            (args.output / f"{target['id']}.json").write_text(json.dumps(row, indent=2) + "\n")
            nav_start = time.perf_counter()
            response = None
            try:
                response = await page.goto(target["url"], wait_until=CONFIG["wait_until"], timeout=20000)
                row["navigation_seconds"] = round(time.perf_counter() - nav_start, 3)
                await asyncio.sleep(CONFIG["settle_seconds"])
            except Exception as exc:
                row["navigation_error"] = {"type": type(exc).__name__, "message": str(exc).splitlines()[0][:400]}
                row["navigation_seconds"] = round(time.perf_counter() - nav_start, 3)
            row["final_url"] = page.url
            if response is not None:
                row["status"] = response.status
                hops = []
                request = response.request
                while request.redirected_from is not None:
                    previous = request.redirected_from
                    prior_response = await previous.response()
                    hops.append({"previous": previous.url, "url": request.url,
                                 "status": prior_response.status if prior_response else None})
                    request = previous
                row["redirect_history"] = list(reversed(hops))
            final_responses = [r for r in row["document_responses"] if r["url"] == page.url]
            if final_responses:
                row["status"] = final_responses[-1]["status"]
            row["redirect_outcome"] = ("followed" if row["redirect_history"] else
                                       "final_url_changed_history_unavailable" if page.url != target["url"] and page.url != "about:blank" else "none")
            try:
                async with asyncio.timeout(3):
                    html = await page.content()
                capture = html.encode("utf-8")
                if len(capture) > CONFIG["max_capture_bytes"]:
                    capture = capture[:CONFIG["max_capture_bytes"]]
                    row["capture_error"] = {"type": "CaptureLimit", "message": "Serialized DOM exceeded 5 MiB; saved prefix only"}
                else:
                    row["capture_complete"] = True
                body = args.output / f"{target['id']}.html"
                body.write_bytes(capture)
                row.update(raw_html_sha256=sha(capture), raw_html_private_path=str(body), 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)
            (args.output / f"{target['id']}.json").write_text(json.dumps(row, indent=2) + "\n")
            await browser.close()
    except Exception as exc:
        row["error"] = {"type": type(exc).__name__, "message": str(exc).splitlines()[0][:400], "classification": "browser_or_cdp_error"}
    finally:
        if proc.poll() is None:
            proc.terminate()
            try:
                await asyncio.to_thread(proc.wait, timeout=1)
            except subprocess.TimeoutExpired:
                proc.kill()
                await asyncio.to_thread(proc.wait)
        log.close()
    row.update(finished_at=utc(), attempt_seconds=round(time.perf_counter() - start, 3),
               browser_exit_code=proc.returncode, 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"
    (args.output / f"{target['id']}.json").write_text(json.dumps(row, indent=2) + "\n")


def pss_tree(root):
    pending, seen, total, read = [root], set(), 0, 0
    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())
            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, ValueError):
            continue
    return total / 1024 if read else None


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("playwright") == "1.63.0"
    os.umask(0o077)
    args.output.mkdir(parents=True, exist_ok=False, mode=0o700)
    binary_version = subprocess.check_output([str(args.binary), "version"], text=True).strip()
    run = {"run_id": args.output.name, "tool": "lightpanda", "binary_version": binary_version,
           "binary_path": str(args.binary), "binary_sha256": sha(args.binary.read_bytes()),
           "python_version": platform.python_version(), "playwright_version": version("playwright"),
           "platform": platform.platform(), "network": "Same local host direct outbound; no configured proxy; different UTC window from wreq",
           "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(), "attempts": [],
           "body_hash_basis": "UTF-8 serialized rendered DOM, not original HTTP response bytes. Full HTML/private browser logs stay in ignored runner/data.",
           "pss_caveat": "100 ms samples can miss brief peaks. Process trees change between reads. Includes Python/Playwright driver and browser startup/teardown; not a pure browser-memory or exact peak measurement.",
           "policy_differences": ["Native Lightpanda redirect cap 10 versus shared wreq cap 3; no public CLI override in this binary.",
                                  "DOM after JavaScript/DOMContentLoaded plus two-second settle versus wreq response source HTML.",
                                  "20-second navigation limit plus a 30-second whole-attempt watchdog versus wreq's 20-second whole transfer.",
                                  "Lightpanda default HTTP transfer timeout 15 seconds and default script/subrequests; no robots fetching, matching wreq.",
                                  "Fresh browser accepts in-session cookies; wreq disables its cookie store.",
                                  "5 MiB limit per response and serialized DOM, versus wreq's final response-body cap.",
                                  "Private-network target/subrequests blocked; no account, proxy, cache directory, or resource opt-ins."]}
    for index, target in enumerate(targets):
        if index:
            await asyncio.sleep(CONFIG["gap_seconds"])
        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:
            proc = await asyncio.create_subprocess_exec(*command, stdout=log, stderr=log, env=env_direct(), start_new_session=True)
            started = time.perf_counter()
            samples = []
            killed = False
            while proc.returncode is None:
                reading = pss_tree(proc.pid)
                if reading is not None:
                    samples.append({"elapsed_seconds": round(time.perf_counter() - started, 3), "pss_mb": round(reading, 3)})
                if time.perf_counter() - started > CONFIG["attempt_timeout_seconds"]:
                    killed = True
                    os.killpg(proc.pid, signal.SIGTERM)
                    await asyncio.sleep(0.2)
                    try:
                        os.killpg(proc.pid, signal.SIGKILL)
                    except ProcessLookupError:
                        pass
                    break
                await asyncio.sleep(CONFIG["pss_sample_interval_seconds"])
            await proc.wait()
        row_path = args.output / f"{target['id']}.json"
        row = json.loads(row_path.read_text()) if row_path.exists() else base_row(target)
        if killed or proc.returncode != 0:
            row["error"] = {"type": "AttemptWatchdog" if killed else "AttemptProcessExit", "message": f"Attempt process exit {proc.returncode}", "classification": "attempt_timeout" if killed else "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)
        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)
        (args.output / "results.json").write_text(json.dumps(run, indent=2) + "\n")
        if row["navigation_started_at"] is None:
            raise RuntimeError(f"Setup blocker before target navigation: {row['error']}")
        print(f"{target['id']}: HTTP {row['status']} | {row['classification']} | {row['html_bytes']} DOM bytes | PSS {row['peak_sampled_tree_pss_mb']} MiB", flush=True)
    run.update(finished_at=utc(), 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")
    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()
    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))
