From c570f7d938f07a9dea0ed11bae726f9af0dcc136 Mon Sep 17 00:00:00 2001 From: DCCONSTRUCTIONS Date: Wed, 2 Sep 2026 17:35:20 +0300 Subject: [PATCH] fix(telemetry): restore worker agent delivery and truthful system status --- .../src/core/system/useWorkerTelemetry.ts | 3 + .../src/core/system/workerTelemetry.ts | 14 ++++- .../system/ComputeModulesWorkspace.tsx | 33 ++++++++--- .../test/systemTelemetry.test.mjs | 5 ++ deploy/telemetry-plane/compose.source.yaml | 9 +++ deploy/telemetry-plane/compose.yaml | 2 + .../telemetry-plane/normalizer/normalizer.py | 55 +++++++++++++------ .../Get-NdcMissionCorePipelineTelemetry.ps1 | 47 +++++++++++++++- .../telegraf/mission-core-linux.conf.tmpl | 4 ++ .../telegraf/mission-core-windows.conf.tmpl | 9 +++ src/k1link/web/compute_contour_network.py | 9 ++- src/k1link/web/system_telemetry_api.py | 34 +++++++++++- tests/test_compute_contour_network.py | 15 +++++ tests/test_system_telemetry_api.py | 30 ++++++++++ tests/test_telemetry_normalizer_contract.py | 16 ++++++ 15 files changed, 254 insertions(+), 31 deletions(-) create mode 100644 deploy/telemetry-plane/compose.source.yaml diff --git a/apps/control-station/src/core/system/useWorkerTelemetry.ts b/apps/control-station/src/core/system/useWorkerTelemetry.ts index 6047289..6fcb6ed 100644 --- a/apps/control-station/src/core/system/useWorkerTelemetry.ts +++ b/apps/control-station/src/core/system/useWorkerTelemetry.ts @@ -38,6 +38,8 @@ export function useWorkerTelemetry( setError(null); return; } + setTelemetry(null); + setError(null); const controller = new AbortController(); const stopPolling = startSequentialPolling(async () => { setLoading(true); @@ -48,6 +50,7 @@ export function useWorkerTelemetry( setError(null); } catch (reason: unknown) { if (controller.signal.aborted) return; + setTelemetry(null); setError( reason instanceof Error ? reason.message diff --git a/apps/control-station/src/core/system/workerTelemetry.ts b/apps/control-station/src/core/system/workerTelemetry.ts index 8b8db9f..e6cdf75 100644 --- a/apps/control-station/src/core/system/workerTelemetry.ts +++ b/apps/control-station/src/core/system/workerTelemetry.ts @@ -114,6 +114,9 @@ export interface WorkerTelemetry { power_watts?: number; temperature_celsius?: number; memory_used_percent?: number; + sm_clock_mhz?: number | null; + memory_clock_mhz?: number | null; + driver_version?: string | null; } | null; triton: { ready: boolean; @@ -131,6 +134,13 @@ export interface WorkerTelemetry { completed_runs: number | null; failed_runs: number | null; model_load_seconds: number | null; + profile_name?: string | null; + input_state?: string | null; + wait_reason?: string | null; + live_children?: number | null; + input_pauses?: number | null; + pending_bundles?: number | null; + buffer_bytes?: number | null; stages: WorkerPipelineStage[]; }; network: { @@ -189,7 +199,9 @@ export async function fetchWorkerTelemetry( `/api/v1/system/contours/${encodeURIComponent(contourId)}/telemetry?history=90`, { method: "GET", - signal, + signal: signal + ? AbortSignal.any([signal, AbortSignal.timeout(5_000)]) + : AbortSignal.timeout(5_000), }, ), "missioncore.worker-telemetry/v1", diff --git a/apps/control-station/src/workspaces/system/ComputeModulesWorkspace.tsx b/apps/control-station/src/workspaces/system/ComputeModulesWorkspace.tsx index 33753f5..58c83ca 100644 --- a/apps/control-station/src/workspaces/system/ComputeModulesWorkspace.tsx +++ b/apps/control-station/src/workspaces/system/ComputeModulesWorkspace.tsx @@ -17,9 +17,14 @@ import { } from "../../core/system/useWorkerTelemetry"; function pipelineStateLabel(state: string): string { - if (state === "busy") return "Выполняет задачу"; + if (state === "busy" || state === "running") return "Выполняет задачу"; + if (state === "waiting") return "Ожидает восстановления данных"; + if (state === "synchronizing") return "Синхронизирует поток"; + if (state === "starting") return "Загружает профиль"; + if (state === "stopped" || state === "cancelled") return "Профиль остановлен"; + if (state === "failed") return "Ошибка профиля"; if (state === "ready") return "Готов к задаче"; - return "Нет live-состояния"; + return "Нет свежих данных о задаче"; } function formatResourcePair( @@ -116,7 +121,7 @@ export function ComputeModulesWorkspace() { {!legacyDiagnostic && !agentTelemetry ? ( - Агентный data-plane ещё не опубликовал нормализованный срез этого контура. + Ожидаем свежую телеметрию выбранного узла. ) : null} @@ -187,6 +192,8 @@ export function ComputeModulesWorkspace() { ? `${node.gpu.temperature_celsius} °C` : "—" } +
Частота GPU
{node?.gpu?.sm_clock_mhz == null ? "—" : `${node.gpu.sm_clock_mhz} МГц`}
+
Частота памяти GPU
{node?.gpu?.memory_clock_mhz == null ? "—" : `${node.gpu.memory_clock_mhz} МГц`}
Мощность
{ typeof node?.gpu?.power_watts === "number" ? `${node.gpu.power_watts.toFixed(1)} Вт` @@ -214,7 +221,7 @@ export function ComputeModulesWorkspace() {
PROCESSING RUNTIME

Контейнеры Mission Core

-

Два ограниченных runtime: inference server и прикладной perception pipeline.

+

Состояние контейнера и готовность вычислительного профиля проверяются отдельно.

{node?.triton.ready ? "Triton ready" : "Triton недоступен"} @@ -233,19 +240,31 @@ export function ComputeModulesWorkspace() { ТЕКУЩАЯ ЗАДАЧА

{pipelineStateLabel(telemetry?.pipeline.service_state ?? "unavailable")}

+ {telemetry?.pipeline.profile_name ?? "Perception"} + {" · "} {telemetry?.pipeline.active_request_id ? `Run ${telemetry.pipeline.active_request_id}` - : "Очередь свободна; модели остаются загруженными в persistent worker."} + : telemetry?.pipeline.service_state === "ready" + ? "Готов к приёму данных." + : "Готовность моделей не подтверждена."}

- {telemetry?.pipeline.service_state ?? "unavailable"} + {pipelineStateLabel(telemetry?.pipeline.service_state ?? "unavailable")} + {telemetry?.pipeline.profile_name ? ( +
+
Процессы профиля{telemetry.pipeline.live_children ?? "—"}
+
Ожиданий потока{telemetry.pipeline.input_pauses ?? "—"}
+
Пакетов в очереди{telemetry.pipeline.pending_bundles ?? "—"}
+
Входные буферы{formatBytes(telemetry.pipeline.buffer_bytes)}
+
+ ) : null}
Завершено запусков{telemetry?.pipeline.completed_runs ?? "—"}
Ошибок запусков{telemetry?.pipeline.failed_runs ?? "—"}
diff --git a/apps/control-station/test/systemTelemetry.test.mjs b/apps/control-station/test/systemTelemetry.test.mjs index 2a9d44c..e5bf8db 100644 --- a/apps/control-station/test/systemTelemetry.test.mjs +++ b/apps/control-station/test/systemTelemetry.test.mjs @@ -124,6 +124,11 @@ test("Compute-contour telemetry remains a bounded system feature slice", async ( assert.match(pollIntervalContract, /MIN_TELEMETRY_POLL_INTERVAL_SECONDS\s*=\s*1/); assert.match(telemetryPolling, /normalizeWorkerTelemetryPollMilliseconds/); assert.match(telemetryPolling, /startSequentialPolling/); + assert.match(telemetryPolling, /catch \(reason: unknown\)[\s\S]*?setTelemetry\(null\)/); + assert.match(computeWorkspace, /sm_clock_mhz/); + assert.match(computeWorkspace, /memory_clock_mhz/); + assert.match(computeWorkspace, /Ожидает восстановления данных/); + assert.doesNotMatch(computeWorkspace, /Очередь свободна; модели остаются загруженными/); assert.doesNotMatch(telemetryPolling, /if\s*\(\s*!enabled\s*\|\|\s*loading\s*\)/); assert.match(pollingScheduler, /A transient failure must not stop polling/); assert.match(styles, /system-telemetry\.css/); diff --git a/deploy/telemetry-plane/compose.source.yaml b/deploy/telemetry-plane/compose.source.yaml new file mode 100644 index 0000000..fa565e6 --- /dev/null +++ b/deploy/telemetry-plane/compose.source.yaml @@ -0,0 +1,9 @@ +# Local, read-only source override for repairing the existing normalizer without +# rebuilding an image on a memory-constrained operator machine. No extra service. +services: + normalizer: + volumes: + - type: bind + source: ${MISSIONCORE_NORMALIZER_SOURCE:?Set the absolute normalizer.py path} + target: /app/normalizer.py + read_only: true diff --git a/deploy/telemetry-plane/compose.yaml b/deploy/telemetry-plane/compose.yaml index bb74d8b..e190853 100644 --- a/deploy/telemetry-plane/compose.yaml +++ b/deploy/telemetry-plane/compose.yaml @@ -49,6 +49,8 @@ services: timescale: image: timescale/timescaledb-ha:pg16.14-ts2.28.2-all-oss + # The HA image auto-tunes for the Docker VM, not this 1 GiB container. + command: ["postgres", "-c", "shared_buffers=128MB", "-c", "work_mem=4MB", "-c", "maintenance_work_mem=64MB", "-c", "effective_cache_size=512MB", "-c", "max_parallel_workers_per_gather=0", "-c", "max_connections=32"] container_name: ndc-mission-core-telemetry-timescaledb restart: unless-stopped security_opt: diff --git a/deploy/telemetry-plane/normalizer/normalizer.py b/deploy/telemetry-plane/normalizer/normalizer.py index 357aa3b..98bea50 100644 --- a/deploy/telemetry-plane/normalizer/normalizer.py +++ b/deploy/telemetry-plane/normalizer/normalizer.py @@ -31,7 +31,8 @@ SAFE_IDENTIFIER: Final = re.compile(r"^[a-z0-9](?:[a-z0-9-]{0,62}[a-z0-9])?$") MAX_PAYLOAD_BYTES: Final = 1024 * 1024 MAX_TAGS: Final = 256 MAX_SERIES_KEY_BYTES: Final = 4096 -RETENTION_INTERVAL_SECONDS: Final = 86_400 +RETENTION_INTERVAL_SECONDS: Final = 60 +RETENTION_BATCH_ROWS: Final = 1000 ALLOWED_TELEMETRY_TAGS: Final = frozenset( { "agent_id", @@ -434,6 +435,39 @@ def _start_query_server( return server +def _retention_batch(connection: Any) -> int: + """Preserve the existing 30-day policy without an unbounded startup DELETE.""" + with connection.transaction(), connection.cursor() as cursor: + cursor.execute("SET LOCAL statement_timeout = '2s'") + cursor.execute("SET LOCAL lock_timeout = '250ms'") + cursor.execute( + """ + WITH expired AS ( + SELECT tableoid, ctid FROM contour_telemetry_samples + WHERE observed_at < CURRENT_TIMESTAMP - INTERVAL '30 days' + ORDER BY observed_at LIMIT %s + ) + DELETE FROM contour_telemetry_samples AS samples USING expired + WHERE samples.tableoid = expired.tableoid AND samples.ctid = expired.ctid + """, + (RETENTION_BATCH_ROWS,), + ) + return cursor.rowcount + + +def _retention_loop(dsn: str) -> None: + import psycopg # type: ignore[import-not-found] + + while True: + # Neither a retention timeout nor database recovery blocks MQTT I/O. + time.sleep(RETENTION_INTERVAL_SECONDS) + try: + with psycopg.connect(dsn, connect_timeout=3) as connection: + _retention_batch(connection) + except psycopg.Error: + print("telemetry retention deferred", flush=True) + + def main() -> None: import paho.mqtt.client as mqtt import psycopg # type: ignore[import-not-found] @@ -447,7 +481,9 @@ def main() -> None: password = _required("MISSIONCORE_MQTT_PASSWORD") runtime_health = TelemetryRuntimeHealth() connection = psycopg.connect(dsn, autocommit=True) - last_retention_monotonic = 0.0 + threading.Thread( + target=_retention_loop, args=(dsn,), name="telemetry-retention", daemon=True + ).start() runtime_health.set_database(True) _start_query_server(dsn, runtime_health) client = mqtt.Client( @@ -480,7 +516,7 @@ def main() -> None: runtime_health.set_mqtt(False, f"MQTT disconnected: {reason_code}") def on_message(_client: Any, _userdata: object, message: Any) -> None: - nonlocal connection, last_retention_monotonic + nonlocal connection try: row = _normalize(message.topic, message.payload) except (ValueError, UnicodeDecodeError, json.JSONDecodeError) as exc: @@ -489,19 +525,6 @@ def main() -> None: for attempt in range(2): try: with connection.cursor() as cursor: - now_monotonic = time.monotonic() - if ( - now_monotonic - last_retention_monotonic - >= RETENTION_INTERVAL_SECONDS - ): - cursor.execute( - """ - DELETE FROM contour_telemetry_samples - WHERE observed_at - < CURRENT_TIMESTAMP - INTERVAL '30 days' - """ - ) - last_retention_monotonic = now_monotonic cursor.execute( """ INSERT INTO contour_telemetry_samples ( diff --git a/deploy/telemetry-plane/telegraf/Get-NdcMissionCorePipelineTelemetry.ps1 b/deploy/telemetry-plane/telegraf/Get-NdcMissionCorePipelineTelemetry.ps1 index 7094895..f352d16 100644 --- a/deploy/telemetry-plane/telegraf/Get-NdcMissionCorePipelineTelemetry.ps1 +++ b/deploy/telemetry-plane/telegraf/Get-NdcMissionCorePipelineTelemetry.ps1 @@ -1,7 +1,8 @@ [CmdletBinding()] param( [string]$ContainerName = "ndc-mission-core-perception-worker", - [string]$HealthUrl = "http://127.0.0.1:18020/health" + [string]$HealthUrl = "http://127.0.0.1:18020/health", + [string]$RuntimeSnapshot = "C:\ProgramData\NDC\MissionCore\telemetry-agent\perception\current.json" ) $ErrorActionPreference = "Stop" @@ -18,6 +19,46 @@ $stageIds = @( ) $health = $null $collectorState = "unavailable" +# A fresh full-profile export supersedes the historical service. Presence of an +# expired/broken export is NOT permission to relabel a legacy model as current. +if (Test-Path -LiteralPath $RuntimeSnapshot -PathType Leaf) { + $runtime = $null + try { + $item = Get-Item -LiteralPath $RuntimeSnapshot + $age = ([DateTime]::UtcNow - $item.LastWriteTimeUtc).TotalSeconds + if ($item.Length -le 8192 -and $age -ge -1 -and $age -le 5) { + $stream = [IO.File]::Open($RuntimeSnapshot, 'Open', 'Read', 'ReadWrite') + try { + $bytes = New-Object byte[] 8193 + $count = $stream.Read($bytes, 0, $bytes.Length) + if ($count -le 8192) { + $runtime = [Text.Encoding]::UTF8.GetString($bytes, 0, $count) | ConvertFrom-Json + } + } finally { $stream.Dispose() } + if ($runtime.schema_version -ne 'missioncore.perception-runtime-observation/v1') { + $runtime = $null + } + } + } catch { $runtime = $null } + $samples = foreach ($stageId in $stageIds) { + $sample = [ordered]@{ + stage_id = $stageId + stage_state = 'unavailable' + service_state = 'unavailable' + collector_state = 'unavailable' + } + if ($runtime) { + $sample.collector_state = 'live' + foreach ($name in @('service_state', 'profile_name', 'active_request_id', + 'input_state', 'wait_reason', 'live_children', 'input_pauses', 'pending_bundles', 'buffer_bytes')) { + $sample[$name] = $runtime.$name + } + } + [pscustomobject]$sample + } + @($samples) | ConvertTo-Json -Compress -Depth 4 + exit 0 +} try { $healthJson = docker exec $ContainerName python3 -c ` "import urllib.request;print(urllib.request.urlopen('$HealthUrl',timeout=2).read().decode())" ` @@ -68,13 +109,13 @@ $completedRuns = if ($health -and $null -ne $health.completed_runs) { [int64]$health.completed_runs } else { - [int64]0 + [int64]-1 } $failedRuns = if ($health -and $null -ne $health.failed_runs) { [int64]$health.failed_runs } else { - [int64]0 + [int64]-1 } $modelLoadSeconds = if ($health -and $null -ne $health.model_load_seconds) { [double]$health.model_load_seconds diff --git a/deploy/telemetry-plane/telegraf/mission-core-linux.conf.tmpl b/deploy/telemetry-plane/telegraf/mission-core-linux.conf.tmpl index 633ade1..68aa1c6 100644 --- a/deploy/telemetry-plane/telegraf/mission-core-linux.conf.tmpl +++ b/deploy/telemetry-plane/telegraf/mission-core-linux.conf.tmpl @@ -1,5 +1,8 @@ [agent] interval = "${MISSIONCORE_TELEMETRY_INTERVAL}" + flush_interval = "${MISSIONCORE_TELEMETRY_INTERVAL}" + metric_batch_size = 200 + metric_buffer_limit = 2000 round_interval = true omit_hostname = false @@ -26,6 +29,7 @@ total = true [[outputs.mqtt]] + startup_error_behavior = "retry" servers = ["tcp://${MISSIONCORE_MQTT_HOST}:${MISSIONCORE_MQTT_PORT}"] topic = "mission-core/v1/contours/${MISSIONCORE_CONTOUR_ID}/agents/${MISSIONCORE_AGENT_ID}/host" username = "${MISSIONCORE_MQTT_USERNAME}" diff --git a/deploy/telemetry-plane/telegraf/mission-core-windows.conf.tmpl b/deploy/telemetry-plane/telegraf/mission-core-windows.conf.tmpl index 13ba712..89acea8 100644 --- a/deploy/telemetry-plane/telegraf/mission-core-windows.conf.tmpl +++ b/deploy/telemetry-plane/telegraf/mission-core-windows.conf.tmpl @@ -1,5 +1,8 @@ [agent] interval = "${MISSIONCORE_TELEMETRY_INTERVAL}" + flush_interval = "${MISSIONCORE_TELEMETRY_INTERVAL}" + metric_batch_size = 200 + metric_buffer_limit = 2000 round_interval = true omit_hostname = false @@ -30,6 +33,7 @@ [[inputs.docker]] endpoint = "npipe:////./pipe/docker_engine" container_name_include = [] + container_state_include = ["running", "exited", "restarting", "paused"] [[inputs.http_response]] urls = ["http://127.0.0.1:8000/v2/health/ready"] @@ -59,6 +63,9 @@ "current_stage", "active_request_id", "collector_state", + "profile_name", + "input_state", + "wait_reason", ] [[inputs.tail]] @@ -73,6 +80,7 @@ data_type = "string" [[outputs.mqtt]] + startup_error_behavior = "retry" servers = ["tcp://${MISSIONCORE_MQTT_HOST}:${MISSIONCORE_MQTT_PORT}"] topic = "mission-core/v1/contours/${MISSIONCORE_CONTOUR_ID}/agents/${MISSIONCORE_AGENT_ID}/host" username = "${MISSIONCORE_MQTT_USERNAME}" @@ -83,6 +91,7 @@ namedrop = ["missioncore_pipeline", "missioncore_pipeline_event"] [[outputs.mqtt]] + startup_error_behavior = "retry" servers = ["tcp://${MISSIONCORE_MQTT_HOST}:${MISSIONCORE_MQTT_PORT}"] topic = "mission-core/v1/contours/${MISSIONCORE_CONTOUR_ID}/agents/${MISSIONCORE_AGENT_ID}/pipeline" username = "${MISSIONCORE_MQTT_USERNAME}" diff --git a/src/k1link/web/compute_contour_network.py b/src/k1link/web/compute_contour_network.py index 0c21ca0..1095788 100644 --- a/src/k1link/web/compute_contour_network.py +++ b/src/k1link/web/compute_contour_network.py @@ -486,7 +486,7 @@ def _run_compose_broker(telemetry_plane_root: Path) -> None: ] for suffix, timeout in ( (["config", "--quiet"], 20), - (["up", "-d", "--no-build", "broker"], 45), + (["up", "-d", "--no-deps", "--no-build", "--force-recreate", "broker"], 45), ): try: completed = subprocess.run( @@ -525,11 +525,17 @@ def apply_broker_network( target.mqtt_bind_address, ) changed = after != before + repaired = False mode = environment_path.stat().st_mode try: if changed: _atomic_private_write(environment_path, after, mode=mode) _run_compose_broker(telemetry_plane_root) + elif not _tcp_reachable(target.mqtt_bind_address, target.mqtt_port, timeout=3): + # An unchanged file does not prove Docker still publishes the port + # after a host/network restart. This is an explicit Apply, not a GET. + _run_compose_broker(telemetry_plane_root) + repaired = True if not _tcp_reachable( target.mqtt_bind_address, target.mqtt_port, @@ -553,6 +559,7 @@ def apply_broker_network( "contour_id": target.contour_id, "target": "broker", "changed": changed, + "repaired": repaired, "applied_at_utc": _utc_now(), "ready": True, } diff --git a/src/k1link/web/system_telemetry_api.py b/src/k1link/web/system_telemetry_api.py index fd3ba40..4bc630c 100644 --- a/src/k1link/web/system_telemetry_api.py +++ b/src/k1link/web/system_telemetry_api.py @@ -478,7 +478,7 @@ def _agent_raw_document(document: dict[str, Any]) -> dict[str, Any]: perception: dict[str, Any] = raw["perception"] pipeline_stage_metrics: dict[str, dict[str, Any]] = {} pipeline_active_stages: set[str] = set() - for sample in samples: + for sample in sorted(samples, key=lambda item: str(item.get("observed_at_utc", ""))): measurement = sample.get("measurement") payload, fields, tags = _metric_payload(sample) if sample.get("kind") == "pipeline": @@ -497,6 +497,9 @@ def _agent_raw_document(document: dict[str, Any]) -> dict[str, Any]: "service_state", "current_stage", "active_request_id", + "profile_name", + "input_state", + "wait_reason", ): value = fields.get(name) if isinstance(value, str): @@ -511,6 +514,8 @@ def _agent_raw_document(document: dict[str, Any]) -> dict[str, Any]: if model_load_seconds is not None: perception["model_load_seconds"] = model_load_seconds perception["collector_state"] = fields.get("collector_state") + for name in ("live_children", "input_pauses", "pending_bundles", "buffer_bytes"): + perception[name] = _number(fields.get(name)) elif sample.get("source_schema") == PIPELINE_TELEMETRY_SCHEMA: # Native lifecycle rows are immutable run evidence. The current # service snapshot remains owned by the periodic health sample, @@ -593,6 +598,9 @@ def _agent_raw_document(document: dict[str, Any]) -> dict[str, Any]: "memory_total_mib": memory_total, "power_watts": _number(fields.get("power_draw")), "temperature_celsius": _number(fields.get("temperature_gpu")), + "sm_clock_mhz": _number(fields.get("clocks_current_sm")), + "memory_clock_mhz": _number(fields.get("clocks_current_memory")), + "driver_version": fields.get("driver_version"), } elif measurement == "missioncore_triton_health": raw["triton"] = { @@ -637,7 +645,15 @@ def _agent_raw_document(document: dict[str, Any]) -> dict[str, Any]: "image": tags.get("container_image"), }, ) - if measurement == "docker_container_cpu": + if measurement == "docker_container_status": + container_states[container_name]["state"]["Status"] = ( + tags.get("container_status") or "unavailable" + ) + elif measurement == "docker_container_health": + container_states[container_name]["state"]["Health"] = { + "Status": fields.get("health_status") or "unknown", + } + elif measurement == "docker_container_cpu": usage = _number(fields.get("usage_percent")) if usage is not None: stats["CPUPerc"] = f"{usage:.3f}%" @@ -906,7 +922,12 @@ def _pipeline_document(raw: dict[str, Any]) -> dict[str, Any]: ) def stage_state(stage_id: str) -> str: - if not perception: + if not perception or perception.get("state") == "unavailable": + return "unavailable" + if perception.get("state") in ("waiting", "synchronizing", "starting"): + return "waiting" + # A full-profile heartbeat is not a measurement of each individual stage. + if perception.get("profile_name"): return "unavailable" if not busy: return "ready" @@ -970,6 +991,13 @@ def _pipeline_document(raw: dict[str, Any]) -> dict[str, Any]: else None ), "model_load_seconds": _number(perception.get("model_load_seconds")), + "profile_name": perception.get("profile_name"), + "input_state": perception.get("input_state"), + "wait_reason": perception.get("wait_reason"), + "live_children": _number(perception.get("live_children")), + "input_pauses": _number(perception.get("input_pauses")), + "pending_bundles": _number(perception.get("pending_bundles")), + "buffer_bytes": _number(perception.get("buffer_bytes")), "stages": [stage_document(stage_id, label) for stage_id, label in stages], } diff --git a/tests/test_compute_contour_network.py b/tests/test_compute_contour_network.py index 2aa0ea6..99cf042 100644 --- a/tests/test_compute_contour_network.py +++ b/tests/test_compute_contour_network.py @@ -63,6 +63,21 @@ def test_broker_apply_preserves_private_environment_and_changes_only_bind( assert result["ready"] is True +def test_broker_apply_repairs_missing_listener_without_changing_secrets(tmp_path, monkeypatch): + before = "MISSIONCORE_MQTT_BIND_ADDRESS=192.168.68.56\nSECRET=keep\n" + (tmp_path / ".env").write_text(before) + (tmp_path / "compose.yaml").write_text("services: {}\n") + reachable = iter((False, True)) + calls = [] + monkeypatch.setattr(network, "_address_belongs_to_host", lambda _: True) + monkeypatch.setattr(network, "_tcp_reachable", lambda *a, **k: next(reachable)) + monkeypatch.setattr(network, "_run_compose_broker", lambda path: calls.append(path)) + result = network.apply_broker_network(_target(), tmp_path) + assert not result["changed"] and result["repaired"] and result["ready"] + assert calls == [tmp_path] + assert (tmp_path / ".env").read_text() == before + + def test_broker_apply_rolls_back_environment_when_listener_does_not_open( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, diff --git a/tests/test_system_telemetry_api.py b/tests/test_system_telemetry_api.py index 82ac9a5..fd90b20 100644 --- a/tests/test_system_telemetry_api.py +++ b/tests/test_system_telemetry_api.py @@ -447,6 +447,36 @@ def test_compute_contour_maps_to_worker_identity_without_singleton_defaults() -> assert profile.address == "192.0.2.25" +def test_gpu_clocks_and_explicit_container_status_are_not_invented() -> None: + raw = _agent_raw_document({"samples": [ + {"measurement": "nvidia_smi", "payload": {"fields": { + "clocks_current_sm": 210, "clocks_current_memory": 405, + }}}, + {"measurement": "docker_container_status", "payload": { + "tags": {"container_name": "sentinel-frigate", "container_status": "exited"}, + "fields": {}, + }}, + ]}) + assert raw["gpu"]["sm_clock_mhz"] == 210 + assert raw["gpu"]["memory_clock_mhz"] == 405 + assert raw["container_states"]["sentinel-frigate"]["state"]["Status"] == "exited" + + +def test_unavailable_pipeline_is_not_ready_and_profile_wait_survives(tmp_path) -> None: + probe = _probe() + probe["raw"]["perception"] = {"state": "unavailable", "collector_state": "unavailable"} + service = WorkerTelemetryService(WorkerProfileStore(tmp_path), lambda _: probe, cache_seconds=0) + assert all(s["state"] == "unavailable" for s in service.snapshot(1)["pipeline"]["stages"]) + probe["raw"]["perception"] = { + "state": "waiting", "profile_name": "K1 DDRNet", "input_pauses": 1, + "live_children": 4, "buffer_bytes": 1024, + } + profile = service.snapshot(1)["pipeline"] + assert profile["service_state"] == "waiting" + assert profile["live_children"] == 4 and profile["buffer_bytes"] == 1024 + assert all(s["state"] == "waiting" for s in profile["stages"]) + + def test_profile_apply_fails_closed_on_wrong_node_and_keeps_old_profile( tmp_path: Path, ) -> None: diff --git a/tests/test_telemetry_normalizer_contract.py b/tests/test_telemetry_normalizer_contract.py index 4cd9ee1..2877f21 100644 --- a/tests/test_telemetry_normalizer_contract.py +++ b/tests/test_telemetry_normalizer_contract.py @@ -126,6 +126,22 @@ def test_normalizer_enforces_oss_compatible_bounded_retention() -> None: assert "DELETE FROM contour_telemetry_samples" in NORMALIZER_SOURCE assert "INTERVAL '30 days'" in NORMALIZER_SOURCE assert "RETENTION_INTERVAL_SECONDS" in NORMALIZER_SOURCE + assert "ORDER BY observed_at LIMIT %s" in NORMALIZER_SOURCE + assert "statement_timeout = '2s'" in NORMALIZER_SOURCE + assert "samples.tableoid = expired.tableoid" in NORMALIZER_SOURCE + assert "target=_retention_loop" in NORMALIZER_SOURCE + callback = NORMALIZER_SOURCE.split(" def on_message(", 1)[1] + assert "DELETE FROM" not in callback + + +def test_telegraf_network_recovery_and_database_memory_are_bounded() -> None: + compose = yaml.safe_load(COMPOSE_PATH.read_text()) + assert "shared_buffers=128MB" in compose["services"]["timescale"]["command"] + for platform, outputs in (("windows", 2), ("linux", 1)): + template = (REPOSITORY_ROOT / "deploy/telemetry-plane/telegraf" / + f"mission-core-{platform}.conf.tmpl").read_text() + assert template.count('startup_error_behavior = "retry"') == outputs + assert "metric_buffer_limit = 2000" in template def test_normalizer_rejects_unschematized_pipeline_payload() -> None: