diff --git a/scripts/benchmark/run_sync_worker_ablation.py b/scripts/benchmark/run_sync_worker_ablation.py new file mode 100644 index 0000000..57181c5 --- /dev/null +++ b/scripts/benchmark/run_sync_worker_ablation.py @@ -0,0 +1,667 @@ +#!/usr/bin/env python3 +"""Run ours RP sync-worker-count ablation on a remote host. + +The experiment intentionally mirrors the Feature #100 ours-pp-object-cache +profile and changes only --parallel-max-repo-sync-workers-global. +""" + +from __future__ import annotations + +import argparse +import csv +import hashlib +import html +import json +import os +import re +import shlex +import subprocess +import sys +import time +from datetime import datetime, timezone +from pathlib import Path +from typing import Any + + +PRODUCT_ROOT = Path(__file__).resolve().parents[2] +DEFAULT_REMOTE = "root@47.77.204.233" +DEFAULT_WORKERS = "1,2,4,8" +DEFAULT_INTERVAL_SECS = 0 +DEFAULT_RUNS_PER_WORKER = 5 + +RIRS = ["afrinic", "apnic", "arin", "lacnic", "ripe"] +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", +} + + +REMOTE_RUNNER = r"""#!/usr/bin/env bash +set -euo pipefail +mode="$1" +state_dir="$2" +run_dir="$3" +shift 3 +mkdir -p "$run_dir" +if [[ "$mode" == "snapshot" ]]; then + rm -rf "$state_dir" +fi +mkdir -p "$state_dir" "$(dirname "$run_dir")" +date -u +"%Y-%m-%dT%H:%M:%SZ" > "$run_dir/start-time.txt" +printf '%q ' "$@" > "$run_dir/command.txt" +printf '\n' >> "$run_dir/command.txt" +set +e +/usr/bin/time -v -o "$run_dir/process-time.txt" "$@" > "$run_dir/stdout.log" 2> "$run_dir/stderr.log" +status=$? +set -e +date -u +"%Y-%m-%dT%H:%M:%SZ" > "$run_dir/end-time.txt" +printf '%s\n' "$status" > "$run_dir/exit-code.txt" +exit "$status" +""" + + +REMOTE_METRICS = r"""import csv +import json +import os +import re +import subprocess +import sys +from pathlib import Path + +run_dir = Path(os.environ["RUN_DIR"]) +state_dir = Path(os.environ["STATE_DIR"]) +workers = int(os.environ["WORKERS"]) +mode = os.environ["MODE"] +run_index = int(os.environ["RUN_INDEX"]) +scheduled_epoch = float(os.environ["SCHEDULED_EPOCH"]) +started_epoch = float(os.environ["STARTED_EPOCH"]) + + +def read_text(path): + try: + return Path(path).read_text(errors="replace") + except FileNotFoundError: + return "" + + +def read_json(path): + try: + with Path(path).open() as fh: + return json.load(fh) + except FileNotFoundError: + return None + except json.JSONDecodeError as exc: + return {"_json_error": str(exc)} + + +def parse_time(path): + text = read_text(path) + out = {} + match = re.search(r"Elapsed \(wall clock\) time.*: (?:(\d+):)?(\d+):(\d+(?:\.\d+)?)", text) + if match: + hours = int(match.group(1) or 0) + minutes = int(match.group(2)) + seconds = float(match.group(3)) + out["wall_ms_timev"] = int(round(((hours * 60 + minutes) * 60 + seconds) * 1000)) + for name, key in [ + ("Maximum resident set size \\(kbytes\\)", "max_rss_kb"), + ("Percent of CPU this job got", "cpu_percent_text"), + ("Major \\(requiring I/O\\) page faults", "major_faults"), + ("Minor \\(reclaiming a frame\\) page faults", "minor_faults"), + ]: + found = re.search(name + r":\s+(.+)", text) + if found: + value = found.group(1).strip() + if key.endswith("_faults") or key == "max_rss_kb": + try: + out[key] = int(value) + except ValueError: + out[key] = value + else: + out[key] = value + return out + + +def count_csv(path): + rows = 0 + unique = set() + try: + with Path(path).open(newline="") as fh: + reader = csv.reader(fh) + header = next(reader, None) + for row in reader: + if not row: + continue + rows += 1 + unique.add(tuple(row)) + except FileNotFoundError: + return {"rows": None, "unique_rows": None} + return {"rows": rows, "unique_rows": len(unique)} + + +def dir_bytes(path): + path = Path(path) + if not path.exists(): + return 0 + total = 0 + for item in path.rglob("*"): + try: + if item.is_file(): + total += item.stat().st_size + except OSError: + pass + return total + + +def file_bytes(path): + try: + return Path(path).stat().st_size + except FileNotFoundError: + return None + + +stage = read_json(run_dir / "stage-timing.json") or {} +process_time = parse_time(run_dir / "process-time.txt") +exit_text = read_text(run_dir / "exit-code.txt").strip() +start_text = read_text(run_dir / "start-time.txt").strip() +end_text = read_text(run_dir / "end-time.txt").strip() + +df = subprocess.run(["df", "-P", "/"], text=True, capture_output=True) +df_line = df.stdout.strip().splitlines()[-1].split() if df.returncode == 0 and len(df.stdout.strip().splitlines()) >= 2 else [] +df_root = {} +if len(df_line) >= 5: + df_root = { + "filesystem": df_line[0], + "blocks_1k": int(df_line[1]), + "used_1k": int(df_line[2]), + "available_1k": int(df_line[3]), + "use_percent": int(df_line[4].rstrip("%")), + } + +metrics = { + "workers": workers, + "mode": mode, + "run_index": run_index, + "scheduled_epoch": scheduled_epoch, + "started_epoch": started_epoch, + "schedule_lag_ms": int(round((started_epoch - scheduled_epoch) * 1000)), + "start_time_utc": start_text, + "end_time_utc": end_text, + "exit_code": int(exit_text) if exit_text.isdigit() else None, + "run_dir": str(run_dir), + "state_dir": str(state_dir), + "stage": stage, + "process_time": process_time, + "vrps": count_csv(run_dir / "vrps.csv"), + "vaps": count_csv(run_dir / "vaps.csv"), + "artifact_bytes": { + "report_json": file_bytes(run_dir / "report.json"), + "result_ccr": file_bytes(run_dir / "result.ccr"), + "input_cir": file_bytes(run_dir / "input.cir"), + "vrps_csv": file_bytes(run_dir / "vrps.csv"), + "vaps_csv": file_bytes(run_dir / "vaps.csv"), + }, + "state_bytes": dir_bytes(state_dir), + "run_bytes": dir_bytes(run_dir), + "df_root": df_root, +} +print(json.dumps(metrics, ensure_ascii=False, sort_keys=True)) +""" + + +def utc_stamp() -> str: + return datetime.now(timezone.utc).strftime("%Y%m%dT%H%M%SZ") + + +def run_cmd(cmd: list[str], *, cwd: Path | None = None, check: bool = True, input_text: str | None = None) -> subprocess.CompletedProcess[str]: + result = subprocess.run(cmd, cwd=cwd, input=input_text, text=True, stdout=subprocess.PIPE, stderr=subprocess.PIPE) + if check and result.returncode != 0: + sys.stderr.write(result.stdout) + sys.stderr.write(result.stderr) + raise SystemExit(result.returncode) + return result + + +def ssh(remote: str, command: str, *, check: bool = True, input_text: str | None = None) -> subprocess.CompletedProcess[str]: + return run_cmd(["ssh", "-o", "BatchMode=yes", "-o", "ConnectTimeout=10", remote, command], check=check, input_text=input_text) + + +def scp_to(remote: str, local: Path, remote_path: str) -> None: + ssh(remote, "mkdir -p " + shlex.quote(str(Path(remote_path).parent))) + run_cmd(["scp", "-q", str(local), f"{remote}:{remote_path}"]) + + +def parse_workers(value: str) -> list[int]: + workers = [] + for part in value.split(","): + part = part.strip() + if not part: + continue + worker = int(part) + if worker <= 0: + raise argparse.ArgumentTypeError("worker count must be positive") + workers.append(worker) + if len(set(workers)) != len(workers): + raise argparse.ArgumentTypeError("worker counts must be unique") + return workers + + +def parse_args() -> argparse.Namespace: + parser = argparse.ArgumentParser(description="Run ours RP sync worker ablation.") + parser.add_argument("--remote", default=DEFAULT_REMOTE) + parser.add_argument("--workers", type=parse_workers, default=parse_workers(DEFAULT_WORKERS)) + parser.add_argument("--interval-secs", type=int, default=DEFAULT_INTERVAL_SECS) + parser.add_argument("--runs-per-worker", type=int, default=DEFAULT_RUNS_PER_WORKER) + parser.add_argument("--remote-root") + parser.add_argument("--out-root", type=Path, default=Path("specs/develop/20260714/feature115_sync_worker_ablation_runs")) + parser.add_argument("--binary", type=Path) + parser.add_argument("--skip-build", action="store_true") + parser.add_argument("--dry-run", action="store_true") + return parser.parse_args() + + +def ensure_binary(args: argparse.Namespace) -> Path: + if args.binary: + binary = args.binary.resolve() + if not binary.exists(): + raise SystemExit(f"binary does not exist: {binary}") + return binary + binary = PRODUCT_ROOT / "target" / "release" / "rpki" + if not args.skip_build: + print("building release rpki binary", flush=True) + run_cmd(["cargo", "build", "--release", "--bin", "rpki"], cwd=PRODUCT_ROOT) + if not binary.exists(): + raise SystemExit(f"release binary not found: {binary}") + return binary + + +def fixed_args(remote_bin: str, state_dir: str, run_dir: str, workers: int) -> list[str]: + args = [ + remote_bin, + "--db", f"{state_dir}/work-db", + "--repo-bytes-db", f"{state_dir}/repo-bytes.db", + "--rsync-mirror-root", f"{state_dir}/rsync-mirror", + "--rsync-scope", "module-root", + "--report-json", f"{run_dir}/report.json", + "--report-json-compact", + "--ccr-out", f"{run_dir}/result.ccr", + "--cir-enable", + "--cir-out", f"{run_dir}/input.cir", + "--vrps-csv-out", f"{run_dir}/vrps.csv", + "--vaps-csv-out", f"{run_dir}/vaps.csv", + "--compare-view-trust-anchor", "all5", + "--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", + "--parallel-max-repo-sync-workers-global", str(workers), + "--enable-publication-point-validation-cache", + "--enable-roa-validation-cache", + "--enable-child-certificate-validation-cache", + ] + for rir in RIRS: + url = TAL_URLS[rir] + args.extend(["--tal-url", url, "--cir-tal-uri", url]) + return args + + +def plan_runs(workers: list[int], runs_per_worker: int, interval_secs: int, remote_root: str, remote_bin: str) -> list[dict[str, Any]]: + if runs_per_worker < 2: + raise SystemExit("--runs-per-worker must be at least 2") + planned = [] + ordinal = 0 + for worker in workers: + case_root = f"{remote_root}/cases/workers-{worker:02d}" + state_dir = f"{case_root}/state" + for run_index in range(1, runs_per_worker + 1): + mode = "snapshot" if run_index == 1 else "delta" + run_dir = f"{case_root}/runs/run_{run_index:04d}_{mode}" + planned.append({ + "ordinal": ordinal, + "workers": worker, + "mode": mode, + "run_index": run_index, + "state_dir": state_dir, + "run_dir": run_dir, + "interval_secs": interval_secs, + "command": fixed_args(remote_bin, state_dir, run_dir, worker), + }) + ordinal += 1 + return planned + + +def write_json(path: Path, data: Any) -> None: + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(json.dumps(data, indent=2, ensure_ascii=False, sort_keys=True) + "\n") + + +def write_text(path: Path, text: str) -> None: + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(text) + + +def prepare_remote(remote: str, remote_root: str, binary: Path) -> str: + remote_bin = f"{remote_root}/bin/rpki" + preflight = f""" +set -euo pipefail +mkdir -p {shlex.quote(remote_root)}/bin {shlex.quote(remote_root)}/env +systemctl disable --now rpki-client.timer >/dev/null 2>&1 || true +systemctl stop rpki-client.service >/dev/null 2>&1 || true +pkill -x rpki-client >/dev/null 2>&1 || true +pkill -x routinator >/dev/null 2>&1 || true +pkill -x fort >/dev/null 2>&1 || true +pkill -x rpki >/dev/null 2>&1 || true +{{ + date -u +"%Y-%m-%dT%H:%M:%SZ" + hostname + uname -a + nproc + df -h / + free -h || true + ps -eo pid,ppid,comm,%cpu,%mem,args --sort=-%cpu | head -30 + systemctl is-active rpki-client.timer 2>/dev/null || true +}} > {shlex.quote(remote_root)}/env/preflight.txt +""" + ssh(remote, "bash -s", input_text=preflight) + scp_to(remote, binary, remote_bin) + ssh(remote, "chmod +x " + shlex.quote(remote_bin)) + return remote_bin + + +def run_remote_plan(remote: str, plan: dict[str, Any], scheduled_epoch: float) -> dict[str, Any]: + now = time.time() + if now < scheduled_epoch: + wait_secs = int(round(scheduled_epoch - now)) + print(f"waiting {wait_secs}s before workers={plan['workers']} {plan['mode']}", flush=True) + time.sleep(scheduled_epoch - now) + started_epoch = time.time() + command = shlex.join(["bash", "-s", "--", plan["mode"], plan["state_dir"], plan["run_dir"], *plan["command"]]) + print(f"start workers={plan['workers']} {plan['mode']} run={plan['run_index']} at {datetime.now(timezone.utc).isoformat()}", flush=True) + result = ssh(remote, command, input_text=REMOTE_RUNNER, check=False) + if result.returncode != 0: + sys.stderr.write(result.stdout) + sys.stderr.write(result.stderr) + metrics_script = "\n".join([ + f"export RUN_DIR={shlex.quote(plan['run_dir'])}", + f"export STATE_DIR={shlex.quote(plan['state_dir'])}", + f"export WORKERS={plan['workers']}", + f"export MODE={shlex.quote(plan['mode'])}", + f"export RUN_INDEX={plan['run_index']}", + f"export SCHEDULED_EPOCH={scheduled_epoch}", + f"export STARTED_EPOCH={started_epoch}", + "python3 - <<'PY'", + REMOTE_METRICS, + "PY", + ]) + metrics_result = ssh(remote, "bash -s", input_text=metrics_script, check=True) + metrics = json.loads(metrics_result.stdout) + metrics["ssh_return_code"] = result.returncode + metrics["command_digest"] = hashlib.sha256("\0".join(plan["command"]).encode()).hexdigest() + ssh( + remote, + "cat > " + shlex.quote(plan["run_dir"] + "/metrics.json"), + input_text=json.dumps(metrics, ensure_ascii=False, sort_keys=True) + "\n", + ) + if result.returncode != 0: + raise RuntimeError(f"remote run failed: workers={plan['workers']} mode={plan['mode']} rc={result.returncode}") + validate_metrics(metrics) + print(f"done workers={plan['workers']} {plan['mode']} wall={metrics.get('wall_ms')}ms vrps={metrics['vrps'].get('rows')} vaps={metrics['vaps'].get('rows')}", flush=True) + return metrics + + +def validate_metrics(metrics: dict[str, Any]) -> None: + stage = metrics.get("stage") or {} + expected = { + "enable_transport_request_prefetch": False, + "enable_publication_point_validation_cache": True, + "enable_roa_validation_cache": True, + "enable_child_certificate_validation_cache": True, + } + for key, value in expected.items(): + if stage.get(key) is not value: + raise RuntimeError(f"unexpected {key}: {stage.get(key)!r}, expected {value!r}") + if metrics["vrps"].get("rows") in (None, 0): + raise RuntimeError("VRP CSV is missing or empty") + if metrics["vaps"].get("rows") is None: + raise RuntimeError("VAP CSV is missing") + df_root = metrics.get("df_root") or {} + if df_root.get("use_percent", 0) >= 90: + raise RuntimeError(f"remote disk is above threshold: {df_root}") + stage_total = stage.get("total_ms") + timev_total = metrics.get("process_time", {}).get("wall_ms_timev") + metrics["wall_ms"] = stage_total if isinstance(stage_total, int) else timev_total + + +def append_metric(out_root: Path, metrics: dict[str, Any]) -> None: + path = out_root / "per_run_metrics.jsonl" + with path.open("a") as fh: + fh.write(json.dumps(metrics, ensure_ascii=False, sort_keys=True) + "\n") + + +def load_metrics(out_root: Path) -> list[dict[str, Any]]: + path = out_root / "per_run_metrics.jsonl" + if not path.exists(): + return [] + return [json.loads(line) for line in path.read_text().splitlines() if line.strip()] + + +def metric_row(metric: dict[str, Any]) -> dict[str, Any]: + stage = metric.get("stage") or {} + roa_cache = stage.get("roa_validation_cache") or {} + return { + "workers": metric["workers"], + "mode": metric["mode"], + "run_index": metric["run_index"], + "start_time_utc": metric.get("start_time_utc"), + "wall_ms": metric.get("wall_ms"), + "timev_wall_ms": metric.get("process_time", {}).get("wall_ms_timev"), + "max_rss_kb": metric.get("process_time", {}).get("max_rss_kb"), + "cpu_percent": metric.get("process_time", {}).get("cpu_percent_text"), + "validation_ms": stage.get("validation_ms"), + "repo_sync_ms_total": stage.get("repo_sync_ms_total"), + "rrdp_download_ms_total": stage.get("rrdp_download_ms_total"), + "rsync_download_ms_total": stage.get("rsync_download_ms_total"), + "download_bytes_total": stage.get("download_bytes_total"), + "publication_points": stage.get("publication_points"), + "pp_cache_hit_ratio": stage.get("publication_point_cache_hit_ratio"), + "roa_hit_roas": roa_cache.get("hit_roas"), + "roa_miss_roas": roa_cache.get("miss_roas"), + "vrps_rows": metric.get("vrps", {}).get("rows"), + "vrps_unique_rows": metric.get("vrps", {}).get("unique_rows"), + "vaps_rows": metric.get("vaps", {}).get("rows"), + "vaps_unique_rows": metric.get("vaps", {}).get("unique_rows"), + "state_bytes": metric.get("state_bytes"), + "run_bytes": metric.get("run_bytes"), + "remote_run_dir": metric.get("run_dir"), + } + + +def write_csv_report(out_root: Path, metrics: list[dict[str, Any]]) -> None: + rows = [metric_row(item) for item in metrics] + if not rows: + return + with (out_root / "per_run_metrics.csv").open("w", newline="") as fh: + writer = csv.DictWriter(fh, fieldnames=list(rows[0])) + writer.writeheader() + writer.writerows(rows) + + +def fmt_ms(value: Any) -> str: + if value is None: + return "-" + try: + return f"{float(value) / 1000:.2f}s" + except (TypeError, ValueError): + return str(value) + + +def fmt_bytes(value: Any) -> str: + if value is None: + return "-" + try: + value = float(value) + except (TypeError, ValueError): + return str(value) + for unit in ["B", "KiB", "MiB", "GiB"]: + if value < 1024 or unit == "GiB": + return f"{value:.1f} {unit}" + value /= 1024 + return str(value) + + +def write_markdown_report(out_root: Path, metrics: list[dict[str, Any]], provenance: dict[str, Any]) -> None: + rows = [metric_row(item) for item in metrics] + lines = [ + "# Feature #115 Sync Worker Ablation Summary", + "", + "## Experiment", + "", + f"- Remote: {provenance['remote']}", + f"- Remote root: {provenance['remote_root']}", + f"- Local commit: {provenance['git_commit']}", + f"- Worker counts: {', '.join(str(w) for w in provenance['workers'])}", + f"- Runs per worker: {provenance['runs_per_worker']}", + f"- Interval seconds: {provenance['interval_secs']}", + "- Fixed profile: all5 live TAL URL, module-root rsync scope, PP cache enabled, ROA cache enabled, child-certificate cache enabled, transport prefetch disabled.", + "", + "## Per Run", + "", + "| workers | mode | wall | rss | validation | repo sync | RRDP | rsync | download | PP | VRPs | VAPs |", + "|---:|---|---:|---:|---:|---:|---:|---:|---:|---:|---:|---:|", + ] + for row in rows: + lines.append( + f"| {row['workers']} | {row['mode']} | {fmt_ms(row['wall_ms'])} | " + f"{fmt_bytes((row['max_rss_kb'] or 0) * 1024 if row['max_rss_kb'] else None)} | " + f"{fmt_ms(row['validation_ms'])} | {fmt_ms(row['repo_sync_ms_total'])} | " + f"{fmt_ms(row['rrdp_download_ms_total'])} | {fmt_ms(row['rsync_download_ms_total'])} | " + f"{fmt_bytes(row['download_bytes_total'])} | {row['publication_points']} | " + f"{row['vrps_rows']} | {row['vaps_rows']} |" + ) + lines.extend(["", "## Artifacts", "", f"- Per-run metrics JSONL: {out_root / 'per_run_metrics.jsonl'}", f"- Per-run metrics CSV: {out_root / 'per_run_metrics.csv'}"]) + write_text(out_root / "summary.md", "\n".join(lines) + "\n") + + +def write_html_report(out_root: Path, metrics: list[dict[str, Any]], provenance: dict[str, Any]) -> None: + rows = [metric_row(item) for item in metrics] + if not rows: + return + max_wall = max(float(row["wall_ms"] or 0) for row in rows) or 1 + bars = [] + for row in rows: + width = max(1, int((float(row["wall_ms"] or 0) / max_wall) * 100)) + bars.append( + f"
This report compares only the global repo-sync worker count while holding the Feature #100 ours-cache profile fixed.
+ +| Workers | Mode | Wall | Wall Bar | Max RSS | Validation | Repo Sync | PP | VRPs | VAPs |
|---|
all5 live TAL URL, module-root rsync scope, PP cache enabled, ROA cache enabled, child-certificate cache enabled, transport prefetch disabled.
+Local commit: {html.escape(provenance['git_commit'])}. Full JSONL and CSV evidence are stored beside this report.