From 27979854a73d8a96db591f77eb9449e4a0f59fa5 Mon Sep 17 00:00:00 2001 From: DCCONSTRUCTIONS Date: Thu, 3 Sep 2026 13:15:07 +0300 Subject: [PATCH] feat(observatory): cache sealed camera and TGS replay on Core --- .../laboratory/canonical_rerun_overlay.py | 6 +- src/k1link/observatory/portable_replay.py | 328 ++++++++++++++++++ src/k1link/observatory/portable_tgs_replay.py | 216 ++++++++++++ src/k1link/web/app.py | 17 + src/k1link/web/portable_replay_api.py | 75 ++++ tests/test_portable_replay.py | 267 ++++++++++++++ 6 files changed, 907 insertions(+), 2 deletions(-) create mode 100644 src/k1link/observatory/portable_replay.py create mode 100644 src/k1link/observatory/portable_tgs_replay.py create mode 100644 src/k1link/web/portable_replay_api.py create mode 100644 tests/test_portable_replay.py diff --git a/src/k1link/laboratory/canonical_rerun_overlay.py b/src/k1link/laboratory/canonical_rerun_overlay.py index 023898a..4e8afe6 100644 --- a/src/k1link/laboratory/canonical_rerun_overlay.py +++ b/src/k1link/laboratory/canonical_rerun_overlay.py @@ -15,6 +15,7 @@ import hashlib import io import json import os +import re import subprocess import sys import tempfile @@ -282,8 +283,9 @@ def canonical_lab_replay( or not _is_sha256(base_generation_sha256) or _sha256(base) != base_generation_sha256 or not _artifact_is_regular(overlay) - or not result_id.startswith("lab-v1-vegetation-shadow-") - or len(result_id) != len("lab-v1-vegetation-shadow-") + 64 + or re.fullmatch( + r"(?:lab-v1-vegetation-shadow|m49-tgs-portable-review)-[a-f0-9]{64}", result_id + ) is None or not recording_id or len(recording_id) > 128 ): diff --git a/src/k1link/observatory/portable_replay.py b/src/k1link/observatory/portable_replay.py new file mode 100644 index 0000000..0a16d33 --- /dev/null +++ b/src/k1link/observatory/portable_replay.py @@ -0,0 +1,328 @@ +"""Bounded, cached native replay packaging on Core; no inference or Worker IO.""" + +from __future__ import annotations + +import hashlib +import json +import os +import re +import shutil +import subprocess +import tempfile +import threading +from collections import OrderedDict +from collections.abc import Callable +from pathlib import Path +from typing import Any, cast + +import numpy as np +import rerun as rr + +from k1link.laboratory.canonical_rerun_overlay import ( + APPLICATION_ID, + CanonicalLabOverlayArtifact, + CanonicalLabReplayArtifact, + canonical_lab_replay, + canonical_recording_id, +) +from k1link.observatory.portable_result_view import PortableResultViewService +from k1link.observatory.portable_tgs_replay import ( + PortableReplayError, + load_tgs_data, + log_tgs, + read_json, + stat_identity, + verified_file, +) +from k1link.observatory.source_admission import PORTABLE_SOURCE_DOCUMENT_DIRECTORY +from k1link.sessions.media import RecordedMediaEpoch, RecordedMediaInspector + +RENDERER_VERSION = "portable-tgs-camera-rerun-0.36.3-v1" +_SHA = re.compile(r"^[a-f0-9]{64}$") +_RESULT = re.compile(r"^m49-tgs-portable-review-[a-f0-9]{64}$") +_LOCK = threading.Lock() + + +class PortableReplayService: + def __init__( + self, + *, + view: PortableResultViewService, + data_dir: Path, + media: RecordedMediaInspector, + recording_source: Callable[[str], tuple[Path, str] | None], + ffmpeg_path: Path, + ) -> None: + self.view = view + self.data_dir = data_dir + self.media = media + self.recording_source = recording_source + self.ffmpeg_path = ffmpeg_path + self.cache = data_dir / "observatory-portable-replay-cache" + self._verified: OrderedDict[tuple[Path, str], tuple[int, ...]] = OrderedDict() + + def _key(self, result_id: str, base_sha: str) -> tuple[str, dict[str, Any]]: + if _RESULT.fullmatch(result_id) is None or _SHA.fullmatch(base_sha) is None: + raise PortableReplayError("unsupported replay identity") + view = cast(dict[str, Any], self.view.read(result_id)) + identity = [RENDERER_VERSION, result_id, base_sha, view["artifact_manifest_id"]] + return hashlib.sha256("\0".join(identity).encode()).hexdigest(), view + + def cached(self, result_id: str, base_sha: str) -> CanonicalLabReplayArtifact | None: + key, _ = self._key(result_id, base_sha) + sidecar = self.cache / f"{key}.json" + if not sidecar.exists(): + return None + metadata = read_json(sidecar, 8192) + digest = metadata.get("sha256") + if ( + metadata.get("key") != key + or not isinstance(digest, str) + or _SHA.fullmatch(digest) is None + ): + raise PortableReplayError("replay cache identity changed") + name = metadata.get("file") + if not isinstance(name, str) or not re.fullmatch(r"[a-f0-9]{64}\.replay\.rrd", name): + raise PortableReplayError("replay cache path is unsafe") + path = self.cache / name + # Warm open reads the already sealed derivative, never source media or + # TGS arrays. Hashing is bounded by the one-RRD browser admission cap. + if ( + path.is_symlink() + or not path.is_file() + or path.stat().st_size != metadata["byte_length"] + ): + raise PortableReplayError("replay cache length changed") + if not 4 <= metadata["byte_length"] <= 1024 * 1024 * 1024: + raise PortableReplayError("replay cache exceeds bounds") + stamp = stat_identity(path.stat()) + fingerprint = (path, digest) + if self._verified.get(fingerprint) != stamp: + with path.open("rb") as stream: + if ( + hashlib.file_digest(stream, "sha256").hexdigest() != digest + or stat_identity(path.stat()) != stamp + ): + raise PortableReplayError("replay cache digest changed") + self._verified[fingerprint] = stamp + while len(self._verified) > 32: + self._verified.popitem(last=False) + return CanonicalLabReplayArtifact(path, metadata["byte_length"], digest) + + def prepare(self, result_id: str, base_sha: str) -> CanonicalLabReplayArtifact: + # Serialize packaging on the constrained operator host, including + # concurrent clicks. Waiting readers cannot start a duplicate render. + with _LOCK: + cached = self.cached(result_id, base_sha) + if cached is not None: + return cached + key, view = self._key(result_id, base_sha) + base = self.recording_source(view["source_session_id"]) + if base is None or base[1] != base_sha: + raise PortableReplayError("source recording generation is not ready") + binding = self.view.sessions.get_lab_instance(result_id) + if binding is None: + raise PortableReplayError("result binding disappeared") + source = binding.provenance["source"] + bundle_sha = source["bundle_sha256"] + if not isinstance(bundle_sha, str) or _SHA.fullmatch(bundle_sha) is None: + raise PortableReplayError("source bundle identity is invalid") + bundle_path = self.data_dir / PORTABLE_SOURCE_DOCUMENT_DIRECTORY / f"{bundle_sha}.json" + bundle = read_json(bundle_path, 32 * 1024 * 1024) + verified_file(bundle_path, bundle_sha, bundle_path.stat().st_size) + if bundle["source_session_id"] != view["source_session_id"]: + raise PortableReplayError("source bundle belongs to another session") + _, catalog_sha = self.view.sessions.get_session_with_catalog_snapshot( + view["source_session_id"] + ) + if ( + catalog_sha != source["catalog_sha256"] + or catalog_sha != bundle["source_catalog_sha256"] + ): + raise PortableReplayError("source catalog changed since calculation") + epoch = self._camera(view["source_session_id"], bundle) + data = load_tgs_data(view, self.view.artifacts, source_bundle_sha256=bundle_sha) + starts = [epoch.timeline_start_seconds] + [ + epoch.timeline_start_seconds + part.end_time_seconds for part in epoch.segments[:-1] + ] + if len(starts) != len(data.rows) or any( + abs(start - row["session_seconds"]) > 1e-8 + for start, row in zip(starts, data.rows, strict=True) + ): + raise PortableReplayError("camera and costmap anchor clocks disagree") + recording_id = canonical_recording_id(base[0]) + self.cache.mkdir(mode=0o700, parents=True, exist_ok=True) + if self.cache.is_symlink() or shutil.disk_usage(self.cache).free < 3 * 1024**3: + raise PortableReplayError("replay cache has insufficient safe space") + with tempfile.TemporaryDirectory(prefix=".pack-", dir=self.cache) as temporary: + root = Path(temporary) + proxy = self._video(epoch, root) + output = root / "overlay.rrd" + recording = rr.RecordingStream(APPLICATION_ID, recording_id=recording_id) + try: + recording.set_sinks(rr.FileSink(output, write_footer=True)) + _log_video(recording, proxy, epoch) + log_tgs(recording, data, epoch.timeline_end_seconds) + recording.flush(timeout_sec=180) + finally: + recording.disconnect() + with output.open("rb") as stream: + digest = hashlib.file_digest(stream, "sha256").hexdigest() + artifact = canonical_lab_replay( + base[0], + base_generation_sha256=base_sha, + overlay=CanonicalLabOverlayArtifact(output, output.stat().st_size, digest), + result_id=result_id, + recording_id=recording_id, + cache_root=self.cache, + ) + metadata = { + "key": key, + "file": artifact.path.name, + "sha256": artifact.sha256, + "byte_length": artifact.byte_length, + } + staged = self.cache / f".{key}.json" + staged.write_text(json.dumps(metadata, sort_keys=True), encoding="utf-8") + os.chmod(staged, 0o600) + os.replace(staged, self.cache / f"{key}.json") + return artifact + + def _camera(self, session_id: str, bundle: dict[str, Any]) -> RecordedMediaEpoch: + camera = bundle["camera"] + artifact = self.view.sessions.get_recorded_media(session_id, camera["artifact_id"]) + command = self.view.sessions.prepare_replay(session_id, speed=1.0, loop=False) + media = self.media.restore_prepared(artifact, command) + if ( + media is None + or media.generation_sha256 != camera["generation_sha256"] + or len(media.epochs) != 1 + ): + raise PortableReplayError("exact source camera is not prepared") + epoch, expected = media.epochs[0], camera["epoch"] + actual = { + "ordinal": epoch.ordinal, + "media_type": epoch.media_type, + "init": {"byte_length": epoch.init_byte_length, "sha256": epoch.init_sha256}, + "timeline_start_seconds": epoch.timeline_start_seconds, + "timeline_end_seconds": epoch.timeline_end_seconds, + "segments": [ + { + "sequence": item.sequence, + "byte_length": item.byte_length, + "sha256": item.sha256, + "random_access": item.random_access, + "end_time_seconds": item.end_time_seconds, + } + for item in epoch.segments + ], + } + if actual != expected: + raise PortableReplayError("source camera epoch changed") + return epoch + + def _video(self, epoch: RecordedMediaEpoch, root: Path) -> Path: + files = [(epoch.init_path, epoch.init_sha256, epoch.init_byte_length)] + [ + (part.path, part.sha256, part.byte_length) for part in epoch.segments + ] + video_key = hashlib.sha256( + ( + "camera-proxy-demux-timebase-v2\0" + "\0".join(digest for _, digest, _ in files) + ).encode() + ).hexdigest() + saved = self.cache / f"{video_key}.camera.mp4" + saved_manifest = self.cache / f"{video_key}.camera.json" + if saved_manifest.exists(): + metadata = read_json(saved_manifest, 8192) + return verified_file(saved, metadata["sha256"], metadata["byte_length"]) + if sum(size for _, _, size in files) > 768 * 1024 * 1024: + raise PortableReplayError("camera packaging exceeds the current viewer bound") + source, proxy = root / "source.mp4", root / "camera.mp4" + with source.open("wb") as target: + for path, digest, size in files: + verified_file(path, digest, size) + with path.open("rb") as part: + shutil.copyfileobj(part, target, 1024 * 1024) + completed = subprocess.run( + [ + str(self.ffmpeg_path), + "-nostdin", + "-hide_banner", + "-loglevel", + "error", + "-i", + str(source), + "-an", + "-c:v", + "libx264", + "-threads", + "2", + "-preset", + "veryfast", + "-crf", + "28", + "-g", + "20", + "-keyint_min", + "20", + "-pix_fmt", + "yuv420p", + "-fps_mode", + "passthrough", + "-enc_time_base", + "-1", + "-movflags", + "+faststart", + str(proxy), + ], + capture_output=True, + timeout=900, + check=False, + ) + if completed.returncode or not proxy.is_file() or proxy.stat().st_size > 192 * 1024 * 1024: + raise PortableReplayError("PTS-preserving camera packaging failed") + with proxy.open("rb") as stream: + digest = hashlib.file_digest(stream, "sha256").hexdigest() + metadata = {"sha256": digest, "byte_length": proxy.stat().st_size} + os.chmod(proxy, 0o600) + os.replace(proxy, saved) + staged = root / "camera.json" + staged.write_text(json.dumps(metadata), encoding="utf-8") + os.chmod(staged, 0o600) + os.replace(staged, saved_manifest) + return saved + + +def _log_video(recording: rr.RecordingStream, proxy: Path, epoch: RecordedMediaEpoch) -> None: + video = rr.AssetVideo(path=proxy) + timestamps = video.read_frame_timestamps_nanos() + if ( + not len(timestamps) + or np.any(np.diff(timestamps) <= 0) + or timestamps[0] < 0 + or abs( + float(timestamps[-1]) / 1e9 + - (epoch.timeline_end_seconds - epoch.timeline_start_seconds) + ) + > 2 + ): + raise PortableReplayError( + f"video presentation timestamps changed: count={len(timestamps)}, " + f"first={timestamps[0] if len(timestamps) else None}, " + f"last={timestamps[-1] if len(timestamps) else None}, " + f"nonincreasing={int(np.count_nonzero(np.diff(timestamps) <= 0))}, " + f"expected_duration={epoch.timeline_end_seconds - epoch.timeline_start_seconds}" + ) + # Native decoded PTS, never fragment count / nominal FPS. Rerun latest-at + # holds the prior sample across missing decodable fragments. + recording.log("/perception/camera/image", video, static=True) + origin = round(epoch.timeline_start_seconds * 1e9) + for timestamp in timestamps: + recording.set_time("session_time", duration=np.timedelta64(origin + int(timestamp), "ns")) + recording.log( + "/perception/camera/image", rr.VideoFrameReference(nanoseconds=int(timestamp)) + ) + recording.set_time( + "session_time", duration=np.timedelta64(round(epoch.timeline_end_seconds * 1e9), "ns") + ) + recording.log("/perception/camera/image", rr.Clear(recursive=False)) diff --git a/src/k1link/observatory/portable_tgs_replay.py b/src/k1link/observatory/portable_tgs_replay.py new file mode 100644 index 0000000..f45b96a --- /dev/null +++ b/src/k1link/observatory/portable_tgs_replay.py @@ -0,0 +1,216 @@ +"""Read-only projection of sealed portable TGS evidence onto the source clock. + +This adapter never calls a model, a Worker, or a legacy LAB result. Local +costmap cells are translated by the exact source pose, not by a viewer eye. +""" + +from __future__ import annotations + +import hashlib +import json +import math +import os +from dataclasses import dataclass +from pathlib import Path +from typing import Any + +import numpy as np +import rerun as rr + +from k1link.artifact_gateway import CentralArtifactStore +from k1link.observatory.m49_portable_source import read_m49_source_index +from k1link.observatory.portable_result_contract import canonical_json + +RESULT_SCHEMA = "missioncore.recorded-tgs-costmap-review/v2" +MAX_ARTIFACT_BYTES = 512 * 1024 * 1024 +MAX_FRAMES = 250_000 +MAX_CELLS = 100_000 +STATE_CODES = { + "UNOBSERVED": 0, + "GROUND_SUPPORT": 1, + "NONGROUND_OCCUPIED": 2, + "UNKNOWN_REJECTED": 3, +} +# Evidence colors, not control states: same four categories as the result. +STATE_COLORS = ( + (110, 115, 125, 90), + (108, 182, 115, 130), + (239, 107, 112, 180), + (229, 186, 83, 150), +) + + +class PortableReplayError(ValueError): + """Published evidence cannot safely produce a viewer derivative.""" + + +def stat_identity(value: os.stat_result) -> tuple[int, ...]: + # stat_result tuple equality omits nanosecond precision on macOS. + return (value.st_dev, value.st_ino, value.st_size, value.st_mtime_ns, value.st_ctime_ns) + + +def verified_file(path: Path, digest: str, size: int) -> Path: + if path.is_symlink() or not path.is_file() or not 0 < size <= MAX_ARTIFACT_BYTES: + raise PortableReplayError("replay artifact is outside bounds") + before = path.stat() + if before.st_size != size: + raise PortableReplayError("replay artifact length changed") + with path.open("rb") as stream: + actual = hashlib.file_digest(stream, "sha256").hexdigest() + if actual != digest or stat_identity(before) != stat_identity(path.stat()): + raise PortableReplayError("replay artifact integrity changed") + return path + + +def read_json(path: Path, limit: int = 8 * 1024 * 1024) -> dict[str, Any]: + if path.is_symlink() or not path.is_file() or not 0 < path.stat().st_size <= limit: + raise PortableReplayError("replay document is outside bounds") + value = json.loads(path.read_bytes()) + if not isinstance(value, dict): + raise PortableReplayError("replay document is not an object") + return value + + +@dataclass(frozen=True) +class TgsReplayData: + rows: tuple[dict[str, Any], ...] + centers: np.ndarray[Any, Any] + states: np.ndarray[Any, Any] + z_bounds: np.ndarray[Any, Any] + cell_size_m: float + + @property + def start_seconds(self) -> float: + return float(self.rows[0]["session_seconds"]) + + +def load_tgs_data( + view: dict[str, Any], store: CentralArtifactStore, *, source_bundle_sha256: str | None = None +) -> TgsReplayData: + doc = view["result_document"] + if doc.get("schema_version") != RESULT_SCHEMA or doc.get("result_id") != view["result_id"]: + raise PortableReplayError("result has no supported TGS replay contract") + costmap = doc["costmap"] + count, cells = doc["timeline"]["frame_count"], costmap["cell_count"] + size = costmap["cell_size_m"] + if ( + type(count) is not int + or not 0 < count <= MAX_FRAMES + or type(cells) is not int + or not 0 < cells <= MAX_CELLS + or not isinstance(size, (int, float)) + or isinstance(size, bool) + or not math.isfinite(size) + or not 0 < size <= 10 + or costmap["coordinate_frame"] != "map-gravity-local" + or costmap["state_codes"] != STATE_CODES + ): + raise PortableReplayError("TGS grid contract is invalid") + members = {item["role"]: item for item in view["artifacts"]} + + def artifact(role: str) -> Path: + member = members[role] + return verified_file( + store.object_path(member["sha256"]), member["sha256"], member["byte_length"] + ) + + stage = read_json(artifact("source-stage-manifest")) + if ( + stage["identity_sha256"] != doc["source_stage"]["identity_sha256"] + or members["source-stage-manifest"]["sha256"] != doc["source_stage"]["manifest_sha256"] + or stage["identity"]["source"]["source_session_id"] != view["source_session_id"] + or stage["identity"]["profile"]["profile_sha256"] != doc["profile"]["profile_sha256"] + or hashlib.sha256(canonical_json(stage["identity"])).hexdigest() != stage["identity_sha256"] + or ( + source_bundle_sha256 is not None + and stage["identity"]["source"]["source_bundle_sha256"] != source_bundle_sha256 + ) + ): + raise PortableReplayError("source stage identity changed") + rows = tuple(read_m49_source_index(artifact("source-stage-index"), expected_frame_count=count)) + with artifact("frame-index").open("rb") as stream: + for row in rows: + line = stream.readline(65537) + if len(line) > 65536: + raise PortableReplayError("TGS frame row exceeds bounds") + frame = json.loads(line) + for key in ( + "timeline_frame_index", + "source_frame_index", + "session_seconds", + "sample_available", + "available_slot", + ): + if frame.get(key) != row[key]: + raise PortableReplayError("TGS frame and source clocks disagree") + if stream.read(1): + raise PortableReplayError("TGS frame index has extra rows") + centers = np.load(artifact("costmap-cell-centers"), mmap_mode="r", allow_pickle=False) + states = np.load(artifact("costmap-states"), mmap_mode="r", allow_pickle=False) + bounds = np.load(artifact("costmap-z-bounds"), mmap_mode="r", allow_pickle=False) + if ( + centers.shape != (cells, 2) + or centers.dtype != np.dtype(" 3) + or not np.isnan(z[~observed]).all() + or not np.isfinite(z[observed]).all() + or np.any(z[observed, 1] < z[observed, 0]) + or (not row["sample_available"] and np.any(observed)) + ): + raise PortableReplayError("TGS states contradict support or height evidence") + return TgsReplayData(rows, centers, states, bounds, float(size)) + + +def world_cells( + data: TgsReplayData, index: int, state: int +) -> tuple[np.ndarray[Any, Any], np.ndarray[Any, Any]]: + """Return measured volumes only; unobserved cells have no measured height.""" + row = data.rows[index] + selected = data.states[index] == state + if state == 0 or not row["sample_available"]: + return np.empty((0, 3)), np.empty((0, 3)) + z = data.z_bounds[index, selected] + pose = np.asarray(row["position_map_m"], dtype=np.float64) + centers = np.column_stack((data.centers[selected], z.mean(axis=1))) + pose + # Planar support stays planar: do not invent a physical thickness. + sizes = np.column_stack((np.full((len(z), 2), data.cell_size_m), z[:, 1] - z[:, 0])) + return centers, sizes + + +def log_tgs(recording: rr.RecordingStream, data: TgsReplayData, end_seconds: float) -> None: + if not data.start_seconds < end_seconds or data.rows[-1]["session_seconds"] >= end_seconds: + raise PortableReplayError("TGS coverage is outside the admitted camera epoch") + for index, row in enumerate(data.rows): + recording.set_time( + "session_time", + duration=np.timedelta64(round(float(row["session_seconds"]) * 1e9), "ns"), + ) + # Explicit per-frame clears prevent stale obstacles surviving a missing + # LiDAR sample or a backwards seek. No source-cloud entity is replaced. + for state, name in ( + (1, "ground_support"), + (2, "nonground_occupied"), + (3, "unknown_rejected"), + ): + centers, sizes = world_cells(data, index, state) + recording.log( + f"/world/costmap/{name}", + rr.Boxes3D( + centers=centers, sizes=sizes, colors=STATE_COLORS[state], fill_mode="solid" + ) + if len(centers) + else rr.Clear(recursive=False), + ) + recording.set_time("session_time", duration=np.timedelta64(round(end_seconds * 1e9), "ns")) + recording.log("/world/costmap", rr.Clear(recursive=True)) diff --git a/src/k1link/web/app.py b/src/k1link/web/app.py index bce9e54..452df2c 100644 --- a/src/k1link/web/app.py +++ b/src/k1link/web/app.py @@ -54,6 +54,7 @@ from k1link.observatory.portable_queue_binding import ( PortableQueueBindingError, PortableRecordedQueueBindingService, ) +from k1link.observatory.portable_replay import PortableReplayService from k1link.observatory.portable_result_cache import PortableResultCache from k1link.observatory.portable_result_contract import ( PortableCalculationProfileRegistry, @@ -212,6 +213,7 @@ from k1link.web.plugin_runtime import ( PluginRuntimeUnavailableError, ) from k1link.web.polygon_api import build_polygon_router, configured_polygon_runs_root +from k1link.web.portable_replay_api import build_portable_replay_router from k1link.web.runtime_diagnostics import configure_scanner_diagnostics from k1link.web.runtime_readiness import ( BackgroundReconcilerReadiness, @@ -1053,6 +1055,21 @@ async def device_plugin_events(websocket: WebSocket, plugin_id: str) -> None: await websocket.close(code=1011, reason="Device plugin state stream failed") +if session_artifact_gateway is not None and _ffmpeg is not None: + app.include_router( + build_portable_replay_router( + PortableReplayService( + view=PortableResultViewService( + sessions=session_store, artifacts=session_artifact_gateway.store + ), + data_dir=session_store.data_dir, + media=session_recorded_media_inspector, + recording_source=_canonical_lab_recording_source, + ffmpeg_path=_ffmpeg, + ) + ) + ) + for legacy_router in plugin_environment.legacy_routers: app.include_router(legacy_router) diff --git a/src/k1link/web/portable_replay_api.py b/src/k1link/web/portable_replay_api.py new file mode 100644 index 0000000..6c7bbe9 --- /dev/null +++ b/src/k1link/web/portable_replay_api.py @@ -0,0 +1,75 @@ +"""Same-origin immutable replay derivatives; GET never starts packaging.""" + +from __future__ import annotations + +import logging +from typing import Annotated + +from fastapi import APIRouter, HTTPException, Path, Query, Response +from fastapi.responses import FileResponse + +from k1link.laboratory.canonical_rerun_overlay import CanonicalLabOverlayError +from k1link.observatory.portable_replay import PortableReplayService +from k1link.observatory.portable_result_view import PortableResultViewError + +ResultId = Annotated[str, Path(pattern=r"^m49-tgs-portable-review-[a-f0-9]{64}$")] +BaseSha = Annotated[str, Path(pattern=r"^[a-f0-9]{64}$")] +_LOG = logging.getLogger(__name__) + + +def build_portable_replay_router(service: PortableReplayService) -> APIRouter: + router = APIRouter(tags=["observatory"]) + path = "/api/v1/observatory/portable-results/{result_id}/replays/{base_sha}/recording.rrd" + + @router.head(path) + def prepare(result_id: ResultId, base_sha: BaseSha) -> Response: + try: + artifact = service.prepare(result_id, base_sha) + except ( + ValueError, + OSError, + KeyError, + TypeError, + CanonicalLabOverlayError, + PortableResultViewError, + ) as exc: + _LOG.exception("Portable replay packaging rejected") + raise HTTPException( + 409, "Сохранённый результат не удалось подготовить к просмотру." + ) from exc + return Response( + media_type="application/vnd.rerun.rrd", + headers={ + "Content-Length": str(artifact.byte_length), + "ETag": f'"{artifact.sha256}"', + "X-Rerun-Format": "RRF2", + "Cache-Control": "private, no-store", + }, + ) + + @router.get(path) + def read( + result_id: ResultId, + base_sha: BaseSha, + generation: Annotated[str, Query(pattern=r"^[a-f0-9]{64}$")], + ) -> FileResponse: + try: + artifact = service.cached(result_id, base_sha) + except (ValueError, OSError, KeyError, TypeError, PortableResultViewError) as exc: + raise HTTPException(409, "Кэш результата не прошёл проверку.") from exc + if artifact is None: + raise HTTPException(409, "Сохранённый просмотр ещё не подготовлен.") + if artifact.sha256 != generation: + raise HTTPException(412, "Версия сохранённого просмотра изменилась.") + return FileResponse( + artifact.path, + media_type="application/vnd.rerun.rrd", + headers={ + "ETag": f'"{artifact.sha256}"', + "X-Rerun-Format": "RRF2", + "Cache-Control": "private, max-age=31536000, immutable", + "X-Content-Type-Options": "nosniff", + }, + ) + + return router diff --git a/tests/test_portable_replay.py b/tests/test_portable_replay.py new file mode 100644 index 0000000..2b69968 --- /dev/null +++ b/tests/test_portable_replay.py @@ -0,0 +1,267 @@ +from __future__ import annotations + +import hashlib +from pathlib import Path +from types import SimpleNamespace + +import numpy as np +import pytest +import rerun as rr +from fastapi import FastAPI +from fastapi.testclient import TestClient + +from k1link.laboratory.canonical_rerun_overlay import CanonicalLabReplayArtifact +from k1link.observatory.portable_replay import PortableReplayService +from k1link.observatory.portable_tgs_replay import ( + PortableReplayError, + TgsReplayData, + log_tgs, + verified_file, + world_cells, +) +from k1link.web.portable_replay_api import build_portable_replay_router + +RESULT = "m49-tgs-portable-review-" + "a" * 64 +BASE = "b" * 64 + + +@pytest.mark.parametrize( + "damage", [None, "frame-time", "source-bundle", "nan-height", "state", "missing"] +) +def test_array_admission_checks_clocks_support_and_identity(tmp_path, monkeypatch, damage): + from k1link.artifact_gateway import CentralArtifactStore + from k1link.observatory.portable_result_contract import canonical_json + from k1link.observatory.portable_tgs_replay import STATE_CODES, load_tgs_data + + store = CentralArtifactStore(tmp_path / "store", create=True) + members = [] + + def member(role, payload): + path = tmp_path / role + path.write_bytes(payload) + item = store.publish_file(path) + members.append({"role": role, "sha256": item.sha256, "byte_length": item.byte_length}) + return item.sha256 + + source_row = { + "timeline_frame_index": 0, + "source_frame_index": 1, + "session_seconds": 1.0, + "sample_available": damage != "missing", + "available_slot": 0, + "position_map_m": [10.0, 20.0, 3.0], + } + monkeypatch.setattr( + "k1link.observatory.portable_tgs_replay.read_m49_source_index", + lambda *args, **kwargs: (source_row,), + ) + member("source-stage-index", b"test") + frame = dict(source_row) + if damage == "frame-time": + frame["session_seconds"] = 2.0 + member("frame-index", canonical_json(frame) + b"\n") + identity = { + "source": {"source_session_id": "source", "source_bundle_sha256": BASE}, + "profile": {"profile_sha256": "c" * 64}, + } + stage_sha = hashlib.sha256(canonical_json(identity)).hexdigest() + stage_manifest = member( + "source-stage-manifest", + canonical_json({"identity": identity, "identity_sha256": stage_sha}), + ) + import io + + for role, array in ( + ("costmap-cell-centers", np.array([[1.0, 2.0]], dtype=" TgsReplayData: + return TgsReplayData( + ( + {"session_seconds": 1.0, "sample_available": True, "position_map_m": [10, 20, 3]}, + {"session_seconds": 2.0, "sample_available": False, "position_map_m": None}, + ), + np.array([[1.0, 2.0], [2.0, 3.0]]), + np.array([[2, 0], [0, 0]], dtype="u1"), + np.array([[[-1.0, 1.0], [np.nan, np.nan]], [[np.nan, np.nan], [np.nan, np.nan]]]), + 0.45, + ) + + +def test_world_cells_use_exact_pose_once_and_missing_is_not_free(): + replay = data() + centers, sizes = world_cells(replay, 0, 2) + np.testing.assert_allclose(centers, [[11.0, 22.0, 3.0]]) + np.testing.assert_allclose(sizes, [[0.45, 0.45, 2.0]]) + assert world_cells(replay, 0, 0)[0].shape == (0, 3) + assert world_cells(replay, 1, 2)[0].shape == (0, 3) + + +def test_tgs_native_serialization_clears_missing_and_coverage_end(tmp_path): + path = tmp_path / "scene.rrd" + recording = rr.RecordingStream("nodedc_mission_core_recorded", recording_id="test") + calls = [] + native_log = recording.log + + def capture(entity, value): + calls.append((entity, value)) + native_log(entity, value) + + recording.log = capture + recording.set_sinks(rr.FileSink(path, write_footer=True)) + log_tgs(recording, data(), 2.1) + recording.flush() + recording.disconnect() + assert path.read_bytes().startswith(b"RRF2") + assert sum(isinstance(value, rr.Boxes3D) for _, value in calls) == 1 + assert all(isinstance(value, rr.Clear) for _, value in calls[3:]) + assert calls[-1][0] == "/world/costmap" + assert not any(entity.startswith("/world/points") for entity, _ in calls) + + +def test_coverage_cannot_end_before_last_tgs_sample(): + with pytest.raises(PortableReplayError): + log_tgs(SimpleNamespace(), data(), 2.0) + + +def test_artifact_corruption_and_symlinks_are_rejected(tmp_path): + path = tmp_path / "artifact" + path.write_bytes(b"first") + sha = hashlib.sha256(b"first").hexdigest() + assert verified_file(path, sha, 5) == path + path.write_bytes(b"other") + with pytest.raises(PortableReplayError): + verified_file(path, sha, 5) + link = tmp_path / "link" + link.symlink_to(path) + with pytest.raises(PortableReplayError): + verified_file(link, hashlib.sha256(b"other").hexdigest(), 5) + + +def test_camera_proxy_keeps_demux_time_base_and_reuses_its_own_sealed_cache(tmp_path, monkeypatch): + import subprocess + + source = tmp_path / "init.mp4" + source.write_bytes(b"synthetic") + epoch = SimpleNamespace( + init_path=source, + init_sha256=hashlib.sha256(b"synthetic").hexdigest(), + init_byte_length=9, + segments=[], + ) + service = PortableReplayService( + view=None, + data_dir=tmp_path, + media=None, + recording_source=lambda _: None, + ffmpeg_path=Path("ffmpeg"), + ) + service.cache.mkdir() + staging = tmp_path / "staging" + staging.mkdir() + calls = [] + + def run(command, **kwargs): + calls.append(command) + assert command[command.index("-enc_time_base") + 1] == "-1" + assert command[command.index("-fps_mode") + 1] == "passthrough" + assert "-r" not in command and "-frames:v" not in command + Path(command[-1]).write_bytes(b"synthetic-proxy") + return SimpleNamespace(returncode=0) + + monkeypatch.setattr(subprocess, "run", run) + first = service._video(epoch, staging) + assert first.parent == service.cache + assert service._video(epoch, staging) == first + assert len(calls) == 1 + + +def test_replay_api_get_never_prepares_and_range_is_generation_fenced(tmp_path): + path = tmp_path / "replay.rrd" + path.write_bytes(b"RRF2test") + sha = hashlib.sha256(path.read_bytes()).hexdigest() + artifact = CanonicalLabReplayArtifact(path, 8, sha) + prepared = [] + + def prepare(*args): + prepared.append(args) + return artifact + + app = FastAPI() + app.include_router( + build_portable_replay_router( + SimpleNamespace(prepare=prepare, cached=lambda *args: artifact if prepared else None) + ) + ) + with TestClient(app) as client: + url = f"/api/v1/observatory/portable-results/{RESULT}/replays/{BASE}/recording.rrd" + assert client.get(url, params={"generation": sha}).status_code == 409 + assert not prepared + response = client.head(url) + assert response.status_code == 200 and response.headers["etag"] == f'"{sha}"' + assert response.headers["x-rerun-format"] == "RRF2" + response = client.get(url, params={"generation": sha}, headers={"Range": "bytes=0-3"}) + assert response.status_code == 206 and response.content == b"RRF2" + assert client.get(url, params={"generation": "c" * 64}).status_code == 412 + assert len(prepared) == 1 + + +def test_warm_reopen_does_not_prepare_and_detects_same_size_corruption(tmp_path, monkeypatch): + import json + + service = PortableReplayService( + view=SimpleNamespace(read=lambda _: {"artifact_manifest_id": "d" * 64}), + data_dir=tmp_path, + media=None, + recording_source=lambda _: None, + ffmpeg_path=Path("ffmpeg"), + ) + key, _ = service._key(RESULT, BASE) + service.cache.mkdir() + path = service.cache / ("e" * 64 + ".replay.rrd") + path.write_bytes(b"RRF2test") + sha = hashlib.sha256(path.read_bytes()).hexdigest() + (service.cache / f"{key}.json").write_text( + json.dumps({"key": key, "file": path.name, "sha256": sha, "byte_length": 8}) + ) + assert service.prepare(RESULT, BASE).sha256 == sha + monkeypatch.setattr(hashlib, "file_digest", lambda *_: pytest.fail("warm byte rehash")) + assert service.prepare(RESULT, BASE).sha256 == sha + monkeypatch.undo() + path.write_bytes(b"RRF2fail") + with pytest.raises(PortableReplayError): + service.cached(RESULT, BASE)