From b43a98ed07f879c60e2c82b29846e0022f56d7ac Mon Sep 17 00:00:00 2001 From: ghangz <152254226+ghangz@users.noreply.github.com> Date: Wed, 1 Jul 2026 11:42:29 +0800 Subject: [PATCH 1/2] Add benchmark log summarizer --- tools/benchmark_log_summary.py | 57 ++++++++++++++++++++++++++++++++++ 1 file changed, 57 insertions(+) create mode 100644 tools/benchmark_log_summary.py diff --git a/tools/benchmark_log_summary.py b/tools/benchmark_log_summary.py new file mode 100644 index 0000000..b26c121 --- /dev/null +++ b/tools/benchmark_log_summary.py @@ -0,0 +1,57 @@ +#!/usr/bin/env python3 +"""Parse benchmark or distributed-test logs into a compact JSON summary.""" + +from __future__ import annotations + + +import argparse +import json +import re +from pathlib import Path + + +METRIC_RE = re.compile('(?P[A-Za-z0-9_./-]+).*?(?P[0-9.]+)\\s*(?:us|ms)', re.IGNORECASE) +ERROR_RE = re.compile(r"(error|failed|timeout|traceback|segmentation fault|core dumped)", re.IGNORECASE) + + +def parse_log(path: Path) -> dict[str, object]: + metrics: list[dict[str, object]] = [] + errors: list[dict[str, object]] = [] + for lineno, line in enumerate(path.read_text(encoding="utf-8", errors="replace").splitlines(), start=1): + match = METRIC_RE.search(line) + if match: + item = {"line": lineno, "text": line.strip()} + item.update({k: v for k, v in match.groupdict().items() if v is not None}) + metrics.append(item) + if ERROR_RE.search(line): + errors.append({"line": lineno, "text": line.strip()}) + return {"path": str(path), "metrics": metrics, "errors": errors, "metric_count": len(metrics), "error_count": len(errors)} + + +def self_test() -> None: + sample = Path("_sample_log_for_parser.txt") + sample.write_text("all_reduce size 1024 algbw 11.5 busbw 10.1 latency 2.5 ms throughput 33.0 tokens/s 12.0 GB/s\nERROR timeout\n", encoding="utf-8") + try: + data = parse_log(sample) + assert data["metric_count"] >= 1 + assert data["error_count"] == 1 + print(json.dumps({"ok": True, "metric_count": data["metric_count"]}, ensure_ascii=False)) + finally: + sample.unlink(missing_ok=True) + + +def main() -> int: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("logs", nargs="*", help="Log files to parse.") + parser.add_argument("--self-test", action="store_true") + args = parser.parse_args() + if args.self_test: + self_test() + return 0 + results = [parse_log(Path(p)) for p in args.logs] + print(json.dumps(results, ensure_ascii=False, indent=2)) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) From 091650ae77908b4d465f1ea932d561f3ddfbaade Mon Sep 17 00:00:00 2001 From: ghangz <152254226+ghangz@users.noreply.github.com> Date: Wed, 1 Jul 2026 15:45:02 +0800 Subject: [PATCH 2/2] Stream benchmark log summary input --- tools/benchmark_log_summary.py | 69 ++++++++++++++++++++++++---------- 1 file changed, 49 insertions(+), 20 deletions(-) diff --git a/tools/benchmark_log_summary.py b/tools/benchmark_log_summary.py index b26c121..9d04f13 100644 --- a/tools/benchmark_log_summary.py +++ b/tools/benchmark_log_summary.py @@ -7,37 +7,60 @@ import argparse import json import re +import sys from pathlib import Path +from tempfile import TemporaryDirectory -METRIC_RE = re.compile('(?P[A-Za-z0-9_./-]+).*?(?P[0-9.]+)\\s*(?:us|ms)', re.IGNORECASE) -ERROR_RE = re.compile(r"(error|failed|timeout|traceback|segmentation fault|core dumped)", re.IGNORECASE) +METRIC_RE = re.compile( + r"(?P[A-Za-z0-9_./-]+).*?(?P[0-9.]+)\s*(?:us|ms)", + re.IGNORECASE, +) +ERROR_RE = re.compile( + r"(error|failed|timeout|traceback|segmentation fault|core dumped)", + re.IGNORECASE, +) def parse_log(path: Path) -> dict[str, object]: metrics: list[dict[str, object]] = [] errors: list[dict[str, object]] = [] - for lineno, line in enumerate(path.read_text(encoding="utf-8", errors="replace").splitlines(), start=1): - match = METRIC_RE.search(line) - if match: - item = {"line": lineno, "text": line.strip()} - item.update({k: v for k, v in match.groupdict().items() if v is not None}) - metrics.append(item) - if ERROR_RE.search(line): - errors.append({"line": lineno, "text": line.strip()}) - return {"path": str(path), "metrics": metrics, "errors": errors, "metric_count": len(metrics), "error_count": len(errors)} + with path.open(encoding="utf-8", errors="replace") as handle: + for lineno, line in enumerate(handle, start=1): + match = METRIC_RE.search(line) + if match: + item = {"line": lineno, "text": line.strip()} + item.update( + {k: v for k, v in match.groupdict().items() if v is not None} + ) + metrics.append(item) + if ERROR_RE.search(line): + errors.append({"line": lineno, "text": line.strip()}) + return { + "path": str(path), + "metrics": metrics, + "errors": errors, + "metric_count": len(metrics), + "error_count": len(errors), + } def self_test() -> None: - sample = Path("_sample_log_for_parser.txt") - sample.write_text("all_reduce size 1024 algbw 11.5 busbw 10.1 latency 2.5 ms throughput 33.0 tokens/s 12.0 GB/s\nERROR timeout\n", encoding="utf-8") - try: + with TemporaryDirectory() as tmp_dir: + sample = Path(tmp_dir) / "sample_log_for_parser.txt" + sample.write_text( + "all_reduce size 1024 algbw 11.5 busbw 10.1 latency 2.5 ms " + "throughput 33.0 tokens/s 12.0 GB/s\nERROR timeout\n", + encoding="utf-8", + ) data = parse_log(sample) - assert data["metric_count"] >= 1 - assert data["error_count"] == 1 - print(json.dumps({"ok": True, "metric_count": data["metric_count"]}, ensure_ascii=False)) - finally: - sample.unlink(missing_ok=True) + if data["metric_count"] < 1 or data["error_count"] != 1: + raise RuntimeError(f"self-test failed: {data}") + print( + json.dumps( + {"ok": True, "metric_count": data["metric_count"]}, ensure_ascii=False + ) + ) def main() -> int: @@ -48,7 +71,13 @@ def main() -> int: if args.self_test: self_test() return 0 - results = [parse_log(Path(p)) for p in args.logs] + results = [] + for p in args.logs: + path = Path(p) + if not path.is_file(): + print(f"Error: log file does not exist or is not a file: {p}", file=sys.stderr) + return 2 + results.append(parse_log(path)) print(json.dumps(results, ensure_ascii=False, indent=2)) return 0