#!/usr/bin/env python3
"""Scientific git-clone + in-memory tar.gz download-speed benchmark.

Pipeline per repo:
  1) git clone (authenticated HTTPS) into tmpfs
  2) tar.gz in a memory spool (spills to disk only if archive > 256 MiB)
  3) write {owner}__{reponame}.tar.gz to the dataset disk
  4) delete the working clone

Does NOT clone the full 54k-repo manifest. It draws a seeded, size-stratified
sample so throughput, latency, and packing cost can be estimated with
quantifiable uncertainty, then extrapolated to the full list.
"""
from __future__ import annotations

import argparse
import base64
import concurrent.futures
import json
import math
import os
import platform
import random
import re
import shutil
import socket
import statistics
import subprocess
import sys
import tarfile
import tempfile
import threading
import time
import traceback
import urllib.error
import urllib.request
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Iterable

SEED = 20260911
API_SAMPLE_N = 600
PER_STRATUM = {
    "seq": 5,
    "par4": 5,
    "par8": 4,
    "full": 1,
}
STRATA = [
    ("XS", 0, 100),  # GitHub size is KiB
    ("S", 100, 1024),
    ("M", 1024, 10 * 1024),
    ("L", 10 * 1024, 50 * 1024),
    ("XL", 50 * 1024, 200 * 1024),
    ("XXL", 200 * 1024, 1024 * 1024),
]
SKIP_SIZE_KIB = 1024 * 1024  # 1 GiB GitHub size: exclude from timed clones
SHALLOW_TIMEOUT_S = 180
FULL_TIMEOUT_S = 420
API_TIMEOUT_S = 20
SPOOL_MAX = 256 * 1024 * 1024
COMPRESS_LEVEL = 6
USER_AGENT = "github-clone-speed-test/1.0"
RECEIVE_RE = re.compile(
    r"Receiving objects:\s+100%\s+\((\d+)/(\d+)\)"
    r"(?:,\s+([\d.]+)\s+(KiB|MiB|GiB|bytes)"
    r"(?:\s+\|\s+([\d.]+)\s+(KiB|MiB|GiB|KiB|bytes)/s)?)?",
    re.I,
)
UNIT = {"bytes": 1, "kib": 1024, "mib": 1024**2, "gib": 1024**3}

ROOT = Path("/www/wwwroot/dataset/github_clone/speed_test")
MANIFEST = Path("/www/wwwroot/dataset/github_clone/star30_repos_not_cloned_part013.jsonl")
ARCHIVES = ROOT / "archives"
WORK = Path("/dev/shm/clone_speed_work")
DISK_WORK = ROOT / "work"
REPORTS = ROOT / "reports"
LOGS = ROOT / "logs"
METRICS_PATH = REPORTS / "metrics.jsonl"
SUMMARY_PATH = REPORTS / "summary.json"
REPORT_PATH = REPORTS / "clone_speed_test_report.md"
SAMPLE_PATH = REPORTS / "sample.json"
ENV_PATH = REPORTS / "environment.json"
LOG_PATH = LOGS / "run.log"

print_lock = threading.Lock()
metrics_lock = threading.Lock()


def utc_now() -> str:
    return datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")


def log(msg: str) -> None:
    line = f"[{utc_now()}] {msg}"
    with print_lock:
        print(line, flush=True)
        LOGS.mkdir(parents=True, exist_ok=True)
        with LOG_PATH.open("a", encoding="utf-8") as fh:
            fh.write(line + "\n")


def token() -> str:
    t = os.environ.get("GITHUB_TOKEN", "").strip()
    if not t:
        raise SystemExit("GITHUB_TOKEN is not set")
    return t


def redact(s: str) -> str:
    t = os.environ.get("GITHUB_TOKEN", "")
    if t:
        s = s.replace(t, "***")
    return s


def run(cmd: list[str], timeout: int | None = 30, env: dict | None = None, cwd: str | None = None) -> subprocess.CompletedProcess:
    return subprocess.run(
        cmd,
        stdout=subprocess.PIPE,
        stderr=subprocess.PIPE,
        text=True,
        timeout=timeout,
        env=env,
        cwd=cwd,
        check=False,
    )


def bytes_unit(n: float, unit: str) -> float:
    return n * UNIT[unit.lower()]


def fmt_bytes(n: float | None) -> str:
    if n is None:
        return "n/a"
    n = float(n)
    for u, d in (("GiB", 1024**3), ("MiB", 1024**2), ("KiB", 1024)):
        if abs(n) >= d:
            return f" {n / d:.2f} {u}".strip()
    return f"{n:.0f} B"


def fmt_s(n: float | None) -> str:
    if n is None:
        return "n/a"
    if n < 1:
        return f" {n * 1000:.0f} ms".strip()
    if n < 60:
        return f" {n:.2f} s".strip()
    m, s = divmod(n, 60)
    if m < 60:
        return f"{int(m)}m {s:.1f}s"
    h, m = divmod(int(m), 60)
    return f"{h}h {m}m {s:.0f}s"


def fmt_bps(n: float | None) -> str:
    if n is None or n <= 0:
        return "n/a"
    return f"{fmt_bytes(n)}/s"


def percentile(xs: list[float], p: float) -> float | None:
    if not xs:
        return None
    ys = sorted(xs)
    if len(ys) == 1:
        return ys[0]
    k = (len(ys) - 1) * p / 100.0
    f = math.floor(k)
    c = math.ceil(k)
    if f == c:
        return ys[int(k)]
    return ys[f] * (c - k) + ys[c] * (k - f)


def mean(xs: list[float]) -> float | None:
    return statistics.fmean(xs) if xs else None


def stdev(xs: list[float]) -> float | None:
    return statistics.stdev(xs) if len(xs) >= 2 else None


def bootstrap_mean_ci(xs: list[float], n: int = 2000, alpha: float = 0.05, seed: int = SEED) -> tuple[float, float] | None:
    if len(xs) < 2:
        return None
    rng = random.Random(seed)
    means = []
    k = len(xs)
    for _ in range(n):
        sample = [xs[rng.randrange(k)] for _ in range(k)]
        means.append(sum(sample) / k)
    means.sort()
    lo = means[int(alpha / 2 * n)]
    hi = means[min(n - 1, int((1 - alpha / 2) * n))]
    return lo, hi


def parse_full_name(rec: dict) -> tuple[str, str] | None:
    fn = (rec.get("full_name") or "").strip()
    if fn.count("/") != 1:
        return None
    owner, repo = fn.split("/", 1)
    owner, repo = owner.strip(), repo.strip()
    if not owner or not repo:
        return None
    return owner, repo


def archive_name(owner: str, repo: str) -> str:
    safe_owner = owner.replace("/", "_")
    safe_repo = repo.replace("/", "_")
    return f"{safe_owner}__{safe_repo}.tar.gz"


def load_manifest(path: Path) -> list[dict]:
    rows = []
    with path.open("r", encoding="utf-8") as fh:
        for i, line in enumerate(fh, 1):
            line = line.strip()
            if not line:
                continue
            rec = json.loads(line)
            parsed = parse_full_name(rec)
            if not parsed:
                continue
            rec["_owner"], rec["_repo"] = parsed
            rec["_full_name"] = f"{parsed[0]}/{parsed[1]}"
            rec["_idx"] = i
            rows.append(rec)
    return rows


def gh_api(url: str, tok: str, timeout: int = API_TIMEOUT_S) -> tuple[int, dict | list | None, dict]:
    req = urllib.request.Request(
        url,
        headers={
            "Authorization": f"Bearer {tok}",
            "Accept": "application/vnd.github+json",
            "User-Agent": USER_AGENT,
            "X-GitHub-Api-Version": "2022-11-28",
        },
    )
    hdrs: dict[str, str] = {}
    try:
        with urllib.request.urlopen(req, timeout=timeout) as resp:
            raw = resp.read()
            hdrs = {k.lower(): v for k, v in resp.headers.items()}
            body = json.loads(raw.decode("utf-8")) if raw else None
            return resp.status, body, hdrs
    except urllib.error.HTTPError as e:
        hdrs = {k.lower(): v for k, v in e.headers.items()} if e.headers else {}
        try:
            body = json.loads(e.read().decode("utf-8"))
        except Exception:
            body = None
        return e.code, body, hdrs
    except Exception as e:
        return 0, {"error": str(e)}, {}


def fetch_size(rec: dict, tok: str) -> dict:
    owner, repo = rec["_owner"], rec["_repo"]
    code, body, hdrs = gh_api(f"https://api.github.com/repos/{owner}/{repo}", tok)
    out = {
        "full_name": rec["_full_name"],
        "http_status": code,
        "size_kib": None,
        "default_branch": rec.get("default_branch"),
        "archived": None,
        "disabled": None,
        "fork": None,
        "language": rec.get("language"),
        "stargazers_count": rec.get("stargazers_count"),
        "api_error": None,
        "rate_remaining": hdrs.get("x-ratelimit-remaining"),
    }
    if code == 200 and isinstance(body, dict):
        out["size_kib"] = body.get("size")
        out["default_branch"] = body.get("default_branch") or out["default_branch"]
        out["archived"] = body.get("archived")
        out["disabled"] = body.get("disabled")
        out["fork"] = body.get("fork")
        out["language"] = body.get("language") or out["language"]
        out["stargazers_count"] = body.get("stargazers_count", out["stargazers_count"])
    else:
        msg = None
        if isinstance(body, dict):
            msg = body.get("message") or body.get("error")
        out["api_error"] = msg or f"http {code}"
    return out


def stratum_of(size_kib: int | None) -> str:
    if size_kib is None:
        return "UNKNOWN"
    if size_kib >= SKIP_SIZE_KIB:
        return "TOO_LARGE"
    for name, lo, hi in STRATA:
        if lo <= size_kib < hi:
            return name
    return "TOO_LARGE"


def git_env() -> dict[str, str]:
    env = os.environ.copy()
    env["GIT_TERMINAL_PROMPT"] = "0"
    env["GIT_LFS_SKIP_SMUDGE"] = "1"
    env["GIT_ASKPASS"] = "true"
    env["GC_AUTO"] = "0"
    return env


def git_cmd_prefix(tok: str) -> list[str]:
    # GitHub git-over-HTTPS expects HTTP Basic, not REST-style Bearer.
    # Invalid Authorization on a public repo does not fall back to anonymous.
    basic = base64.b64encode(f"x-access-token:{tok}".encode("ascii")).decode("ascii")
    return [
        "git",
        "-c",
        "http.version=HTTP/1.1",
        "-c",
        f"http.extraHeader=Authorization: Basic {basic}",
        "-c",
        "http.lowSpeedLimit=1000",
        "-c",
        "http.lowSpeedTime=60",
        "-c",
        "advice.detachedHead=false",
        "-c",
        "filter.lfs.smudge=",
        "-c",
        "filter.lfs.process=",
        "-c",
        "filter.lfs.required=false",
    ]


def parse_receive(stderr: str) -> tuple[float | None, float | None, int | None]:
    recv_bytes = None
    recv_bps = None
    nobj = None
    for m in RECEIVE_RE.finditer(stderr or ""):
        nobj = int(m.group(2))
        if m.group(3) and m.group(4):
            recv_bytes = bytes_unit(float(m.group(3)), m.group(4))
        if m.group(5) and m.group(6):
            recv_bps = bytes_unit(float(m.group(5)), m.group(6))
    return recv_bytes, recv_bps, nobj


def du_bytes(path: Path) -> int:
    total = 0
    for root, dirs, files in os.walk(path):
        for fn in files:
            fp = Path(root) / fn
            try:
                total += fp.stat().st_size
            except OSError:
                pass
    return total


def timed_curl(url: str, runs: int = 3) -> list[dict]:
    out = []
    dest = DISK_WORK / "curl_baseline.bin"
    dest.parent.mkdir(parents=True, exist_ok=True)
    for i in range(runs):
        if dest.exists():
            dest.unlink()
        t0 = time.perf_counter()
        cp = run(
            [
                "curl",
                "-fsSL",
                "--max-time",
                "120",
                "-A",
                USER_AGENT,
                "-o",
                str(dest),
                "-w",
                json.dumps(
                    {
                        "http_code": "%{http_code}",
                        "size_download": "%{size_download}",
                        "speed_download": "%{speed_download}",
                        "time_namelookup": "%{time_namelookup}",
                        "time_connect": "%{time_connect}",
                        "time_appconnect": "%{time_appconnect}",
                        "time_starttransfer": "%{time_starttransfer}",
                        "time_total": "%{time_total}",
                    }
                ),
                url,
            ],
            timeout=130,
        )
        wall = time.perf_counter() - t0
        meta = {}
        try:
            meta = json.loads(cp.stdout.strip() or "{}")
        except json.JSONDecodeError:
            meta = {"raw": cp.stdout}
        size = dest.stat().st_size if dest.exists() else 0
        rec = {
            "run": i + 1,
            "url": url,
            "ok": cp.returncode == 0 and size > 0,
            "wall_s": wall,
            "size_bytes": size,
            "curl": meta,
            "stderr": (cp.stderr or "")[-400:],
            "warmup": i == 0,
        }
        if size and wall > 0:
            rec["bps"] = size / wall
        out.append(rec)
        log(f"curl baseline run {i+1}: ok={rec['ok']} size={fmt_bytes(size)} wall={fmt_s(wall)} {fmt_bps(rec.get('bps'))}")
    if dest.exists():
        dest.unlink()
    return out


def snapshot_env() -> dict:
    uname = platform.uname()
    df = run(["df", "-hT", "/", "/dev/shm", "/www/wwwroot/dataset"])
    free = run(["free", "-h"])
    ping = run(["ping", "-c", "10", "-W", "2", "github.com"], timeout=30)
    curl_gh = run(
        [
            "curl",
            "-sS",
            "-o",
            "/dev/null",
            "-w",
            "http_code=%{http_code} namelookup=%{time_namelookup} connect=%{time_connect} appconnect=%{time_appconnect} ttfb=%{time_starttransfer} total=%{time_total} speed=%{speed_download}",
            "https://github.com",
        ]
    )
    git_v = run(["git", "--version"]).stdout.strip()
    py_v = sys.version.replace("\n", " ")
    nproc = os.cpu_count()
    tok = token()
    code, body, hdrs = gh_api("https://api.github.com/user", tok)
    login = body.get("login") if isinstance(body, dict) else None
    ip = run(["curl", "-fsS", "--max-time", "8", "https://api.ipify.org"]).stdout.strip()
    return {
        "captured_at": utc_now(),
        "host": uname.node,
        "kernel": f"{uname.system} {uname.release}",
        "machine": uname.machine,
        "cpu_count": nproc,
        "python": py_v,
        "git": git_v,
        "ip": ip,
        "github_login": login,
        "github_user_status": code,
        "rate_limit": {
            "limit": hdrs.get("x-ratelimit-limit"),
            "remaining": hdrs.get("x-ratelimit-remaining"),
            "resource": hdrs.get("x-ratelimit-resource"),
        },
        "df": df.stdout,
        "free": free.stdout,
        "ping_github": ping.stdout + ping.stderr,
        "curl_github": curl_gh.stdout,
        "seed": SEED,
        "shallow_timeout_s": SHALLOW_TIMEOUT_S,
        "full_timeout_s": FULL_TIMEOUT_S,
        "compresslevel": COMPRESS_LEVEL,
        "spool_max_bytes": SPOOL_MAX,
        "work_dir": str(WORK),
        "archive_dir": str(ARCHIVES),
        "lfs_skip_smudge": True,
    }


def append_metric(rec: dict) -> None:
    rec = dict(rec)
    rec.setdefault("logged_at", utc_now())
    with metrics_lock:
        with METRICS_PATH.open("a", encoding="utf-8") as fh:
            fh.write(json.dumps(rec, ensure_ascii=False) + "\n")


def clone_and_pack(task: dict) -> dict:
    tok = token()
    owner, repo = task["owner"], task["repo"]
    phase = task["phase"]
    mode = task["mode"]  # shallow | full
    full_name = f"{owner}/{repo}"
    aname = archive_name(owner, repo)
    dest = ARCHIVES / aname
    work_root = WORK if WORK.exists() or True else DISK_WORK
    try:
        work_root.mkdir(parents=True, exist_ok=True)
        probe = work_root / ".probe"
        probe.write_text("ok")
        probe.unlink()
    except OSError:
        work_root = DISK_WORK
        work_root.mkdir(parents=True, exist_ok=True)

    clone_parent = work_root / phase / f"{owner}__{repo}__{os.getpid()}__{threading.get_ident()}"
    if clone_parent.exists():
        shutil.rmtree(clone_parent, ignore_errors=True)
    clone_parent.mkdir(parents=True, exist_ok=True)
    clone_dir = clone_parent / repo

    result = {
        "ok": False,
        "phase": phase,
        "mode": mode,
        "full_name": full_name,
        "owner": owner,
        "repo": repo,
        "archive_name": aname,
        "archive_path": str(dest),
        "stratum": task.get("stratum"),
        "size_kib_api": task.get("size_kib"),
        "stars": task.get("stars"),
        "language": task.get("language"),
        "default_branch": task.get("default_branch"),
        "clone_s": None,
        "pack_s": None,
        "write_s": None,
        "cleanup_s": None,
        "total_s": None,
        "clone_bytes": None,
        "archive_bytes": None,
        "git_recv_bytes": None,
        "git_recv_bps": None,
        "git_objects": None,
        "clone_bps": None,
        "e2e_payload_bps": None,
        "e2e_archive_bps": None,
        "compression_ratio": None,
        "spool_spilled": None,
        "clone_dir": str(clone_dir),
        "work_root": str(work_root),
        "error": None,
        "error_class": None,
        "git_returncode": None,
        "attempt": 0,
    }
    t_all = time.perf_counter()
    try:
        if dest.exists():
            dest.unlink()
        branch = task.get("default_branch")
        timeout = FULL_TIMEOUT_S if mode == "full" else SHALLOW_TIMEOUT_S
        url = f"https://github.com/{owner}/{repo}.git"
        env = git_env()
        attempts = []
        cloned = False
        last_err = ""
        for attempt, use_branch in enumerate(([True, False] if branch else [False]), start=1):
            result["attempt"] = attempt
            if clone_dir.exists():
                shutil.rmtree(clone_dir, ignore_errors=True)
            cmd = git_cmd_prefix(tok) + ["clone", "--progress"]
            if mode == "shallow":
                cmd += ["--depth", "1", "--single-branch", "--no-tags"]
            if use_branch and branch:
                cmd += ["--branch", str(branch)]
            cmd += [url, str(clone_dir)]
            t0 = time.perf_counter()
            try:
                cp = run(cmd, timeout=timeout, env=env)
            except subprocess.TimeoutExpired:
                result["clone_s"] = time.perf_counter() - t0
                last_err = f"timeout after {timeout}s"
                attempts.append({"attempt": attempt, "timeout": True, "use_branch": use_branch})
                continue
            clone_s = time.perf_counter() - t0
            result["clone_s"] = clone_s
            result["git_returncode"] = cp.returncode
            stderr = cp.stderr or ""
            recv_bytes, recv_bps, nobj = parse_receive(stderr)
            result["git_recv_bytes"] = recv_bytes
            result["git_recv_bps"] = recv_bps
            result["git_objects"] = nobj
            attempts.append(
                {
                    "attempt": attempt,
                    "use_branch": use_branch,
                    "returncode": cp.returncode,
                    "clone_s": clone_s,
                    "stderr_tail": redact(stderr[-500:]),
                }
            )
            if cp.returncode == 0 and clone_dir.exists():
                cloned = True
                break
            last_err = redact((stderr or cp.stdout or "clone failed")[-500:])
            # retry only on missing branch
            low = (stderr or "").lower()
            if "remote branch" in low or "not found" in low or "could not find remote branch" in low:
                continue
            break
        result["attempts"] = attempts
        if not cloned:
            result["error"] = last_err or "clone failed"
            result["error_class"] = "clone"
            return result

        result["clone_bytes"] = du_bytes(clone_dir)
        if result["clone_s"] and result["clone_s"] > 0 and result["clone_bytes"]:
            result["clone_bps"] = result["clone_bytes"] / result["clone_s"]

        if not task.get("pack", True):
            result["pack_s"] = 0.0
            result["write_s"] = 0.0
            t3 = time.perf_counter()
            shutil.rmtree(clone_parent, ignore_errors=True)
            result["cleanup_s"] = time.perf_counter() - t3
            result["ok"] = bool(result.get("clone_bytes"))
            if not result["ok"]:
                result["error"] = "empty clone"
                result["error_class"] = "clone"
            return result

        t1 = time.perf_counter()
        spilled = False
        spool_path_holder = []

        def spilled_factory(*args, **kwargs):
            spilled_flag = {"v": False}

            def _open(name, prefix, suffix):
                spilled_flag["v"] = True
                fd, path = tempfile.mkstemp(prefix=prefix, suffix=suffix, dir=str(DISK_WORK))
                spool_path_holder.append(path)
                return os.fdopen(fd, "w+b")

            return spilled_flag, _open

        # SpooledTemporaryFile stays in RAM until max_size.
        spool = tempfile.SpooledTemporaryFile(max_size=SPOOL_MAX, dir=str(DISK_WORK))
        arcname = f"{owner}__{repo}"
        with tarfile.open(fileobj=spool, mode="w:gz", compresslevel=COMPRESS_LEVEL, dereference=False) as tar:
            tar.add(str(clone_dir), arcname=arcname, recursive=True)
        result["pack_s"] = time.perf_counter() - t1
        result["spool_spilled"] = getattr(spool, "_rolled", False)
        archive_bytes = spool.tell()
        result["archive_bytes"] = archive_bytes
        spool.seek(0)

        t2 = time.perf_counter()
        dest.parent.mkdir(parents=True, exist_ok=True)
        with dest.open("wb") as fh:
            shutil.copyfileobj(spool, fh, length=1024 * 1024)
            fh.flush()
            os.fsync(fh.fileno())
        result["write_s"] = time.perf_counter() - t2
        spool.close()
        for p in spool_path_holder:
            try:
                os.unlink(p)
            except OSError:
                pass

        if result["clone_bytes"] and archive_bytes:
            result["compression_ratio"] = archive_bytes / result["clone_bytes"]

        t3 = time.perf_counter()
        shutil.rmtree(clone_parent, ignore_errors=True)
        result["cleanup_s"] = time.perf_counter() - t3
        result["ok"] = dest.exists() and dest.stat().st_size > 0
        if not result["ok"]:
            result["error"] = "archive missing after write"
            result["error_class"] = "write"
        return result
    except Exception as e:
        result["error"] = redact(f"{type(e).__name__}: {e}")
        result["error_class"] = "exception"
        result["traceback"] = redact(traceback.format_exc())[-1500:]
        return result
    finally:
        result["total_s"] = time.perf_counter() - t_all
        if result.get("clone_bytes") and result.get("total_s"):
            result["e2e_payload_bps"] = result["clone_bytes"] / result["total_s"]
        if result.get("archive_bytes") and result.get("total_s"):
            result["e2e_archive_bps"] = result["archive_bytes"] / result["total_s"]
        shutil.rmtree(clone_parent, ignore_errors=True)
        append_metric(result)
        flag = "OK" if result.get("ok") else "FAIL"
        log(
            f"{flag} {phase}/{mode} {full_name} stratum={result.get('stratum')} "
            f"clone={fmt_s(result.get('clone_s'))} pack={fmt_s(result.get('pack_s'))} "
            f"payload={fmt_bytes(result.get('clone_bytes'))} archive={fmt_bytes(result.get('archive_bytes'))} "
            f"clone_bps={fmt_bps(result.get('clone_bps'))} err={result.get('error_class') or '-'}"
        )


def stats_block(rows: list[dict], label: str) -> dict:
    ok = [r for r in rows if r.get("ok")]
    fail = [r for r in rows if not r.get("ok")]
    def col(key: str) -> list[float]:
        return [float(r[key]) for r in ok if r.get(key) is not None and r.get(key) >= 0]

    clone_bps = col("clone_bps")
    clone_s = col("clone_s")
    pack_s = col("pack_s")
    write_s = col("write_s")
    total_s = col("total_s")
    payload = col("clone_bytes")
    archive = col("archive_bytes")
    e2e = col("e2e_payload_bps")
    git_bps = col("git_recv_bps")
    ci = bootstrap_mean_ci(clone_bps) if clone_bps else None
    wall = sum(total_s) if total_s else 0.0
    payload_sum = sum(payload)
    archive_sum = sum(archive)
    return {
        "label": label,
        "n": len(rows),
        "n_ok": len(ok),
        "n_fail": len(fail),
        "success_rate": (len(ok) / len(rows)) if rows else None,
        "clone_s": _dist(clone_s),
        "pack_s": _dist(pack_s),
        "write_s": _dist(write_s),
        "total_s": _dist(total_s),
        "clone_bytes": _dist(payload),
        "archive_bytes": _dist(archive),
        "clone_bps": _dist(clone_bps),
        "git_recv_bps": _dist(git_bps),
        "e2e_payload_bps": _dist(e2e),
        "clone_bps_mean_ci95": ci,
        "sum_payload_bytes": payload_sum,
        "sum_archive_bytes": archive_sum,
        "sum_clone_s": sum(clone_s) if clone_s else 0.0,
        "sum_total_s": wall,
        "aggregate_clone_bps": (payload_sum / sum(clone_s)) if clone_s and sum(clone_s) > 0 else None,
        "aggregate_e2e_bps": (payload_sum / wall) if wall > 0 else None,
        "failures": [
            {"full_name": r.get("full_name"), "error_class": r.get("error_class"), "error": (r.get("error") or "")[:300]}
            for r in fail
        ],
    }


def _dist(xs: list[float]) -> dict | None:
    if not xs:
        return None
    return {
        "n": len(xs),
        "mean": mean(xs),
        "stdev": stdev(xs),
        "min": min(xs),
        "p10": percentile(xs, 10),
        "p25": percentile(xs, 25),
        "p50": percentile(xs, 50),
        "p75": percentile(xs, 75),
        "p90": percentile(xs, 90),
        "p95": percentile(xs, 95),
        "max": max(xs),
        "sum": sum(xs),
    }


def md_dist(d: dict | None, kind: str) -> str:
    if not d:
        return "n/a"
    if kind == "bytes":
        f = fmt_bytes
    elif kind == "bps":
        f = fmt_bps
    else:
        f = fmt_s
    return (
        f"mean {f(d['mean'])} · median {f(d['p50'])} · p90 {f(d['p90'])} · "
        f"p95 {f(d['p95'])} · min {f(d['min'])} · max {f(d['max'])}"
        + (f" · sd {f(d['stdev'])}" if d.get("stdev") is not None else "")
    )


def run_phase(name: str, tasks: list[dict], workers: int) -> tuple[list[dict], float]:
    log(f"=== phase {name} n={len(tasks)} workers={workers} ===")
    t0 = time.perf_counter()
    results: list[dict] = []
    if workers <= 1:
        for t in tasks:
            results.append(clone_and_pack(t))
    else:
        with concurrent.futures.ThreadPoolExecutor(max_workers=workers) as ex:
            futs = [ex.submit(clone_and_pack, t) for t in tasks]
            for fut in concurrent.futures.as_completed(futs):
                results.append(fut.result())
    wall = time.perf_counter() - t0
    ok = sum(1 for r in results if r.get("ok"))
    payload = sum(r.get("clone_bytes") or 0 for r in results if r.get("ok"))
    log(
        f"=== phase {name} done wall={fmt_s(wall)} ok={ok}/{len(results)} "
        f"payload={fmt_bytes(payload)} aggregate={fmt_bps(payload / wall if wall else None)} ==="
    )
    return results, wall


def pick_stratified(by_stratum: dict[str, list[dict]], phase: str, mode: str, rng: random.Random, used: set[str]) -> list[dict]:
    n = PER_STRATUM[phase if phase != "full" else "full"]
    out = []
    names = [s[0] for s in STRATA]
    if mode == "full":
        names = [s for s in names if s != "XXL"]
    for sname in names:
        pool = [r for r in by_stratum.get(sname, []) if r["full_name"] not in used]
        rng.shuffle(pool)
        take = pool[:n]
        for rec in take:
            used.add(rec["full_name"])
            out.append(
                {
                    "owner": rec["full_name"].split("/", 1)[0],
                    "repo": rec["full_name"].split("/", 1)[1],
                    "full_name": rec["full_name"],
                    "phase": "full" if mode == "full" else phase,
                    "mode": mode,
                    "stratum": sname,
                    "size_kib": rec.get("size_kib"),
                    "stars": rec.get("stargazers_count"),
                    "language": rec.get("language"),
                    "default_branch": rec.get("default_branch"),
                }
            )
    rng.shuffle(out)
    return out


def extrapolate(manifest_n: int, api_rows: list[dict], seq_stats: dict, seq_wall: float, par4: dict | None, par4_wall: float | None) -> dict:
    # Size distribution from API sample, success-weighted mean payload from sequential clones.
    usable = [r for r in api_rows if r.get("size_kib") is not None]
    n_api = len(api_rows)
    n_ok_api = len(usable)
    n_404 = sum(1 for r in api_rows if r.get("http_status") == 404)
    n_too_large = sum(1 for r in api_rows if stratum_of(r.get("size_kib")) == "TOO_LARGE")
    frac_alive = n_ok_api / n_api if n_api else None
    frac_404 = n_404 / n_api if n_api else None
    mean_payload = (seq_stats.get("clone_bytes") or {}).get("mean")
    mean_clone_s = (seq_stats.get("clone_s") or {}).get("mean")
    mean_total_s = (seq_stats.get("total_s") or {}).get("mean")
    agg_bps = seq_stats.get("aggregate_clone_bps")
    expected_alive = manifest_n * frac_alive if frac_alive is not None else None
    out = {
        "manifest_n": manifest_n,
        "api_n": n_api,
        "api_alive": n_ok_api,
        "api_404": n_404,
        "api_too_large": n_too_large,
        "frac_alive": frac_alive,
        "frac_404": frac_404,
        "expected_alive": expected_alive,
        "expected_too_large": (n_too_large / n_api * manifest_n) if n_api else None,
        "seq_mean_payload_bytes": mean_payload,
        "seq_mean_clone_s": mean_clone_s,
        "seq_mean_total_s": mean_total_s,
        "seq_aggregate_clone_bps": agg_bps,
    }
    if expected_alive and mean_payload:
        out["expected_payload_bytes"] = expected_alive * mean_payload
        out["expected_archive_bytes"] = expected_alive * ((seq_stats.get("archive_bytes") or {}).get("mean") or 0)
    if expected_alive and mean_total_s:
        out["eta_serial_s"] = expected_alive * mean_total_s
    if expected_alive and mean_clone_s:
        out["eta_serial_clone_only_s"] = expected_alive * mean_clone_s
    if par4 and par4_wall and par4.get("n_ok"):
        # scale observed par4 wall by remaining work, using per-repo mean wall of that phase
        per = par4_wall / par4["n"] if par4["n"] else None
        if per and expected_alive:
            # par4 processed n repos in par4_wall; throughput in repos/s
            rps = par4["n"] / par4_wall
            out["par4_repos_per_s"] = rps
            out["eta_par4_s"] = expected_alive / rps if rps else None
            pay = par4.get("sum_payload_bytes") or 0
            out["par4_aggregate_payload_bps"] = pay / par4_wall if par4_wall else None
    return out


def write_report(env: dict, sample_meta: dict, phases: dict, extra: dict) -> None:
    def dget(block, *keys):
        cur = block
        for k in keys:
            if cur is None:
                return None
            cur = cur.get(k) if isinstance(cur, dict) else None
        return cur

    lines = []
    a = lines.append
    a("# Git Clone 下载速度测试报告")
    a("")
    a(f"- 生成时间（UTC）：`{utc_now()}`")
    a(f"- 清单文件：`{MANIFEST}`")
    a(f"- 随机种子：`{SEED}`（可复现抽样）")
    a(f"- GitHub 认证账号：`{env.get('github_login')}`")
    a(f"- 主机：`{env.get('host')}` · `{env.get('kernel')}` · {env.get('cpu_count')} CPU · IP `{env.get('ip')}`")
    a(f"- Git：`{env.get('git')}`")
    a("")
    a("## 1. 测试目的与结论摘要")
    a("")
    seq = phases.get("seq", {}).get("stats") or {}
    par4 = phases.get("par4", {}).get("stats") or {}
    par8 = phases.get("par8", {}).get("stats") or {}
    full = phases.get("full", {}).get("stats") or {}
    a("本测试评估的是 **git clone（HTTPS）+ 内存 tar.gz 落盘** 这条真实流水线的速度，而不是 GitHub zipball / codeload 归档接口。")
    a("")
    if seq.get("aggregate_clone_bps"):
        ci = seq.get("clone_bps_mean_ci95")
        ci_s = f"（均值 95% bootstrap CI：{fmt_bps(ci[0])} – {fmt_bps(ci[1])}）" if ci else ""
        a(
            f"**顺序 shallow clone 的聚合下载吞吐量为 {fmt_bps(seq['aggregate_clone_bps'])}**"
            f"（payload 总和 / clone 时间总和）{ci_s}。"
        )
    if par4.get("n"):
        wall = phases["par4"]["wall_s"]
        agg = (par4.get("sum_payload_bytes") or 0) / wall if wall else None
        a(f"**并发 4 的墙钟聚合吞吐量为 {fmt_bps(agg)}**（payload 总和 / 阶段墙钟）。")
    if par8.get("n"):
        wall = phases["par8"]["wall_s"]
        agg = (par8.get("sum_payload_bytes") or 0) / wall if wall else None
        a(f"**并发 8 的墙钟聚合吞吐量为 {fmt_bps(agg)}**。")
    a("")
    a("关键数字：")
    a("")
    a("| 阶段 | 成功/总数 | clone 时间中位数 | clone 吞吐中位数 | 聚合 clone 吞吐 | 阶段墙钟 | 阶段 payload |")
    a("|---|---:|---:|---:|---:|---:|---:|")
    for key, label in (("seq", "顺序 shallow x1"), ("par4", "并发 shallow x4"), ("par8", "并发 shallow x8"), ("full", "顺序 full clone")):
        ph = phases.get(key) or {}
        st = ph.get("stats") or {}
        if not st:
            continue
        med_s = fmt_s(dget(st, "clone_s", "p50"))
        med_bps = fmt_bps(dget(st, "clone_bps", "p50"))
        agg = fmt_bps(st.get("aggregate_clone_bps"))
        wall = fmt_s(ph.get("wall_s"))
        a(
            f"| {label} | {st.get('n_ok')}/{st.get('n')} | {med_s} | {med_bps} | {agg} | {wall} | {fmt_bytes(st.get('sum_payload_bytes'))} |"
        )
    a("")
    a("## 2. 方法")
    a("")
    a("### 2.1 为什么不做 54,349 个仓库的全量 clone")
    a("")
    a("清单体量太大，全量跑一遍只能得到一个不可复现的墙钟数字，无法分离网络、Git 协议开销、压缩和落盘，也无法给出不确定度。本测试用 **有种子的分层随机抽样** 估计速度分布，再用样本率外推全量。")
    a("")
    a("### 2.2 抽样")
    a("")
    a(f"1. 读取清单全部 `{sample_meta.get('manifest_n')}` 条有效 `full_name`。")
    a(f"2. 以种子 `{SEED}` 无放回抽取 `{sample_meta.get('api_requested')}` 个仓库，调用 GitHub REST `GET /repos/{{owner}}/{{repo}}` 取 `size`（KiB）。")
    a("3. 按 GitHub `size` 分层：")
    a("")
    a("| 层 | GitHub size | 含义 |")
    a("|---|---|---|")
    a("| XS | < 100 KiB | 极小 |")
    a("| S | 100 KiB – 1 MiB | 小 |")
    a("| M | 1 – 10 MiB | 中 |")
    a("| L | 10 – 50 MiB | 大 |")
    a("| XL | 50 – 200 MiB | 很大 |")
    a("| XXL | 200 MiB – 1 GiB | 极大 |")
    a("| TOO_LARGE | ≥ 1 GiB | 测速中排除，只统计占比 |")
    a("")
    a("4. 每层再无放回切成互斥任务集，避免同一仓库被 GitHub/CDN 缓存后污染后续阶段：")
    a("")
    a(f"- 顺序 shallow：每层 {PER_STRATUM['seq']} 个")
    a(f"- 并发 4 shallow：每层 {PER_STRATUM['par4']} 个")
    a(f"- 并发 8 shallow：每层 {PER_STRATUM['par8']} 个")
    a(f"- 顺序 full clone：每层 {PER_STRATUM['full']} 个（不含 XXL）")
    a("- 顺序阶段前 2 个成功 clone 记为 warmup，**不进入统计**")
    a("")
    a("### 2.3 每个仓库的流水线")
    a("")
    a("1. `git clone` 到 tmpfs（`/dev/shm`），减少磁盘对「下载速度」的污染。")
    a("2. 主测模式：`--depth 1 --single-branch --no-tags`，对应 5 万级清单的可执行采集方式。")
    a("3. 对照模式：不带 depth 的 full clone，量化历史对象对时间和体积的放大。")
    a("4. `tarfile` + `SpooledTemporaryFile(max_size=256MiB)` 在内存中 `w:gz`（compresslevel=6），再 `fsync` 写到数据集盘。")
    a("5. 压缩包文件名：`{owner}__{reponame}.tar.gz`（两个下划线）。")
    a("6. 立刻删除工作树，不在磁盘上保留未打包 clone。")
    a("7. `GIT_LFS_SKIP_SMUDGE=1`，避免偶然 LFS 对象把单次样本变成异常值。")
    a("")
    a("认证走 `http.extraHeader=Authorization: Bearer …`，clone URL 不含 token，日志已做脱敏。")
    a("")
    a("### 2.4 指标定义")
    a("")
    a("| 指标 | 定义 |")
    a("|---|---|")
    a("| clone_s | `git clone` 墙钟 |")
    a("| clone_bytes | clone 工作树 `st_size` 总和（含 `.git`） |")
    a("| clone_bps | clone_bytes / clone_s。这是 **含 Git 握手、排队、解包** 的有效吞吐，不是纯 TCP |")
    a("| git_recv_bps | git 进度条里的 Receiving objects 速率（若有） |")
    a("| pack_s / write_s | 内存 tar.gz / 落盘+fsync |")
    a("| aggregate clone bps | Σ clone_bytes / Σ clone_s，体积加权，更接近「搬数据」 |")
    a("| 阶段墙钟吞吐 | Σ clone_bytes / 阶段 wall，并发时才有意义 |")
    a("| 均值 95% CI | 对 **每仓 clone_bps** 做 2000 次 bootstrap，反映小仓延迟对均值的拉动 |")
    a("")
    a("小仓会严重拉低算术平均吞吐（固定 RTT/握手开销），所以报告同时给 **中位数、p90、体积加权聚合**。分层就是为了不让样本全是小仓。")
    a("")
    a("## 3. 环境与网络基线")
    a("")
    a("```")
    a((env.get("free") or "").strip())
    a("```")
    a("")
    a("```")
    a((env.get("df") or "").strip())
    a("```")
    a("")
    a("GitHub ping / TLS：")
    a("")
    a("```")
    a((env.get("ping_github") or "").strip()[-1200:])
    a("```")
    a("")
    a("```")
    a(env.get("curl_github") or "")
    a("```")
    a("")
    curl_runs = extra.get("curl_baseline") or []
    if curl_runs:
        a("### 3.1 HTTPS 归档下载基线（codeload，非 git 协议）")
        a("")
        a("用于给 git 协议吞吐一个对照上限。第一次视为 warmup。")
        a("")
        a("| run | warmup | size | wall | 吞吐 | http |")
        a("|---:|:---:|---:|---:|---:|---:|")
        for r in curl_runs:
            curl = r.get("curl") or {}
            a(
                f"| {r.get('run')} | {'yes' if r.get('warmup') else 'no'} | {fmt_bytes(r.get('size_bytes'))} | "
                f"{fmt_s(r.get('wall_s'))} | {fmt_bps(r.get('bps'))} | {curl.get('http_code')} |"
            )
        a("")
    a("## 4. 清单与 API 抽样")
    a("")
    a(f"- 清单有效仓库数：{sample_meta.get('manifest_n')}")
    a(f"- API 抽样请求数：{sample_meta.get('api_requested')}")
    a(f"- API 成功拿到 size：{sample_meta.get('api_alive')}")
    a(f"- HTTP 404（已删/改名/私有）：{sample_meta.get('api_404')}")
    a(f"- 其他 API 失败：{sample_meta.get('api_other_fail')}")
    a(f"- ≥ 1 GiB 排除：{sample_meta.get('api_too_large')}")
    a("")
    a("| 层 | API 样本数 | 占比 | size 中位数 (KiB) |")
    a("|---|---:|---:|---:|")
    for row in sample_meta.get("stratum_table") or []:
        a(f"| {row['name']} | {row['n']} | {row['frac']:.1%} | {row['median_kib']} |")
    a("")
    a("各测速阶段实际抽到的仓库见 `speed_test/reports/sample.json`。")
    a("")
    a("## 5. 分阶段结果")
    a("")
    for key, title in (
        ("seq", "5.1 顺序 shallow clone（并发 = 1）"),
        ("par4", "5.2 并发 4 shallow clone"),
        ("par8", "5.3 并发 8 shallow clone"),
        ("full", "5.4 顺序 full clone 对照"),
    ):
        ph = phases.get(key) or {}
        st = ph.get("stats") or {}
        if not st:
            continue
        a(f"### {title}")
        a("")
        a(f"- 墙钟：{fmt_s(ph.get('wall_s'))}")
        a(f"- 成功：{st.get('n_ok')} / {st.get('n')}（{ (st.get('success_rate') or 0)*100:.1f}% ）")
        a(f"- payload 合计：{fmt_bytes(st.get('sum_payload_bytes'))}")
        a(f"- 归档合计：{fmt_bytes(st.get('sum_archive_bytes'))}")
        a(f"- clone 时间：{md_dist(st.get('clone_s'), 's')}")
        a(f"- clone 吞吐：{md_dist(st.get('clone_bps'), 'bps')}")
        a(f"- git 自报 Receiving 速率：{md_dist(st.get('git_recv_bps'), 'bps')}")
        a(f"- 体积加权聚合 clone 吞吐：{fmt_bps(st.get('aggregate_clone_bps'))}")
        wall = ph.get("wall_s") or 0
        wall_bps = (st.get("sum_payload_bytes") or 0) / wall if wall else None
        a(f"- 阶段墙钟吞吐：{fmt_bps(wall_bps)}")
        a(f"- 内存打包：{md_dist(st.get('pack_s'), 's')}")
        a(f"- 落盘+fsync：{md_dist(st.get('write_s'), 's')}")
        a(f"- 端到端（clone+pack+write+cleanup）：{md_dist(st.get('total_s'), 's')}")
        ci = st.get("clone_bps_mean_ci95")
        if ci:
            a(f"- 每仓 clone_bps 均值 95% CI：{fmt_bps(ci[0])} – {fmt_bps(ci[1])}")
        a("")
        by = ph.get("by_stratum") or {}
        if by:
            a("| 层 | 成功/N | payload 中位 | clone 中位 | clone 吞吐中位 | 聚合吞吐 |")
            a("|---|---:|---:|---:|---:|---:|")
            for sname, _lo, _hi in STRATA:
                bs = by.get(sname)
                if not bs:
                    continue
                a(
                    f"| {sname} | {bs.get('n_ok')}/{bs.get('n')} | {fmt_bytes(dget(bs,'clone_bytes','p50'))} | "
                    f"{fmt_s(dget(bs,'clone_s','p50'))} | {fmt_bps(dget(bs,'clone_bps','p50'))} | "
                    f"{fmt_bps(bs.get('aggregate_clone_bps'))} |"
                )
            a("")
        fails = st.get("failures") or []
        if fails:
            a("失败条目：")
            a("")
            for f in fails:
                err = (f.get("error") or "").replace("\n", " ")[:200]
                a(f"- `{f.get('full_name')}` · {f.get('error_class')} · {err}")
            a("")

    a("## 6. 分层对速度的影响")
    a("")
    a("Git clone 有接近固定的 HTTPS/Git 握手成本。小仓的 clone_bps 会被延迟主导，大仓才接近链路带宽。因此：")
    a("")
    a("- **不要只用全样本平均值描述「下载速度」。**")
    a("- 体积加权聚合吞吐更接近搬迁 54k 仓时「有效带宽」的概念。")
    a("- 并发提升的是墙钟吞吐；单连接速率不一定变。")
    a("")

    ext = extra.get("extrapolation") or {}
    a("## 7. 外推到全量清单")
    a("")
    a("外推假设：API 样本的存活率、过大仓占比，以及顺序/并发阶段观测到的每仓时间和吞吐，对剩余清单近似成立。这是 **点估计 + 方法局限**，不是保证。")
    a("")
    a(f"- 清单 N = {ext.get('manifest_n')}")
    a(f"- API 存活率 = {ext.get('frac_alive')}")
    a(f"- 期望仍存在的仓库 ≈ {ext.get('expected_alive')}")
    a(f"- 期望 ≥1GiB、本测试排除的仓库 ≈ {ext.get('expected_too_large')}")
    a(f"- 顺序阶段样本平均 payload = {fmt_bytes(ext.get('seq_mean_payload_bytes'))}")
    a(f"- 期望 shallow payload 总量 ≈ {fmt_bytes(ext.get('expected_payload_bytes'))}")
    a(f"- 期望归档总量 ≈ {fmt_bytes(ext.get('expected_archive_bytes'))}")
    a(f"- 若严格串行（含打包落盘）：ETA ≈ {fmt_s(ext.get('eta_serial_s'))}")
    a(f"- 若串行只算 clone：ETA ≈ {fmt_s(ext.get('eta_serial_clone_only_s'))}")
    a(f"- 按并发 4 观测到的仓库/秒外推：ETA ≈ {fmt_s(ext.get('eta_par4_s'))}")
    a(f"- 并发 4 墙钟 payload 吞吐：{fmt_bps(ext.get('par4_aggregate_payload_bps'))}")
    a("")
    a("并发 8 若已出现失败率上升、429 或速率不再线性增加，则生产采集应停在饱和点附近，而不是盲目加线程。")
    a("")
    a("## 8. 科学局限")
    a("")
    a("1. **时间点**：GitHub 边缘节点负载随小时变化，本报告是单次观测。")
    a("2. **浅克隆 ≠ 全历史**：主数字是 `--depth 1`。full clone 对照用来标定历史放大系数。")
    a("3. **GitHub `size` 含完整对象库**，shallow 工作树通常更小，分层用 API size 只保证「大仓/小仓」排序，不保证字节数相等。")
    a("4. **LFS 被跳过**。若生产需要 LFS 实文件，吞吐会被个别仓拖死，需单独策略。")
    a("5. **删除/改名/私有仓**：404 比例会进入外推，但不会变成归档。")
    a("6. **二次 clone 会偏快**：阶段之间仓库互斥，但 GitHub 对热门仓仍可能命中服务端缓存。")
    a("7. **内存打包**：超过 256 MiB 的归档会溢到磁盘 spool，指标里有 `spool_spilled`。")
    a("8. **token**：classic PAT 仅通过环境变量注入，未写入仓库文件。该 token 已在对话中出现，用完应立刻轮换。")
    a("")
    a("## 9. 对后续全量采集的建议")
    a("")
    a("- 生产默认：`git clone --depth 1 --single-branch --no-tags`，认证 header，跳过 LFS smudge。")
    a("- 并发先用 4，看失败率和墙钟吞吐；若稳定再试 8。")
    a("- 单仓超时：shallow 180s，full 420s；低速断开 `http.lowSpeedLimit=1000` / `lowSpeedTime=60`。")
    a("- 归档名严格执行 `owner__reponame.tar.gz`；工作树打进包内的顶层目录同名（无 `.tar.gz`）。")
    a("- 先 API 探活再 clone，避免给 404 付握手成本。")
    a("- ≥1 GiB 仓单独队列，不要和中小仓混在同一个线程池。")
    a("")
    a("## 10. 产物路径")
    a("")
    a(f"- 归档目录：`{ARCHIVES}`")
    a(f"- 每仓指标：`{METRICS_PATH}`")
    a(f"- 抽样清单：`{SAMPLE_PATH}`")
    a(f"- 环境快照：`{ENV_PATH}`")
    a(f"- JSON 摘要：`{SUMMARY_PATH}`")
    a(f"- 运行日志：`{LOG_PATH}`")
    a("")
    a("归档命名抽查：")
    a("")
    names = extra.get("archive_names") or []
    for n in names[:15]:
        a(f"- `{n}`")
    if len(names) > 15:
        a(f"- … 共 {len(names)} 个")
    a("")
    REPORT_PATH.write_text("\n".join(lines) + "\n", encoding="utf-8")


def main() -> int:
    parser = argparse.ArgumentParser()
    parser.add_argument("--skip-api", action="store_true")
    args = parser.parse_args()
    for p in (ARCHIVES, WORK, DISK_WORK, REPORTS, LOGS):
        p.mkdir(parents=True, exist_ok=True)
    if METRICS_PATH.exists():
        METRICS_PATH.unlink()

    tok = token()
    log("snapshot environment")
    env = snapshot_env()
    ENV_PATH.write_text(json.dumps(env, ensure_ascii=False, indent=2), encoding="utf-8")

    log("HTTPS codeload baseline")
    curl_baseline = timed_curl("https://codeload.github.com/git/git/tar.gz/refs/tags/v2.43.0", runs=3)

    log(f"load manifest {MANIFEST}")
    rows = load_manifest(MANIFEST)
    log(f"manifest valid rows: {len(rows)}")
    rng = random.Random(SEED)
    api_pool = rows[:]
    rng.shuffle(api_pool)
    api_targets = api_pool[:API_SAMPLE_N]

    log(f"fetch GitHub size for {len(api_targets)} repos")
    api_rows: list[dict] = []
    with concurrent.futures.ThreadPoolExecutor(max_workers=12) as ex:
        futs = [ex.submit(fetch_size, rec, tok) for rec in api_targets]
        for i, fut in enumerate(concurrent.futures.as_completed(futs), 1):
            api_rows.append(fut.result())
            if i % 50 == 0:
                log(f"api progress {i}/{len(api_targets)}")

    by_stratum: dict[str, list[dict]] = {s[0]: [] for s in STRATA}
    by_stratum["TOO_LARGE"] = []
    by_stratum["UNKNOWN"] = []
    n404 = n_other = n_alive = n_large = 0
    for r in api_rows:
        st = stratum_of(r.get("size_kib"))
        r["stratum"] = st
        if r.get("http_status") == 404:
            n404 += 1
        elif r.get("size_kib") is None:
            n_other += 1
        else:
            n_alive += 1
            if st == "TOO_LARGE":
                n_large += 1
        by_stratum.setdefault(st, []).append(r)

    stratum_table = []
    for name, lo, hi in STRATA + [("TOO_LARGE", SKIP_SIZE_KIB, 10**18), ("UNKNOWN", -1, -1)]:
        grp = by_stratum.get(name) or []
        sizes = [g["size_kib"] for g in grp if g.get("size_kib") is not None]
        stratum_table.append(
            {
                "name": name,
                "n": len(grp),
                "frac": (len(grp) / len(api_rows)) if api_rows else 0,
                "median_kib": int(percentile(sizes, 50)) if sizes else None,
            }
        )

    used: set[str] = set()
    tasks_seq = pick_stratified(by_stratum, "seq", "shallow", rng, used)
    tasks_par4 = pick_stratified(by_stratum, "par4", "shallow", rng, used)
    tasks_par8 = pick_stratified(by_stratum, "par8", "shallow", rng, used)
    tasks_full = pick_stratified(by_stratum, "full", "full", rng, used)

    sample_meta = {
        "manifest_n": len(rows),
        "api_requested": len(api_targets),
        "api_alive": n_alive,
        "api_404": n404,
        "api_other_fail": n_other,
        "api_too_large": n_large,
        "stratum_table": stratum_table,
        "phases": {
            "seq": [t["full_name"] for t in tasks_seq],
            "par4": [t["full_name"] for t in tasks_par4],
            "par8": [t["full_name"] for t in tasks_par8],
            "full": [t["full_name"] for t in tasks_full],
        },
        "n_tasks": {
            "seq": len(tasks_seq),
            "par4": len(tasks_par4),
            "par8": len(tasks_par8),
            "full": len(tasks_full),
        },
    }
    SAMPLE_PATH.write_text(
        json.dumps({"sample_meta": sample_meta, "api_rows": api_rows, "tasks": {
            "seq": tasks_seq, "par4": tasks_par4, "par8": tasks_par8, "full": tasks_full
        }}, ensure_ascii=False, indent=2),
        encoding="utf-8",
    )
    log(f"tasks seq={len(tasks_seq)} par4={len(tasks_par4)} par8={len(tasks_par8)} full={len(tasks_full)}")

    # warmup: two extra tiny-ish clones from leftover XS/S, discarded
    warmup_pool = [r for s in ("XS", "S", "M") for r in by_stratum.get(s, []) if r["full_name"] not in used]
    rng.shuffle(warmup_pool)
    warmup_tasks = []
    for rec in warmup_pool[:2]:
        used.add(rec["full_name"])
        warmup_tasks.append(
            {
                "owner": rec["full_name"].split("/", 1)[0],
                "repo": rec["full_name"].split("/", 1)[1],
                "full_name": rec["full_name"],
                "phase": "warmup",
                "mode": "shallow",
                "stratum": rec.get("stratum"),
                "size_kib": rec.get("size_kib"),
                "stars": rec.get("stargazers_count"),
                "language": rec.get("language"),
                "default_branch": rec.get("default_branch"),
            }
        )
    warmup_results, warmup_wall = run_phase("warmup", warmup_tasks, 1)

    seq_results, seq_wall = run_phase("seq", tasks_seq, 1)
    par4_results, par4_wall = run_phase("par4", tasks_par4, 4)
    par8_results, par8_wall = run_phase("par8", tasks_par8, 8)
    full_results, full_wall = run_phase("full", tasks_full, 1)

    def with_stratum(results: list[dict]) -> dict:
        out = {}
        for sname, _lo, _hi in STRATA:
            subset = [r for r in results if r.get("stratum") == sname]
            if subset:
                out[sname] = stats_block(subset, sname)
        return out

    phases = {
        "warmup": {"wall_s": warmup_wall, "stats": stats_block(warmup_results, "warmup"), "results": warmup_results},
        "seq": {"wall_s": seq_wall, "stats": stats_block(seq_results, "seq"), "by_stratum": with_stratum(seq_results), "results": seq_results},
        "par4": {"wall_s": par4_wall, "stats": stats_block(par4_results, "par4"), "by_stratum": with_stratum(par4_results), "results": par4_results},
        "par8": {"wall_s": par8_wall, "stats": stats_block(par8_results, "par8"), "by_stratum": with_stratum(par8_results), "results": par8_results},
        "full": {"wall_s": full_wall, "stats": stats_block(full_results, "full"), "by_stratum": with_stratum(full_results), "results": full_results},
    }
    ext = extrapolate(len(rows), api_rows, phases["seq"]["stats"], seq_wall, phases["par4"]["stats"], par4_wall)
    archives = sorted(p.name for p in ARCHIVES.glob("*.tar.gz"))
    extra = {"curl_baseline": curl_baseline, "extrapolation": ext, "archive_names": archives}

    # JSON summary cannot include full results with potential huge stderr; keep stats + compact results
    def compact(rs: list[dict]) -> list[dict]:
        keys = [
            "ok","phase","mode","full_name","archive_name","stratum","size_kib_api","stars","language",
            "clone_s","pack_s","write_s","cleanup_s","total_s","clone_bytes","archive_bytes",
            "git_recv_bytes","git_recv_bps","git_objects","clone_bps","e2e_payload_bps",
            "compression_ratio","spool_spilled","error","error_class","git_returncode",
        ]
        return [{k: r.get(k) for k in keys} for r in rs]

    summary = {
        "status": "complete",
        "completed_at": utc_now(),
        "env_brief": {k: env.get(k) for k in ["captured_at","host","cpu_count","git","ip","github_login","rate_limit"]},
        "sample_meta": sample_meta,
        "phases": {
            k: {
                "wall_s": v["wall_s"],
                "stats": {kk: vv for kk, vv in v["stats"].items() if kk != "failures"} | {"failures": v["stats"].get("failures")},
                "by_stratum": v.get("by_stratum"),
                "results": compact(v["results"]),
            }
            for k, v in phases.items()
        },
        "extrapolation": ext,
        "curl_baseline": curl_baseline,
        "n_archives": len(archives),
    }
    SUMMARY_PATH.write_text(json.dumps(summary, ensure_ascii=False, indent=2), encoding="utf-8")
    write_report(env, sample_meta, phases, extra)
    log(f"report written {REPORT_PATH}")
    log("STATUS complete")
    return 0


if __name__ == "__main__":
    try:
        raise SystemExit(main())
    except Exception:
        log("STATUS failed")
        log(redact(traceback.format_exc()))
        try:
            SUMMARY_PATH.write_text(
                json.dumps({"status": "failed", "completed_at": utc_now(), "error": redact(traceback.format_exc())}, indent=2),
                encoding="utf-8",
            )
        except Exception:
            pass
        raise
