Add production benchmark runner for Meeting Lab
This commit is contained in:
+266
-8
@@ -21,6 +21,8 @@ from collections import Counter
|
||||
from dataclasses import asdict
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
from queue import Empty, Queue
|
||||
from threading import Thread
|
||||
from typing import Any, Callable
|
||||
|
||||
import requests
|
||||
@@ -94,6 +96,184 @@ class RunnerError(RuntimeError):
|
||||
self.details = details or {}
|
||||
|
||||
|
||||
class ProgressReporter:
|
||||
"""Reusable progress and timing reporter for benchmark runs."""
|
||||
|
||||
STAGES = (
|
||||
"preflight",
|
||||
"chunking",
|
||||
"normalization",
|
||||
"extraction",
|
||||
"canonicalizer",
|
||||
"semantic_consolidator",
|
||||
"renderer",
|
||||
)
|
||||
|
||||
DISPLAY_NAMES = {
|
||||
"preflight": "Preflight",
|
||||
"chunking": "Chunking",
|
||||
"normalization": "Normalization",
|
||||
"extraction": "Extraction",
|
||||
"canonicalizer": "Canonicalizer",
|
||||
"semantic_consolidator": "Semantic Consolidator",
|
||||
"renderer": "Renderer",
|
||||
}
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
stream: Any = sys.stdout,
|
||||
monotonic: Callable[[], float] = time.perf_counter,
|
||||
wall_clock: Callable[[], datetime] = datetime.now,
|
||||
stages: tuple[str, ...] = STAGES,
|
||||
) -> None:
|
||||
self.stream = stream
|
||||
self.monotonic = monotonic
|
||||
self.wall_clock = wall_clock
|
||||
self.stages = stages
|
||||
self.completed: set[str] = set()
|
||||
self.current_stage: str | None = None
|
||||
self.pipeline_start = self.monotonic()
|
||||
self.stage_starts: dict[str, float] = {}
|
||||
self.stage_wall_starts: dict[str, datetime] = {}
|
||||
self.stage_durations: dict[str, float] = {}
|
||||
|
||||
def stage_name(self, stage: str) -> str:
|
||||
return self.DISPLAY_NAMES.get(stage, stage.replace("_", " ").title())
|
||||
|
||||
def format_hms(self, seconds: float) -> str:
|
||||
seconds_int = max(0, int(round(seconds)))
|
||||
hours, remainder = divmod(seconds_int, 3600)
|
||||
minutes, secs = divmod(remainder, 60)
|
||||
return f"{hours:02d}:{minutes:02d}:{secs:02d}"
|
||||
|
||||
def format_mmss(self, seconds: float) -> str:
|
||||
seconds_int = max(0, int(round(seconds)))
|
||||
minutes, secs = divmod(seconds_int, 60)
|
||||
return f"{minutes:02d}:{secs:02d}"
|
||||
|
||||
def format_duration_seconds(self, seconds: float) -> str:
|
||||
return f"{seconds:.1f} s"
|
||||
|
||||
def extraction_eta(self, completed_chunks: int, total_chunks: int, elapsed: float) -> tuple[float | None, float | None]:
|
||||
if completed_chunks <= 0 or total_chunks <= 0:
|
||||
return None, None
|
||||
average = elapsed / completed_chunks
|
||||
remaining = average * max(total_chunks - completed_chunks, 0)
|
||||
return average, remaining
|
||||
|
||||
def pipeline_percent(self) -> int:
|
||||
if not self.stages:
|
||||
return 100
|
||||
completed = len(self.completed)
|
||||
if self.current_stage and self.current_stage not in self.completed:
|
||||
completed += 0.5
|
||||
return min(100, int((completed / len(self.stages)) * 100))
|
||||
|
||||
def progress_bar(self, percent: int, width: int = 20) -> str:
|
||||
filled = int(width * percent / 100)
|
||||
return "[" + "#" * filled + "-" * (width - filled) + f"] {percent}%"
|
||||
|
||||
def render_pipeline_progress(self) -> str:
|
||||
lines = [
|
||||
"=" * 60,
|
||||
"Meeting Lab Benchmark",
|
||||
"",
|
||||
self.progress_bar(self.pipeline_percent()),
|
||||
"",
|
||||
]
|
||||
for stage in self.stages:
|
||||
marker = "x" if stage in self.completed else ">" if stage == self.current_stage else " "
|
||||
lines.append(f"[{marker}] {self.stage_name(stage)}")
|
||||
lines.extend(
|
||||
[
|
||||
"",
|
||||
f"Elapsed: {self.format_hms(self.monotonic() - self.pipeline_start)}",
|
||||
"=" * 60,
|
||||
]
|
||||
)
|
||||
return "\n".join(lines)
|
||||
|
||||
def print_pipeline_progress(self) -> None:
|
||||
print(self.render_pipeline_progress(), file=self.stream, flush=True)
|
||||
|
||||
def start_stage(self, stage: str) -> None:
|
||||
self.current_stage = stage
|
||||
self.stage_starts[stage] = self.monotonic()
|
||||
self.stage_wall_starts[stage] = self.wall_clock()
|
||||
self.print_pipeline_progress()
|
||||
|
||||
def finish_stage(self, stage: str, status: str) -> float:
|
||||
now = self.monotonic()
|
||||
elapsed = now - self.stage_starts.get(stage, now)
|
||||
self.stage_durations[stage] = elapsed
|
||||
if status == "passed":
|
||||
self.completed.add(stage)
|
||||
if self.current_stage == stage:
|
||||
self.current_stage = None
|
||||
self.print_pipeline_progress()
|
||||
return elapsed
|
||||
|
||||
def format_extraction_progress(self, completed_chunks: int, total_chunks: int, elapsed: float) -> str:
|
||||
percent = int((completed_chunks / total_chunks) * 100) if total_chunks else 100
|
||||
average, remaining = self.extraction_eta(completed_chunks, total_chunks, elapsed)
|
||||
average_text = self.format_mmss(average) if average is not None else "--:--"
|
||||
remaining_text = "~" + self.format_hms(remaining) if remaining is not None else "unknown"
|
||||
return "\n".join(
|
||||
[
|
||||
"Extraction",
|
||||
f"Chunk {completed_chunks} / {total_chunks} ({percent}%)",
|
||||
"",
|
||||
f"Elapsed: {self.format_hms(elapsed)}",
|
||||
f"Average: {average_text} / chunk",
|
||||
f"Remaining: {remaining_text}",
|
||||
]
|
||||
)
|
||||
|
||||
def print_extraction_progress(self, completed_chunks: int, total_chunks: int, elapsed: float) -> None:
|
||||
print(
|
||||
self.format_extraction_progress(completed_chunks, total_chunks, elapsed),
|
||||
file=self.stream,
|
||||
flush=True,
|
||||
)
|
||||
|
||||
def format_stage_heartbeat(self, stage: str, elapsed: float) -> str:
|
||||
started = self.stage_wall_starts.get(stage)
|
||||
started_text = started.isoformat(timespec="seconds") if started else "unknown"
|
||||
return "\n".join(
|
||||
[
|
||||
self.stage_name(stage),
|
||||
"",
|
||||
"Running...",
|
||||
f"Started: {started_text}",
|
||||
f"Elapsed: {self.format_mmss(elapsed)}",
|
||||
]
|
||||
)
|
||||
|
||||
def print_stage_heartbeat(self, stage: str) -> None:
|
||||
started = self.stage_starts.get(stage, self.monotonic())
|
||||
print(
|
||||
self.format_stage_heartbeat(stage, self.monotonic() - started),
|
||||
file=self.stream,
|
||||
flush=True,
|
||||
)
|
||||
|
||||
def format_stage_timing_summary(self, timings: dict[str, Any]) -> str:
|
||||
rows = []
|
||||
for stage in self.stages:
|
||||
data = timings.get(stage)
|
||||
if not isinstance(data, dict):
|
||||
continue
|
||||
seconds = data.get("runtime_seconds")
|
||||
if isinstance(seconds, (int, float)):
|
||||
rows.append((self.stage_name(stage), float(seconds)))
|
||||
total = sum(seconds for _stage, seconds in rows)
|
||||
lines = ["Stage Timing Summary", "", f"{'Stage':<24} Duration", ""]
|
||||
lines.extend(f"{stage:<24} {self.format_duration_seconds(seconds):>10}" for stage, seconds in rows)
|
||||
lines.extend(["", f"{'Total':<24} {self.format_duration_seconds(total):>10}"])
|
||||
return "\n".join(lines)
|
||||
|
||||
|
||||
def parse_args(argv: list[str] | None = None) -> argparse.Namespace:
|
||||
parser = argparse.ArgumentParser(
|
||||
description="Run the Meeting Lab production benchmark pipeline."
|
||||
@@ -284,22 +464,66 @@ def total_ram_bytes() -> int | None:
|
||||
return int(psutil.virtual_memory().total)
|
||||
|
||||
|
||||
def run_stage(name: str, timings: dict[str, Any], func: Callable[[], Any]) -> Any:
|
||||
def run_stage(
|
||||
name: str,
|
||||
timings: dict[str, Any],
|
||||
func: Callable[[], Any],
|
||||
reporter: ProgressReporter,
|
||||
*,
|
||||
heartbeat: bool = False,
|
||||
heartbeat_interval: int = DEFAULT_PROGRESS_INTERVAL,
|
||||
) -> Any:
|
||||
reporter.start_stage(name)
|
||||
print(f"[{datetime.now().isoformat(timespec='seconds')}] Starting {name}", flush=True)
|
||||
start = time.perf_counter()
|
||||
try:
|
||||
result = func()
|
||||
if heartbeat:
|
||||
result = run_with_heartbeat(
|
||||
name,
|
||||
func,
|
||||
reporter=reporter,
|
||||
heartbeat_interval=heartbeat_interval,
|
||||
)
|
||||
else:
|
||||
result = func()
|
||||
except Exception:
|
||||
elapsed = time.perf_counter() - start
|
||||
elapsed = reporter.finish_stage(name, "failed")
|
||||
timings[name] = {"runtime_seconds": round(elapsed, 3), "status": "failed"}
|
||||
print(f"[{datetime.now().isoformat(timespec='seconds')}] FAILED {name} after {elapsed:.3f}s", flush=True)
|
||||
raise
|
||||
elapsed = time.perf_counter() - start
|
||||
elapsed = reporter.finish_stage(name, "passed")
|
||||
timings[name] = {"runtime_seconds": round(elapsed, 3), "status": "passed"}
|
||||
print(f"[{datetime.now().isoformat(timespec='seconds')}] Finished {name} in {elapsed:.3f}s", flush=True)
|
||||
return result
|
||||
|
||||
|
||||
def run_with_heartbeat(
|
||||
name: str,
|
||||
func: Callable[[], Any],
|
||||
*,
|
||||
reporter: ProgressReporter,
|
||||
heartbeat_interval: int,
|
||||
) -> Any:
|
||||
results: Queue[tuple[str, Any]] = Queue()
|
||||
|
||||
def worker() -> None:
|
||||
try:
|
||||
results.put(("result", func()))
|
||||
except BaseException as exc: # noqa: BLE001 - forwarded to caller.
|
||||
results.put(("error", exc))
|
||||
|
||||
thread = Thread(target=worker, daemon=True)
|
||||
thread.start()
|
||||
while True:
|
||||
try:
|
||||
status, payload = results.get(timeout=max(heartbeat_interval, 1))
|
||||
except Empty:
|
||||
reporter.print_stage_heartbeat(name)
|
||||
continue
|
||||
if status == "error":
|
||||
raise payload
|
||||
return payload
|
||||
|
||||
|
||||
def write_normalized_chunk(input_path: Path, output_path: Path, changes_path: Path) -> dict[str, Any]:
|
||||
source = input_path.read_text(encoding="utf-8-sig")
|
||||
blocks = split_blocks(source)
|
||||
@@ -338,10 +562,21 @@ def run_extraction_stage(
|
||||
num_predict: int | None,
|
||||
num_ctx: int,
|
||||
meeting_context: MeetingContext,
|
||||
progress_reporter: ProgressReporter | None = None,
|
||||
) -> list[Path]:
|
||||
output_dir.mkdir(parents=True, exist_ok=True)
|
||||
output_paths = []
|
||||
for chunk_path in find_normalized_chunks(normalized_dir):
|
||||
chunk_paths = find_normalized_chunks(normalized_dir)
|
||||
total_chunks = len(chunk_paths)
|
||||
start = progress_reporter.monotonic() if progress_reporter else time.perf_counter()
|
||||
|
||||
def elapsed() -> float:
|
||||
clock = progress_reporter.monotonic if progress_reporter else time.perf_counter
|
||||
return clock() - start
|
||||
|
||||
for chunk_index, chunk_path in enumerate(chunk_paths, start=1):
|
||||
if progress_reporter:
|
||||
progress_reporter.print_extraction_progress(chunk_index - 1, total_chunks, elapsed())
|
||||
output_path = extraction_output_path(chunk_path, output_dir)
|
||||
raw_path = output_path.with_suffix(".raw.txt")
|
||||
metadata_path = output_path.with_name(output_path.stem + "_metadata.json")
|
||||
@@ -366,6 +601,8 @@ def run_extraction_stage(
|
||||
extraction["context"] = meeting_context.provenance()
|
||||
output_path.write_text(json.dumps(extraction, ensure_ascii=False, indent=2) + "\n", encoding="utf-8")
|
||||
output_paths.append(output_path)
|
||||
if progress_reporter:
|
||||
progress_reporter.print_extraction_progress(chunk_index, total_chunks, elapsed())
|
||||
return output_paths
|
||||
|
||||
|
||||
@@ -539,6 +776,7 @@ def write_benchmark_report(path: Path, metadata: dict[str, Any]) -> None:
|
||||
semantic = metadata.get("semantic_consolidator", {})
|
||||
timings = metadata.get("runtime_seconds_by_stage", {})
|
||||
ollama = metadata.get("preflight", {}).get("ollama", {})
|
||||
reporter = ProgressReporter()
|
||||
lines = [
|
||||
"# Meeting Lab Benchmark Report",
|
||||
"",
|
||||
@@ -568,6 +806,12 @@ def write_benchmark_report(path: Path, metadata: dict[str, Any]) -> None:
|
||||
lines.append(f"| {stage} | {data.get('runtime_seconds')} | {data.get('status')} |")
|
||||
lines.extend(
|
||||
[
|
||||
"",
|
||||
"## Stage Timing Summary",
|
||||
"",
|
||||
"```text",
|
||||
reporter.format_stage_timing_summary(timings),
|
||||
"```",
|
||||
"",
|
||||
f"- Total wall-clock runtime: {metadata.get('total_runtime_seconds')} seconds",
|
||||
f"- Average extraction time per chunk: {metadata.get('average_extraction_seconds_per_chunk')} seconds",
|
||||
@@ -628,6 +872,7 @@ def run_pipeline(args: argparse.Namespace) -> tuple[int, Path, Path, Path | None
|
||||
protocol_path: Path | None = None
|
||||
timings: dict[str, Any] = {}
|
||||
start = time.perf_counter()
|
||||
reporter = ProgressReporter()
|
||||
metadata: dict[str, Any] = {
|
||||
"status": "running",
|
||||
"input": str(args.input),
|
||||
@@ -657,6 +902,7 @@ def run_pipeline(args: argparse.Namespace) -> tuple[int, Path, Path, Path | None
|
||||
"preflight",
|
||||
timings,
|
||||
lambda: preflight(args, output_dir),
|
||||
reporter,
|
||||
)
|
||||
metadata["preflight"] = preflight_data
|
||||
|
||||
@@ -675,6 +921,7 @@ def run_pipeline(args: argparse.Namespace) -> tuple[int, Path, Path, Path | None
|
||||
chunks_dir,
|
||||
args.input.name,
|
||||
),
|
||||
reporter,
|
||||
)
|
||||
metadata["chunk_distribution"] = chunk_distribution(chunks_dir)
|
||||
|
||||
@@ -691,7 +938,7 @@ def run_pipeline(args: argparse.Namespace) -> tuple[int, Path, Path, Path | None
|
||||
)
|
||||
shutil.copyfile(chunks_dir / "manifest.json", normalized_dir / "manifest.json")
|
||||
|
||||
run_stage("normalization", timings, normalize_all)
|
||||
run_stage("normalization", timings, normalize_all, reporter)
|
||||
|
||||
extractions_dir = output_dir / "extractions"
|
||||
extraction_paths = run_stage(
|
||||
@@ -707,7 +954,9 @@ def run_pipeline(args: argparse.Namespace) -> tuple[int, Path, Path, Path | None
|
||||
num_predict=args.num_predict,
|
||||
num_ctx=args.num_ctx,
|
||||
meeting_context=meeting_context,
|
||||
progress_reporter=reporter,
|
||||
),
|
||||
reporter,
|
||||
)
|
||||
metadata["extraction_files"] = len(extraction_paths)
|
||||
metadata["extraction_counts"] = category_counts_from_extractions(extractions_dir)
|
||||
@@ -721,6 +970,9 @@ def run_pipeline(args: argparse.Namespace) -> tuple[int, Path, Path, Path | None
|
||||
merge_duplicates=True,
|
||||
meeting_context_path=args.context,
|
||||
),
|
||||
reporter,
|
||||
heartbeat=True,
|
||||
heartbeat_interval=args.progress_interval,
|
||||
)
|
||||
write_canonicalized(canonicalized, canonicalized_path)
|
||||
metadata["canonicalizer_stats"] = canonicalized.get("stats")
|
||||
@@ -736,8 +988,11 @@ def run_pipeline(args: argparse.Namespace) -> tuple[int, Path, Path, Path | None
|
||||
timeout=args.timeout,
|
||||
num_ctx=args.num_ctx,
|
||||
num_predict=args.consolidator_num_predict,
|
||||
progress_interval=args.progress_interval,
|
||||
progress_interval=max(args.timeout, args.progress_interval) + 1,
|
||||
),
|
||||
reporter,
|
||||
heartbeat=True,
|
||||
heartbeat_interval=args.progress_interval,
|
||||
)
|
||||
metadata["semantic_consolidator"] = semantic_metadata
|
||||
|
||||
@@ -754,6 +1009,9 @@ def run_pipeline(args: argparse.Namespace) -> tuple[int, Path, Path, Path | None
|
||||
num_predict=args.renderer_num_predict,
|
||||
think=False,
|
||||
),
|
||||
reporter,
|
||||
heartbeat=True,
|
||||
heartbeat_interval=args.progress_interval,
|
||||
)
|
||||
metadata["renderer"] = {
|
||||
"valid": renderer_metadata.get("valid"),
|
||||
|
||||
Reference in New Issue
Block a user