#!/usr/bin/env python3
from __future__ import annotations

import argparse
import csv
import gzip
import json
import os
import re
import shutil
import signal
import subprocess
import sys
import time
from pathlib import Path
from typing import Any

RIR_TAL = {
    "afrinic": "afrinic.tal",
    "apnic": "apnic-rfc7730-https.tal",
    "arin": "arin.tal",
    "lacnic": "lacnic.tal",
    "ripe": "ripe-ncc.tal",
}

TAL_URLS = {
    "afrinic": "https://rpki.afrinic.net/tal/afrinic.tal",
    "apnic": "https://tal.apnic.net/apnic.tal",
    "arin": "https://www.arin.net/resources/manage/rpki/arin.tal",
    "lacnic": "https://www.lacnic.net/innovaportal/file/4983/1/lacnic.tal",
    "ripe": "https://tal.rpki.ripe.net/ripe-ncc.tal",
}

PROFILES = {
    "ours-no-cache": "ours-rp",
    "ours-pp-object-cache": "ours-rp",
    "routinator-latest-release": "routinator",
    "rpki-client-latest-release": "rpki-client",
    "fort-latest-release": "fort",
}

ALIASES = {
    "ours": "ours-pp-object-cache",
    "ours-rp": "ours-pp-object-cache",
    "ours-cached": "ours-pp-object-cache",
    "ours-baseline": "ours-no-cache",
    "ours-no-cache": "ours-no-cache",
    "ours-pp-object-cache": "ours-pp-object-cache",
    "routinator": "routinator-latest-release",
    "routinator-latest-release": "routinator-latest-release",
    "rpki-client": "rpki-client-latest-release",
    "rpki-client-latest-release": "rpki-client-latest-release",
    "fort": "fort-latest-release",
    "fort-latest-release": "fort-latest-release",
}

DEFAULT_CONFIG = {
    "root": "/root/rpki_4rp_ctl",
    "bundle": "",
}


def utc_now() -> str:
    return time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime())


def utc_id() -> str:
    return time.strftime("%Y%m%dT%H%M%SZ", time.gmtime())


def script_root() -> Path:
    return Path(__file__).resolve().parent


def ctl_root() -> Path:
    return script_root()


def state_dir() -> Path:
    return ctl_root() / "state"


def logs_dir() -> Path:
    return ctl_root() / "logs"


def runs_dir() -> Path:
    return ctl_root() / "runs"


def config_path() -> Path:
    return ctl_root() / "config.json"


def current_path() -> Path:
    return state_dir() / "current.json"


def pid_path() -> Path:
    return state_dir() / "runner.pid"


def pgid_path() -> Path:
    return state_dir() / "runner.pgid"


def ensure_dirs() -> None:
    for path in [state_dir(), logs_dir(), runs_dir()]:
        path.mkdir(parents=True, exist_ok=True)


def read_json(path: Path, default: Any) -> Any:
    if not path.exists():
        return default
    try:
        return json.loads(path.read_text(encoding="utf-8"))
    except Exception:
        return default


def write_json(path: Path, payload: Any) -> None:
    path.parent.mkdir(parents=True, exist_ok=True)
    tmp = path.with_suffix(path.suffix + ".tmp")
    tmp.write_text(json.dumps(payload, ensure_ascii=False, indent=2, sort_keys=True) + "\n", encoding="utf-8")
    tmp.replace(path)


def load_config() -> dict[str, Any]:
    cfg = dict(DEFAULT_CONFIG)
    cfg.update(read_json(config_path(), {}))
    return cfg


def save_config(cfg: dict[str, Any]) -> None:
    write_json(config_path(), cfg)


def pid_alive(pid: int) -> bool:
    if pid <= 0:
        return False
    try:
        os.kill(pid, 0)
    except ProcessLookupError:
        return False
    except PermissionError:
        return True
    return True


def read_pid(path: Path) -> int:
    try:
        return int(path.read_text(encoding="utf-8").strip())
    except Exception:
        return 0


def normalize_profile(value: str) -> str:
    key = value.strip()
    if key not in ALIASES:
        raise SystemExit(f"unknown rp/profile: {value}; use one of: {', '.join(sorted(ALIASES))}")
    return ALIASES[key]


def parse_rirs(value: str) -> list[str]:
    if value in ("all5", "all"):
        return ["afrinic", "apnic", "arin", "lacnic", "ripe"]
    rirs = [part.strip().lower() for part in value.split(",") if part.strip()]
    bad = [rir for rir in rirs if rir not in RIR_TAL]
    if bad:
        raise SystemExit(f"unknown RIR(s): {', '.join(bad)}")
    if not rirs:
        raise SystemExit("empty --rirs")
    return rirs


def parse_interval(value: str) -> int:
    text = value.strip().lower()
    if text in ("0", "0s", "now"):
        return 0
    match = re.fullmatch(r"(\d+)([smh]?)", text)
    if not match:
        raise SystemExit(f"invalid interval: {value}")
    n = int(match.group(1))
    unit = match.group(2) or "s"
    return n * {"s": 1, "m": 60, "h": 3600}[unit]


def command_string(argv: list[str]) -> str:
    import shlex
    return " ".join(shlex.quote(x) for x in argv)


def check_bundle(root: Path) -> None:
    required = ["bin/rpki", "bin/routinator", "bin/rpki-client", "bin/fort", "fixtures/tal", "scripts/normalize_and_collect.py", "scripts/write_group_report.py", "versions.json"]
    missing = [item for item in required if not (root / item).exists()]
    if missing:
        raise SystemExit(f"runtime bundle incomplete at {root}: missing {missing}")


def install_from_bundle(bundle: Path) -> None:
    check_bundle(bundle)
    for name in ["bin", "lib", "fixtures", "scripts"]:
        src = bundle / name
        dst = ctl_root() / name
        if dst.exists():
            if dst.is_dir():
                shutil.rmtree(dst)
            else:
                dst.unlink()
        shutil.copytree(src, dst)
    shutil.copy2(bundle / "versions.json", ctl_root() / "versions.json")
    for path in [ctl_root() / "bin", ctl_root() / "scripts"]:
        for item in path.iterdir():
            if item.is_file():
                item.chmod(item.stat().st_mode | 0o111)


def cmd_init(args: argparse.Namespace) -> None:
    ensure_dirs()
    bundle = Path(args.bundle).resolve()
    install_from_bundle(bundle)
    cfg = load_config()
    cfg["bundle"] = str(bundle)
    cfg["initializedAtUtc"] = utc_now()
    save_config(cfg)
    print(f"initialized root={ctl_root()}")
    print(f"bundle={bundle}")


def timed_run(run_dir: Path, argv: list[str], env: dict[str, str] | None = None) -> int:
    run_dir.mkdir(parents=True, exist_ok=True)
    (run_dir / "start-utc.txt").write_text(utc_now() + "\n", encoding="utf-8")
    stdout = (run_dir / "stdout.log").open("wb")
    stderr = (run_dir / "stderr.log").open("wb")
    time_path = run_dir / "process-time.txt"
    full = ["/usr/bin/time", "-v", "-o", str(time_path), "--", *argv]
    try:
        proc = subprocess.run(full, stdout=stdout, stderr=stderr, env=env, check=False)
        code = proc.returncode
    finally:
        stdout.close()
        stderr.close()
    (run_dir / "end-utc.txt").write_text(utc_now() + "\n", encoding="utf-8")
    (run_dir / "exit-code.txt").write_text(str(code) + "\n", encoding="utf-8")
    return code


def clean_state_for_snapshot(profile_root: Path) -> None:
    state = profile_root / "state"
    if state.exists():
        shutil.rmtree(state)
    state.mkdir(parents=True, exist_ok=True)


def prepare_state(profile_root: Path, snapshot: bool) -> Path:
    state = profile_root / "state"
    if snapshot:
        clean_state_for_snapshot(profile_root)
    else:
        state.mkdir(parents=True, exist_ok=True)
    return state


def build_ours_args(profile: str, run_root: Path, run_dir: Path, state: Path, rirs: list[str]) -> list[str]:
    argv = [
        str(ctl_root() / "bin" / "rpki"),
        "--db", str(state / "work-db"),
        "--repo-bytes-db", str(state / "repo-bytes.db"),
        "--rsync-mirror-root", str(state / "rsync-mirror"),
        "--rsync-scope", "module-root",
        "--report-json", str(run_dir / "report.json"),
        "--report-json-compact",
        "--ccr-out", str(run_dir / "result.ccr"),
        "--cir-enable",
        "--cir-out", str(run_dir / "input.cir"),
        "--vrps-csv-out", str(run_dir / "vrps.csv"),
        "--vaps-csv-out", str(run_dir / "vaps.csv"),
        "--compare-view-trust-anchor", "all5" if len(rirs) > 1 else rirs[0],
        "--parallel-phase2-ready-batch-size", "256",
        "--parallel-phase2-ready-batch-wall-time-budget-ms", "100",
        "--parallel-phase2-result-drain-batch-size", "2048",
        "--parallel-phase2-finalize-batch-size", "256",
        "--parallel-phase2-finalize-batch-wall-time-budget-ms", "100",
    ]
    for rir in rirs:
        argv += ["--tal-url", TAL_URLS[rir]]
    for rir in rirs:
        argv += ["--cir-tal-uri", TAL_URLS[rir]]
    if profile == "ours-pp-object-cache":
        argv += [
            "--enable-publication-point-validation-cache",
            "--enable-roa-validation-cache",
            "--enable-child-certificate-validation-cache",
        ]
    return argv


def build_routinator_args(run_dir: Path, state: Path, rirs: list[str], snapshot: bool) -> list[str]:
    extra_tals = state / "extra-tals"
    repo = state / "repository"
    extra_tals.mkdir(parents=True, exist_ok=True)
    repo.mkdir(parents=True, exist_ok=True)
    for rir in rirs:
        shutil.copy2(ctl_root() / "fixtures" / "tal" / RIR_TAL[rir], extra_tals / RIR_TAL[rir])
    argv = [
        str(ctl_root() / "bin" / "routinator"),
        "-r", str(repo),
        "--no-rir-tals",
        "--extra-tals-dir", str(extra_tals),
        "--enable-aspa",
    ]
    if snapshot:
        argv.append("--fresh")
    argv += ["vrps", "-f", "jsonext", "-o", str(run_dir / "routinator.json")]
    return argv


def ensure_rpki_client_user() -> None:
    subprocess.run(["id", "-u", "_rpki-client"], stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, check=False)
    if subprocess.run(["id", "-u", "_rpki-client"], stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, check=False).returncode != 0:
        subprocess.run(["useradd", "-r", "-M", "-s", "/usr/sbin/nologin", "_rpki-client"], check=False)


def build_rpki_client_args(run_dir: Path, state: Path, rirs: list[str]) -> tuple[list[str], dict[str, str]]:
    ensure_rpki_client_user()
    cache = state / "cache"
    cache.mkdir(parents=True, exist_ok=True)
    run_dir.mkdir(parents=True, exist_ok=True)
    skiplist = state / "skiplist"
    skiplist.touch(exist_ok=True)
    for path in [state, cache, run_dir]:
        path.chmod(0o755)
    subprocess.run(["chown", "-R", "_rpki-client:_rpki-client", str(state), str(run_dir)], check=False)
    subprocess.run(["chmod", "-R", "u+rwX,go+rX", str(state), str(run_dir)], check=False)
    tal_args: list[str] = []
    for rir in rirs:
        tal_args += ["-t", str(ctl_root() / "fixtures" / "tal" / RIR_TAL[rir])]
    argv = [
        str(ctl_root() / "bin" / "rpki-client"),
        "-p", "4",
        "-j",
        "-c",
        "-S", str(skiplist),
        "-d", str(cache),
        *tal_args,
        str(run_dir),
    ]
    env = os.environ.copy()
    env["LD_LIBRARY_PATH"] = str(ctl_root() / "lib") + ":" + env.get("LD_LIBRARY_PATH", "")
    return argv, env


def build_fort_args(run_dir: Path, state: Path, rirs: list[str]) -> tuple[list[str], dict[str, str]]:
    repo = state / "repository"
    tal_dir = state / "tal"
    repo.mkdir(parents=True, exist_ok=True)
    tal_dir.mkdir(parents=True, exist_ok=True)
    for old_tal in tal_dir.glob("*.tal"):
        old_tal.unlink()
    for rir in rirs:
        shutil.copy2(ctl_root() / "fixtures" / "tal" / RIR_TAL[rir], tal_dir / RIR_TAL[rir])
    argv = [
        str(ctl_root() / "bin" / "fort"),
        "--mode", "standalone",
        "--tal", str(tal_dir),
        "--local-repository", str(repo),
        "--output.roa", str(run_dir / "vrps.csv"),
        "--output.format", "csv",
        "--log.output", "console",
        "--log.level", "warning",
        "--thread-pool.validation.max", "5",
        "--http.connect-timeout", "30",
        "--http.transfer-timeout", "180",
        "--rsync.strategy", "root-except-ta",
    ]
    env = os.environ.copy()
    env["MALLOC_ARENA_MAX"] = "2"
    env["LD_LIBRARY_PATH"] = str(ctl_root() / "lib") + ":" + env.get("LD_LIBRARY_PATH", "")
    return argv, env


def normalize_one(run_root: Path, profile: str, run_seq: int) -> None:
    script = ctl_root() / "scripts" / "normalize_and_collect.py"
    if not script.exists():
        return
    subprocess.run(["python3", str(script), "--root", str(run_root), "--profile", profile, "--run-seq", str(run_seq)], check=False)


def write_group_report(run_root: Path, interval_secs: int, runs: int) -> None:
    script = ctl_root() / "scripts" / "write_group_report.py"
    if not script.exists():
        return
    delta_count = max(0, runs - 1)
    interval_minutes = max(0, interval_secs // 60)
    subprocess.run([
        "python3", str(script),
        "--root", str(run_root),
        "--interval-minutes", str(interval_minutes),
        "--delta-count", str(delta_count),
        "--mode", "manual-rpctl",
    ], check=False)


def parse_latest_meta(run_root: Path, profile: str, run_seq: int) -> dict[str, Any]:
    meta = run_root / "profiles" / profile / "runs" / f"run_{run_seq:04d}" / "run-meta.json"
    return read_json(meta, {})


def update_current(**fields: Any) -> None:
    cur = read_json(current_path(), {})
    cur.update(fields)
    cur["updatedAtUtc"] = utc_now()
    write_json(current_path(), cur)


def is_zero_exit(value: Any) -> bool:
    return str(value) in {"0", "0.0"}


def completed_run_count(profile: str, run_root: Path) -> int:
    profile_runs = run_root / "profiles" / profile / "runs"
    if not profile_runs.exists():
        return 0
    count = 0
    for expected_seq, run_dir in enumerate(sorted(profile_runs.glob("run_*")), start=1):
        try:
            seq = int(run_dir.name.replace("run_", ""))
        except ValueError:
            raise RuntimeError(f"invalid run directory name: {run_dir}")
        if seq != expected_seq:
            raise RuntimeError(f"non-contiguous run directories under {profile_runs}: expected run_{expected_seq:04d}, got {run_dir.name}")
        exit_code = read_text_strip(run_dir / "exit-code.txt")
        if not is_zero_exit(exit_code):
            raise RuntimeError(f"run is not successful: {run_dir} exit_code={exit_code or '<missing>'}")
        count = seq
    return count


def parse_entry_rirs(entry: dict[str, Any]) -> list[str]:
    value = entry.get("rirs") or ""
    if isinstance(value, list):
        return parse_rirs(",".join(str(item) for item in value))
    return parse_rirs(str(value))


def run_profile(profile: str, run_id: str, runs: int, interval_secs: int, rirs: list[str], append: bool = False) -> int:
    ensure_dirs()
    rp = PROFILES[profile]
    run_root = runs_dir() / profile / run_id
    profile_root = run_root / "profiles" / profile
    (run_root / "logs").mkdir(parents=True, exist_ok=True)
    (run_root / "reports").mkdir(parents=True, exist_ok=True)
    (run_root / "versions.json").write_text((ctl_root() / "versions.json").read_text(encoding="utf-8"), encoding="utf-8")
    now = utc_now()
    if append:
        run_info = read_json(run_root / "run-info.json", {})
        if not run_info:
            raise RuntimeError(f"cannot continue run without run-info.json: {run_root}")
        if run_info.get("state") != "success" or not is_zero_exit(run_info.get("exitCode")):
            raise RuntimeError(f"can only continue a successful run: run_id={run_id} state={run_info.get('state')} exitCode={run_info.get('exitCode')}")
        existing_runs = completed_run_count(profile, run_root)
        if existing_runs < 1:
            raise RuntimeError(f"cannot continue run without successful child runs: {run_root}")
        first_run_seq = existing_runs + 1
        last_run_seq = existing_runs + runs
        run_info.update({
            "runsRequested": last_run_seq,
            "intervalSecs": interval_secs,
            "state": "running",
            "exitCode": "",
            "lastContinueAtUtc": now,
            "lastContinueFromRunSeq": first_run_seq,
            "lastContinueRunsRequested": runs,
        })
    else:
        first_run_seq = 1
        last_run_seq = runs
        run_info = {
            "runId": run_id,
            "profile": profile,
            "rp": rp,
            "rirs": rirs,
            "runsRequested": runs,
            "intervalSecs": interval_secs,
            "cacheProfile": cache_profile_name(profile),
            "startedAtUtc": now,
            "runRoot": str(run_root),
            "stateRoot": str(profile_root / "state"),
        }
    write_json(run_root / "run-info.json", run_info)
    status = {
        "state": "running",
        "rp": rp,
        "profile": profile,
        "pid": os.getpid(),
        "pgid": os.getpgrp(),
        "runId": run_id,
        "runRoot": str(run_root),
        "intervalSecs": interval_secs,
        "runsRequested": last_run_seq,
        "continueMode": append,
        "firstRunSeq": first_run_seq,
        "lastRunSeqRequested": last_run_seq,
        "rirs": rirs,
        "startedAtUtc": now,
        "currentRun": None,
    }
    write_json(current_path(), status)
    exit_code = 0
    last_start = 0.0
    completed_seq = first_run_seq - 1
    for run_seq in range(first_run_seq, last_run_seq + 1):
        if run_seq > first_run_seq and interval_secs > 0:
            target = last_start + interval_secs
            sleep_seconds = max(0.0, target - time.time())
            update_current(state="waiting", nextRunSeq=run_seq, nextRunAtEpoch=target, sleepSeconds=round(sleep_seconds, 3))
            if sleep_seconds > 0:
                time.sleep(sleep_seconds)
        last_start = time.time()
        snapshot = run_seq == 1
        mode = "snapshot" if snapshot else "delta"
        run_dir = profile_root / "runs" / f"run_{run_seq:04d}"
        if run_dir.exists() and any(run_dir.iterdir()):
            raise RuntimeError(f"refusing to overwrite existing run directory: {run_dir}")
        state = prepare_state(profile_root, snapshot)
        update_current(state="running", currentRun=run_seq, currentMode=mode, currentRunDir=str(run_dir), lastRunStartedAtUtc=utc_now())
        env = None
        if rp == "ours-rp":
            argv = build_ours_args(profile, run_root, run_dir, state, rirs)
        elif rp == "routinator":
            argv = build_routinator_args(run_dir, state, rirs, snapshot)
        elif rp == "rpki-client":
            argv, env = build_rpki_client_args(run_dir, state, rirs)
        elif rp == "fort":
            argv, env = build_fort_args(run_dir, state, rirs)
        else:
            raise RuntimeError(rp)
        run_dir.mkdir(parents=True, exist_ok=True)
        (run_dir / "command.txt").write_text(command_string(argv) + "\n", encoding="utf-8")
        code = timed_run(run_dir, argv, env=env)
        normalize_one(run_root, profile, run_seq)
        meta = parse_latest_meta(run_root, profile, run_seq)
        update_current(
            state="running",
            lastRunSeq=run_seq,
            lastRunExitCode=code,
            lastRunWallMs=meta.get("timing", {}).get("wallMs"),
            lastRunMaxRssKb=meta.get("timing", {}).get("maxRssKb"),
            lastRunVrps=meta.get("counts", {}).get("vrps"),
            lastRunVaps=meta.get("counts", {}).get("vaps"),
            lastRunDir=str(run_dir),
        )
        completed_seq = run_seq
        if code != 0:
            exit_code = code
            break
    write_group_report(run_root, interval_secs, completed_seq)
    final_state = "success" if exit_code == 0 else "failed"
    update_current(state=final_state, finishedAtUtc=utc_now(), exitCode=exit_code)
    run_info = read_json(run_root / "run-info.json", {})
    run_info.update({
        "finishedAtUtc": utc_now(),
        "state": final_state,
        "exitCode": exit_code,
        "completedRuns": completed_seq,
    })
    write_json(run_root / "run-info.json", run_info)
    return exit_code


def cmd_runner(args: argparse.Namespace) -> None:
    code = run_profile(args.profile, args.run_id, args.runs, args.interval_secs, args.rirs.split(","), append=args.append)
    sys.exit(code)


def current_live_status() -> dict[str, Any]:
    cur = read_json(current_path(), {})
    pid = read_pid(pid_path()) or int(cur.get("pid") or 0)
    if not cur and not pid:
        return {}
    alive = pid_alive(pid)
    if cur.get("state") in {"starting", "running", "waiting"} and not alive:
        cur = dict(cur)
        cur["state"] = "exited"
        cur["pidAlive"] = False
        cur["observedExitAtUtc"] = utc_now()
        write_json(current_path(), cur)
    else:
        cur["pidAlive"] = alive
    if pid:
        cur["pid"] = pid
    return cur


def cmd_status(args: argparse.Namespace) -> None:
    ensure_dirs()
    cur = current_live_status()
    if args.json:
        print(json.dumps(cur, ensure_ascii=False, indent=2, sort_keys=True))
        return
    if not cur:
        print("state=idle")
        return
    keys = ["state", "profile", "rp", "pid", "pidAlive", "runId", "runRoot", "currentRun", "currentMode", "lastRunSeq", "lastRunExitCode", "lastRunWallMs", "lastRunMaxRssKb", "lastRunVrps", "lastRunVaps", "lastRunDir"]
    for key in keys:
        if key in cur:
            print(f"{key}={cur[key]}")


def cmd_list(args: argparse.Namespace) -> None:
    print("profile\trp\taliases")
    reverse: dict[str, list[str]] = {p: [] for p in PROFILES}
    for alias, profile in ALIASES.items():
        reverse.setdefault(profile, []).append(alias)
    for profile, rp in PROFILES.items():
        print(f"{profile}\t{rp}\t{','.join(sorted(reverse.get(profile, [])))}")


def cmd_paths(args: argparse.Namespace) -> None:
    cfg = load_config()
    print(f"root={ctl_root()}")
    print(f"config={config_path()}")
    print(f"state={state_dir()}")
    print(f"runs={runs_dir()}")
    print(f"logs={logs_dir()}")
    print(f"bundle={cfg.get('bundle','')}")
    cur = current_live_status()
    if cur.get("runRoot"):
        print(f"current_run_root={cur['runRoot']}")
    if cur.get("lastRunDir"):
        print(f"last_run_dir={cur['lastRunDir']}")


def refuse_if_running() -> None:
    cur = current_live_status()
    if cur.get("state") in {"starting", "running", "waiting"} and cur.get("pidAlive"):
        raise SystemExit(f"another rpctl run is active: profile={cur.get('profile')} pid={cur.get('pid')} runRoot={cur.get('runRoot')}")


def cmd_start(args: argparse.Namespace) -> None:
    ensure_dirs()
    refuse_if_running()
    profile = normalize_profile(args.rp)
    rirs = parse_rirs(args.rirs)
    interval_secs = parse_interval(args.interval)
    if args.runs < 1:
        raise SystemExit("--runs must be >= 1")
    if not (ctl_root() / "versions.json").exists():
        raise SystemExit("rpctl is not initialized; run ./rpctl init --bundle <bundle-root> first")
    run_id = args.run_id or f"{utc_id()}_{profile}_{args.runs}runs_{interval_secs}s"
    runner_args = [
        sys.executable,
        str(Path(__file__).resolve()),
        "_runner",
        "--profile", profile,
        "--run-id", run_id,
        "--runs", str(args.runs),
        "--interval-secs", str(interval_secs),
        "--rirs", ",".join(rirs),
    ]
    nohup = logs_dir() / f"runner-{run_id}.log"
    out = nohup.open("ab")
    proc = subprocess.Popen(runner_args, stdout=out, stderr=subprocess.STDOUT, start_new_session=True, cwd=str(ctl_root()))
    pid_path().write_text(str(proc.pid) + "\n", encoding="utf-8")
    try:
        pgid = os.getpgid(proc.pid)
    except Exception:
        pgid = proc.pid
    pgid_path().write_text(str(pgid) + "\n", encoding="utf-8")
    write_json(current_path(), {
        "state": "starting",
        "profile": profile,
        "rp": PROFILES[profile],
        "pid": proc.pid,
        "pgid": pgid,
        "runId": run_id,
        "intervalSecs": interval_secs,
        "runsRequested": args.runs,
        "rirs": rirs,
        "runnerLog": str(nohup),
        "startedAtUtc": utc_now(),
    })
    print(f"started profile={profile} pid={proc.pid} pgid={pgid} run_id={run_id}")
    print(f"log={nohup}")


def cmd_continue(args: argparse.Namespace) -> None:
    ensure_dirs()
    refuse_if_running()
    if args.runs < 1:
        raise SystemExit("--runs must be >= 1")
    if not (ctl_root() / "versions.json").exists():
        raise SystemExit("rpctl is not initialized; run ./rpctl init --bundle <bundle-root> first")
    entry = find_run_entry(args.rp, args.run_id)
    if entry["state"] != "success":
        raise SystemExit(f"can only continue a successful run: run_id={args.run_id} state={entry['state']}")
    profile = entry["profile"]
    rirs = parse_entry_rirs(entry)
    interval_secs = parse_interval(args.interval)
    runner_args = [
        sys.executable,
        str(Path(__file__).resolve()),
        "_runner",
        "--profile", profile,
        "--run-id", args.run_id,
        "--runs", str(args.runs),
        "--interval-secs", str(interval_secs),
        "--rirs", ",".join(rirs),
        "--append",
    ]
    nohup = logs_dir() / f"runner-{args.run_id}-continue-{utc_id()}.log"
    out = nohup.open("ab")
    proc = subprocess.Popen(runner_args, stdout=out, stderr=subprocess.STDOUT, start_new_session=True, cwd=str(ctl_root()))
    pid_path().write_text(str(proc.pid) + "\n", encoding="utf-8")
    try:
        pgid = os.getpgid(proc.pid)
    except Exception:
        pgid = proc.pid
    pgid_path().write_text(str(pgid) + "\n", encoding="utf-8")
    previous_runs = int(entry.get("runCount") or 0)
    write_json(current_path(), {
        "state": "starting",
        "profile": profile,
        "rp": PROFILES[profile],
        "pid": proc.pid,
        "pgid": pgid,
        "runId": args.run_id,
        "intervalSecs": interval_secs,
        "runsRequested": previous_runs + args.runs,
        "continueMode": True,
        "continueRunsRequested": args.runs,
        "previousRunCount": previous_runs,
        "rirs": rirs,
        "runnerLog": str(nohup),
        "startedAtUtc": utc_now(),
    })
    print(f"continued profile={profile} pid={proc.pid} pgid={pgid} run_id={args.run_id} add_runs={args.runs}")
    print(f"log={nohup}")


def cmd_stop(args: argparse.Namespace) -> None:
    ensure_dirs()
    cur = current_live_status()
    pid = int(cur.get("pid") or read_pid(pid_path()) or 0)
    pgid = int(cur.get("pgid") or read_pid(pgid_path()) or pid or 0)
    if not pid:
        print("no runner pid")
        return
    if not pid_alive(pid):
        update_current(state="exited", pidAlive=False)
        print(f"pid {pid} is not alive")
        return
    sig = signal.SIGKILL if args.force else signal.SIGTERM
    try:
        os.killpg(pgid, sig)
        print(f"sent {sig.name} to pgid={pgid}")
    except ProcessLookupError:
        try:
            os.kill(pid, sig)
            print(f"sent {sig.name} to pid={pid}")
        except ProcessLookupError:
            print(f"pid {pid} disappeared")
    deadline = time.time() + (2 if args.force else 10)
    while time.time() < deadline:
        if not pid_alive(pid):
            break
        time.sleep(0.2)
    alive = pid_alive(pid)
    update_current(state="stopped" if not alive else "stopping", pidAlive=alive, stoppedAtUtc=utc_now())
    print(f"pidAlive={alive}")


def cmd_tail(args: argparse.Namespace) -> None:
    cur = current_live_status()
    candidates = []
    runner_log = cur.get("runnerLog")
    if runner_log:
        candidates.append(Path(runner_log))
    run_root = cur.get("runRoot")
    if run_root:
        candidates.append(Path(run_root) / "logs" / "driver-progress.log")
    candidates.extend(sorted(logs_dir().glob("runner-*.log"), key=lambda p: p.stat().st_mtime if p.exists() else 0, reverse=True))
    path = next((p for p in candidates if p.is_file()), None)
    if path is None:
        print("no log found")
        return
    lines = path.read_text(encoding="utf-8", errors="replace").splitlines()
    for line in lines[-args.lines:]:
        print(line)


def load_csv_rows(path: Path) -> list[dict[str, str]]:
    if not path.exists():
        return []
    try:
        with path.open("r", encoding="utf-8", errors="replace", newline="") as handle:
            return list(csv.DictReader(handle))
    except Exception:
        return []


def first_non_empty(*values: Any) -> str:
    for value in values:
        if value is not None and str(value) != "":
            return str(value)
    return ""


def display_value(value: Any) -> str:
    text = first_non_empty(value)
    return text if text else "-"


def read_text_strip(path: Path) -> str:
    try:
        return path.read_text(encoding="utf-8", errors="replace").strip()
    except Exception:
        return ""


def parse_command_summary(command: str) -> str:
    if not command:
        return ""
    import shlex
    try:
        tokens = shlex.split(command)
    except ValueError:
        tokens = command.split()
    token_set = set(tokens)
    keys = []
    for flag in [
        "--enable-publication-point-validation-cache",
        "--enable-roa-validation-cache",
        "--enable-child-certificate-validation-cache",
        "--enable-aspa",
        "-p",
        "-j",
        "-c",
    ]:
        if flag in token_set:
            keys.append(flag)
    for idx, token in enumerate(tokens):
        if token == "--rsync-scope" and idx + 1 < len(tokens):
            keys.append(f"--rsync-scope={tokens[idx + 1]}")
        elif token.startswith("--rsync-scope="):
            keys.append(token)
    return " ".join(keys)


def parse_process_time(path: Path) -> dict[str, int]:
    if not path.exists():
        return {}
    result: dict[str, int] = {}
    for line in path.read_text(encoding="utf-8", errors="replace").splitlines():
        if "Elapsed (wall clock) time" in line:
            raw = line.rsplit("):", 1)[-1].strip() if "):" in line else line.rsplit(":", 1)[-1].strip()
            try:
                result["wallMs"] = elapsed_to_ms(raw)
            except Exception:
                pass
        elif "Maximum resident set size" in line:
            try:
                result["maxRssKb"] = int(line.rsplit(":", 1)[1].strip() or 0)
            except Exception:
                pass
    return result


def elapsed_to_ms(raw: str) -> int:
    raw = raw.strip()
    days = 0
    if "-" in raw:
        day_text, raw = raw.split("-", 1)
        days = int(day_text)
    parts = raw.split(":")
    if len(parts) == 3:
        hours, minutes, seconds = parts
    elif len(parts) == 2:
        hours, minutes, seconds = "0", parts[0], parts[1]
    else:
        hours, minutes, seconds = "0", "0", parts[0]
    return int(round((days * 86400 + int(hours) * 3600 + int(minutes) * 60 + float(seconds)) * 1000))


def human_ms(value: Any) -> str:
    text = first_non_empty(value)
    if not text:
        return "-"
    try:
        ms = int(float(text))
    except ValueError:
        return text
    if ms >= 60_000:
        return f"{ms / 60_000:.1f}m"
    if ms >= 1000:
        return f"{ms / 1000:.1f}s"
    return f"{ms}ms"


def compact_utc(value: Any) -> str:
    text = first_non_empty(value)
    if not text:
        return "-"
    match = re.fullmatch(r"(\d{4})-(\d{2})-(\d{2})T(\d{2}:\d{2}:\d{2})Z", text)
    if match:
        return f"{match.group(2)}-{match.group(3)} {match.group(4)}"
    return text


def compact_rirs(value: Any) -> str:
    text = first_non_empty(value)
    if not text:
        return "-"
    rirs = [part.strip() for part in text.split(",") if part.strip()]
    if rirs == ["afrinic", "apnic", "arin", "lacnic", "ripe"]:
        return "all5"
    return ",".join(rirs) if rirs else text


def cache_profile_name(profile: str) -> str:
    if profile == "ours-pp-object-cache":
        return "pp+roa+crt"
    if profile == "ours-no-cache":
        return "off"
    return "n/a"


def format_run_params(entry: dict[str, Any]) -> str:
    parts = []
    runs_requested = first_non_empty(entry.get("runsRequested"))
    interval = first_non_empty(entry.get("intervalSecs"))
    rirs = compact_rirs(entry.get("rirs"))
    cache = first_non_empty(entry.get("cacheProfile"))
    if runs_requested or interval:
        interval_text = f"{interval}s" if interval else "?"
        parts.append(f"{runs_requested or '?'}@{interval_text}")
    if rirs != "-":
        parts.append(rirs)
    if cache:
        parts.append(cache)
    return " ".join(parts)


def profile_label(profile: str) -> str:
    labels = {
        "ours-pp-object-cache": "ours-cache",
        "ours-no-cache": "ours-off",
        "routinator-latest-release": "routinator",
        "rpki-client-latest-release": "rpki-client",
        "fort-latest-release": "fort",
    }
    return labels.get(profile, profile)


def count_gzip_lines(path: Path) -> str:
    if not path.exists():
        return ""
    try:
        with gzip.open(path, "rt", encoding="utf-8", errors="replace") as handle:
            return str(sum(1 for line in handle if line.strip()))
    except Exception:
        return ""


def count_csv_data_rows(path: Path) -> str:
    if not path.exists():
        return ""
    try:
        with path.open("r", encoding="utf-8", errors="replace", newline="") as handle:
            return str(sum(1 for _ in csv.DictReader(handle)))
    except Exception:
        return ""


def count_vrp_json(path: Path) -> str:
    if not path.exists():
        return ""
    try:
        data = json.loads(path.read_text(encoding="utf-8", errors="replace"))
    except Exception:
        return ""
    if not isinstance(data, dict):
        return ""
    for key in ("roas", "routeOrigins", "valid_roas"):
        items = data.get(key)
        if isinstance(items, list):
            return str(len(items))
    return ""


def count_vap_json(path: Path) -> str:
    if not path.exists():
        return ""
    try:
        data = json.loads(path.read_text(encoding="utf-8", errors="replace"))
    except Exception:
        return ""
    if not isinstance(data, dict):
        return ""
    for key in ("aspas", "aspaAssertions", "vaps"):
        items = data.get(key)
        if isinstance(items, list):
            return str(len(items))
    return ""


def infer_output_counts(profile: str, run_dir: Path) -> tuple[str, str]:
    normalized_vrps = count_gzip_lines(run_dir / "normalized-vrps.txt.gz")
    normalized_vaps = count_gzip_lines(run_dir / "normalized-vaps.txt.gz")
    if normalized_vrps or normalized_vaps:
        return normalized_vrps, normalized_vaps
    rp = PROFILES.get(profile, profile)
    if rp == "ours-rp":
        return count_csv_data_rows(run_dir / "vrps.csv"), count_csv_data_rows(run_dir / "vaps.csv")
    if rp == "routinator":
        return count_vrp_json(run_dir / "routinator.json"), count_vap_json(run_dir / "routinator.json")
    if rp == "rpki-client":
        vrps = count_csv_data_rows(run_dir / "csv") or count_vrp_json(run_dir / "json")
        vaps = count_vap_json(run_dir / "json")
        return vrps, vaps
    if rp == "fort":
        return count_csv_data_rows(run_dir / "vrps.csv"), "0" if (run_dir / "vrps.csv").exists() else ""
    return "", ""


def discover_run_entries() -> list[dict[str, Any]]:
    ensure_dirs()
    current = current_live_status()
    entries: list[dict[str, Any]] = []
    for profile_dir in sorted([p for p in runs_dir().iterdir() if p.is_dir()]):
        profile = profile_dir.name
        for run_root in sorted([p for p in profile_dir.iterdir() if p.is_dir()]):
            run_id = run_root.name
            run_info = read_json(run_root / "run-info.json", {})
            reports_csv = run_root / "reports" / "per_run_metrics.csv"
            rows = load_csv_rows(reports_csv)
            profile_run_root = run_root / "profiles" / profile / "runs"
            run_dirs = sorted(profile_run_root.glob("run_*")) if profile_run_root.exists() else []
            started = first_non_empty(run_info.get("startedAtUtc"))
            ended = first_non_empty(run_info.get("finishedAtUtc"))
            exit_codes: list[str] = []
            walls: list[int] = []
            max_rss_values: list[int] = []
            vrps = ""
            vaps = ""
            run_count = 0
            failed_runs: list[str] = []
            if rows:
                run_count = len(rows)
                started = first_non_empty(rows[0].get("start_utc"), rows[0].get("startUtc"))
                ended = first_non_empty(rows[-1].get("end_utc"), rows[-1].get("endUtc"))
                for row in rows:
                    seq = first_non_empty(row.get("run_seq"), row.get("runSeq"))
                    exit_code = first_non_empty(row.get("exit_code"), row.get("exitCode"))
                    if exit_code:
                        exit_codes.append(exit_code)
                        if exit_code not in {"0", "0.0"}:
                            failed_runs.append(seq or "?")
                    try:
                        walls.append(int(float(first_non_empty(row.get("wall_ms"), row.get("wallMs")) or 0)))
                    except ValueError:
                        pass
                    try:
                        max_rss_values.append(int(float(first_non_empty(row.get("max_rss_kb"), row.get("maxRssKb")) or 0)))
                    except ValueError:
                        pass
                vrps = first_non_empty(rows[-1].get("vrps"), rows[-1].get("VRPs"))
                vaps = first_non_empty(rows[-1].get("vaps"), rows[-1].get("VAPs"))
                if run_dirs and (not vrps or not vaps):
                    inferred_vrps, inferred_vaps = infer_output_counts(profile, run_dirs[-1])
                    vrps = first_non_empty(vrps, inferred_vrps)
                    vaps = first_non_empty(vaps, inferred_vaps)
            else:
                run_count = len(run_dirs)
                for run_dir in run_dirs:
                    seq = run_dir.name.replace("run_", "")
                    started = started or read_text_strip(run_dir / "start-utc.txt")
                    end = read_text_strip(run_dir / "end-utc.txt")
                    ended = end or ended
                    exit_code = read_text_strip(run_dir / "exit-code.txt")
                    if exit_code:
                        exit_codes.append(exit_code)
                        if exit_code != "0":
                            failed_runs.append(seq)
                    meta = read_json(run_dir / "run-meta.json", {})
                    if meta:
                        try:
                            walls.append(int(meta.get("timing", {}).get("wallMs") or 0))
                        except Exception:
                            pass
                        try:
                            max_rss_values.append(int(meta.get("timing", {}).get("maxRssKb") or 0))
                        except Exception:
                            pass
                        vrps = str(meta.get("counts", {}).get("vrps", vrps))
                        vaps = str(meta.get("counts", {}).get("vaps", vaps))
                    else:
                        timing = parse_process_time(run_dir / "process-time.txt")
                        if timing.get("wallMs"):
                            walls.append(timing["wallMs"])
                        if timing.get("maxRssKb"):
                            max_rss_values.append(timing["maxRssKb"])
                    inferred_vrps, inferred_vaps = infer_output_counts(profile, run_dir)
                    vrps = first_non_empty(inferred_vrps, vrps)
                    vaps = first_non_empty(inferred_vaps, vaps)
            command = ""
            first_cmd = run_root / "profiles" / profile / "runs" / "run_0001" / "command.txt"
            command = read_text_strip(first_cmd)
            state = "unknown"
            if current.get("runRoot") == str(run_root) and current.get("state") in {"starting", "running", "waiting", "stopping", "stopped", "success", "failed", "exited"}:
                state = str(current.get("state"))
            elif failed_runs:
                state = "failed"
            elif exit_codes and all(code in {"0", "0.0"} for code in exit_codes):
                state = "success"
            elif run_dirs and not ended:
                state = "incomplete"
            elif run_root.exists():
                state = "partial"
            interval = first_non_empty(run_info.get("intervalSecs"))
            runs_requested = first_non_empty(run_info.get("runsRequested"))
            rirs = ",".join(run_info.get("rirs", [])) if isinstance(run_info.get("rirs"), list) else first_non_empty(run_info.get("rirs"))
            cache_profile = first_non_empty(run_info.get("cacheProfile"), cache_profile_name(profile))
            if current.get("runRoot") == str(run_root):
                interval = first_non_empty(current.get("intervalSecs"), interval)
                runs_requested = first_non_empty(current.get("runsRequested"), runs_requested)
                current_rirs = current.get("rirs")
                if isinstance(current_rirs, list):
                    rirs = ",".join(current_rirs)
            entries.append({
                "profile": profile,
                "rp": PROFILES.get(profile, profile),
                "runId": run_id,
                "state": state,
                "startedAt": started,
                "endedAt": ended,
                "runCount": run_count,
                "runsRequested": runs_requested,
                "intervalSecs": interval,
                "rirs": rirs,
                "cacheProfile": cache_profile,
                "exitCodes": ",".join(exit_codes),
                "failedRuns": ",".join(failed_runs),
                "wallMsTotal": sum(walls) if walls else "",
                "wallMsMax": max(walls) if walls else "",
                "maxRssKb": max(max_rss_values) if max_rss_values else "",
                "vrps": vrps,
                "vaps": vaps,
                "runRoot": str(run_root),
                "commandSummary": parse_command_summary(command),
                "command": command,
            })
    entries.sort(key=lambda item: (item.get("startedAt") or "", item.get("runRoot") or ""), reverse=True)
    return entries


def select_show_entries(args: argparse.Namespace) -> list[dict[str, Any]]:
    entries = discover_run_entries()
    if args.rp:
        profile = normalize_profile(args.rp)
        entries = [entry for entry in entries if entry["profile"] == profile]
    if args.failed:
        entries = [entry for entry in entries if entry["state"] in {"failed", "incomplete", "partial", "exited"} or entry.get("failedRuns")]
    if args.limit and args.limit > 0:
        entries = entries[: args.limit]
    return entries


def print_table(rows: list[list[Any]], headers: list[str]) -> None:
    data = [[display_value(cell) for cell in row] for row in rows]
    widths = [len(header) for header in headers]
    for row in data:
        for idx, cell in enumerate(row):
            widths[idx] = max(widths[idx], min(len(cell), 80))
    def trim(cell: str, width: int) -> str:
        return cell if len(cell) <= width else cell[: max(0, width - 1)] + "…"
    print("  ".join(header.ljust(widths[idx]) for idx, header in enumerate(headers)))
    print("  ".join("-" * widths[idx] for idx in range(len(headers))))
    for row in data:
        print("  ".join(trim(cell, widths[idx]).ljust(widths[idx]) for idx, cell in enumerate(row)))


def cmd_show(args: argparse.Namespace) -> None:
    entries = select_show_entries(args)
    if args.json:
        print(json.dumps(entries, ensure_ascii=False, indent=2, sort_keys=True))
        return
    headers = ["profile", "run_id", "state", "runs", "params", "start", "end", "failed", "wall", "vrps", "vaps"]
    rows = [[
        profile_label(e["profile"]),
        e["runId"],
        e["state"],
        e["runCount"],
        format_run_params(e),
        compact_utc(e["startedAt"]),
        compact_utc(e["endedAt"]),
        e["failedRuns"],
        human_ms(e["wallMsTotal"]),
        e["vrps"],
        e["vaps"],
    ] for e in entries]
    if args.paths:
        headers.append("path")
        for row, entry in zip(rows, entries):
            row.append(entry["runRoot"])
    if not rows:
        print("no runs found")
        return
    print_table(rows, headers)


def find_run_entry(profile_or_alias: str | None, run_id: str) -> dict[str, Any]:
    profile = normalize_profile(profile_or_alias) if profile_or_alias else None
    entries = discover_run_entries()
    matches = [entry for entry in entries if entry["runId"] == run_id and (profile is None or entry["profile"] == profile)]
    if not matches:
        raise SystemExit(f"run not found: run_id={run_id}" + (f" profile={profile}" if profile else ""))
    if len(matches) > 1:
        profiles = ", ".join(entry["profile"] for entry in matches)
        raise SystemExit(f"run_id is ambiguous; specify --rp. matches: {profiles}")
    return matches[0]


def is_running_entry(entry: dict[str, Any]) -> bool:
    run_root = str(Path(str(entry.get("runRoot"))).resolve())
    cur = current_live_status()
    if cur.get("runRoot"):
        try:
            current_root = Path(str(cur.get("runRoot"))).resolve()
            target_root = Path(run_root)
            if (
                current_root == target_root
                and cur.get("state") in {"starting", "running", "waiting", "stopping"}
                and bool(cur.get("pidAlive"))
            ):
                return True
        except Exception:
            pass
    try:
        proc = subprocess.run(
            ["ps", "-eo", "pid=,args="],
            text=True,
            stdout=subprocess.PIPE,
            stderr=subprocess.DEVNULL,
            check=False,
        )
    except Exception:
        return False
    for line in proc.stdout.splitlines():
        if run_root in line and " rpctl delete " not in line:
            return True
    return False


def safe_run_root(path: Path) -> Path:
    resolved = path.resolve()
    base = runs_dir().resolve()
    if resolved == base or base not in resolved.parents:
        raise SystemExit(f"refusing to delete path outside runs dir: {resolved}")
    if len(resolved.relative_to(base).parts) < 2:
        raise SystemExit(f"refusing to delete non-run root path: {resolved}")
    return resolved


def matching_runner_logs(run_id: str) -> list[Path]:
    prefix = f"runner-{run_id}"
    return sorted(path for path in logs_dir().glob("runner-*.log") if path.name.startswith(prefix))


def clear_current_if_points_to(run_root: Path) -> None:
    cur = current_live_status()
    if not cur.get("runRoot"):
        return
    try:
        if Path(str(cur.get("runRoot"))).resolve() != run_root.resolve():
            return
    except Exception:
        return
    for path in [current_path(), pid_path(), pgid_path()]:
        try:
            path.unlink()
        except FileNotFoundError:
            pass


def cmd_delete(args: argparse.Namespace) -> None:
    entry = find_run_entry(args.rp, args.run_id)
    if is_running_entry(entry):
        raise SystemExit(f"refusing to delete active run: run_id={args.run_id}. stop it first.")
    run_root = safe_run_root(Path(entry["runRoot"]))
    logs = matching_runner_logs(args.run_id)
    if args.dry_run:
        print(f"run_root={run_root}")
        for log in logs:
            print(f"log={log}")
        return
    if not args.yes:
        raise SystemExit(f"refusing to delete without --yes. target={run_root}")
    if run_root.exists():
        shutil.rmtree(run_root)
    deleted_logs = 0
    for log in logs:
        try:
            log.unlink()
            deleted_logs += 1
        except FileNotFoundError:
            pass
    clear_current_if_points_to(run_root)
    print(f"deleted run_id={args.run_id}")
    print(f"run_root={run_root}")
    print(f"runner_logs_deleted={deleted_logs}")


def run_rows_for_entry(entry: dict[str, Any]) -> list[dict[str, str]]:
    return load_csv_rows(Path(entry["runRoot"]) / "reports" / "per_run_metrics.csv")


def fallback_run_details(entry: dict[str, Any]) -> list[dict[str, Any]]:
    profile = entry["profile"]
    root = Path(entry["runRoot"]) / "profiles" / profile / "runs"
    details = []
    for run_dir in sorted(root.glob("run_*")):
        meta = read_json(run_dir / "run-meta.json", {})
        timing = meta.get("timing", {}) if meta else parse_process_time(run_dir / "process-time.txt")
        inferred_vrps, inferred_vaps = infer_output_counts(profile, run_dir)
        details.append({
            "run_seq": run_dir.name.replace("run_", ""),
            "sync_mode": meta.get("syncMode") or ("snapshot" if run_dir.name.endswith("0001") else "delta"),
            "start_utc": read_text_strip(run_dir / "start-utc.txt"),
            "end_utc": read_text_strip(run_dir / "end-utc.txt"),
            "exit_code": read_text_strip(run_dir / "exit-code.txt"),
            "wall_ms": timing.get("wallMs", ""),
            "max_rss_kb": timing.get("maxRssKb", ""),
            "vrps": first_non_empty(meta.get("counts", {}).get("vrps") if meta else "", inferred_vrps),
            "vaps": first_non_empty(meta.get("counts", {}).get("vaps") if meta else "", inferred_vaps),
            "remote_run_dir": str(run_dir),
        })
    return details


def cmd_view(args: argparse.Namespace) -> None:
    entry = find_run_entry(args.rp, args.run_id)
    rows = run_rows_for_entry(entry) or fallback_run_details(entry)
    payload = {"summary": entry, "runs": rows}
    if args.json:
        print(json.dumps(payload, ensure_ascii=False, indent=2, sort_keys=True))
        return
    print(f"run_id={display_value(entry['runId'])}")
    print(f"profile={display_value(entry['profile'])}")
    print(f"rp={display_value(entry['rp'])}")
    print(f"state={display_value(entry['state'])}")
    print(f"start={display_value(entry['startedAt'])}")
    print(f"end={display_value(entry['endedAt'])}")
    print(f"run_count={display_value(entry['runCount'])}")
    print(f"runs_requested={display_value(entry['runsRequested'])}")
    print(f"interval_secs={display_value(entry['intervalSecs'])}")
    print(f"rirs={display_value(entry['rirs'])}")
    print(f"cache_profile={display_value(entry['cacheProfile'])}")
    print(f"failed_runs={display_value(entry['failedRuns'])}")
    print(f"wall_total_ms={display_value(entry['wallMsTotal'])}")
    print(f"wall_max_ms={display_value(entry['wallMsMax'])}")
    print(f"max_rss_kb={display_value(entry['maxRssKb'])}")
    print(f"vrps={display_value(entry['vrps'])}")
    print(f"vaps={display_value(entry['vaps'])}")
    print(f"run_root={display_value(entry['runRoot'])}")
    print(f"command_summary={display_value(entry['commandSummary'])}")
    if args.command and entry.get("command"):
        print("command=" + entry["command"])
    if rows:
        print("\nper-run:")
        headers = ["run", "mode", "exit", "start", "end", "wall_ms", "rss_kb", "vrps", "vaps", "path"]
        table_rows = []
        for row in rows:
            table_rows.append([
                first_non_empty(row.get("run_seq"), row.get("runSeq")),
                first_non_empty(row.get("sync_mode"), row.get("syncMode")),
                first_non_empty(row.get("exit_code"), row.get("exitCode")),
                first_non_empty(row.get("start_utc"), row.get("startUtc")),
                first_non_empty(row.get("end_utc"), row.get("endUtc")),
                first_non_empty(row.get("wall_ms"), row.get("wallMs")),
                first_non_empty(row.get("max_rss_kb"), row.get("maxRssKb")),
                first_non_empty(row.get("vrps"), row.get("VRPs")),
                first_non_empty(row.get("vaps"), row.get("VAPs")),
                first_non_empty(row.get("remote_run_dir"), row.get("remoteRunDir")),
            ])
        print_table(table_rows, headers)


def build_parser() -> argparse.ArgumentParser:
    formatter = argparse.RawDescriptionHelpFormatter
    parser = argparse.ArgumentParser(
        description="Control one local four-RP performance run on this server.",
        formatter_class=formatter,
        epilog="""
RP/profile naming:
  ours / ours-rp / ours-cached
      Alias of ours-pp-object-cache. Starts ours RP with PP cache + ROA cache
      + child certificate cache enabled.

  ours-no-cache / ours-baseline
      Starts ours RP without the three cache acceleration flags. This is the
      baseline/no-cache profile.

  routinator
      Alias of routinator-latest-release. Uses the archived Routinator release
      binary and enables ASPA output when supported.

  rpki-client
      Alias of rpki-client-latest-release. Uses archived official rpki-client
      release binary with the same parameters as the #100 experiment.

  fort
      Alias of fort-latest-release. Uses archived FORT validator release binary.

Cache difference for ours RP:
  --rp ours adds:
    --enable-publication-point-validation-cache
    --enable-roa-validation-cache
    --enable-child-certificate-validation-cache

  --rp ours-no-cache does not add those flags.

Typical usage:
  ./rpctl init --bundle /root/rpki_4rp_ctl_runtime_bundle
  ./rpctl list
  ./rpctl start --rp ours --interval 10m --runs 6 --rirs all5
  ./rpctl start --rp ours-no-cache --interval 10m --runs 6 --rirs all5
  ./rpctl status
  ./rpctl show
  ./rpctl show --failed
  ./rpctl view --run-id 20260709T081122Z_ours-pp-object-cache_2runs_600s
  ./rpctl delete --run-id old_run_id --dry-run
  ./rpctl tail -n 120
  ./rpctl stop

Data model:
  - Run data stays under ./runs/<profile>/<run-id>/.
  - The first run is snapshot; later runs are delta and reuse the profile state.
  - By default only one RP/profile can run at a time.
""",
    )
    sub = parser.add_subparsers(dest="cmd", required=True, metavar="COMMAND")

    p = sub.add_parser(
        "init",
        help="initialize runtime files from a #100 runtime bundle",
        description="Initialize rpctl by copying the archived runtime bundle into this local control directory.",
        formatter_class=formatter,
        epilog="""
The bundle must contain:
  bin/rpki
  bin/routinator
  bin/rpki-client
  bin/fort
  lib/
  fixtures/tal/
  scripts/normalize_and_collect.py
  scripts/write_group_report.py
  versions.json

Recommended source on 47.77.237.41:
  /root/rpki_4rp_ctl_runtime_bundle

Example:
  ./rpctl init --bundle /root/rpki_4rp_ctl_runtime_bundle
""",
    )
    p.add_argument("--bundle", required=True, help="Path to the archived #100 runtime bundle root.")
    p.set_defaults(func=cmd_init)

    p = sub.add_parser(
        "list",
        help="list supported RPs/profiles and aliases",
        description="Show supported profile names, underlying RP name, and accepted aliases.",
        formatter_class=formatter,
        epilog="""
Use this command when unsure which value to pass to --rp.

Important ours RP profiles:
  ours-pp-object-cache  cache-enabled profile
  ours-no-cache         no-cache baseline profile
""",
    )
    p.set_defaults(func=cmd_list)

    p = sub.add_parser(
        "start",
        help="start one RP/profile in background",
        description="Start a single RP/profile run loop in the background on this server.",
        formatter_class=formatter,
        epilog="""
Run semantics:
  --runs 1 means one snapshot run.
  --runs 6 means one snapshot run plus five delta runs.
  --interval controls target start-to-start spacing between adjacent runs.
  The first run always resets that profile's state directory.
  Delta runs reuse the same profile state directory.

RP/profile examples:
  ./rpctl start --rp ours --interval 10m --runs 6 --rirs all5
  ./rpctl start --rp ours-no-cache --interval 10m --runs 6 --rirs all5
  ./rpctl start --rp routinator --interval 15m --runs 2 --rirs apnic
  ./rpctl start --rp rpki-client --interval 0 --runs 1 --rirs apnic,arin
  ./rpctl start --rp fort --interval 0 --runs 1 --rirs all5

Cache mapping:
  --rp ours         enables PP cache + ROA cache + child certificate cache.
  --rp ours-no-cache disables those cache flags.

Safety:
  rpctl refuses to start a second RP while another rpctl runner is alive.
""",
    )
    p.add_argument(
        "--rp",
        required=True,
        help="RP/profile or alias. Common values: ours, ours-no-cache, routinator, rpki-client, fort.",
    )
    p.add_argument(
        "--interval",
        default="10m",
        help="Target start-to-start interval between runs. Supports seconds/minutes/hours, e.g. 0, 30s, 10m, 1h. Default: 10m.",
    )
    p.add_argument(
        "--runs",
        type=int,
        default=6,
        help="Total run count. 1=snapshot only; 6=snapshot+5 deltas. Default: 6.",
    )
    p.add_argument(
        "--rirs",
        default="all5",
        help="RIR set: all5/all, or comma list from afrinic,apnic,arin,lacnic,ripe. Default: all5.",
    )
    p.add_argument(
        "--run-id",
        default="",
        help="Optional run id used in ./runs/<profile>/<run-id>. Default: generated UTC id.",
    )
    p.set_defaults(func=cmd_start)

    p = sub.add_parser(
        "continue",
        help="append more runs to a successful historical run id",
        description="Continue an existing successful run id by appending additional runs without resetting RP state.",
        formatter_class=formatter,
        epilog="""
Run semantics:
  --run-id must refer to a historical run whose existing child runs all exited 0.
  The command appends run_000N directories under the same run root.
  rpctl does not clean profile state and does not force snapshot-only flags.
  Each RP decides delta/snapshot behavior from its persisted state and native implementation.
  The original RIR set is reused; continue does not allow changing scope mid-run.

Examples:
  ./rpctl continue --run-id 20260709T081122Z_ours-pp-object-cache_2runs_600s --runs 3 --interval 10m
  ./rpctl continue --rp routinator --run-id m18_237_routinator_all5_snapshot_delta_20260710T002726Z --runs 1 --interval 0
""",
    )
    p.add_argument("--run-id", required=True, help="Existing successful run id to append to.")
    p.add_argument("--rp", default="", help="Optional RP/profile filter if the run id exists under multiple profiles.")
    p.add_argument(
        "--interval",
        default="10m",
        help="Target start-to-start interval between newly appended runs. The first appended run starts immediately. Default: 10m.",
    )
    p.add_argument("--runs", type=int, default=1, help="Number of additional runs to append. Default: 1.")
    p.set_defaults(func=cmd_continue)

    p = sub.add_parser(
        "stop",
        help="stop active run",
        description="Stop the currently active rpctl runner process group.",
        formatter_class=formatter,
        epilog="""
Default stop sends SIGTERM to the runner process group and waits briefly.
Use --force to send SIGKILL when the RP does not exit.

Stopping does not delete run data.

Examples:
  ./rpctl stop
  ./rpctl stop --force
""",
    )
    p.add_argument("--force", action="store_true", help="Send SIGKILL instead of SIGTERM.")
    p.set_defaults(func=cmd_stop)

    p = sub.add_parser(
        "status",
        help="show current status",
        description="Show current or last rpctl run status from local state files and PID liveness.",
        formatter_class=formatter,
        epilog="""
Human output includes:
  state, profile, rp, pid, pidAlive, runId, runRoot,
  currentRun/currentMode, lastRun wall/RSS/counts if available.

States:
  idle      no state file and no runner pid
  starting  runner just launched
  running   currently executing an RP run
  waiting   waiting for next interval
  stopped   stopped by rpctl stop
  failed    runner finished with non-zero exit code
  success   all requested runs completed successfully
  exited    state said running, but PID is gone

Examples:
  ./rpctl status
  ./rpctl status --json
""",
    )
    p.add_argument("--json", action="store_true", help="Print machine-readable JSON.")
    p.set_defaults(func=cmd_status)

    p = sub.add_parser(
        "paths",
        help="show important paths",
        description="Print local control, state, log, bundle and run data paths.",
        formatter_class=formatter,
        epilog="""
Useful paths:
  root       rpctl installation root
  config     init metadata
  state      current.json, runner.pid, runner.pgid
  runs       all generated run data
  logs       runner nohup logs
  bundle     original bundle path recorded by init

Example:
  ./rpctl paths
""",
    )
    p.set_defaults(func=cmd_paths)

    p = sub.add_parser(
        "tail",
        help="show runner log tail",
        description="Print the tail of the latest/current rpctl runner log.",
        formatter_class=formatter,
        epilog="""
This shows rpctl runner errors and Python stack traces if the wrapper fails.
Per-RP stdout/stderr are stored inside the current run directory.
Use ./rpctl paths to find the latest run directory.

Examples:
  ./rpctl tail
  ./rpctl tail -n 200
""",
    )
    p.add_argument("-n", "--lines", type=int, default=80, help="Number of log lines to print. Default: 80.")
    p.set_defaults(func=cmd_tail)

    p = sub.add_parser(
        "show",
        help="list historical run ids and summaries",
        description="List historical rpctl run roots with status, start/end time, parameters and product counts.",
        formatter_class=formatter,
        epilog="""
Show scans ./runs/<profile>/<run-id>/ and does not require the run to be active.
It can identify failed/partial runs from exit codes, missing end time, or current status.

Examples:
  ./rpctl show
  ./rpctl show --failed
  ./rpctl show --rp ours --limit 20
  ./rpctl show --paths
  ./rpctl show --json
""",
    )
    p.add_argument("--rp", default="", help="Optional RP/profile filter, e.g. ours, ours-no-cache, routinator, rpki-client, fort.")
    p.add_argument("--failed", action="store_true", help="Only show failed, partial, incomplete or exited runs.")
    p.add_argument("--limit", type=int, default=50, help="Maximum rows to show. 0 means no limit. Default: 50.")
    p.add_argument("--paths", action="store_true", help="Also show full run root paths. Hidden by default to keep the table readable.")
    p.add_argument("--json", action="store_true", help="Print machine-readable JSON.")
    p.set_defaults(func=cmd_show)

    p = sub.add_parser(
        "view",
        help="show one historical run summary",
        description="Show summary and per-run details for a specific historical run id.",
        formatter_class=formatter,
        epilog="""
Use show first to find the run id. If the same run id exists under multiple
profiles, add --rp to disambiguate.

Examples:
  ./rpctl view --run-id 20260709T081122Z_ours-pp-object-cache_2runs_600s
  ./rpctl view --rp ours --run-id 20260709T081122Z_ours-pp-object-cache_2runs_600s
  ./rpctl view --run-id m3_ours_apnic_smoke --command
  ./rpctl view --run-id m3_ours_apnic_smoke --json
""",
    )
    p.add_argument("--run-id", required=True, help="Run id under ./runs/<profile>/<run-id>.")
    p.add_argument("--rp", default="", help="Optional RP/profile filter to disambiguate duplicated run ids.")
    p.add_argument("--command", action="store_true", help="Also print the first run command line.")
    p.add_argument("--json", action="store_true", help="Print machine-readable JSON.")
    p.set_defaults(func=cmd_view)

    p = sub.add_parser(
        "delete",
        help="delete one non-running run id and its local data",
        description="Delete a historical run root and matching runner logs. Active/running runs are refused.",
        formatter_class=formatter,
        epilog="""
Safety rules:
  --run-id must point to a non-running run.
  If the same run id exists under multiple profiles, add --rp.
  Use --dry-run first to print the exact run root and matching runner logs.
  Actual deletion requires --yes.

Examples:
  ./rpctl delete --run-id 20260710T025439Z_ours-pp-object-cache_4runs_180s --dry-run
  ./rpctl delete --run-id 20260710T025439Z_ours-pp-object-cache_4runs_180s --yes
  ./rpctl delete --rp ours --run-id duplicated_run_id --yes
""",
    )
    p.add_argument("--run-id", required=True, help="Historical non-running run id to delete.")
    p.add_argument("--rp", default="", help="Optional RP/profile filter to disambiguate duplicated run ids.")
    p.add_argument("--dry-run", action="store_true", help="Only print the paths that would be deleted.")
    p.add_argument("--yes", action="store_true", help="Confirm deletion.")
    p.set_defaults(func=cmd_delete)

    p = sub.add_parser("_runner", help=argparse.SUPPRESS)
    p.add_argument("--profile", required=True)
    p.add_argument("--run-id", required=True)
    p.add_argument("--runs", type=int, required=True)
    p.add_argument("--interval-secs", type=int, required=True)
    p.add_argument("--rirs", required=True)
    p.add_argument("--append", action="store_true")
    p.set_defaults(func=cmd_runner)
    return parser

def main() -> None:
    args = build_parser().parse_args()
    args.func(args)


if __name__ == "__main__":
    main()
