feat(telemetry): add bounded native pipeline lifecycle

This commit is contained in:
DCCONSTRUCTIONS
2026-07-29 12:32:52 +03:00
parent d8ccc9345a
commit 636db089ae
13 changed files with 830 additions and 13 deletions
@@ -7,7 +7,13 @@ param(
[string]$PersistentOutputRoot = "D:\NDC_MISSIONCORE\runtime\derived\.perception-persistent-publish",
[ValidateRange(1024, 65535)] [int]$PersistentPort = 18020,
[ValidateRange(5, 3600)] [int]$MaximumDurationSeconds = 20,
[ValidateRange(1, 1000)] [int]$FreeGiBFloor = 360
[ValidateRange(1, 1000)] [int]$FreeGiBFloor = 360,
[ValidatePattern("^[a-z0-9](?:[a-z0-9-]{0,62}[a-z0-9])?$")]
[string]$ContourId = "worker-006",
[ValidatePattern("^[a-z0-9](?:[a-z0-9-]{0,62}[a-z0-9])?$")]
[string]$AgentId = "worker-006",
[ValidatePattern("^[A-Za-z0-9][A-Za-z0-9._-]{0,127}$")]
[string]$NodeId = "DESKTOP-OPJ8J04"
)
$ErrorActionPreference = "Stop"
@@ -55,6 +61,11 @@ $request = @{
output_name = $outputName
token = $token
max_duration_seconds = $MaximumDurationSeconds
telemetry = @{
contour_id = $ContourId
agent_id = $AgentId
node_id = $NodeId
}
} | ConvertTo-Json -Compress
$token = $null
$client = (
@@ -22,7 +22,7 @@ import time
from collections import Counter, deque
from collections.abc import Iterator
from contextlib import contextmanager, suppress
from dataclasses import dataclass
from dataclasses import dataclass, replace
from datetime import UTC, datetime
from pathlib import Path
from typing import Any
@@ -106,6 +106,21 @@ from k1link.compute.live_perception import (
from k1link.data_plane import DecodedPointCloudView, DecodedPoseView
from k1link.device_plugins.xgrids_k1.protocol.normalizer import normalize_k1_message
try:
from k1link.compute.pipeline_telemetry import (
JsonlPipelineTelemetrySink,
PipelineTelemetryEmitter,
PipelineTelemetryIdentity,
)
except ModuleNotFoundError as exc:
if exc.name != "k1link.compute.pipeline_telemetry":
raise
from pipeline_telemetry import ( # type: ignore[no-redef]
JsonlPipelineTelemetrySink,
PipelineTelemetryEmitter,
PipelineTelemetryIdentity,
)
PROFILE_SCHEMA = "missioncore.e15-shadow-inference-profile/v1"
PROJECTION_SCHEMA = "missioncore.e15-live-projection-pack/v1"
WORKER_PACKAGE_SCHEMA = "missioncore.e15-worker-package/v1"
@@ -790,19 +805,71 @@ class _RuntimeTelemetry:
return summary
class _FailureIsolatingPipelineTelemetrySink:
"""Keep telemetry failures visible without changing compute control flow."""
def __init__(self, sink: Any) -> None:
self._sink = sink
self._lock = threading.Lock()
self._publish_failures = 0
self._last_error_type: str | None = None
def publish(self, topic: str, payload: bytes) -> None:
try:
self._sink.publish(topic, payload)
except Exception as exc:
with self._lock:
self._publish_failures += 1
self._last_error_type = type(exc).__name__
print(
json.dumps(
{
"event": "pipeline-telemetry-publish-failed",
"error_type": type(exc).__name__,
},
sort_keys=True,
),
flush=True,
)
def status(self) -> dict[str, Any]:
with self._lock:
return {
"transport": "jsonl-tail",
"ready": self._publish_failures == 0,
"publish_failures": self._publish_failures,
"last_error_type": self._last_error_type,
}
class _StageExecutionTelemetry:
"""Measure named pipeline spans without pretending they are OS processes."""
def __init__(self, stage_ids: tuple[str, ...] = PIPELINE_STAGE_IDS) -> None:
def __init__(
self,
stage_ids: tuple[str, ...] = PIPELINE_STAGE_IDS,
*,
identity: PipelineTelemetryIdentity | None = None,
sink: _FailureIsolatingPipelineTelemetrySink | None = None,
) -> None:
if not stage_ids or len(stage_ids) != len(set(stage_ids)):
raise RuntimeError("pipeline stage identities are invalid")
if (identity is None) != (sink is None):
raise RuntimeError("pipeline telemetry identity and sink must be paired")
self._stage_ids = stage_ids
self._identity = identity
self._sink = sink
self._lock = threading.Lock()
self._next_token = 0
self._active: dict[int, tuple[str, float, int | None]] = {}
self._elapsed_seconds = dict.fromkeys(stage_ids, 0.0)
self._activations = dict.fromkeys(stage_ids, 0)
self._last_frame_index: int | None = None
self._first_frame_by_stage: dict[str, int | None] = {}
self._last_frame_by_stage: dict[str, int | None] = {}
self._native_started: set[str] = set()
self._native_failed: set[str] = set()
self._native_finalized = False
@contextmanager
def measure(
@@ -813,21 +880,116 @@ class _StageExecutionTelemetry:
if stage_id not in self._elapsed_seconds:
raise RuntimeError(f"unknown pipeline stage: {stage_id}")
started = time.perf_counter()
emit_started = False
with self._lock:
self._next_token += 1
token = self._next_token
self._active[token] = (stage_id, started, frame_index)
self._activations[stage_id] += 1
self._last_frame_by_stage[stage_id] = frame_index
if stage_id not in self._first_frame_by_stage:
self._first_frame_by_stage[stage_id] = frame_index
if self.native_events_enabled and stage_id not in self._native_started:
self._native_started.add(stage_id)
emit_started = True
if frame_index is not None:
self._last_frame_index = frame_index
if emit_started:
self._emit_native(
stage_id=stage_id,
state="started",
frame_index=frame_index,
activation_count=1,
)
failure: BaseException | None = None
try:
yield
except BaseException as exc:
failure = exc
raise
finally:
finished = time.perf_counter()
failure_event: tuple[float, int, int | None] | None = None
with self._lock:
active = self._active.pop(token, None)
if active is not None:
self._elapsed_seconds[stage_id] += max(0.0, finished - active[1])
if (
failure is not None
and self.native_events_enabled
and stage_id not in self._native_failed
):
self._native_failed.add(stage_id)
failure_event = (
self._elapsed_seconds[stage_id],
self._activations[stage_id],
self._last_frame_by_stage.get(stage_id),
)
if failure_event is not None:
elapsed_seconds, activations, last_frame_index = failure_event
self._emit_native(
stage_id=stage_id,
state="failed",
frame_index=last_frame_index,
duration_ms=elapsed_seconds * 1000,
activation_count=activations,
error_type=type(failure).__name__,
)
@property
def native_events_enabled(self) -> bool:
return self._identity is not None and self._sink is not None
def finalize_native(self) -> None:
"""Emit one aggregate terminal event for each stage used by the run."""
rows: list[tuple[str, float, int, int | None]] = []
with self._lock:
if not self.native_events_enabled or self._native_finalized:
return
self._native_finalized = True
rows = [
(
stage_id,
self._elapsed_seconds[stage_id],
self._activations[stage_id],
self._last_frame_by_stage.get(stage_id),
)
for stage_id in self._stage_ids
if stage_id in self._native_started
and stage_id not in self._native_failed
]
for stage_id, elapsed_seconds, activations, frame_index in rows:
self._emit_native(
stage_id=stage_id,
state="completed",
frame_index=frame_index,
duration_ms=elapsed_seconds * 1000,
activation_count=activations,
)
def _emit_native(
self,
*,
stage_id: str,
state: str,
frame_index: int | None,
activation_count: int,
duration_ms: float | None = None,
error_type: str | None = None,
) -> None:
if self._identity is None or self._sink is None:
return
PipelineTelemetryEmitter(
identity=replace(self._identity, frame_index=frame_index),
sink=self._sink,
).stage_event(
stage_id,
state,
duration_ms=duration_ms,
activation_count=activation_count,
error_type=error_type,
)
def snapshot(self) -> dict[str, Any]:
now = time.perf_counter()
@@ -2172,6 +2334,44 @@ def _persistent_run_arguments(
return argparse.Namespace(**values)
def _persistent_run_telemetry_identity(
common: dict[str, Any],
request: dict[str, Any],
) -> PipelineTelemetryIdentity | None:
telemetry = request.get("telemetry")
if telemetry is None:
return None
if not isinstance(telemetry, dict) or set(telemetry) != {
"contour_id",
"agent_id",
"node_id",
}:
raise RuntimeError("persistent worker telemetry identity is invalid")
request_id = request.get("request_id")
if not isinstance(request_id, str):
raise RuntimeError("persistent worker telemetry request identity is invalid")
stability = common.get("stability")
lab_id = (
stability["profile_id"]
if isinstance(stability, dict) and isinstance(stability.get("profile_id"), str)
else "lab-e15-shadow-inference-v1"
)
method_id = (
INLINE_TEMPORAL_PIPELINE_ID if isinstance(stability, dict) else PIPELINE_ID
)
return PipelineTelemetryIdentity(
contour_id=telemetry.get("contour_id"),
agent_id=telemetry.get("agent_id"),
node_id=telemetry.get("node_id"),
lab_id=lab_id,
run_id=request_id,
request_id=request_id,
source_id=common["live"]["source"]["source_id"],
source_package_id=common["worker_package"]["package_id"],
method_id=method_id,
)
def serve(args: argparse.Namespace) -> int:
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
@@ -2183,6 +2383,9 @@ def serve(args: argparse.Namespace) -> int:
raise RuntimeError("LAB E15 PyAV version changed")
output_root = args.output_root.resolve(strict=True)
_assert_disk(output_root, args.free_bytes_floor)
telemetry_sink = _FailureIsolatingPipelineTelemetrySink(
JsonlPipelineTelemetrySink(output_root / "pipeline-telemetry.jsonl")
)
load_started = time.perf_counter()
loaded = _load_models(args, common)
model_load_seconds = time.perf_counter() - load_started
@@ -2193,6 +2396,7 @@ def serve(args: argparse.Namespace) -> int:
"failed_runs": 0,
"active_request_id": None,
"active_frame_index": None,
"last_run_outcome": None,
"_stage_telemetry": _StageExecutionTelemetry(),
}
@@ -2237,6 +2441,13 @@ def serve(args: argparse.Namespace) -> int:
"current_stage": stage_snapshot["current_stage"],
"active_stages": stage_snapshot["active_stages"],
"stage_metrics": stage_snapshot["stages"],
"last_run_outcome": state["last_run_outcome"],
"native_pipeline_telemetry": {
**telemetry_sink.status(),
"active_run_events_enabled": (
state["_stage_telemetry"].native_events_enabled
),
},
"authority": common["live"]["authority"],
"gpu": torch.cuda.get_device_name(),
},
@@ -2259,16 +2470,51 @@ def serve(args: argparse.Namespace) -> int:
return
state["busy"] = True
request_id = "invalid"
run_identity: PipelineTelemetryIdentity | None = None
run_emitter: PipelineTelemetryEmitter | None = None
run_started: float | None = None
terminal_event_emitted = False
try:
document = json.loads(self.rfile.read(length))
if not isinstance(document, dict):
raise RuntimeError("persistent worker request is not an object")
request_id = str(document.get("request_id", "invalid"))
run_args = _persistent_run_arguments(args, document)
run_identity = _persistent_run_telemetry_identity(common, document)
document["token"] = None
state["active_request_id"] = request_id
state["_stage_telemetry"] = _StageExecutionTelemetry()
state["_stage_telemetry"] = _StageExecutionTelemetry(
identity=run_identity,
sink=telemetry_sink if run_identity is not None else None,
)
run_emitter = (
PipelineTelemetryEmitter(
identity=run_identity,
sink=telemetry_sink,
)
if run_identity is not None
else None
)
run_started = time.perf_counter()
if run_emitter is not None:
run_emitter.run("started")
exit_code = run(run_args, loaded, state)
state["_stage_telemetry"].finalize_native()
duration_ms = max(0.0, (time.perf_counter() - run_started) * 1000)
if run_emitter is not None:
run_emitter.run(
"completed",
duration_ms=duration_ms,
exit_code=exit_code,
)
terminal_event_emitted = True
state["last_run_outcome"] = {
"request_id": request_id,
"state": "completed",
"duration_ms": round(duration_ms, 6),
"exit_code": exit_code,
"error_type": None,
}
state["completed_runs"] += 1
self._send(
200,
@@ -2280,6 +2526,31 @@ def serve(args: argparse.Namespace) -> int:
},
)
except Exception as exc:
state["_stage_telemetry"].finalize_native()
duration_ms = (
max(0.0, (time.perf_counter() - run_started) * 1000)
if run_started is not None
else None
)
if (
run_emitter is not None
and duration_ms is not None
and not terminal_event_emitted
):
run_emitter.run(
"failed",
duration_ms=duration_ms,
error_type=type(exc).__name__,
)
state["last_run_outcome"] = {
"request_id": request_id,
"state": "failed",
"duration_ms": (
round(duration_ms, 6) if duration_ms is not None else None
),
"exit_code": None,
"error_type": type(exc).__name__,
}
state["failed_runs"] += 1
print(
json.dumps(