Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions changelog.d/us-staging-process-telemetry.added.md
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
US fiscal-refresh staging telemetry now records how a build is running, not only which stage it is in. Every stage event and the progress document carry a `resources` snapshot (CPU user and system seconds including finished child processes, current and peak RSS). A heartbeat refreshes `heartbeat_at` and `resources` every 60 s, so a run killed by the operating system shows as stale within minutes instead of staying `running`. Target compilation (one base pass plus one income-tax pass per requested JCT family, over the same household batches) and each post-export scoring sweep report `work` in the progress document: units done, total, unit and elapsed seconds, so a single snapshot gives a rate. A failed run's event records `failure_class` (`gate_refused`, `terminated`, `interrupted`, `out_of_memory`, `refused` or `error`), `failed_during` and `elapsed_seconds`; a completed run records `elapsed_seconds`. The run manifest records an `identity`: git commit, runtime package versions, platform, CPU count, memory, the thread-pool default, and the command line as option names plus a digest, never its values. All of this is opt-in on `StagingTelemetry` (`record_resources`, `heartbeat_seconds`, `record_outcome`, `record_identity()`, `work_progress()`); other callers and the version 1 contract fixtures are unchanged.
226 changes: 202 additions & 24 deletions packages/microcosm-build/src/microcosm/build/staging.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,10 @@

import json
import math
import os
import platform
import re
import resource
import shutil
import sys
import threading
Expand All @@ -21,6 +25,7 @@
from pathlib import Path
from typing import Any

from microcosm.build.stage_profile import current_rss_bytes
from microcosm.build.staging_storage import (
BestEffortUploadSession,
HuggingFaceDatasetStorage,
Expand Down Expand Up @@ -152,6 +157,73 @@ def restage_run(
return written


def resource_snapshot() -> dict[str, Any]:
"""CPU time and memory of this process and its finished children.

CPU seconds include waited-for child processes, so a stage that shells out
still counts. ``rss_bytes`` is the current resident set; it is omitted when
the platform will not report it.
"""

own = resource.getrusage(resource.RUSAGE_SELF)
children = resource.getrusage(resource.RUSAGE_CHILDREN)
peak = int(own.ru_maxrss)
snapshot: dict[str, Any] = {
"cpu_user_seconds": round(own.ru_utime + children.ru_utime, 3),
"cpu_system_seconds": round(own.ru_stime + children.ru_stime, 3),
# ru_maxrss is bytes on macOS and KiB on Linux.
"peak_rss_bytes": peak if sys.platform == "darwin" else peak * 1024,
}
try:
snapshot["rss_bytes"] = current_rss_bytes()
except Exception:
pass
return snapshot


def host_identity() -> dict[str, Any]:
"""The machine a run is on, without its name or any path."""

try:
memory_bytes = os.sysconf("SC_PAGE_SIZE") * os.sysconf("SC_PHYS_PAGES")
except (OSError, ValueError, AttributeError):
memory_bytes = None
return {
"platform": platform.platform(),
"machine": platform.machine(),
"python": platform.python_version(),
"cpu_count": os.cpu_count(),
"memory_bytes": memory_bytes,
}


_GATE_REFUSAL = re.compile(r"\bgates? (?:failed|refused)\b|\brefus", re.IGNORECASE)


def classify_failure(error: BaseException) -> str:
"""A coarse class for a failed run, for counting failures by kind.

``gate_refused``: a release gate or register refused the candidate.
``terminated``: SIGTERM (a supervisor or budget stop). ``interrupted``:
Ctrl-C. ``out_of_memory``: Python ran out of memory (an operating-system
kill leaves no record at all; the heartbeat going stale shows it).
``refused``: the run stopped itself with a message (a pin or input
refusal). ``error``: anything else, usually a code or data defect.
"""

if type(error).__name__ == "BuildTerminatedError":
return "terminated"
if isinstance(error, KeyboardInterrupt):
return "interrupted"
if isinstance(error, MemoryError):
return "out_of_memory"
if _GATE_REFUSAL.search(str(error)):
return "gate_refused"
if isinstance(error, SystemExit):
return "refused"
return "error"


def _write_json(path: Path, payload: dict[str, Any]) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text(json.dumps(_jsonable(payload), indent=1, allow_nan=False))
Expand Down Expand Up @@ -184,6 +256,16 @@ class StagingTelemetry:
its telemetry local, records why in the run manifest, and builds
as usual; ``restage_run`` uploads the folder later. A check that
cannot reach the Hub changes nothing: uploads stay best-effort.
record_resources: Add a ``resources`` snapshot (CPU seconds, current
and peak RSS) to every stage event and to the progress document.
heartbeat_seconds: When set, a worker refreshes ``heartbeat_at`` and
``resources`` in the progress document this often, so a run that
stops without a final event (an operating-system kill) is visible
within minutes. Heartbeat uploads go through the background
uploader only, so this needs ``background_uploads``.
record_outcome: Add ``failure_class``, ``failed_during`` and
``elapsed_seconds`` to the ``failed`` event, and
``elapsed_seconds`` to the ``complete`` event.
"""

run_id: str
Expand All @@ -197,6 +279,9 @@ class StagingTelemetry:
background_uploads: bool = False
final_upload_timeout_seconds: float = 120.0
check_write_access: bool = False
record_resources: bool = False
heartbeat_seconds: float | None = None
record_outcome: bool = False
started_at: str = field(default_factory=_now)

def __post_init__(self) -> None:
Expand Down Expand Up @@ -231,6 +316,11 @@ def __post_init__(self) -> None:
self._pending_artifacts: list[tuple[Path, str]] = []
self._upload_requested = threading.Event()
self._uploads_closing = False
self._started_monotonic = time.monotonic()
self._identity: dict[str, Any] | None = None
self._work: dict[str, Any] | None = None
self._heartbeat_stop = threading.Event()
self._heartbeat_thread: threading.Thread | None = None
self._upload_thread: threading.Thread | None = None
if self._upload_session is not None and self.background_uploads:
self._upload_thread = threading.Thread(
Expand All @@ -250,6 +340,13 @@ def __post_init__(self) -> None:
}
self._write_run_manifest()
self.stage("created", message="Staging run created.")
if self.heartbeat_seconds is not None and self.heartbeat_seconds > 0:
self._heartbeat_thread = threading.Thread(
target=self._heartbeat_worker,
name=f"staging-heartbeat-{self.run_id}",
daemon=True,
)
self._heartbeat_thread.start()

def _check_write_access(self) -> dict[str, Any] | None:
"""Keep the run local, loudly, when no credential can write the repo.
Expand Down Expand Up @@ -505,6 +602,7 @@ def _write_run_manifest(self) -> None:
if self._delivery_check
else {}
),
**({"identity": self._identity} if self._identity else {}),
},
)

Expand Down Expand Up @@ -541,26 +639,34 @@ def stage(
**details: Any,
) -> None:
updated_at = _now()
self._progress.update(
{
"status": status,
"stage": stage,
"message": message,
"updated_at": updated_at,
"details": _jsonable(details),
}
)
self._write_progress()
self._append_event(
{
resources = resource_snapshot() if self.record_resources else None
with self._io_lock:
if self._work is not None and self._work.get("stage") != stage:
self._work = None
self._progress.pop("work", None)
self._progress.update(
{
"status": status,
"stage": stage,
"message": message,
"updated_at": updated_at,
"details": _jsonable(details),
}
)
if resources is not None:
self._progress["resources"] = resources
self._write_progress()
event: dict[str, Any] = {
"time": updated_at,
"type": "stage",
"status": status,
"stage": stage,
"message": message,
"details": details,
}
)
if resources is not None:
event["resources"] = resources
self._append_event(event)
self._maybe_upload(force=force_upload)

def calibration_progress(self, event: dict[str, object]) -> None:
Expand All @@ -576,17 +682,20 @@ def calibration_progress(self, event: dict[str, object]) -> None:
"l0_lambda": event.get("l0_lambda"),
"time": _now(),
}
self._calibration_events.append(_jsonable(row))
self._progress.update(
{
"status": "running",
"stage": "calibrating",
"updated_at": row["time"],
"calibration": row,
}
)
self._write_progress()
self._write_calibration_progress()
# The heartbeat thread also writes the progress document, so every
# change to it happens under the write lock.
with self._io_lock:
self._calibration_events.append(_jsonable(row))
self._progress.update(
{
"status": "running",
"stage": "calibrating",
"updated_at": row["time"],
"calibration": row,
}
)
self._write_progress()
self._write_calibration_progress()
self._maybe_upload()

def attach_artifact(
Expand Down Expand Up @@ -615,22 +724,91 @@ def attach_artifact(
self._pending_artifacts.append((local, path_in_repo))
self._maybe_upload(force=force_upload)

def record_identity(self, **identity: Any) -> None:
"""Record what produced this run (commit, host, versions) in its manifest."""

self._identity = _jsonable(identity)
self._write_run_manifest()

def work_progress(
self, done: int, total: int, *, unit: str, **details: Any
) -> None:
"""Report how far the current stage is through its work units.

Written to the progress document only (no event), so a stage with
thousands of units adds no events. ``elapsed_seconds`` counts from the
stage's first report, so one snapshot gives a rate.
"""

now = time.monotonic()
with self._io_lock:
stage = self._progress.get("stage")
if self._work is None or self._work.get("stage") != stage:
self._work = {"stage": stage, "started": now}
self._progress["work"] = {
"stage": stage,
"done": int(done),
"total": int(total),
"unit": unit,
"elapsed_seconds": round(now - self._work["started"], 3),
"updated_at": _now(),
"details": _jsonable(details),
}
self._write_progress()
self._maybe_upload()

def _heartbeat_worker(self) -> None:
interval = float(self.heartbeat_seconds or 0)
while not self._heartbeat_stop.wait(interval):
try:
with self._io_lock:
self._progress["heartbeat_at"] = _now()
self._progress["resources"] = resource_snapshot()
self._write_progress()
if self._upload_thread is not None:
self._maybe_upload()
except Exception as error: # pragma: no cover - defensive
print(f"warning: staging heartbeat failed: {error}", file=sys.stderr)

def _stop_heartbeat(self) -> None:
self._heartbeat_stop.set()
thread = self._heartbeat_thread
if thread is not None and thread is not threading.current_thread():
thread.join(timeout=5)

def _elapsed_seconds(self) -> float:
return round(time.monotonic() - self._started_monotonic, 3)

def fail(self, error: BaseException) -> None:
self._stop_heartbeat()
outcome: dict[str, Any] = {}
if self.record_outcome:
outcome = {
"failure_class": classify_failure(error),
"failed_during": self._progress.get("stage"),
"elapsed_seconds": self._elapsed_seconds(),
}
self.stage(
"failed",
status="failed",
message=str(error),
force_upload=True,
error_type=type(error).__name__,
**outcome,
traceback=traceback.format_exc(),
)
self._finish_uploads()

def complete(self) -> None:
self._stop_heartbeat()
outcome = (
{"elapsed_seconds": self._elapsed_seconds()} if self.record_outcome else {}
)
self.stage(
"complete",
status="passed",
message="Staging run completed.",
force_upload=True,
**outcome,
)
self._finish_uploads()
Loading