Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 42 additions & 0 deletions scripts/aiperf_runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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"
Expand Down
3 changes: 3 additions & 0 deletions scripts/benchmark_routing_algorithms.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
44 changes: 44 additions & 0 deletions tests/test_aiperf_runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,12 +4,15 @@
import json
import sys
import time
from pathlib import Path

import pytest

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(
Expand Down Expand Up @@ -99,6 +102,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: Path, error_type: str, log_message: str, recovers: bool
) -> None:
Comment thread
coderabbitai[bot] marked this conversation as resolved.
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
Expand Down
7 changes: 6 additions & 1 deletion tests/test_routing_performance_report.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@

import csv
import json
from collections.abc import Sequence
from dataclasses import replace
from pathlib import Path

Expand Down Expand Up @@ -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
_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)
Expand Down
Loading