From 96d82957a22ddaf91dda922811b6e1956c03e430 Mon Sep 17 00:00:00 2001 From: yuyr Date: Fri, 10 Jul 2026 12:39:24 +0800 Subject: [PATCH] =?UTF-8?q?20260710=20=E5=A2=9E=E5=8A=A0=E5=9B=9BRP?= =?UTF-8?q?=E6=9C=AC=E5=9C=B0=E6=8E=A7=E5=88=B6=E5=8F=B0rpctl?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- scripts/rpctl/README.md | 103 +++ scripts/rpctl/rpctl | 1741 +++++++++++++++++++++++++++++++++++++++ 2 files changed, 1844 insertions(+) create mode 100644 scripts/rpctl/README.md create mode 100755 scripts/rpctl/rpctl diff --git a/scripts/rpctl/README.md b/scripts/rpctl/README.md new file mode 100644 index 0000000..c0d607b --- /dev/null +++ b/scripts/rpctl/README.md @@ -0,0 +1,103 @@ +# rpctl + +`rpctl` 是四款 RP 本地控制 CLI,用于在实验服务器本地独立启动、停止和观察单个 RP/profile 的 all5 run。 + +## 部署 + +```bash +mkdir -p /root/rpki_4rp_ctl +rsync -a rpki_2/rpki/scripts/rpctl/ root@47.77.237.41:/root/rpki_4rp_ctl/ +rsync -a specs/develop/20260707_3/feature100_four_rp_interval_archive/runtime_bundle/ root@47.77.237.41:/root/rpki_4rp_ctl_runtime_bundle/ +ssh root@47.77.237.41 'cd /root/rpki_4rp_ctl && ./rpctl init --bundle /root/rpki_4rp_ctl_runtime_bundle' +``` + +`--bundle` 目录需要包含 #100 runtime bundle 的 `bin/`、`lib/`、`fixtures/`、`scripts/` 和 `versions.json`。不要使用早期 precheck 目录作为标准 bundle,因为它可能缺少 FORT 运行所需动态库。 + +## 常用命令 + +```bash +./rpctl list +./rpctl status +./rpctl start --rp ours-pp-object-cache --interval 10m --runs 6 --rirs all5 +./rpctl start --rp routinator --interval 0 --runs 1 --rirs apnic +./rpctl continue --run-id --runs 3 --interval 10m +./rpctl tail +./rpctl show +./rpctl show --failed +./rpctl view --run-id 20260709T081122Z_ours-pp-object-cache_2runs_600s +./rpctl delete --run-id --dry-run +./rpctl paths +./rpctl stop +./rpctl stop --force +``` + +## Profile + +| 输入别名 | 实际 profile | +|---|---| +| `ours`, `ours-rp`, `ours-cached` | `ours-pp-object-cache` | +| `ours-baseline` | `ours-no-cache` | +| `routinator` | `routinator-latest-release` | +| `rpki-client` | `rpki-client-latest-release` | +| `fort` | `fort-latest-release` | + +## 数据目录 + +默认部署根目录就是 `rpctl` 所在目录,例如 `/root/rpki_4rp_ctl`。 + +```text +/root/rpki_4rp_ctl/ + bin/ lib/ fixtures/ scripts/ + state/current.json + logs/runner-*.log + runs///profiles//runs/run_0001/ +``` + +每个 run 目录包含 `stdout.log`、`stderr.log`、`process-time.txt`、`exit-code.txt`、`run-meta.json` 以及各 RP 产物。 + +## 历史查看与删除 + +```bash +./rpctl show +./rpctl show --failed +./rpctl show --rp ours --limit 20 +./rpctl show --json +``` + +`show` 会扫描 `runs///`,展示历史 `run_id`、状态、start/end 时间、失败 run 列表、运行参数、产物数量和目录路径。 + +默认表格为了可读性不显示完整路径;需要路径时使用: + +```bash +./rpctl show --paths +``` + +```bash +./rpctl view --run-id +./rpctl view --rp ours --run-id +./rpctl view --run-id --command +./rpctl view --run-id --json +``` + +`view` 用于查看指定历史 run 的摘要和每一轮子 run 明细,包括运行耗时、max RSS、VRP/VAP 数量、同步模式、exit code 和目录路径。 + +```bash +./rpctl delete --run-id --dry-run +./rpctl delete --run-id --yes +./rpctl delete --rp ours --run-id --yes +``` + +`delete` 只能删除非 running 的 run。建议先用 `--dry-run` 确认路径;真实删除必须显式加 `--yes`。删除内容包括 `runs///` 和匹配的 `logs/runner-*.log`。 + +如果历史 run 是中途停止的 incomplete run,且没有 `run-meta.json`、normalized 文件或原始输出文件,则 VRP/VAP、wall、RSS 会显示为 `-`,表示该 run 没有可统计产物。 + +## 续跑已有成功 run + +```bash +./rpctl continue --run-id --runs 3 --interval 10m +./rpctl continue --rp routinator --run-id --runs 1 --interval 0 +``` + +`continue` 只接受已经成功结束的历史 run id,并在原目录下追加 `run_000N`。它不会清理 profile state,也不会给 RP 强制传 snapshot-only 参数;后续 run 是 delta 还是 snapshot 由各 RP 根据自身持久状态决定。 + +续跑会复用原 run 的 RIR 集合,避免在同一个 state 目录里混用不同 scope。若同一个 `run-id` 在多个 profile 下存在,需要加 `--rp` 消歧。 diff --git a/scripts/rpctl/rpctl b/scripts/rpctl/rpctl new file mode 100755 index 0000000..9d206ac --- /dev/null +++ b/scripts/rpctl/rpctl @@ -0,0 +1,1741 @@ +#!/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 ''}") + 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 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 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///. + - 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//. 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/// 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//.") + 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()