From c977364d8c0fc0aea0324a0439cc3ffa27da45ba Mon Sep 17 00:00:00 2001 From: PavelMakarchuk Date: Thu, 1 Oct 2026 13:36:34 -0400 Subject: [PATCH] Record process telemetry for US release builds Every stage event carries CPU and RSS; a 60 s heartbeat keeps a silent stage visibly alive and makes an operating-system kill visible within minutes; target compilation and post-export scoring report batches done of the planned total; failures are classed and timed; the run manifest records the commit, runtime, host and a digest of the command line. All opt-in on StagingTelemetry and turned on by the fiscal-refresh release. Co-Authored-By: Claude Opus 5.5 --- .../us-staging-process-telemetry.added.md | 1 + .../src/microcosm/build/staging.py | 226 ++++++++++++++++-- .../tests/engine_free/shared/test_staging.py | 132 ++++++++++ .../us/test_us_fiscal_refresh_builder.py | 67 ++++++ tools/build_us_fiscal_refresh_release.py | 96 +++++++- 5 files changed, 497 insertions(+), 25 deletions(-) create mode 100644 changelog.d/us-staging-process-telemetry.added.md diff --git a/changelog.d/us-staging-process-telemetry.added.md b/changelog.d/us-staging-process-telemetry.added.md new file mode 100644 index 000000000..d889072f3 --- /dev/null +++ b/changelog.d/us-staging-process-telemetry.added.md @@ -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. diff --git a/packages/microcosm-build/src/microcosm/build/staging.py b/packages/microcosm-build/src/microcosm/build/staging.py index 96125ea28..2b531a597 100644 --- a/packages/microcosm-build/src/microcosm/build/staging.py +++ b/packages/microcosm-build/src/microcosm/build/staging.py @@ -11,6 +11,10 @@ import json import math +import os +import platform +import re +import resource import shutil import sys import threading @@ -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, @@ -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)) @@ -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 @@ -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: @@ -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( @@ -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. @@ -505,6 +602,7 @@ def _write_run_manifest(self) -> None: if self._delivery_check else {} ), + **({"identity": self._identity} if self._identity else {}), }, ) @@ -541,18 +639,24 @@ 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, @@ -560,7 +664,9 @@ def 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: @@ -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( @@ -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() diff --git a/packages/microcosm-build/tests/engine_free/shared/test_staging.py b/packages/microcosm-build/tests/engine_free/shared/test_staging.py index c0c630f1e..98a0215c9 100644 --- a/packages/microcosm-build/tests/engine_free/shared/test_staging.py +++ b/packages/microcosm-build/tests/engine_free/shared/test_staging.py @@ -1,5 +1,6 @@ import json import threading +import time import microcosm.build.staging as staging_module from microcosm.build.staging import StagingTelemetry @@ -429,3 +430,134 @@ def test_restage_uploads_a_local_run_without_moving_the_pointer(tmp_path): update_index=False, ) assert "runs.json" not in written + + +def test_resources_are_recorded_on_events_and_progress(tmp_path) -> None: + telemetry = StagingTelemetry( + run_id="run-r", + candidate_release_id="populace-us-2024-r", + run_dir=tmp_path / "run-r", + record_resources=True, + ) + telemetry.stage("target_compilation") + + events = [ + json.loads(line) + for line in (tmp_path / "run-r" / "events.ndjson").read_text().splitlines() + ] + resources = events[-1]["resources"] + assert resources["cpu_user_seconds"] >= 0 + assert resources["peak_rss_bytes"] > 0 + progress = json.loads((tmp_path / "run-r" / "progress.json").read_text()) + assert set(progress["resources"]) >= {"cpu_user_seconds", "peak_rss_bytes"} + + +def test_work_progress_reports_a_rate_and_clears_on_the_next_stage(tmp_path) -> None: + telemetry = StagingTelemetry( + run_id="run-w", + candidate_release_id="populace-us-2024-w", + run_dir=tmp_path / "run-w", + ) + telemetry.stage("target_compilation") + telemetry.work_progress(10, 2124, unit="engine batch", pass_name="base") + + progress = json.loads((tmp_path / "run-w" / "progress.json").read_text()) + assert progress["work"]["stage"] == "target_compilation" + assert (progress["work"]["done"], progress["work"]["total"]) == (10, 2124) + assert progress["work"]["elapsed_seconds"] >= 0 + assert progress["work"]["details"] == {"pass_name": "base"} + # Work reports add no events. + events = (tmp_path / "run-w" / "events.ndjson").read_text().splitlines() + assert [json.loads(line)["stage"] for line in events] == [ + "created", + "target_compilation", + ] + + telemetry.stage("calibrating") + progress = json.loads((tmp_path / "run-w" / "progress.json").read_text()) + assert "work" not in progress + + +def test_heartbeat_keeps_a_silent_stage_visibly_alive(tmp_path) -> None: + telemetry = StagingTelemetry( + run_id="run-h", + candidate_release_id="populace-us-2024-h", + run_dir=tmp_path / "run-h", + heartbeat_seconds=0.05, + ) + telemetry.stage("target_compilation") + progress_path = tmp_path / "run-h" / "progress.json" + deadline = time.monotonic() + 5 + while "heartbeat_at" not in json.loads(progress_path.read_text()): + assert time.monotonic() < deadline, "no heartbeat within 5 s" + time.sleep(0.02) + progress = json.loads(progress_path.read_text()) + assert progress["stage"] == "target_compilation" + assert "resources" in progress + + telemetry.complete() + assert telemetry._heartbeat_thread is not None + assert not telemetry._heartbeat_thread.is_alive() + + +def test_recorded_outcome_classes_the_failure(tmp_path) -> None: + telemetry = StagingTelemetry( + run_id="run-f", + candidate_release_id="populace-us-2024-f", + run_dir=tmp_path / "run-f", + record_outcome=True, + ) + telemetry.stage("release_gates") + telemetry.fail(RuntimeError("Release gates failed: QRF tail concentration")) + + events = [ + json.loads(line) + for line in (tmp_path / "run-f" / "events.ndjson").read_text().splitlines() + ] + details = events[-1]["details"] + assert details["failure_class"] == "gate_refused" + assert details["failed_during"] == "release_gates" + assert details["elapsed_seconds"] >= 0 + + +def test_outcome_fields_are_off_by_default(tmp_path) -> None: + telemetry = StagingTelemetry( + run_id="run-d", + candidate_release_id="populace-us-2024-d", + run_dir=tmp_path / "run-d", + ) + telemetry.complete() + + last = (tmp_path / "run-d" / "events.ndjson").read_text().splitlines()[-1] + assert json.loads(last)["details"] == {} + assert "resources" not in json.loads(last) + + +def test_failures_are_classed() -> None: + class BuildTerminatedError(SystemExit): + pass + + assert staging_module.classify_failure(BuildTerminatedError(143)) == "terminated" + assert staging_module.classify_failure(KeyboardInterrupt()) == "interrupted" + assert staging_module.classify_failure(MemoryError()) == "out_of_memory" + assert ( + staging_module.classify_failure(RuntimeError("Release gates failed: x")) + == "gate_refused" + ) + assert ( + staging_module.classify_failure(SystemExit("The feed pin does not match.")) + == "refused" + ) + assert staging_module.classify_failure(ValueError("bad column")) == "error" + + +def test_identity_is_recorded_in_the_run_manifest(tmp_path) -> None: + telemetry = StagingTelemetry( + run_id="run-i", + candidate_release_id="populace-us-2024-i", + run_dir=tmp_path / "run-i", + ) + telemetry.record_identity(git_commit="abc123", host={"cpu_count": 18}) + + manifest = json.loads((tmp_path / "run-i" / "run_manifest.json").read_text()) + assert manifest["identity"] == {"git_commit": "abc123", "host": {"cpu_count": 18}} diff --git a/packages/microcosm-build/tests/engine_free/us/test_us_fiscal_refresh_builder.py b/packages/microcosm-build/tests/engine_free/us/test_us_fiscal_refresh_builder.py index 58d174ac8..1661928a8 100644 --- a/packages/microcosm-build/tests/engine_free/us/test_us_fiscal_refresh_builder.py +++ b/packages/microcosm-build/tests/engine_free/us/test_us_fiscal_refresh_builder.py @@ -2047,6 +2047,9 @@ class RecordingTelemetry: def __init__(self, **kwargs): constructed.update(kwargs) + def record_identity(self, **identity): + constructed["identity"] = identity + monkeypatch.setattr(builder, "StagingTelemetry", RecordingTelemetry) telemetry = builder._staging_telemetry( args, release_root=tmp_path, release_id="populace-us-2024-k20000-fixture" @@ -2061,6 +2064,70 @@ def __init__(self, **kwargs): assert constructed["check_write_access"] is True +def test_release_staging_records_process_telemetry(tmp_path, monkeypatch) -> None: + builder = _load_builder_module() + constructed: dict[str, object] = {} + + class RecordingTelemetry: + def __init__(self, **kwargs): + constructed.update(kwargs) + + def record_identity(self, **identity): + constructed["identity"] = identity + + monkeypatch.setattr(builder, "StagingTelemetry", RecordingTelemetry) + monkeypatch.setattr( + builder.sys, + "argv", + ["build_us_fiscal_refresh_release.py", "--base-h5", "/Users/someone/base.h5"], + ) + args = SimpleNamespace( + no_staging=False, + staging_dir=tmp_path / "stage", + staging_repo_id=None, + staging_run_id=None, + staging_prefix=builder.DEFAULT_STAGING_PREFIX, + staging_upload_interval_seconds=60.0, + ) + builder._staging_telemetry(args, release_root=tmp_path, release_id="rel-1") + + assert constructed["record_resources"] is True + assert constructed["record_outcome"] is True + assert constructed["heartbeat_seconds"] == builder.STAGING_HEARTBEAT_SECONDS + identity = constructed["identity"] + # Option names and a digest, never argument values (local paths). + assert identity["options"] == ["--base-h5"] + assert "/Users/someone" not in json.dumps(identity) + assert len(identity["argv_sha256"]) == 64 + assert identity["host"]["cpu_count"] >= 1 + + +def test_work_counter_reports_progress_to_the_active_run(monkeypatch) -> None: + builder = _load_builder_module() + reports: list[tuple[int, int, str, dict]] = [] + + class Telemetry: + def work_progress(self, done, total, *, unit, **details): + reports.append((done, total, unit, details)) + + monkeypatch.setattr(builder, "_ACTIVE_TELEMETRY", Telemetry()) + counter = builder._WorkCounter(unit="engine batch", total=4) + counter.advance(pass_name="base") + counter.advance(3, pass_name="charitable", cached=True) + counter.advance(pass_name="overflow") + + assert [(done, total) for done, total, _, _ in reports] == [(1, 4), (4, 4), (4, 4)] + assert reports[1][3] == {"pass_name": "charitable", "cached": True} + + +def test_work_counter_is_silent_without_telemetry(monkeypatch) -> None: + builder = _load_builder_module() + monkeypatch.setattr(builder, "_ACTIVE_TELEMETRY", None) + counter = builder._WorkCounter(unit="engine batch", total=2) + counter.advance() + assert counter.done == 1 + + def test_builder_pool_release_identity_is_manifest_authenticated() -> None: builder = _load_builder_module() diff --git a/tools/build_us_fiscal_refresh_release.py b/tools/build_us_fiscal_refresh_release.py index 6c3b95348..dd56fe588 100644 --- a/tools/build_us_fiscal_refresh_release.py +++ b/tools/build_us_fiscal_refresh_release.py @@ -68,7 +68,11 @@ ) from microcosm.build.ledger_artifact import load_ledger_consumer_artifact from microcosm.build.source_runtime import SourceRuntimeConfig, run_source_stage -from microcosm.build.staging import DEFAULT_STAGING_PREFIX, StagingTelemetry +from microcosm.build.staging import ( + DEFAULT_STAGING_PREFIX, + StagingTelemetry, + host_identity, +) from microcosm.build.us_runtime import ( ASEC_2023_WEEKS_UNEMPLOYED_SOURCE_SHA256, CONGRESSIONAL_DISTRICT_VINTAGE_CROSSWALK_SHA256_ATTR, @@ -5171,6 +5175,8 @@ def _score( parts: dict[PostExportKey, list[tuple[np.ndarray, np.ndarray]]] = { key: [] for key in keys } + # Each sweep is reported on its own: the count restarts per sweep. + sweep = _WorkCounter(unit="scoring batch", total=len(batches)) for batch_frame in batches: with _automatic_gc_suspended(): simulation = self._construct(batch_frame, reform_system=reform_system) @@ -5197,6 +5203,7 @@ def _score( release_engine_simulation(simulation) del simulation, engine _collect_batch_garbage() + sweep.advance(sweep_label=label, keys=len(keys)) _collect_family_garbage() return { key: _concatenate_post_export_parts(key_parts) @@ -5426,6 +5433,7 @@ def _reform_household_income_tax( reform_income_tax[household_positions] = batch_income_tax del batch_income_tax, reformed, reformed_dataset, batch_frame _collect_batch_garbage() + _advance_work(pass_name=reform_spec.measure, batch=batch, batches=len(batches)) del reform_system _collect_family_garbage() return reform_income_tax @@ -6786,6 +6794,7 @@ def _materialize_base_simulation_columns( columns[column] = pool_values del batch_columns, batch_frame _collect_batch_garbage() + _advance_work(pass_name="base", batch=batch, batches=len(batches)) _collect_family_garbage() assert column_order is not None @@ -6833,6 +6842,17 @@ def _materialize_target_frame( _assert_no_formula_owned_columns(base_frame) system = CountryTaxBenefitSystem() n_households = base_frame.n("household") + # One base pass plus one income-tax pass per requested JCT family, each + # over the same household batches. + global _ACTIVE_WORK + requested_measures = {spec.measure for spec in target_specs} + passes = 1 + sum( + spec.measure in requested_measures for spec in US_JCT_TAX_EXPENDITURE_REFORMS + ) + batch_count = len( + tuple(_household_position_batches(n_households, maximum_microsim_batch_size)) + ) + _ACTIVE_WORK = _WorkCounter(unit="engine batch", total=passes * batch_count) # The base simulation uses the JCT reform loop's household partition. base_columns, base_simulation_batching = _materialize_base_simulation_columns( base_frame, @@ -6891,6 +6911,7 @@ def _materialize_target_frame( if cached is not None: reform_income_tax, cache_digest, cache_path = cached cache_stats["hits"] = int(cache_stats["hits"]) + 1 + _advance_work(batch_count, pass_name=reform_spec.measure, cached=True) cache_entry = { "measure": reform_spec.measure, "neutralized_variable": reform_spec.neutralized_variable, @@ -6979,6 +7000,7 @@ def _materialize_target_frame( "jct_reform_families_simulated": jct_reform_families_simulated, }, } + _ACTIVE_WORK = None return ( target_frame, registry, @@ -11414,6 +11436,40 @@ def _assert_exact_k_original_pool_alignment( _ACTIVE_TELEMETRY: StagingTelemetry | None = None +class _WorkCounter: + """Engine batches done out of a stage's planned total. + + Reported to the staging run's progress document, so a stage that runs for + hours (target compilation, post-export scoring) shows how far it is and + at what rate. Telemetry only: nothing in the build reads it. + """ + + def __init__(self, *, unit: str, total: int) -> None: + self.unit = unit + self.total = max(int(total), 0) + self.done = 0 + + def advance(self, units: int = 1, **details: object) -> None: + self.done = ( + min(self.done + units, self.total) if self.total else self.done + units + ) + telemetry = _ACTIVE_TELEMETRY + if telemetry is None: + return + try: + telemetry.work_progress(self.done, self.total, unit=self.unit, **details) + except Exception as error: # pragma: no cover - telemetry is best-effort + print(f"warning: could not report work progress: {error}", file=sys.stderr) + + +_ACTIVE_WORK: _WorkCounter | None = None + + +def _advance_work(units: int = 1, **details: object) -> None: + if _ACTIVE_WORK is not None: + _ACTIVE_WORK.advance(units, **details) + + class _ReleaseDryRun: """One ``--dry-run-gates-report`` run. @@ -12064,10 +12120,48 @@ def _staging_telemetry( background_uploads=True, # No write token: warn, keep telemetry local, build as usual. check_write_access=True, + # Process telemetry: CPU and memory at every stage, a heartbeat so a + # killed run is visible within minutes, and classed failures. + record_resources=True, + heartbeat_seconds=STAGING_HEARTBEAT_SECONDS, + record_outcome=True, ) + _ACTIVE_TELEMETRY.record_identity(**_staging_run_identity()) return _ACTIVE_TELEMETRY +STAGING_HEARTBEAT_SECONDS = 60.0 + + +def _staging_run_identity() -> dict[str, object]: + """What produced this run, for comparing runs like with like. + + The command line is recorded as its option names and a digest, never its + values: staging documents are served publicly and argument values are + local paths. + """ + + argv = list(sys.argv[1:]) + identity: dict[str, object] = { + "tool": Path(sys.argv[0]).name if sys.argv else None, + "options": sorted( + {arg.split("=", 1)[0] for arg in argv if arg.startswith("--")} + ), + "argv_sha256": hashlib.sha256("\0".join(argv).encode()).hexdigest(), + "host": host_identity(), + "thread_pool_default": _THREAD_POOL_DEFAULT, + } + try: + identity["git_commit"] = _git_output("rev-parse", "HEAD") + except (OSError, subprocess.SubprocessError): + identity["git_commit"] = None + try: + identity["runtime"] = _runtime_versions() + except Exception as error: # pragma: no cover - defensive + identity["runtime_error"] = type(error).__name__ + return identity + + class _TerminalBatchTelemetry: """Turn terminal-batch telemetry crashes into release-gate failures.