diff --git a/pyproject.toml b/pyproject.toml index 9b9b82a..7319a96 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,4 +4,5 @@ version = "0.1.0" requires-python = ">=3.11" dependencies = [ "PyYAML>=6.0", + "requests>=2.31", ] diff --git a/scripts/run_meeting.py b/scripts/run_meeting.py index d33d70a..4a7844b 100644 --- a/scripts/run_meeting.py +++ b/scripts/run_meeting.py @@ -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"),