"""Append-only workflow status reporting for external process supervisors.""" from __future__ import annotations import json import threading import time import uuid from datetime import datetime, timezone from pathlib import Path from typing import Any, Literal, TextIO TerminalStatus = Literal["succeeded", "failed", "cancelled"] class StatusReporter: """Write JSON Lines status events and periodic heartbeats.""" def __init__( self, path: str | Path = "report.txt", *, heartbeat_interval: float = 30.0, ) -> None: if heartbeat_interval <= 0: raise ValueError("heartbeat_interval must be positive") self._path = Path(path) self._heartbeat_interval = heartbeat_interval self._run_id = uuid.uuid4().hex self._started_at = time.monotonic() self._step: str | None = None self._step_index: int | None = None self._steps_total: int | None = None self._lock = threading.Lock() self._stop_event = threading.Event() self._thread: threading.Thread | None = None self._file: TextIO | None = None self._finished = False @staticmethod def _timestamp() -> str: return datetime.now(timezone.utc).isoformat(timespec="seconds").replace( "+00:00", "Z", ) def _write(self, status: str, **fields: Any) -> None: record = { "timestamp": self._timestamp(), "run_id": self._run_id, "status": status, **fields, } line = json.dumps(record, ensure_ascii=False, separators=(",", ":")) with self._lock: if self._file is None: raise RuntimeError("StatusReporter is not running") self._file.write(line + "\n") self._file.flush() def start(self) -> None: """Open the report and start emitting heartbeats.""" if self._file is not None or self._finished: raise RuntimeError("StatusReporter cannot be started again") self._file = self._path.open("a", encoding="utf-8", buffering=1) self._write("started") self._thread = threading.Thread( target=self._heartbeat_loop, name="status-reporter", daemon=True, ) self._thread.start() def _heartbeat_loop(self) -> None: while not self._stop_event.wait(self._heartbeat_interval): with self._lock: step = self._step step_index = self._step_index steps_total = self._steps_total fields: dict[str, Any] = { "elapsed_seconds": round(time.monotonic() - self._started_at), } if step is not None: fields.update({ "step": step, "step_index": step_index, "steps_total": steps_total, }) self._write("heartbeat", **fields) def step_started(self, step: str, index: int, total: int) -> None: with self._lock: self._step = step self._step_index = index self._steps_total = total self._write( "running", step=step, step_index=index, steps_total=total, ) def step_completed(self, step: str, index: int, total: int) -> None: self._write( "step_completed", step=step, step_index=index, steps_total=total, ) def finish(self, status: TerminalStatus, **fields: Any) -> None: """Write the terminal event and close the report.""" if self._finished: return self._finished = True self._stop_event.set() if self._thread is not None: self._thread.join() self._write( status, elapsed_seconds=round(time.monotonic() - self._started_at), **fields, ) assert self._file is not None self._file.close() self._file = None