From 1bbec8884472a1e8b52aa3019f65125197610b11 Mon Sep 17 00:00:00 2001 From: Alex Fournier Date: Fri, 18 Sep 2026 00:38:07 +0200 Subject: [PATCH 1/2] fix(benchmark): report expected cancellation runs Signed-off-by: Alex Fournier --- scripts/aiperf_runner.py | 42 +++++++++++++++++++++++ scripts/benchmark_routing_algorithms.py | 3 ++ tests/test_aiperf_runner.py | 43 ++++++++++++++++++++++++ tests/test_routing_performance_report.py | 2 +- 4 files changed, 89 insertions(+), 1 deletion(-) diff --git a/scripts/aiperf_runner.py b/scripts/aiperf_runner.py index f83165df1..a7ebe2465 100644 --- a/scripts/aiperf_runner.py +++ b/scripts/aiperf_runner.py @@ -31,6 +31,7 @@ ("request_count", "avg"), ("request_throughput", "avg"), ) +_ALL_REQUESTS_FAILED = "inference request(s) failed; no successful responses were collected." def validate_aiperf_version(binary: str) -> None: @@ -98,11 +99,48 @@ def _stop_process_group(process: subprocess.Popen[bytes]) -> None: raise RuntimeError(f"could not reap AIPerf process {process.pid}") from error +def _recover_all_timeout_export( + log_path: Path, artifact_dir: Path, expected_timeout_count: int +) -> Path | None: + """Build the summary AIPerf omits when every expected request times out.""" + try: + expected_log = f"All {expected_timeout_count} {_ALL_REQUESTS_FAILED}" + with log_path.open(encoding="utf-8", errors="replace") as log: + if not any(expected_log in line for line in log): + return None + timeout_count = 0 + with (artifact_dir / "profile_export.jsonl").open(encoding="utf-8") as records: + for line in records: + record = json.loads(line) + error = record.get("error") if isinstance(record, dict) else None + if not isinstance(error, dict) or error.get("type") != "TimeoutError": + return None + timeout_count += 1 + except (OSError, json.JSONDecodeError): + return None + if timeout_count != expected_timeout_count: + return None + export_path = artifact_dir / "profile_export_aiperf.json" + summary = { + "aiperf_version": SUPPORTED_AIPERF_VERSION, + "error_request_count": {"unit": "requests", "avg": timeout_count}, + "request_count": {"unit": "requests", "avg": 0}, + "request_throughput": {"unit": "requests/sec", "avg": 0.0}, + } + export_path.write_text( + f"{json.dumps(summary, indent=2)}\n", + encoding="utf-8", + ) + return export_path + + def run_profile( command: Sequence[str], log_path: Path, artifact_dir: Path, timeout_seconds: int, + *, + expected_timeout_count: int | None = None, ) -> Path: """Run one bounded AIPerf process and return its verified export.""" artifact_dir.parent.mkdir(parents=True, exist_ok=True) @@ -130,6 +168,10 @@ def run_profile( finally: if not stopped and (process.poll() is None or _process_group_exists(process.pid)): _stop_process_group(process) + if returncode == 1 and expected_timeout_count is not None: + recovered = _recover_all_timeout_export(log_path, artifact_dir, expected_timeout_count) + if recovered is not None: + return recovered if returncode != 0: raise RuntimeError(f"AIPerf failed with status {returncode}; see {log_path}") export_path = artifact_dir / "profile_export_aiperf.json" diff --git a/scripts/benchmark_routing_algorithms.py b/scripts/benchmark_routing_algorithms.py index 9a0a8baa7..35ba9fd4a 100755 --- a/scripts/benchmark_routing_algorithms.py +++ b/scripts/benchmark_routing_algorithms.py @@ -1317,6 +1317,9 @@ def run_aiperf( output_dir / f"{run_label}.log", trial_root, aiperf_timeout_seconds(config, scenario, profile, concurrency), + expected_timeout_count=( + config.request_count if scenario.id == "client-cancellation" else None + ), ) ) if len(exports) == 1: diff --git a/tests/test_aiperf_runner.py b/tests/test_aiperf_runner.py index 71f5fafa7..92dccab2b 100644 --- a/tests/test_aiperf_runner.py +++ b/tests/test_aiperf_runner.py @@ -10,6 +10,8 @@ import scripts.aiperf_runner as aiperf_runner from scripts.aiperf_runner import aggregate_exports, run_profile, validate_aiperf_version +_ALL_FAILED = "All 2 inference request(s) failed; no successful responses were collected." + def _write_stubborn_worker(worker) -> None: worker.write_text( @@ -99,6 +101,47 @@ def test_run_profile_stops_workers_after_the_leader_exits(tmp_path, monkeypatch) _assert_heartbeat_stopped(heartbeat) +@pytest.mark.parametrize( + ("error_type", "log_message", "recovers"), + [ + ("TimeoutError", _ALL_FAILED, True), + ("ClientConnectorError", _ALL_FAILED, False), + ("TimeoutError", "worker crashed", False), + ], +) +def test_run_profile_recovers_only_expected_timeouts( + tmp_path, error_type, log_message, recovers +) -> None: + artifact_dir = tmp_path / "artifacts" + artifact_dir.mkdir(parents=True) + record = json.dumps({"error": {"type": error_type}}) + (artifact_dir / "profile_export.jsonl").write_text(f"{record}\n{record}\n") + command = [sys.executable, "-c", f"print({log_message!r}); raise SystemExit(1)"] + + if recovers: + export = run_profile( + command, + tmp_path / "aiperf.log", + artifact_dir, + timeout_seconds=10, + expected_timeout_count=2, + ) + summary = json.loads(export.read_text()) + assert summary["error_request_count"]["avg"] == 2 + assert summary["request_count"]["avg"] == 0 + assert summary["request_throughput"]["avg"] == 0.0 + return + + with pytest.raises(RuntimeError, match="AIPerf failed with status 1"): + run_profile( + command, + tmp_path / "aiperf.log", + artifact_dir, + timeout_seconds=10, + expected_timeout_count=2, + ) + + def test_process_group_probe_ignores_an_unowned_reused_group(monkeypatch) -> None: def deny_signal(_process_group, _signal) -> None: raise PermissionError diff --git a/tests/test_routing_performance_report.py b/tests/test_routing_performance_report.py index 445357cc1..2acd5f36a 100644 --- a/tests/test_routing_performance_report.py +++ b/tests/test_routing_performance_report.py @@ -359,7 +359,7 @@ def test_aiperf_cells_use_disjoint_artifacts(tmp_path, monkeypatch) -> None: observed: list[tuple[Path, Path]] = [] def fake_run_profile( - _command, log_path: Path, artifact_dir: Path, _timeout_seconds: int + _command, log_path: Path, artifact_dir: Path, _timeout_seconds: int, **_kwargs ) -> Path: observed.append((log_path, artifact_dir)) artifact_dir.mkdir(parents=True) From a10fa035eedee021a7d1d3f6fc67e1dacb82e07a Mon Sep 17 00:00:00 2001 From: Alex Fournier Date: Sun, 27 Sep 2026 11:50:26 -0400 Subject: [PATCH 2/2] test(benchmark): annotate cancellation tests Signed-off-by: Alex Fournier --- tests/test_aiperf_runner.py | 3 ++- tests/test_routing_performance_report.py | 7 ++++++- 2 files changed, 8 insertions(+), 2 deletions(-) diff --git a/tests/test_aiperf_runner.py b/tests/test_aiperf_runner.py index 92dccab2b..e5d3bb4b4 100644 --- a/tests/test_aiperf_runner.py +++ b/tests/test_aiperf_runner.py @@ -4,6 +4,7 @@ import json import sys import time +from pathlib import Path import pytest @@ -110,7 +111,7 @@ def test_run_profile_stops_workers_after_the_leader_exits(tmp_path, monkeypatch) ], ) def test_run_profile_recovers_only_expected_timeouts( - tmp_path, error_type, log_message, recovers + tmp_path: Path, error_type: str, log_message: str, recovers: bool ) -> None: artifact_dir = tmp_path / "artifacts" artifact_dir.mkdir(parents=True) diff --git a/tests/test_routing_performance_report.py b/tests/test_routing_performance_report.py index 2acd5f36a..e58f00bea 100644 --- a/tests/test_routing_performance_report.py +++ b/tests/test_routing_performance_report.py @@ -3,6 +3,7 @@ import csv import json +from collections.abc import Sequence from dataclasses import replace from pathlib import Path @@ -359,7 +360,11 @@ def test_aiperf_cells_use_disjoint_artifacts(tmp_path, monkeypatch) -> None: observed: list[tuple[Path, Path]] = [] def fake_run_profile( - _command, log_path: Path, artifact_dir: Path, _timeout_seconds: int, **_kwargs + _command: Sequence[str], + log_path: Path, + artifact_dir: Path, + _timeout_seconds: int, + **_kwargs: object, ) -> Path: observed.append((log_path, artifact_dir)) artifact_dir.mkdir(parents=True)