129 lines
4.0 KiB
Python
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
|