Files
2026-linux-sumka/reporting.py
2026-09-24 18:30:05 +03:00

129 lines
4.0 KiB
Python

"""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