"""One isolated direct curl_cffi Session.get per frozen URL; private captures only.
Requires Python >=3.11 and curl_cffi==0.16.3. No Scrapling, browser or JavaScript.
"""
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
from urllib.parse import urlsplit
from curl_cffi import Curl, CurlOpt
from curl_cffi.requests import Session

CONFIG = {'mode': 'Direct Session.get HTTP only', 'proxy_mode': 'none', 'login': False, 'captcha_solving': False, 'concurrency': 1, 'gap_seconds': 2, 'application_retries': 0, 'fresh_process_and_session_per_target': True, 'impersonate': 'chrome150', 'default_headers': True, 'custom_headers': {}, 'google_referer_added': False, 'trust_env': False, 'curl_proxy': 'empty string; explicitly disables libcurl proxy use', 'follow_redirects': True, 'max_redirects': 3, 'timeout': 20, 'attempt_timeout_seconds': 30, 'slice_timeout_seconds': 1000, 'max_capture_bytes': 5242880, 'network_response_byte_cap': None, 'cookies': 'Fresh curl session; redirect cookies accepted; discard_cookies=False; session closed after one get; no imported or cross-target state', 'resources': 'HTTP response only; no browser, JavaScript or resource loads', 'verify': True, 'http_version': 'v2 (HTTP/2 negotiation with HTTP/1 fallback; no HTTP/3)', 'cache': None, 'pss_sample_interval_seconds': 0.1, 'min_host_available_mb': 3072, 'max_attempt_tree_pss_mb': 1024, 'pss_scope': 'Isolated Python attempt and descendants, including imports, curl session/request/capture/checks and close; excludes batch parent', 'private_debug': 'HEADER_IN and HEADER_OUT only; no TLS/data body dump; cookie values remain private'}

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(),
            "request_started_at": None, "status": None, "request_outcome": "not_started",
            "redirect_history": [], "redirect_outcome": "unobserved", "document_responses": [],
            "explicit_request_headers": {}, "response_metadata": {}, "error": None,
            "navigation_error": None, "capture_error": None, "screenshot_error": None,
            "screenshot_private_path": None, "screenshot_sha256": None,
            "capture_complete": False, "raw_response_sha256": None, "raw_response_private_path": None,
            "response_bytes": 0, "downloaded_body_bytes": 0, "body_complete": False, "navigation_seconds": None, "attempt_seconds": None}


async def single(args, target):
    row, capture = base_row(target), b""
    row_path = args.output / f"{target['id']}.json"
    def save():
        row_path.write_text(json.dumps(row, indent=2) + "\n")
    start = time.perf_counter()
    events = []
    def debug(kind, data):
        if kind in (1, 2):
            events.append({"kind": "response_header" if kind == 1 else "request_headers", "data": data.decode("latin1")})
    response = None
    with Session(impersonate="chrome150", default_headers=True, headers={}, cookies={},
                 trust_env=False, proxies={}, retry=0, discard_cookies=False,
                 allow_redirects=True, max_redirects=3, timeout=20, verify=True,
                 http_version="v2", cache=None, raise_for_status=False,
                 curl_options={CurlOpt.PROXY: b"", CurlOpt.VERBOSE: 1, CurlOpt.DEBUGFUNCTION: debug}) as session:
        row["request_started_at"] = utc()
        save()
        request_start = time.perf_counter()
        try:
            response = session.get(target["url"])
            row["request_outcome"] = "returned"
            row["capture_complete"] = True
        except Exception as exc:
            message = str(exc).splitlines()[0][:500]
            code = getattr(exc, "code", None)
            category = "redirect_limit_or_loop" if code == 47 else "timeout" if code == 28 or "timed out" in message.lower() else "transport_error"
            row["request_outcome"] = "error"
            row["error"] = {"type": type(exc).__name__, "message": message, "classification": category, "curl_code": code}
            response = getattr(exc, "response", None)
        row["latency_seconds"] = round(time.perf_counter() - request_start, 3)
        transport = {"header_events": events, "status": None, "final_url": None, "redirect_history": []}
        if response is not None:
            row["status"] = response.status_code or None
            row["final_url"] = str(response.url) or None
            capture = response.content
            row["downloaded_body_bytes"] = len(capture)
            history = list(response.history)
            row["redirect_history"] = [{"previous": h.url, "url": history[i+1].url if i+1 < len(history) else row["final_url"], "status": h.status_code} for i,h in enumerate(history)]
            row["redirect_count"] = response.redirect_count
            row["redirect_outcome"] = "followed" if history else "none" if row["final_url"] == target["url"] else "final_url_changed_history_unavailable"
            if row["error"] and row["error"]["classification"] == "redirect_limit_or_loop":
                row["redirect_outcome"] = "limit_or_loop_error"
            row["response_metadata"] = {"content_type": response.headers.get("content-type"), "content_encoding": response.headers.get("content-encoding"), "http_version_enum": response.http_version, "cookie_values": "private/not published"}
            transport.update(status=row["status"], final_url=row["final_url"], redirect_history=row["redirect_history"], response_headers=list(response.headers.multi_items()))
            if len(capture) > CONFIG["max_capture_bytes"]:
                capture = capture[:CONFIG["max_capture_bytes"]]
                row["capture_complete"] = False
                row["capture_error"] = {"type": "BodyCaptureLimit", "message": "Downloaded body exceeded 5 MiB; saved prefix only. Transfer is bounded by time/RAM, not bytes."}
            dest = args.output / f"{target['id']}.body"
            dest.write_bytes(capture)
            row.update(raw_response_private_path=str(dest), raw_response_sha256=sha(capture), response_bytes=len(capture))
    log = args.output / f"{target['id']}-transport.json"
    log.write_text(json.dumps(transport, indent=2) + "\n")
    row["transport_log_sha256"] = sha(log.read_bytes())
    row["observed_request_headers"] = []
    for event in events:
        if event["kind"] != "request_headers": continue
        lines = event["data"].splitlines()
        row["observed_request_headers"].append({"request_line": lines[0] if lines else None, "headers": [{"name": line.split(":",1)[0], "value": line.split(":",1)[1].strip()} for line in lines[1:] if ":" in line and line.split(":",1)[0].lower() not in ("cookie", "authorization", "proxy-authorization")], "cookie_values": "private/not published"})
    row.update(finished_at=utc(), attempt_seconds=round(time.perf_counter()-start,3), checks=evaluate(capture,target))
    row["empty_content"] = row["checks"]["text_chars"] == 0
    row["body_complete"] = row["capture_complete"]
    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):
    # Signal only recorded worker/descendant identities; never other host processes.
    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("curl_cffi") == "0.16.3"
    if args.target:
        targets = [t for t in targets if t["id"] == args.target]
        assert len(targets) == 1
    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": "curl_cffi", "tool_version": version("curl_cffi"),
           "curl_cffi_version": version("curl_cffi"), "requirements_sha256": sha((Path(__file__).parent/"requirements.txt").read_bytes()), "documentation": json.loads((Path(__file__).parent/"documentation.json").read_text()), "mode": CONFIG["mode"],
           "python_version": platform.python_version(), "platform": platform.platform(),
           "network": "Same local host direct outbound; http_proxy/https_proxy/all_proxy removed case-insensitively; no configured proxies. 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": [],
           "mode_selection": "Documented direct Session.get with chrome150 and default headers, no custom headers/Google Referer, no wrapper/selector or JS; one unchanged mode for all targets.",
           "body_hash_basis": "Original client-decompressed response bytes; no re-encoding or selector transformation. Full private bodies/logs in ignored runner/data.",
           "latency_basis": "Immediately before session.get through response return or exception, including body transfer and redirects. Excludes session construction/close, body saving and frozen checks; no Scrapling selector construction.",
           "pss_caveat": "100 ms sampled isolated Python worker and descendants, from spawn through imports, session/request/capture/checks/close and exit; excludes batch parent. Peaks may be missed; different browser/wrapper workload.",
           "baseline_sha256": {name: sha((Path(__file__).parent.parent/folder/"results.json").read_bytes()) for name,folder in {"botasaurus":"botasaurus-public-30-20260929","camoufox":"camoufox-public-30-20260929","lightpanda":"lightpanda-public-30-20260929","patchright":"patchright-public-30-20260929","scrapling":"scrapling-public-30-20260929","seleniumbase":"seleniumbase-public-30-20260929","wreq":"wreq-public-30-redirects-20260929"}.items()},
           "policy_differences": ["Direct curl_cffi 0.16.3 chrome150 with default impersonation headers; no added Google Referer or custom headers. Scrapling 0.4.15 uses the same curl_cffi/profile, adds stealthy headers including Google Referer, and constructs a selector.","Fresh curl session per process with redirect cookies accepted; no imported or cross-target cookies. wreq disables cookie storage.","No browser/JavaScript/resources. Normal redirects cap 3 as in wreq and Scrapling. Installed release exposes Response.history statuses/URLs without intermediate bodies.","verify=True; proxy environment removed and libcurl PROXY explicitly empty; trust_env=False retained. HTTP/2 negotiation with HTTP/1 fallback, no HTTP/3, no retry or cache.","20-second transfer timeout and 30-second whole-worker watchdog, 3 GiB host RAM floor and 1 GiB worker-tree PSS guard; 5 MiB saved-body cap is not a network byte cap.","Different UTC windows and header/cookie/resource policies; one observation per target establishes no reliability estimate, feature cause or global winner."]}

    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")
    slice_started = time.perf_counter()
    for index, target in enumerate(targets):
        if time.perf_counter() - slice_started > CONFIG["slice_timeout_seconds"]:
            run["blocker"] = "Whole-slice deadline before next target"
            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()), "--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, 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)
        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']}.body"
        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())
        survivors = tree_members(proc.pid)
        row["cleanup_surviving_owned_pids"] = [pid for pid,birth in survivors.items() if owned.get(pid)==birth]
        row["driver_log_sha256"] = sha((args.output / f"{target['id']}-driver.log").read_bytes())
        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, 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['response_bytes']} body bytes | PSS {row['peak_sampled_tree_pss_mb']} MiB", flush=True)
        if row["request_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("--manifest", type=Path, required=True)
    parser.add_argument("--output", type=Path, required=True)
    parser.add_argument("--single")
    parser.add_argument("--target", help="One separate guarded public canary; manifest remains frozen")
    args = parser.parse_args()
    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))
