"""Fresh sequential Crawlee CheerioCrawler workers; frozen checks; private captures."""
import argparse, asyncio, hashlib, json, os, platform, re, shlex, signal, subprocess, sys, time
from collections import Counter
from datetime import datetime, timezone
from html.parser import HTMLParser
from pathlib import Path
from urllib.parse import urlsplit
CONFIG = {'mode':'CheerioCrawler single-page direct HTTP / GotScrapingHttpClient', 'proxy_mode':'none', 'login':False, 'captcha_solving':False, 'concurrency':1, 'gap_seconds':2, 'application_retries':0, 'maxRequestRetries':0, 'maxSessionRotations':0, 'retryOnBlocked':False, 'useSessionPool':False, 'persistCookiesPerSession':False, 'respectRobotsTxtFile':False, 'maxRequestsPerCrawl':1, 'maxRedirects':3, 'got_retry_limit':0, 'ignoreSslErrors':False, 'http2':True, 'useHeaderGenerator':True, 'custom_headers':{}, 'cookies':'No cookie jar/session pool; no imported, cross-target or redirect cookie storage', 'resources':'HTTP body and Cheerio parsing only; no browser/JS/subresources/enqueueLinks', 'navigationTimeoutSecs':20, 'requestHandlerTimeoutSecs':5, 'attempt_timeout_seconds':35, 'slice_timeout_seconds':1200, 'max_capture_bytes':5242880, 'pss_sample_interval_seconds':0.1, 'min_host_available_mb':3072, 'max_attempt_tree_pss_mb':1024, 'pss_scope':'Isolated Node worker and descendants from spawn/imports through crawler startup, request/body/Cheerio parsing/capture and exit; excludes Python batch parent and frozen checks', 'storage':'Fresh process; nonpersistent memory storage; own private target directory', 'headers':'Stock got-scraping header generator; actual outgoing hook headers recorded per request; no pinned Chrome version/profile', 'tls':'Stock got-scraping TLS customization with Node/OpenSSL; certificate verification enabled for page request; no exact Chrome TLS identity claim'}
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}


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 json.loads((args.package_root/"node_modules/@crawlee/cheerio/package.json").read_text())["version"] == "3.18.2"
    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)
    here=Path(__file__).parent
    versions={name:json.loads((args.package_root/'node_modules'/name/'package.json').read_text())['version'] for name in ['@crawlee/cheerio','@crawlee/http','@crawlee/basic','@crawlee/core','got-scraping','got','cheerio','header-generator']}
    run={'run_id':args.output.name,'tool':'Crawlee','tool_version':'3.18.2','node_version':subprocess.check_output(['node','--version']).decode().strip(),'installed_versions':versions,'python_version':platform.python_version(),'platform':platform.platform(),'mode':CONFIG['mode'],'config':CONFIG,'started_at':utc(),'attempts':[], 'host_ram_preflight':initial_ram,
      'manifest_id':manifest['id'],'manifest_sha256':sha(manifest_bytes),'script_sha256':sha(here.joinpath('run.py').read_bytes()),'worker_sha256':sha(here.joinpath('worker.mjs').read_bytes()),'lock_sha256':sha(args.package_root.joinpath('package-lock.json').read_bytes()),'package_sha256':sha(args.package_root.joinpath('package.json').read_bytes()),'executed_command':shlex.join([sys.executable,*sys.argv]),
      'network':'Same local host, direct outbound; proxy environment removed case-insensitively, no proxyConfiguration/proxyUrl; different UTC window',
      'latency_basis':'preNavigationHook before got request through requestHandler entry (including redirects, transfer and Cheerio parsing), or failure-handler/run completion; excludes imports/startup, storage startup, save/checks and teardown',
      'body_hash_basis':'Original client-decompressed got stream bytes, before Crawlee charset conversion/Cheerio serialization. Frozen UTF-8 replacement decoder unchanged.',
      'pss_caveat':'100 ms sampled Node worker and descendants, includes imports/storage/HTTP/Cheerio/capture/exit, excludes Python batch parent/frozen checking. Samples can miss brief peaks or changing processes; no crawl-scale memory claim.',
      'preserved_evidence_sha256':json.loads(here.joinpath('preserved-evidence.json').read_text()),
      'documentation':json.loads(here.joinpath('documentation.json').read_text())}

    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=['node',str(Path(__file__).parent/'worker.mjs'),target['url'],str(args.output),target['id']]
        env=env_direct();env.update(CRAWLEE_PACKAGE_ROOT=str(args.package_root),CRAWLEE_PERSIST_STORAGE='false',CRAWLEE_STORAGE_DIR=str(args.output/(target['id']+'-storage')))

        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, 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["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["cleanup_surviving_owned_pids"] = [pid for pid,birth in owned.items() if tree_members(pid).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("--package-root", type=Path, required=True)
    parser.add_argument("--target", help="One separate guarded public canary; manifest remains frozen")
    args = parser.parse_args()
    asyncio.run(batch(args))
