Files
NODEDC_MISSION_CORE/tests/test_session_api.py
T

2888 lines
106 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
from __future__ import annotations
import asyncio
import hashlib
import json
import stat
import threading
import time
from collections.abc import Callable
from dataclasses import replace
from pathlib import Path
from typing import Any
import pytest
from fastapi import APIRouter, FastAPI, HTTPException, Request, Response
from fastapi.responses import FileResponse
from fastapi.routing import APIRoute
from fastapi.testclient import TestClient
import k1link.sessions.media as recorded_media_module
import k1link.web.session_api as session_api_module
from k1link.compute import RecordedPerceptionOverlayArtifact, RecordedPerceptionVideo
from k1link.device_plugins.xgrids_k1 import xgrids_k1_archive_source
from k1link.device_plugins.xgrids_k1.mqtt.capture import FRAME_HEADER, RAW_MAGIC
from k1link.sessions import (
LabReplayCapability,
MaterializedRecording,
RecordedMediaInspector,
RecordedMediaManifest,
ReplayCommand,
SessionIntegrityError,
SessionNotFoundError,
SessionRecordingMaterializer,
SessionRecordingPreparationManager,
SessionStore,
)
from k1link.sessions.media import _manifest_generation_sha256
from k1link.viewer.recorded import recorded_blueprint_rrd
from k1link.viewer.rerun_bridge import RerunSceneSettings
from k1link.web.camera_archive import CameraArchiveWriter
from k1link.web.session_api import (
LayoutPutRequest,
RecordedBlueprintRequest,
RecordedPerceptionRequest,
RecordedPointColorsRequest,
ReplayRequest,
_recorded_media_manifest_document,
build_session_router,
)
def lab_method() -> dict[str, object]:
return {
"schema_version": "missioncore.laboratory-method/v1",
"completeness": "complete",
"execution_class": "deterministic",
"pipeline_id": "test-pipeline/v1",
"components": [
{
"kind": "algorithm",
"name": "test algorithm",
"version": "v1",
"role": "contract fixture",
"identity_sha256": "9" * 64,
}
],
}
def make_legacy_session(sessions_root: Path, session_id: str) -> Path:
session = sessions_root / session_id
capture = session / "captures" / "mqtt_live"
capture.mkdir(parents=True)
topic_counts = {
"lixel/application/report/lio_pcl": 1,
"lixel/application/report/lio_pose": 1,
}
raw = bytearray(RAW_MAGIC)
metadata: list[dict[str, object]] = []
for sequence, topic in enumerate(topic_counts, start=1):
topic_bytes = topic.encode()
payload = b"fixture"
frame_offset = len(raw)
raw.extend(FRAME_HEADER.pack(len(topic_bytes), len(payload)))
raw.extend(topic_bytes)
raw.extend(payload)
metadata.append(
{
"record_type": "message",
"sequence": sequence,
"received_at_epoch_ns": 1_000_000_000 + sequence - 1,
"received_monotonic_ns": 2_000_000_000 + sequence - 1,
"topic": topic,
"payload_bytes": len(payload),
"raw_frame_offset": frame_offset,
"raw_payload_offset": frame_offset + FRAME_HEADER.size + len(topic_bytes),
"raw_frame_bytes": FRAME_HEADER.size + len(topic_bytes) + len(payload),
}
)
raw_path = capture / "mqtt.raw.k1mqtt"
raw_path.write_bytes(raw)
metadata_path = capture / "mqtt.metadata.jsonl"
metadata_path.write_text(
"".join(json.dumps(record) + "\n" for record in metadata),
encoding="utf-8",
)
(capture / "mqtt.summary.json").write_text(
json.dumps(
{
"created_at_utc": "2026-07-16T20:56:32.699Z",
"completed_at_utc": "2026-07-16T21:20:43.018Z",
"capture_elapsed_seconds": 1440.1,
"stop_reason": "external_stop",
"error": None,
"message_count": 2,
"raw_bytes": len(raw),
"topic_counts": topic_counts,
"artifact_hashes": {
"raw_sha256": hashlib.sha256(raw).hexdigest(),
"metadata_jsonl_sha256": hashlib.sha256(metadata_path.read_bytes()).hexdigest(),
},
}
),
encoding="utf-8",
)
(session / "manifest.redacted.json").write_text("{}", encoding="utf-8")
return session
def make_recorded_h264_fixture(
*,
timescale: int = 1_000,
sample_duration: int = 500,
base_decode_time: int = 0,
trex_default_sample_flags: int = 0,
tfhd_default_sample_flags: int | None = None,
first_sample_flags: int | None = None,
sample_flags: int | None = None,
fragment_sample_count: int = 1,
) -> tuple[bytes, bytes]:
def box(box_type: bytes, payload: bytes = b"") -> bytes:
return (8 + len(payload)).to_bytes(4, "big") + box_type + payload
def full_box(
box_type: bytes,
payload: bytes = b"",
*,
flags: int = 0,
) -> bytes:
return box(box_type, bytes([0]) + flags.to_bytes(3, "big") + payload)
track_id = 1
tkhd = full_box(
b"tkhd",
b"\x00" * 8 + track_id.to_bytes(4, "big") + b"\x00" * 4,
)
mdhd = full_box(
b"mdhd",
b"\x00" * 8 + timescale.to_bytes(4, "big") + b"\x00" * 4,
)
hdlr = full_box(b"hdlr", b"\x00" * 4 + b"vide")
trak = box(b"trak", tkhd + box(b"mdia", mdhd + hdlr))
trex = full_box(
b"trex",
track_id.to_bytes(4, "big")
+ (1).to_bytes(4, "big")
+ sample_duration.to_bytes(4, "big")
+ b"\x00" * 4
+ trex_default_sample_flags.to_bytes(4, "big"),
)
avcc = box(b"avcC", b"\x01\x64\x00\x28")
init = box(b"ftyp", b"isom") + box(b"moov", trak + box(b"mvex", trex) + avcc)
tfhd_flags = 0x020000
tfhd_payload = track_id.to_bytes(4, "big")
if tfhd_default_sample_flags is not None:
tfhd_flags |= 0x000020
tfhd_payload += tfhd_default_sample_flags.to_bytes(4, "big")
tfhd = full_box(b"tfhd", tfhd_payload, flags=tfhd_flags)
tfdt = full_box(b"tfdt", base_decode_time.to_bytes(4, "big"))
trun_flags = 0
trun_payload = fragment_sample_count.to_bytes(4, "big")
if first_sample_flags is not None:
trun_flags |= 0x000004
trun_payload += first_sample_flags.to_bytes(4, "big")
if sample_flags is not None:
trun_flags |= 0x000400
trun_payload += sample_flags.to_bytes(4, "big") * fragment_sample_count
trun = full_box(
b"trun",
trun_payload,
flags=trun_flags,
)
fragment = box(b"moof", box(b"traf", tfhd + tfdt + trun)) + box(b"mdat", b"frame")
return init, fragment
def endpoint(router: APIRouter, path: str, method: str) -> Callable[..., Any]:
routes = router.routes
return next(
route.endpoint
for route in routes
if isinstance(route, APIRoute) and route.path == path and method in route.methods
)
def assert_no_local_paths(value: object, root: Path) -> None:
serialized = json.dumps(value, ensure_ascii=False, default=str)
assert str(root) not in serialized
assert "mqtt.raw.k1mqtt" not in serialized
async def render_file_response(
response: Response,
*,
range_header: bytes | None,
) -> tuple[dict[str, Any], bytes]:
body_iterator = getattr(response, "body_iterator", None)
if body_iterator is not None:
chunks = [chunk async for chunk in body_iterator]
return {
"status": response.status_code,
"headers": response.raw_headers,
}, b"".join(chunk.encode("utf-8") if isinstance(chunk, str) else chunk for chunk in chunks)
sent: list[dict[str, Any]] = []
request_delivered = False
async def receive() -> dict[str, Any]:
nonlocal request_delivered
if request_delivered:
return {"type": "http.disconnect"}
request_delivered = True
return {"type": "http.request", "body": b"", "more_body": False}
async def send(message: dict[str, Any]) -> None:
sent.append(message)
scope: dict[str, Any] = {
"type": "http",
"asgi": {"version": "3.0"},
"http_version": "1.1",
"method": "GET",
"scheme": "http",
"path": "/recording.rrd",
"raw_path": b"/recording.rrd",
"query_string": b"",
"headers": [] if range_header is None else [(b"range", range_header)],
"client": ("127.0.0.1", 1),
"server": ("127.0.0.1", 80),
}
await response(scope, receive, send)
start = next(message for message in sent if message["type"] == "http.response.start")
body = b"".join(
message.get("body", b"") for message in sent if message["type"] == "http.response.body"
)
return start, body
def test_session_router_lists_details_and_dispatches_opaque_replay(tmp_path: Path) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
observed: list[ReplayCommand] = []
async def launch(command: ReplayCommand) -> dict[str, Any]:
observed.append(command)
return {"source_mode": "replay", "phase": "replay"}
router = build_session_router(store, replay_launcher=launch)
list_route = endpoint(router, "/api/v1/observation-sessions", "GET")
detail_route = endpoint(router, "/api/v1/observation-sessions/{session_id}", "GET")
replay_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/replay",
"POST",
)
listing = list_route(limit=20, cursor=None)
detail = detail_route(session_id=session.name)
replay = asyncio.run(replay_route(session_id=session.name, request=None))
assert set(listing) == {"items"}
assert listing["items"][0] == {
"id": session.name,
"label": session.name,
"started_at_utc": "2026-07-16T20:56:32.699Z",
"completed_at_utc": "2026-07-16T21:20:43.018Z",
"status": "ready",
"modalities": ["point-cloud", "trajectory"],
"duration_seconds": 1440.1,
"replayable": True,
}
assert detail["timeline"]["seekable"] is True
assert detail["modalities"] == ["point-cloud", "trajectory"]
assert replay["launch"] == {
"kind": "plugin-action",
"plugin_id": "nodedc.device.xgrids-lixelkity-k1",
"action_id": "stream.start-replay",
"session_id": session.name,
"speed": 1.0,
"loop": False,
"dispatched": True,
}
assert replay["runtime_state"]["source_mode"] == "replay"
assert observed[0].primary_artifact.path.name == "mqtt.raw.k1mqtt"
for value in (listing, detail, replay):
assert_no_local_paths(value, repository)
def test_session_router_exposes_immutable_lab_provenance(tmp_path: Path) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
source = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
binding = store.publish_lab_instance(
session_id="lab-e21-d0201712",
source_session_id=source.name,
display_name="LAB E21 · Real-time 1× · d0201712",
lab_id="LAB E21",
result_kind="e21-realtime-envelope",
result_id="e10-integrated-perception-" + "a" * 64,
source_result_id="e21-realtime-envelope-" + "b" * 64,
config_sha256="c" * 64,
run_created_at_utc="2026-07-23T15:55:15.548Z",
provenance={
"source_payloads_mutated": False,
"method": lab_method(),
},
)
router = build_session_router(store)
list_route = endpoint(router, "/api/v1/observation-sessions", "GET")
detail_route = endpoint(router, "/api/v1/observation-sessions/{session_id}", "GET")
listing = list_route(limit=20, cursor=None, scope="all")
item = next(value for value in listing["items"] if value["id"] == binding.session_id)
source_listing = list_route(limit=20, cursor=None, scope="source")
laboratory_listing = list_route(limit=20, cursor=None, scope="laboratory")
detail = detail_route(session_id=binding.session_id)
assert item["lab"] == binding.as_dict()
assert [value["id"] for value in source_listing["items"]] == [source.name]
assert [value["id"] for value in laboratory_listing["items"]] == [binding.session_id]
assert detail["lab"] == binding.as_dict()
assert item["lab"]["source_session_id"] == source.name
assert item["lab"]["provenance"]["source_payloads_mutated"] is False
assert_no_local_paths((item, detail), repository)
def test_session_router_rolls_capability_projections_out_only_in_v2(
tmp_path: Path,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
source = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
legacy = store.publish_lab_instance(
session_id="lab-e21-legacy",
source_session_id=source.name,
display_name="LAB E21 · legacy",
lab_id="LAB E21",
result_kind="e21-realtime-envelope",
result_id="e21-realtime-envelope-" + "1" * 64,
run_created_at_utc="2026-07-23T15:55:15.548Z",
provenance={"method": lab_method()},
)
capability = LabReplayCapability(
schema_version="missioncore.observation-lab-replay-capability/v1",
kind="canonical-recorded-rerun",
viewer_profile="recorded-session",
timeline="session_time",
activation="explicit",
commands_enabled=False,
)
canonical = store.publish_lab_instance(
session_id="lab-v1-vegetation-shadow-" + "8" * 64,
source_session_id=source.name,
display_name="RAV004 · recorded",
lab_id="LAB V1",
result_kind="recorded-perception-qualification",
result_id="lab-v1-vegetation-shadow-" + "8" * 64,
run_created_at_utc="2026-08-29T18:05:11.329061+00:00",
replay_capability=capability,
provenance={
"replay_capability": capability.as_dict(),
"method": lab_method(),
},
)
router = build_session_router(store)
list_route = endpoint(router, "/api/v1/observation-sessions", "GET")
default_items = list_route(limit=20, cursor=None, scope="all")["items"]
assert {item["id"] for item in default_items} == {source.name, legacy.session_id}
default_legacy = next(item for item in default_items if item["id"] == legacy.session_id)
assert "replay_capability" not in default_legacy["lab"]
v1_labs = list_route(
limit=20,
cursor=None,
scope="laboratory",
lab_contract="v1",
)["items"]
assert [item["id"] for item in v1_labs] == [legacy.session_id]
v2_labs = list_route(
limit=20,
cursor=None,
scope="laboratory",
lab_contract="v2",
)["items"]
by_id = {item["id"]: item for item in v2_labs}
assert set(by_id) == {legacy.session_id, canonical.session_id}
assert by_id[legacy.session_id]["lab"]["replay_capability"] is None
assert by_id[canonical.session_id]["lab"]["replay_capability"] == capability.as_dict()
application = FastAPI()
application.include_router(router)
assert TestClient(application).get(
"/api/v1/observation-sessions?lab_contract=v3"
).status_code == 422
def test_delete_session_removes_evidence_and_cache_but_refuses_an_open_recording(
tmp_path: Path,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
def export_recording(source: Path, destination: Path) -> dict[str, object]:
payload = b"deletable-recording"
destination.write_bytes(payload)
return {
"source_sha256": hashlib.sha256(source.read_bytes()).hexdigest(),
"rrd_sha256": hashlib.sha256(payload).hexdigest(),
"rrd_bytes": len(payload),
"timeline": "session_time",
"timeline_start_ns": 0,
"timeline_end_ns": 1,
}
materializer = SessionRecordingMaterializer(
store.data_dir,
exporter=export_recording,
)
recording, release = materializer.materialize_pinned(store.prepare_replay(session.name))
router = build_session_router(store, recording_materializer=materializer)
delete_route = endpoint(router, "/api/v1/observation-sessions/{session_id}", "DELETE")
with pytest.raises(HTTPException) as conflict:
asyncio.run(delete_route(session_id=session.name))
assert conflict.value.status_code == 409
assert session.is_dir()
assert recording.path.is_file()
release()
response = asyncio.run(delete_route(session_id=session.name))
repeated_response = asyncio.run(delete_route(session_id=session.name))
assert response.status_code == 204
assert repeated_response.status_code == 204
assert not session.exists()
assert not recording.path.parent.exists()
with pytest.raises(SessionNotFoundError):
store.get_session(session.name)
def test_completed_recording_get_does_not_hold_delete_for_launch_lease(
tmp_path: Path,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
payload = b"recording-opened-and-closed"
def export_recording(source: Path, destination: Path) -> dict[str, object]:
destination.write_bytes(payload)
return {
"source_sha256": hashlib.sha256(source.read_bytes()).hexdigest(),
"rrd_sha256": hashlib.sha256(payload).hexdigest(),
"rrd_bytes": len(payload),
"timeline": "session_time",
"timeline_start_ns": 0,
"timeline_end_ns": 1_000_000_000,
}
materializer = SessionRecordingMaterializer(store.data_dir, exporter=export_recording)
command = store.prepare_replay(session.name)
materializer.materialize(command)
manager = SessionRecordingPreparationManager(materializer)
resolved = manager.resolve_cached(command)
assert resolved is not None and resolved.recording is not None
router = build_session_router(
store,
recording_materializer=materializer,
recording_preparation_manager=manager,
)
replay_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/replay",
"POST",
)
recording_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/recording.rrd",
"GET",
)
delete_route = endpoint(router, "/api/v1/observation-sessions/{session_id}", "DELETE")
try:
replay = asyncio.run(replay_route(session_id=session.name, request=None))
generation = replay["launch"]["sha256"]
file_response = asyncio.run(
recording_route(
session_id=session.name,
generation=generation,
if_match=None,
if_none_match=None,
)
)
_, body = asyncio.run(render_file_response(file_response, range_header=None))
assert body == payload
deleted = asyncio.run(delete_route(session_id=session.name))
assert deleted.status_code == 204
assert not session.exists()
finally:
manager.close()
def test_canonical_lab_spatial_frame_uses_ready_immutable_recording(
tmp_path: Path,
monkeypatch: pytest.MonkeyPatch,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
payload = b"sealed-spatial-recording"
def export_recording(source: Path, destination: Path) -> dict[str, object]:
destination.write_bytes(payload)
return {
"source_sha256": hashlib.sha256(source.read_bytes()).hexdigest(),
"rrd_sha256": hashlib.sha256(payload).hexdigest(),
"rrd_bytes": len(payload),
"timeline": "session_time",
"timeline_start_ns": 0,
"timeline_end_ns": 1_000_000_000,
}
materializer = SessionRecordingMaterializer(store.data_dir, exporter=export_recording)
command = store.prepare_replay(session.name)
recording = materializer.materialize(command)
manager = SessionRecordingPreparationManager(materializer)
resolved = manager.resolve_cached(command)
assert resolved is not None and resolved.recording is not None
generation = hashlib.sha256(payload).hexdigest()
expected = {
"schema_version": "missioncore.canonical-recorded-lab-spatial-frame/v3",
"target_time_ns": 500_000_000,
"source_time_ns": 499_000_000,
"pose_time_ns": 499_000_000,
"trajectory_time_ns": 490_000_000,
"coordinate_frame": "body-ground",
"sensor_height": {
"meters": 0.32,
"source": "local-source-cloud-ground-quantile-median",
"sample_count": 20,
"mad_m": 0.03,
"authority": "visual-derived",
},
"spatial_profile": {
"profile_id": "source-paced-ground-v3",
"local_slam_history_seconds": 5.0,
"local_slam_radius_m": 30.0,
"local_slam_voxel_size_m": 0.12,
"local_slam_point_limit": 27000,
},
"body_frame": {
"origin_map_xyz_m": [0.0, 0.0, 0.0],
"sensor_origin_map_xyz_m": [0.0, 0.0, 0.32],
"basis_map_from_body": [[1.0, 0.0, 0.0], [0.0, 1.0, 0.0], [0.0, 0.0, 1.0]],
},
"source_point_count": 1,
"source_points_body_xyz_m": [[1.0, 2.0, 3.0]],
"local_slam_source_frame_count": 1,
"local_slam_source_point_count": 1,
"local_slam_point_count": 1,
"local_slam_body_xyz_m": [[0.0, 0.0, 0.0]],
}
def spatial_frame(path: Path, sha256: str, time_ns: int) -> dict[str, object]:
assert path == recording.path
assert sha256 == generation
assert time_ns == 500_000_000
return expected
monkeypatch.setattr(session_api_module, "canonical_lab_spatial_frame", spatial_frame)
router = build_session_router(
store,
recording_materializer=materializer,
recording_preparation_manager=manager,
)
spatial_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/canonical-lab/spatial-frame",
"GET",
)
try:
response = asyncio.run(spatial_route(
session_id=session.name,
generation=generation,
time_ns=500_000_000,
profile="source-paced-ground-v3",
))
assert json.loads(response.body) == expected
assert response.headers["etag"] == (
f'"{generation}:source-paced-ground-v3:499000000"'
)
assert response.headers["cache-control"].endswith("immutable")
finally:
manager.close()
def test_session_router_returns_seekable_recording_and_serves_byte_ranges(
tmp_path: Path,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
payload = b"RRD-seekable-recording"
export_calls = 0
def export_recording(source: Path, destination: Path) -> dict[str, object]:
nonlocal export_calls
export_calls += 1
destination.write_bytes(payload)
return {
"source_sha256": hashlib.sha256(source.read_bytes()).hexdigest(),
"rrd_sha256": hashlib.sha256(payload).hexdigest(),
"rrd_bytes": len(payload),
"timeline": "session_time",
"timeline_start_ns": 0,
"timeline_end_ns": 2_500_000_000,
}
materializer = SessionRecordingMaterializer(
store.data_dir,
exporter=export_recording,
)
router = build_session_router(
store,
recording_materializer=materializer,
allow_synchronous_recording_fallback=True,
)
replay_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/replay",
"POST",
)
recording_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/recording.rrd",
"GET",
)
replay = asyncio.run(
replay_route(
session_id=session.name,
request=ReplayRequest(speed=2.0, loop=True),
)
)
assert replay == {
"schema_version": "missioncore.observation-session-replay/v2",
"launch": {
"kind": "rerun-recording",
"session_id": session.name,
"source_url": (f"/api/v1/observation-sessions/{session.name}/recording.rrd"),
"viewer_source_url": (
f"/api/v1/observation-sessions/{session.name}/recording.rrd"
f"?generation={hashlib.sha256(payload).hexdigest()}"
),
"media_type": "application/vnd.rerun.rrd",
"timeline": "session_time",
"timeline_start_seconds": 0.0,
"timeline_end_seconds": 2.5,
"seekable": True,
"byte_length": len(payload),
"sha256": hashlib.sha256(payload).hexdigest(),
"playback": {"speed": 2.0, "loop": True},
"media_sources": [],
},
}
assert_no_local_paths(replay, repository)
file_response = asyncio.run(recording_route(session_id=session.name))
response_start, response_body = asyncio.run(
render_file_response(file_response, range_header=b"bytes=4-11")
)
headers = {
key.decode("latin-1"): value.decode("latin-1") for key, value in response_start["headers"]
}
assert response_start["status"] == 206
assert response_body == payload[4:12]
assert headers["accept-ranges"] == "bytes"
assert headers["content-range"] == f"bytes 4-11/{len(payload)}"
assert headers["content-type"] == "application/vnd.rerun.rrd"
assert headers["cache-control"] == "private, no-cache, no-transform"
assert str(repository) not in json.dumps(headers)
not_modified = asyncio.run(
recording_route(
session_id=session.name,
if_none_match=f'W/{headers["etag"]}, "unrelated"',
)
)
assert not_modified.status_code == 304
assert not_modified.body == b""
assert not_modified.headers["etag"] == headers["etag"]
assert not_modified.headers["cache-control"] == "private, no-cache, no-transform"
assert export_calls == 1
def test_cold_replay_returns_quick_202_then_status_returns_ready_launch(
tmp_path: Path,
monkeypatch: pytest.MonkeyPatch,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
started = threading.Event()
release = threading.Event()
calls = 0
payload = b"background-recording"
def export_recording(source: Path, destination: Path) -> dict[str, object]:
nonlocal calls
calls += 1
started.set()
assert release.wait(timeout=2)
destination.write_bytes(payload)
return {
"source_sha256": hashlib.sha256(source.read_bytes()).hexdigest(),
"rrd_sha256": hashlib.sha256(payload).hexdigest(),
"rrd_bytes": len(payload),
"timeline": "session_time",
"timeline_start_ns": 0,
"timeline_end_ns": 2_500_000_000,
}
materializer = SessionRecordingMaterializer(
store.data_dir,
exporter=export_recording,
)
manager = SessionRecordingPreparationManager(materializer)
def forbidden_cold_cache_resolution(*_args: object, **_kwargs: object) -> None:
raise AssertionError("HTTP must not perform cold cache validation")
monkeypatch.setattr(manager, "resolve_cached", forbidden_cold_cache_resolution)
monkeypatch.setattr(manager, "resolve_cached_pinned", forbidden_cold_cache_resolution)
router = build_session_router(
store,
recording_materializer=materializer,
recording_preparation_manager=manager,
)
replay_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/replay",
"POST",
)
status_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/recording-preparation",
"GET",
)
list_route = endpoint(router, "/api/v1/observation-sessions", "GET")
recording_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/recording.rrd",
"GET",
)
try:
with pytest.raises(HTTPException) as missing_generation:
asyncio.run(recording_route(session_id=session.name, if_none_match=None))
assert missing_generation.value.status_code == 428
cold_file_before = time.monotonic()
cold_file = asyncio.run(
recording_route(
session_id=session.name,
if_match='"sha256:previous-generation"',
if_none_match=None,
)
)
assert time.monotonic() - cold_file_before < 0.1
assert cold_file.status_code == 202
cold_preparation_id = json.loads(cold_file.body)["preparation"]["preparation_id"]
before = time.monotonic()
pending = asyncio.run(
replay_route(
session_id=session.name,
request=ReplayRequest(speed=2.0, loop=True),
)
)
assert time.monotonic() - before < 0.1
assert pending.status_code == 202
document = json.loads(pending.body)
preparation = document["preparation"]
assert document["schema_version"] == "missioncore.observation-session-preparation/v1"
assert preparation["preparation_id"] == cold_preparation_id
assert preparation["state"] in {"queued", "validating", "exporting"}
assert preparation["status_url"].endswith("/recording-preparation")
assert pending.headers["location"] == preparation["status_url"]
assert pending.headers["retry-after"] == "1"
preparation_etag = pending.headers["etag"]
assert preparation_etag == f'"{preparation["preparation_id"]}"'
assert started.wait(timeout=1)
catalog_preparation = list_route(limit=20, cursor=None)["items"][0]["preparation"]
assert catalog_preparation["preparation_id"] == preparation["preparation_id"]
assert catalog_preparation["state"] in {"validating", "exporting"}
duplicate = asyncio.run(
replay_route(
session_id=session.name,
request=ReplayRequest(speed=2.0, loop=True),
)
)
assert (
json.loads(duplicate.body)["preparation"]["preparation_id"]
== preparation["preparation_id"]
)
release.set()
deadline = time.monotonic() + 2
ready: dict[str, Any] | None = None
while time.monotonic() < deadline:
candidate = asyncio.run(
status_route(
session_id=session.name,
if_match=preparation_etag,
)
)
if candidate.status_code == 200:
assert candidate.headers["etag"] == preparation_etag
ready = json.loads(candidate.body)
break
time.sleep(0.005)
assert ready is not None
assert ready["launch"]["source_url"].endswith("/recording.rrd")
assert ready["launch"]["viewer_source_url"] == (
f"{ready['launch']['source_url']}?generation={ready['launch']['sha256']}"
)
assert ready["launch"]["playback"] == {"speed": 1.0, "loop": False}
recording_etag = f'"sha256:{ready["launch"]["sha256"]}"'
with pytest.raises(HTTPException) as stale_generation:
asyncio.run(
recording_route(
session_id=session.name,
if_match='"sha256:stale"',
if_none_match=None,
)
)
assert stale_generation.value.status_code == 412
with pytest.raises(HTTPException) as stale_query_generation:
asyncio.run(
recording_route(
session_id=session.name,
generation="0" * 64,
if_none_match=None,
)
)
assert stale_query_generation.value.status_code == 412
with pytest.raises(HTTPException) as malformed_query_generation:
asyncio.run(
recording_route(
session_id=session.name,
generation="NOT-A-CANONICAL-DIGEST",
if_none_match=None,
)
)
assert malformed_query_generation.value.status_code == 412
recording_response = asyncio.run(
recording_route(
session_id=session.name,
if_match=recording_etag,
if_none_match=None,
)
)
recording_start, recording_body = asyncio.run(
render_file_response(recording_response, range_header=None)
)
recording_headers = {
key.decode("latin-1"): value.decode("latin-1")
for key, value in recording_start["headers"]
}
assert recording_start["status"] == 200
assert recording_body == payload
assert recording_headers["etag"] == recording_etag
assert recording_headers["content-length"] == str(len(payload))
assert recording_headers["content-type"] == "application/vnd.rerun.rrd"
assert recording_headers["cache-control"] == "private, no-cache, no-transform"
viewer_response = asyncio.run(
recording_route(
session_id=session.name,
generation=ready["launch"]["sha256"],
if_none_match=None,
)
)
viewer_start, viewer_body = asyncio.run(
render_file_response(viewer_response, range_header=None)
)
viewer_headers = {
key.decode("latin-1"): value.decode("latin-1") for key, value in viewer_start["headers"]
}
assert viewer_start["status"] == 200
assert viewer_body == payload
assert viewer_headers["etag"] == recording_etag
assert viewer_headers["content-length"] == str(len(payload))
assert viewer_headers["content-type"] == "application/vnd.rerun.rrd"
assert viewer_headers["cache-control"] == (
"private, max-age=31536000, immutable, no-transform"
)
viewer_not_modified = asyncio.run(
recording_route(
session_id=session.name,
generation=ready["launch"]["sha256"],
if_none_match=recording_etag,
)
)
assert viewer_not_modified.status_code == 304
assert viewer_not_modified.headers["etag"] == recording_etag
assert viewer_not_modified.headers["cache-control"] == (
"private, max-age=31536000, immutable, no-transform"
)
assert calls == 1
finally:
release.set()
manager.close()
def test_replay_post_retries_terminal_camera_finalization_failure(
tmp_path: Path,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
payload = b"recording-ready-before-camera-archive"
export_calls = 0
finalization_calls = 0
def export_recording(source: Path, destination: Path) -> dict[str, object]:
nonlocal export_calls
export_calls += 1
destination.write_bytes(payload)
return {
"source_sha256": hashlib.sha256(source.read_bytes()).hexdigest(),
"rrd_sha256": hashlib.sha256(payload).hexdigest(),
"rrd_bytes": len(payload),
"timeline": "session_time",
"timeline_start_ns": 0,
"timeline_end_ns": 2_500_000_000,
}
def prepare_recorded_media(
_command: ReplayCommand,
_recording: MaterializedRecording,
) -> tuple[RecordedMediaManifest, ...]:
nonlocal finalization_calls
finalization_calls += 1
if finalization_calls == 1:
raise SessionIntegrityError("camera archive is still sealing")
return ()
materializer = SessionRecordingMaterializer(
store.data_dir,
exporter=export_recording,
)
manager = SessionRecordingPreparationManager(
materializer,
ready_preparer=prepare_recorded_media,
)
router = build_session_router(
store,
recording_materializer=materializer,
recording_preparation_manager=manager,
)
replay_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/replay",
"POST",
)
try:
first = asyncio.run(replay_route(session_id=session.name, request=None))
assert first.status_code == 202
deadline = time.monotonic() + 2
failed = manager.status(session.name)
while failed is None or failed.state != "failed":
assert time.monotonic() < deadline
time.sleep(0.005)
failed = manager.status(session.name)
failed_preparation_id = failed.preparation_id
retry = asyncio.run(replay_route(session_id=session.name, request=None))
assert retry.status_code == 202
retry_document = json.loads(retry.body)
assert retry_document["preparation"]["preparation_id"] != failed_preparation_id
deadline = time.monotonic() + 2
ready = manager.status(session.name)
while ready is None or ready.state != "ready":
assert time.monotonic() < deadline
time.sleep(0.005)
ready = manager.status(session.name)
launch = asyncio.run(replay_route(session_id=session.name, request=None))
assert launch["launch"]["session_id"] == session.name
assert export_calls == 1
assert finalization_calls == 2
finally:
manager.close()
def test_catalog_read_never_prepares_a_cold_historical_session(
tmp_path: Path,
monkeypatch: pytest.MonkeyPatch,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
calls = 0
def exporter(_source: Path, _destination: Path) -> dict[str, object]:
nonlocal calls
calls += 1
raise AssertionError("catalog reads must not invoke the exporter")
materializer = SessionRecordingMaterializer(store.data_dir, exporter=exporter)
manager = SessionRecordingPreparationManager(materializer)
def forbidden_restore(*_args: object, **_kwargs: object) -> None:
raise AssertionError("catalog reads must not restore published recordings")
monkeypatch.setattr(manager, "restore_published", forbidden_restore)
router = build_session_router(
store,
recording_materializer=materializer,
recording_preparation_manager=manager,
)
list_route = endpoint(router, "/api/v1/observation-sessions", "GET")
try:
item = list_route(limit=20, cursor=None)["items"][0]
time.sleep(0.02)
assert item["id"] == session.name
assert item["preparation"] is None
assert manager.status(session.name) is None
assert calls == 0
finally:
manager.close()
def test_recording_response_releases_pin_once_when_asgi_send_fails(
tmp_path: Path,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
recording_path = tmp_path / "pinned.rrd"
payload = b"RRD-pinned-response"
recording_path.write_bytes(payload)
recording = MaterializedRecording(
session_id=session.name,
path=recording_path,
media_type="application/vnd.rerun.rrd",
byte_length=len(payload),
sha256=hashlib.sha256(payload).hexdigest(),
source_sha256="source-digest",
timeline="session_time",
timeline_start_ns=0,
timeline_end_ns=1_000_000_000,
)
class PinnedMaterializer:
def __init__(self) -> None:
self.release_calls = 0
def __call__(self, _command: ReplayCommand) -> MaterializedRecording:
return recording
def materialize_pinned(
self,
_command: ReplayCommand,
) -> tuple[MaterializedRecording, Callable[[], None]]:
return recording, self.release
def release(self) -> None:
self.release_calls += 1
materializer = PinnedMaterializer()
router = build_session_router(
store,
recording_materializer=materializer,
allow_synchronous_recording_fallback=True,
)
recording_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/recording.rrd",
"GET",
)
response = asyncio.run(recording_route(session_id=session.name))
async def send_failure() -> None:
async def receive() -> dict[str, Any]:
return {"type": "http.request", "body": b"", "more_body": False}
async def send(message: dict[str, Any]) -> None:
if message["type"] == "http.response.body":
raise RuntimeError("synthetic ASGI send failure")
scope: dict[str, Any] = {
"type": "http",
"asgi": {"version": "3.0"},
"http_version": "1.1",
"method": "GET",
"scheme": "http",
"path": "/recording.rrd",
"raw_path": b"/recording.rrd",
"query_string": b"",
"headers": [],
"client": ("127.0.0.1", 1),
"server": ("127.0.0.1", 80),
}
await response(scope, receive, send)
with pytest.raises(RuntimeError, match="synthetic ASGI send failure"):
asyncio.run(send_failure())
assert materializer.release_calls == 1
# A response must never double-release the same cache lease, even if an
# ASGI host accidentally invokes the response object again after failure.
with pytest.raises(RuntimeError, match="synthetic ASGI send failure"):
asyncio.run(send_failure())
assert materializer.release_calls == 1
def test_preparation_status_and_cancel_are_bound_to_exact_job_etag(tmp_path: Path) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
started = threading.Event()
release = threading.Event()
def exporter(source: Path, destination: Path) -> dict[str, object]:
started.set()
assert release.wait(timeout=2)
payload = b"recording"
destination.write_bytes(payload)
return {
"source_sha256": hashlib.sha256(source.read_bytes()).hexdigest(),
"rrd_sha256": hashlib.sha256(payload).hexdigest(),
"rrd_bytes": len(payload),
"timeline": "session_time",
"timeline_start_ns": 0,
"timeline_end_ns": 1,
}
materializer = SessionRecordingMaterializer(store.data_dir, exporter=exporter)
manager = SessionRecordingPreparationManager(materializer)
router = build_session_router(
store,
recording_materializer=materializer,
recording_preparation_manager=manager,
)
replay_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/replay",
"POST",
)
status_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/recording-preparation",
"GET",
)
cancel_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/recording-preparation",
"DELETE",
)
try:
pending = asyncio.run(replay_route(session_id=session.name, request=None))
assert pending.status_code == 202
assert started.wait(timeout=1)
etag = pending.headers["etag"]
with pytest.raises(HTTPException) as stale_status:
asyncio.run(status_route(session_id=session.name, if_match='"stale-job"'))
assert stale_status.value.status_code == 412
with pytest.raises(HTTPException) as stale_cancel:
cancel_route(session_id=session.name, if_match='"stale-job"')
assert stale_cancel.value.status_code == 412
current = manager.status(session.name)
assert current is not None
assert f'"{current.preparation_id}"' == etag
assert current.state in {"validating", "exporting"}
assert cancel_route(session_id=session.name, if_match=etag).status_code == 204
release.set()
deadline = time.monotonic() + 2
while time.monotonic() < deadline:
current = manager.status(session.name)
if current is not None and current.state == "cancelled":
break
time.sleep(0.005)
assert current is not None and current.state == "cancelled"
finally:
release.set()
manager.close()
def test_stale_status_poll_does_not_create_a_replacement_job(tmp_path: Path) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
calls = 0
def exporter(_source: Path, _destination: Path) -> dict[str, object]:
nonlocal calls
calls += 1
raise AssertionError("a stale status poll must not enqueue work")
materializer = SessionRecordingMaterializer(store.data_dir, exporter=exporter)
manager = SessionRecordingPreparationManager(materializer)
router = build_session_router(
store,
recording_materializer=materializer,
recording_preparation_manager=manager,
)
status_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/recording-preparation",
"GET",
)
try:
with pytest.raises(HTTPException) as stale:
asyncio.run(
status_route(
session_id=session.name,
if_match='"no-longer-current"',
)
)
assert stale.value.status_code == 412
assert manager.status(session.name) is None
assert calls == 0
finally:
manager.close()
def test_ready_post_uses_each_callers_playback_without_mutating_shared_job(
tmp_path: Path,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
calls = 0
def exporter(source: Path, destination: Path) -> dict[str, object]:
nonlocal calls
calls += 1
payload = b"recording"
destination.write_bytes(payload)
return {
"source_sha256": hashlib.sha256(source.read_bytes()).hexdigest(),
"rrd_sha256": hashlib.sha256(payload).hexdigest(),
"rrd_bytes": len(payload),
"timeline": "session_time",
"timeline_start_ns": 0,
"timeline_end_ns": 1,
}
materializer = SessionRecordingMaterializer(store.data_dir, exporter=exporter)
manager = SessionRecordingPreparationManager(materializer)
command = store.prepare_replay(session.name)
manager.enqueue(command)
deadline = time.monotonic() + 2
while time.monotonic() < deadline:
ready = manager.status(session.name)
if ready is not None and ready.state == "ready":
break
time.sleep(0.005)
router = build_session_router(
store,
recording_materializer=materializer,
recording_preparation_manager=manager,
)
replay_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/replay",
"POST",
)
try:
first = asyncio.run(
replay_route(
session_id=session.name,
request=ReplayRequest(speed=2.0, loop=True),
)
)
second = asyncio.run(
replay_route(
session_id=session.name,
request=ReplayRequest(speed=7.0, loop=False),
)
)
assert first["launch"]["playback"] == {"speed": 2.0, "loop": True}
assert second["launch"]["playback"] == {"speed": 7.0, "loop": False}
shared = manager.status(session.name)
assert shared is not None
assert shared.command.speed == 1.0
assert shared.command.loop is False
assert calls == 1
finally:
manager.close()
def test_production_router_never_materializes_recording_inline(tmp_path: Path) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
calls = 0
def exporter(_source: Path, _destination: Path) -> dict[str, object]:
nonlocal calls
calls += 1
raise AssertionError("public route must not export inline")
materializer = SessionRecordingMaterializer(store.data_dir, exporter=exporter)
router = build_session_router(store, recording_materializer=materializer)
replay_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/replay",
"POST",
)
recording_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/recording.rrd",
"GET",
)
with pytest.raises(HTTPException) as replay_error:
asyncio.run(replay_route(session_id=session.name, request=None))
with pytest.raises(HTTPException) as recording_error:
asyncio.run(recording_route(session_id=session.name, if_none_match=None))
assert replay_error.value.status_code == 503
assert recording_error.value.status_code == 503
assert calls == 0
def test_recorded_blueprint_endpoint_is_small_strict_and_session_scoped(
tmp_path: Path,
monkeypatch: pytest.MonkeyPatch,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
observed_settings: list[RerunSceneSettings] = []
observed_kwargs: list[dict[str, Any]] = []
def capture_settings(settings: RerunSceneSettings, **kwargs: Any) -> bytes:
observed_settings.append(settings)
observed_kwargs.append(kwargs)
return recorded_blueprint_rrd(settings, **kwargs)
monkeypatch.setattr(
"k1link.web.session_api.recorded_blueprint_rrd",
capture_settings,
)
router = build_session_router(store)
blueprint_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/blueprint.rrd",
"POST",
)
response = asyncio.run(
blueprint_route(
session_id=session.name,
request=RecordedBlueprintRequest(
application_id="nodedc_mission_core_recorded",
recording_id="recording-001",
blueprint_session_id="a" * 32,
accumulation_seconds=18.5,
show_points=False,
show_trajectory=True,
show_grid=False,
point_size=6.25,
color_mode="height",
palette="custom",
custom_color="#112233",
active_view="perception",
view_reset_generation=1,
unified_perception=True,
show_detections_2d=True,
show_segmentation=True,
show_cuboids_3d=True,
follow_trajectory=True,
),
)
)
assert response.media_type == "application/vnd.rerun.rrd"
assert response.headers["cache-control"] == "no-store"
assert response.headers["x-content-type-options"] == "nosniff"
assert response.body.startswith(b"RRF2")
assert len(response.body) < 250_000
assert_no_local_paths(response.headers, repository)
assert len(observed_settings) == 1
assert observed_settings[0].show_points is False
assert observed_settings[0].show_trajectory is True
assert observed_settings[0].color_mode == "height"
assert observed_kwargs[0]["blueprint_session_id"] == "a" * 32
assert observed_kwargs[0]["active_view"] == "perception"
assert observed_kwargs[0]["view_reset_generation"] == 1
assert observed_kwargs[0]["unified_perception"] is True
assert observed_kwargs[0]["show_detections_2d"] is True
assert observed_kwargs[0]["show_segmentation"] is True
assert observed_kwargs[0]["show_cuboids_3d"] is True
assert observed_kwargs[0]["follow_trajectory"] is True
with pytest.raises(ValueError):
RecordedBlueprintRequest.model_validate(
{
"application_id": "nodedc_mission_core_recorded",
"recording_id": "recording-001",
"accumulation_seconds": 12,
"show_grid": True,
"source_url": "https://outside.invalid/recording.rrd",
}
)
with pytest.raises(HTTPException) as missing:
asyncio.run(
blueprint_route(
session_id="missing-session",
request=RecordedBlueprintRequest(
application_id="nodedc_mission_core_recorded",
recording_id="recording-001",
blueprint_session_id="a" * 32,
accumulation_seconds=12,
show_points=True,
show_trajectory=True,
show_grid=True,
),
)
)
assert missing.value.status_code == 404
def test_recorded_perception_endpoint_returns_one_complete_optional_overlay(
tmp_path: Path,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
calls: list[tuple[str, str, str]] = []
class Provider:
def render(
self,
session_id: str,
*,
application_id: str,
recording_id: str,
) -> bytes:
calls.append((session_id, application_id, recording_id))
return b"RRF2perception"
def status(self, session_id: str, *, recording_id: str) -> dict[str, object]:
assert session_id == session.name
assert recording_id == "recording-001"
return {
"state": "preparing",
"phase": "artifact-validation",
"elapsed_seconds": 2.5,
"byte_length": None,
}
router = build_session_router(store, perception_overlay_provider=Provider())
perception_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/perception.rrd",
"POST",
)
response = asyncio.run(
perception_route(
session_id=session.name,
request=RecordedPerceptionRequest(
application_id="nodedc_mission_core_recorded",
recording_id="recording-001",
),
)
)
assert response.body == b"RRF2perception"
assert response.media_type == "application/vnd.rerun.rrd"
assert response.headers["content-length"] == str(len(response.body))
assert calls == [(session.name, "nodedc_mission_core_recorded", "recording-001")]
status_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/perception/status",
"GET",
)
status_response = Response()
status = asyncio.run(
status_route(
session_id=session.name,
response=status_response,
recording_id="recording-001",
)
)
assert status == {
"state": "preparing",
"phase": "artifact-validation",
"elapsed_seconds": 2.5,
"byte_length": None,
}
assert status_response.headers["cache-control"] == "no-store"
empty_router = build_session_router(store)
empty_route = endpoint(
empty_router,
"/api/v1/observation-sessions/{session_id}/perception.rrd",
"POST",
)
empty = asyncio.run(
empty_route(
session_id=session.name,
request=RecordedPerceptionRequest(
application_id="nodedc_mission_core_recorded",
recording_id="recording-001",
),
)
)
assert empty.status_code == 204
def test_recorded_perception_endpoint_streams_validated_file_artifact(
tmp_path: Path,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
payload = b"RRF2file-backed-perception"
overlay = tmp_path / "sealed-perception.rrd"
overlay.write_bytes(payload)
digest = hashlib.sha256(payload).hexdigest()
calls: list[tuple[str, str, str]] = []
class Provider:
def materialize(
self,
session_id: str,
*,
application_id: str,
recording_id: str,
) -> RecordedPerceptionOverlayArtifact:
calls.append((session_id, application_id, recording_id))
return RecordedPerceptionOverlayArtifact(
path=overlay,
byte_length=len(payload),
sha256=digest,
)
def render(
self,
_session_id: str,
*,
application_id: str,
recording_id: str,
) -> bytes:
raise AssertionError(f"heap rendering must not run: {application_id}/{recording_id}")
router = build_session_router(store, perception_overlay_provider=Provider())
perception_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/perception.rrd",
"POST",
)
response = asyncio.run(
perception_route(
session_id=session.name,
request=RecordedPerceptionRequest(
application_id="nodedc_mission_core_recorded",
recording_id="recording-001",
),
)
)
assert isinstance(response, FileResponse)
assert Path(response.path) == overlay
assert response.media_type == "application/vnd.rerun.rrd"
assert response.headers["content-length"] == str(len(payload))
assert response.headers["etag"] == f'"sha256:{digest}"'
assert response.headers["content-disposition"] == 'inline; filename="perception.rrd"'
assert str(overlay) not in str(response.headers)
assert calls == [(session.name, "nodedc_mission_core_recorded", "recording-001")]
def test_recorded_point_color_endpoint_is_strict_and_forwards_confined_command(
tmp_path: Path,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
calls: list[tuple[str, str, str, str, str, str]] = []
class Provider:
def __call__(
self,
command: ReplayCommand,
*,
application_id: str,
recording_id: str,
color_mode: str,
palette: str,
custom_color: str,
) -> bytes:
calls.append(
(
command.session_id,
application_id,
recording_id,
color_mode,
palette,
custom_color,
)
)
return b"RRF2colors"
router = build_session_router(
store,
point_color_renderers={"nodedc.device.xgrids-lixelkity-k1": Provider()},
)
color_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/point-colors.rrd",
"POST",
)
response = asyncio.run(
color_route(
session_id=session.name,
request=RecordedPointColorsRequest(
application_id="nodedc_mission_core_recorded",
recording_id="recording-001",
color_mode="distance",
palette="viridis",
custom_color="#112233",
),
)
)
assert response.body == b"RRF2colors"
assert response.media_type == "application/vnd.rerun.rrd"
assert response.headers["content-length"] == str(len(response.body))
assert len(calls) == 1
assert calls[0] == (
session.name,
"nodedc_mission_core_recorded",
"recording-001",
"distance",
"viridis",
"#112233",
)
unavailable_router = build_session_router(store)
unavailable_route = endpoint(
unavailable_router,
"/api/v1/observation-sessions/{session_id}/point-colors.rrd",
"POST",
)
with pytest.raises(HTTPException) as unavailable:
asyncio.run(
unavailable_route(
session_id=session.name,
request=RecordedPointColorsRequest(
application_id="nodedc_mission_core_recorded",
recording_id="recording-001",
color_mode="intensity",
palette="turbo",
custom_color="#35d7c1",
),
)
)
assert unavailable.value.status_code == 409
def test_session_router_exposes_opaque_recorded_media_manifest_and_ranges(
tmp_path: Path,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
writer = CameraArchiveWriter(session, "sensor.camera.private-left", 9)
init, segment = make_recorded_h264_fixture()
writer.append("init", init, host_epoch_ns=1_100_000_000, host_monotonic_ns=2_100_000_000)
writer.append(
"media",
segment,
host_epoch_ns=1_250_000_000,
host_monotonic_ns=2_250_000_000,
)
writer.close()
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
payload = b"recording"
def export_recording(source: Path, destination: Path) -> dict[str, object]:
destination.write_bytes(payload)
return {
"source_sha256": hashlib.sha256(source.read_bytes()).hexdigest(),
"rrd_sha256": hashlib.sha256(payload).hexdigest(),
"rrd_bytes": len(payload),
"timeline": "session_time",
"timeline_start_ns": 0,
"timeline_end_ns": 5_000_000_000,
}
router = build_session_router(
store,
recording_materializer=SessionRecordingMaterializer(
store.data_dir,
exporter=export_recording,
),
allow_synchronous_recording_fallback=True,
)
replay_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/replay",
"POST",
)
manifest_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/media/{artifact_id}/manifest",
"GET",
)
init_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/media/{artifact_id}/"
"epochs/{epoch_ordinal}/init.mp4",
"GET",
)
segment_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/media/{artifact_id}/"
"epochs/{epoch_ordinal}/segments/{segment_sequence}.m4s",
"GET",
)
stream_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/media/{artifact_id}/"
"epochs/{epoch_ordinal}/recording.mp4",
"GET",
)
replay = asyncio.run(replay_route(session_id=session.name, request=None))
source = replay["launch"]["media_sources"][0]
assert source == {
"id": source["id"],
"label": "Записанная камера 1",
"modality": "video",
"manifest_url": source["manifest_url"],
"manifest_generation_sha256": source["manifest_generation_sha256"],
"byte_length": len(init) + len(segment),
"media_type": "video/mp4",
"timeline_start_seconds": 0.0,
"timeline_end_seconds": 0.5,
"seekable": True,
"synchronization": "host-arrival-best-effort",
}
assert source["id"].startswith("recorded.camera.")
assert "sensor.camera.private-left" not in json.dumps(replay)
artifact_id = source["manifest_url"].split("/")[-2]
with pytest.raises(HTTPException) as missing_manifest_generation:
manifest_route(session_id=session.name, artifact_id=artifact_id)
assert missing_manifest_generation.value.status_code == 428
with pytest.raises(HTTPException) as stale_manifest_generation:
manifest_route(
session_id=session.name,
artifact_id=artifact_id,
if_match='"sha256:stale"',
)
assert stale_manifest_generation.value.status_code == 412
manifest_response = manifest_route(
session_id=session.name,
artifact_id=artifact_id,
if_match=f'"sha256:{source["manifest_generation_sha256"]}"',
)
manifest = json.loads(manifest_response.body)
assert manifest["schema_version"] == "missioncore.observation-recorded-media/v4"
assert manifest["source_id"] == source["id"]
assert manifest["generation_sha256"] == source["manifest_generation_sha256"]
assert manifest["byte_length"] == source["byte_length"]
assert manifest["timeline_start_seconds"] == source["timeline_start_seconds"]
assert manifest["timeline_end_seconds"] == source["timeline_end_seconds"]
assert manifest_response.headers["etag"] == (f'"sha256:{manifest["generation_sha256"]}"')
generation = manifest["generation_sha256"]
assert manifest["epochs"] == [
{
"ordinal": 1,
"timeline_start_seconds": 0.0,
"timeline_end_seconds": 0.5,
"media_type": 'video/mp4; codecs="avc1.640028"',
"byte_length": len(init) + len(segment),
"segment_count": 1,
"random_access_sequences": [1],
"segment_end_times_seconds": [0.5],
"stream_url": (
f"{source['manifest_url'].removesuffix('/manifest')}/epochs/1/recording.mp4"
f"?generation={generation}"
),
}
]
assert "sensor.camera.private-left" not in json.dumps(manifest)
assert_no_local_paths(manifest, repository)
request = Request({"type": "http", "method": "GET", "headers": []})
with pytest.raises(HTTPException) as stale_stream_generation:
stream_route(
request=request,
session_id=session.name,
artifact_id=artifact_id,
epoch_ordinal=1,
generation="0" * 64,
range_header="bytes=0-5",
)
assert stale_stream_generation.value.status_code == 412
stream_response = stream_route(
request=request,
session_id=session.name,
artifact_id=artifact_id,
epoch_ordinal=1,
generation=generation,
range_header=f"bytes={len(init) - 2}-{len(init) + 3}",
)
stream_start, stream_body = asyncio.run(
render_file_response(stream_response, range_header=None)
)
stream_headers = {
key.decode("latin-1"): value.decode("latin-1") for key, value in stream_start["headers"]
}
assert stream_start["status"] == 206
assert stream_body == init[-2:] + segment[:4]
assert stream_headers["content-range"] == (
f"bytes {len(init) - 2}-{len(init) + 3}/{len(init) + len(segment)}"
)
assert stream_headers["content-length"] == "6"
assert stream_headers["accept-ranges"] == "bytes"
assert stream_headers["etag"] == f'"generation:{generation}:epoch:1"'
with pytest.raises(HTTPException) as missing_init_generation:
init_route(
session_id=session.name,
artifact_id=artifact_id,
epoch_ordinal=1,
range_header="bytes=4-7",
)
assert missing_init_generation.value.status_code == 428
init_sha256 = hashlib.sha256(init).hexdigest()
segment_sha256 = hashlib.sha256(segment).hexdigest()
init_response = init_route(
session_id=session.name,
artifact_id=artifact_id,
epoch_ordinal=1,
if_match=f'"sha256:{init_sha256}"',
range_header="bytes=4-7",
)
init_start, init_body = asyncio.run(
render_file_response(init_response, range_header=b"bytes=4-7")
)
assert init_start["status"] == 206
assert init_body == init[4:8]
generation_bound_init = init_route(
session_id=session.name,
artifact_id=artifact_id,
epoch_ordinal=1,
generation=generation,
range_header="bytes=0-3",
)
generation_init_start, generation_init_body = asyncio.run(
render_file_response(generation_bound_init, range_header=b"bytes=0-3")
)
assert generation_init_start["status"] == 206
assert generation_init_body == init[:4]
with pytest.raises(HTTPException) as stale_segment_generation:
segment_route(
session_id=session.name,
artifact_id=artifact_id,
epoch_ordinal=1,
segment_sequence=1,
if_match='"sha256:stale"',
range_header=None,
)
assert stale_segment_generation.value.status_code == 412
segment_response = segment_route(
session_id=session.name,
artifact_id=artifact_id,
epoch_ordinal=1,
segment_sequence=1,
if_match=f'"sha256:{segment_sha256}"',
range_header="bytes=0-5",
)
segment_start, segment_body = asyncio.run(
render_file_response(segment_response, range_header=b"bytes=0-5")
)
headers = {
key.decode("latin-1"): value.decode("latin-1") for key, value in segment_start["headers"]
}
assert segment_start["status"] == 206
assert segment_body == segment[:6]
assert headers["cache-control"] == ("private, max-age=31536000, immutable, no-transform")
assert headers["etag"] == f'"sha256:{segment_sha256}"'
assert headers["content-length"] == "6"
assert headers["x-content-type-options"] == "nosniff"
generation_bound_segment = segment_route(
session_id=session.name,
artifact_id=artifact_id,
epoch_ordinal=1,
segment_sequence=1,
generation=generation,
range_header="bytes=0-3",
)
generation_segment_start, generation_segment_body = asyncio.run(
render_file_response(generation_bound_segment, range_header=b"bytes=0-3")
)
assert generation_segment_start["status"] == 206
assert generation_segment_body == segment[:4]
with pytest.raises(HTTPException) as stale_fragment_generation:
segment_route(
session_id=session.name,
artifact_id=artifact_id,
epoch_ordinal=1,
segment_sequence=1,
generation="0" * 64,
range_header=None,
)
assert stale_fragment_generation.value.status_code == 412
original_segment = writer.segments_dir / "1.m4s"
outside = tmp_path / "outside.m4s"
outside.write_bytes(segment)
original_segment.unlink()
original_segment.symlink_to(outside)
with pytest.raises(HTTPException) as confined:
segment_route(
session_id=session.name,
artifact_id=artifact_id,
epoch_ordinal=1,
segment_sequence=1,
if_match=f'"sha256:{segment_sha256}"',
range_header=None,
)
assert confined.value.status_code == 409
def test_recorded_media_persists_and_exposes_random_access_sequences(
tmp_path: Path,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
init, first = make_recorded_h264_fixture()
_, second = make_recorded_h264_fixture(
base_decode_time=500,
tfhd_default_sample_flags=0x00010000,
)
_, third = make_recorded_h264_fixture(
base_decode_time=1_000,
tfhd_default_sample_flags=0x00010000,
first_sample_flags=0,
)
_, fourth = make_recorded_h264_fixture(
base_decode_time=1_500,
sample_flags=0x00010000,
)
writer = CameraArchiveWriter(session, "sensor.camera.private-left", 1)
writer.append("init", init)
for sequence, fragment in enumerate((first, second, third, fourth), start=1):
writer.append(
"media",
fragment,
host_epoch_ns=1_000_000_000 + sequence * 500_000_000,
host_monotonic_ns=2_000_000_000 + sequence * 500_000_000,
)
writer.close()
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
command = store.prepare_replay(session.name)
artifact = store.list_recorded_media(session.name)[0]
cache_root = tmp_path / "prepared-media"
manifest = RecordedMediaInspector(cache_root).inspect(artifact, command)
assert [segment.random_access for segment in manifest.epochs[0].segments] == [
True,
False,
True,
False,
]
document = _recorded_media_manifest_document(manifest)
assert document["schema_version"] == "missioncore.observation-recorded-media/v4"
assert document["epochs"][0]["segment_count"] == 4
assert document["epochs"][0]["random_access_sequences"] == [1, 3]
assert document["epochs"][0]["segment_end_times_seconds"] == [0.5, 1.0, 1.5, 2.0]
assert document["epochs"][0]["segment_end_times_seconds"][-1] == (
manifest.epochs[0].timeline_end_seconds - manifest.epochs[0].timeline_start_seconds
)
sidecar = json.loads(next(cache_root.glob("*.json")).read_text(encoding="utf-8"))
assert sidecar["schema_version"] == "missioncore.recorded-media-preparation/v3"
assert sidecar["manifest"]["schema_version"] == ("missioncore.observation-recorded-media/v4")
assert [
segment["random_access"] for segment in sidecar["manifest"]["epochs"][0]["segments"]
] == [True, False, True, False]
assert [
segment["end_time_seconds"] for segment in sidecar["manifest"]["epochs"][0]["segments"]
] == [0.5, 1.0, 1.5, 2.0]
restarted = RecordedMediaInspector(cache_root).inspect(artifact, command)
assert [segment.end_time_seconds for segment in restarted.epochs[0].segments] == [
0.5,
1.0,
1.5,
2.0,
]
def test_recorded_media_rejects_ambiguous_trun_sample_flags() -> None:
init, fragment = make_recorded_h264_fixture(
first_sample_flags=0,
sample_flags=0,
)
timing = recorded_media_module._mp4_video_timing(
init,
recorded_media_module._Mp4ParseBudget(),
)
with pytest.raises(SessionIntegrityError, match="sample flags are ambiguous"):
recorded_media_module._mp4_video_fragment_timing(
fragment,
timing,
recorded_media_module._Mp4ParseBudget(),
)
def test_recorded_media_rejects_multi_sample_video_fragment() -> None:
init, fragment = make_recorded_h264_fixture(fragment_sample_count=2)
timing = recorded_media_module._mp4_video_timing(
init,
recorded_media_module._Mp4ParseBudget(),
)
with pytest.raises(SessionIntegrityError, match="exactly one sample"):
recorded_media_module._mp4_video_fragment_timing(
fragment,
timing,
recorded_media_module._Mp4ParseBudget(),
)
def test_recorded_media_rejects_codec_epoch_without_initial_random_access(
tmp_path: Path,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
init, non_random_access = make_recorded_h264_fixture(sample_flags=0x00010000)
writer = CameraArchiveWriter(session, "sensor.camera.private-left", 1)
writer.append("init", init)
writer.append(
"media",
non_random_access,
host_epoch_ns=1_500_000_000,
host_monotonic_ns=2_500_000_000,
)
writer.close()
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
with pytest.raises(SessionIntegrityError, match="does not begin with a random-access"):
RecordedMediaInspector().inspect(
store.list_recorded_media(session.name)[0],
store.prepare_replay(session.name),
)
def test_replay_exposes_generation_bound_perception_video_with_native_ranges(
tmp_path: Path,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260720T065719Z_viewer_live")
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
recording_payload = b"recording"
def export_recording(source: Path, destination: Path) -> dict[str, object]:
destination.write_bytes(recording_payload)
return {
"source_sha256": hashlib.sha256(source.read_bytes()).hexdigest(),
"rrd_sha256": hashlib.sha256(recording_payload).hexdigest(),
"rrd_bytes": len(recording_payload),
"timeline": "session_time",
"timeline_start_ns": 0,
"timeline_end_ns": 5_000_000_000,
}
video_payload = b"sealed-panoptic-video"
video_path = tmp_path / "perception.mp4"
video_path.write_bytes(video_payload)
video_sha256 = hashlib.sha256(video_payload).hexdigest()
result_id = f"result-{'d' * 64}"
video = RecordedPerceptionVideo(
result_id=result_id,
session_id=session.name,
source_id="sensor.camera.right",
public_source_id="recorded.perception.right",
label="Сегментация · камера right",
path=video_path,
media_type='video/mp4; codecs="avc1.640028"',
byte_length=len(video_payload),
sha256=video_sha256,
timeline_start_seconds=0.5,
timeline_end_seconds=4.5,
)
class Provider:
def __init__(self) -> None:
self.calls: list[tuple[str, str | None]] = []
def video(
self,
session_id: str,
requested_result_id: str | None = None,
) -> RecordedPerceptionVideo | None:
self.calls.append((session_id, requested_result_id))
if session_id != video.session_id:
return None
if requested_result_id is not None and requested_result_id != video.result_id:
return None
return video
provider = Provider()
router = build_session_router(
store,
recording_materializer=SessionRecordingMaterializer(
store.data_dir,
exporter=export_recording,
),
perception_media_provider=provider,
allow_synchronous_recording_fallback=True,
)
replay_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/replay",
"POST",
)
manifest_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/perception-media/{result_id}/manifest",
"GET",
)
stream_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/perception-media/{result_id}/recording.mp4",
"GET",
)
replay = asyncio.run(replay_route(session_id=session.name, request=None))
# Perception is loaded through the on-demand RRD overlay and no longer
# blocks replay launch by hashing a duplicate pre-rendered video.
assert replay["launch"]["media_sources"] == []
assert provider.calls == []
with pytest.raises(HTTPException) as missing_generation:
manifest_route(session_id=session.name, result_id=result_id)
assert missing_generation.value.status_code == 428
manifest_response = manifest_route(
session_id=session.name,
result_id=result_id,
if_match=f'"sha256:{video_sha256}"',
)
assert provider.calls == [
(session.name, result_id),
(session.name, result_id),
]
manifest = json.loads(manifest_response.body)
assert manifest == {
"schema_version": "missioncore.observation-recorded-media/v3",
"source_id": "recorded.perception.right",
"generation_sha256": video_sha256,
"byte_length": len(video_payload),
"timeline_start_seconds": 0.5,
"timeline_end_seconds": 4.5,
"synchronization": "host-arrival-best-effort",
"epochs": [
{
"ordinal": 1,
"timeline_start_seconds": 0.5,
"timeline_end_seconds": 4.5,
"media_type": 'video/mp4; codecs="avc1.640028"',
"byte_length": len(video_payload),
"stream_url": (
f"/api/v1/observation-sessions/{session.name}/perception-media/"
f"{result_id}/recording.mp4?generation={video_sha256}"
),
}
],
}
request = Request({"type": "http", "method": "GET", "headers": []})
with pytest.raises(HTTPException) as stale_generation:
stream_route(
request=request,
session_id=session.name,
result_id=result_id,
generation="0" * 64,
range_header=None,
)
assert stale_generation.value.status_code == 412
response = stream_route(
request=request,
session_id=session.name,
result_id=result_id,
generation=video_sha256,
range_header="bytes=7-14",
)
start, body = asyncio.run(render_file_response(response, range_header=b"bytes=7-14"))
headers = {key.decode("latin-1"): value.decode("latin-1") for key, value in start["headers"]}
assert start["status"] == 206
assert body == video_payload[7:15]
assert headers["content-range"] == f"bytes 7-14/{len(video_payload)}"
assert headers["accept-ranges"] == "bytes"
assert headers["etag"] == f'"sha256:{video_sha256}"'
def test_ready_launch_uses_background_prepared_camera_manifest_without_rescan(
tmp_path: Path,
monkeypatch: pytest.MonkeyPatch,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
writer = CameraArchiveWriter(session, "sensor.camera.private-left", 1)
init, segment = make_recorded_h264_fixture()
writer.append(
"init",
init,
host_epoch_ns=1_100_000_000,
host_monotonic_ns=2_100_000_000,
)
writer.append(
"media",
segment,
host_epoch_ns=1_250_000_000,
host_monotonic_ns=2_250_000_000,
)
writer.close()
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
def exporter(source: Path, destination: Path) -> dict[str, object]:
payload = b"prepared-recording"
destination.write_bytes(payload)
return {
"source_sha256": hashlib.sha256(source.read_bytes()).hexdigest(),
"rrd_sha256": hashlib.sha256(payload).hexdigest(),
"rrd_bytes": len(payload),
"timeline": "session_time",
"timeline_start_ns": 0,
"timeline_end_ns": 5_000_000_000,
}
materializer = SessionRecordingMaterializer(store.data_dir, exporter=exporter)
inspector = RecordedMediaInspector()
def prepare_media(
command: ReplayCommand,
_recording: MaterializedRecording,
) -> tuple[Any, ...]:
return tuple(
inspector.inspect(artifact, command)
for artifact in store.list_recorded_media(command.session_id)
)
manager = SessionRecordingPreparationManager(
materializer,
ready_preparer=prepare_media,
)
try:
manager.enqueue(store.prepare_replay(session.name))
deadline = time.monotonic() + 2
snapshot = manager.status(session.name)
while (snapshot is None or snapshot.state != "ready") and time.monotonic() < deadline:
time.sleep(0.005)
snapshot = manager.status(session.name)
assert snapshot is not None and snapshot.state == "ready"
assert snapshot.recorded_media is not None and len(snapshot.recorded_media) == 1
def forbidden_rescan(*_args: object, **_kwargs: object) -> None:
raise AssertionError("ready HTTP path must not rescan RRD or camera indexes")
monkeypatch.setattr(inspector, "inspect", forbidden_rescan)
monkeypatch.setattr(manager, "resolve_cached", forbidden_rescan)
monkeypatch.setattr(manager, "resolve_cached_pinned", forbidden_rescan)
router = build_session_router(
store,
recording_materializer=materializer,
recording_preparation_manager=manager,
media_inspector=inspector,
)
replay_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/replay",
"POST",
)
manifest_route = endpoint(
router,
"/api/v1/observation-sessions/{session_id}/media/{artifact_id}/manifest",
"GET",
)
replay = asyncio.run(replay_route(session_id=session.name, request=None))
source = replay["launch"]["media_sources"][0]
artifact_id = source["manifest_url"].split("/")[-2]
manifest_response = manifest_route(
session_id=session.name,
artifact_id=artifact_id,
if_match=f'"sha256:{source["manifest_generation_sha256"]}"',
)
manifest = json.loads(manifest_response.body)
assert manifest["generation_sha256"] == source["manifest_generation_sha256"]
finally:
manager.close()
def test_recorded_media_epoch_ends_are_monotonic_generation_bound_and_not_rrd_padded(
tmp_path: Path,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
init, segment = make_recorded_h264_fixture(sample_duration=250)
_, following_segment = make_recorded_h264_fixture(
sample_duration=250,
base_decode_time=250,
)
first = CameraArchiveWriter(session, "sensor.camera.private-left", 1)
first.append("init", init, host_epoch_ns=1_100_000_000, host_monotonic_ns=2_100_000_000)
first.append(
"media",
segment,
host_epoch_ns=1_250_000_000,
host_monotonic_ns=2_250_000_000,
)
first.close()
second = CameraArchiveWriter(session, "sensor.camera.private-left", 2)
second.append("init", init, host_epoch_ns=1_800_000_000, host_monotonic_ns=2_800_000_000)
second.append(
"media",
segment,
host_epoch_ns=1_800_000_000,
host_monotonic_ns=2_800_000_000,
)
second.append(
"media",
following_segment,
host_epoch_ns=2_000_000_000,
host_monotonic_ns=3_000_000_000,
)
second.close()
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
command = store.prepare_replay(session.name)
artifact = store.list_recorded_media(session.name)[0]
manifest = RecordedMediaInspector().inspect(artifact, command)
assert [epoch.timeline_start_seconds for epoch in manifest.epochs] == [0.0, 0.55]
assert [epoch.timeline_end_seconds for epoch in manifest.epochs] == [0.25, 1.05]
assert manifest.timeline_end_seconds == 1.05
assert manifest.byte_length == 2 * len(init) + 3 * len(segment)
original_generation = manifest.generation_sha256
changed_epochs = (
*manifest.epochs[:-1],
replace(
manifest.epochs[-1],
timeline_end_seconds=manifest.epochs[-1].timeline_end_seconds + 0.1,
),
)
changed_generation = _manifest_generation_sha256(
public_source_id=manifest.public_source_id,
artifact_id=manifest.artifact_id,
synchronization=manifest.synchronization,
epochs=changed_epochs,
)
assert changed_generation != original_generation
def test_recorded_media_rejects_non_monotonic_segments_and_overlapping_epochs(
tmp_path: Path,
) -> None:
init, segment = make_recorded_h264_fixture(sample_duration=500)
non_monotonic_root = tmp_path / "non-monotonic"
non_monotonic_session = make_legacy_session(
non_monotonic_root / "sessions",
"20260716T205632Z_viewer_live",
)
writer = CameraArchiveWriter(non_monotonic_session, "sensor.camera.private-left", 1)
writer.append("init", init)
writer.append(
"media",
segment,
host_epoch_ns=1_600_000_000,
host_monotonic_ns=2_600_000_000,
)
writer.append(
"media",
segment,
host_epoch_ns=1_500_000_000,
host_monotonic_ns=2_500_000_000,
)
writer.close()
store = SessionStore(non_monotonic_root, data_dir=tmp_path / "non-monotonic-data")
store.reconcile_archive(xgrids_k1_archive_source(non_monotonic_root / "sessions"))
with pytest.raises(SessionIntegrityError, match="strictly monotonic"):
RecordedMediaInspector().inspect(
store.list_recorded_media(non_monotonic_session.name)[0],
store.prepare_replay(non_monotonic_session.name),
)
overlap_root = tmp_path / "overlap"
overlap_session = make_legacy_session(
overlap_root / "sessions",
"20260716T205633Z_viewer_live",
)
first = CameraArchiveWriter(overlap_session, "sensor.camera.private-left", 1)
first.append("init", init)
first.append(
"media",
segment,
host_epoch_ns=1_250_000_000,
host_monotonic_ns=2_250_000_000,
)
first.close()
second = CameraArchiveWriter(overlap_session, "sensor.camera.private-left", 2)
second.append("init", init)
second.append(
"media",
segment,
host_epoch_ns=1_500_000_000,
host_monotonic_ns=2_500_000_000,
)
second.close()
store = SessionStore(overlap_root, data_dir=tmp_path / "overlap-data")
store.reconcile_archive(xgrids_k1_archive_source(overlap_root / "sessions"))
with pytest.raises(SessionIntegrityError, match="overlap"):
RecordedMediaInspector().inspect(
store.list_recorded_media(overlap_session.name)[0],
store.prepare_replay(overlap_session.name),
)
def test_recorded_media_rejects_discontinuous_tfdt_decode_timeline(
tmp_path: Path,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
init, first_segment = make_recorded_h264_fixture(sample_duration=500)
_, discontinuous_segment = make_recorded_h264_fixture(
sample_duration=500,
base_decode_time=1_000,
)
writer = CameraArchiveWriter(session, "sensor.camera.private-left", 1)
writer.append("init", init)
writer.append(
"media",
first_segment,
host_epoch_ns=1_250_000_000,
host_monotonic_ns=2_250_000_000,
)
writer.append(
"media",
discontinuous_segment,
host_epoch_ns=1_750_000_000,
host_monotonic_ns=2_750_000_000,
)
writer.close()
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
with pytest.raises(SessionIntegrityError, match="decode timeline is discontinuous"):
RecordedMediaInspector().inspect(
store.list_recorded_media(session.name)[0],
store.prepare_replay(session.name),
)
def test_background_preparation_fails_closed_on_unparseable_recorded_media(
tmp_path: Path,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
writer = CameraArchiveWriter(session, "sensor.camera.private-left", 1)
writer.append("init", b"not-an-iso-bmff-init")
writer.append("media", b"not-an-iso-bmff-fragment")
writer.close()
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
def exporter(source: Path, destination: Path) -> dict[str, object]:
return {
"source_sha256": hashlib.sha256(source.read_bytes()).hexdigest(),
"rrd_sha256": hashlib.sha256(b"recording").hexdigest(),
"rrd_bytes": destination.write_bytes(b"recording"),
"timeline": "session_time",
"timeline_start_ns": 0,
"timeline_end_ns": 5_000_000_000,
}
inspector = RecordedMediaInspector()
def prepare_media(
command: ReplayCommand,
_recording: MaterializedRecording,
) -> tuple[Any, ...]:
return tuple(
inspector.inspect(artifact, command)
for artifact in store.list_recorded_media(command.session_id)
)
manager = SessionRecordingPreparationManager(
SessionRecordingMaterializer(store.data_dir, exporter=exporter),
ready_preparer=prepare_media,
)
try:
manager.enqueue(store.prepare_replay(session.name))
deadline = time.monotonic() + 2
snapshot = manager.status(session.name)
while snapshot is not None and snapshot.state not in {"failed", "ready"}:
assert time.monotonic() < deadline
time.sleep(0.005)
snapshot = manager.status(session.name)
assert snapshot is not None
assert snapshot.state == "failed"
assert snapshot.recording is None
assert snapshot.recorded_media is None
finally:
manager.close()
def test_background_preparation_rejects_camera_outside_rrd_timeline(
tmp_path: Path,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
init, segment = make_recorded_h264_fixture(sample_duration=500)
writer = CameraArchiveWriter(session, "sensor.camera.private-left", 1)
writer.append("init", init)
writer.append(
"media",
segment,
host_epoch_ns=1_500_000_000,
host_monotonic_ns=2_500_000_000,
)
writer.close()
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
def exporter(source: Path, destination: Path) -> dict[str, object]:
payload = b"short-spatial-recording"
destination.write_bytes(payload)
return {
"source_sha256": hashlib.sha256(source.read_bytes()).hexdigest(),
"rrd_sha256": hashlib.sha256(payload).hexdigest(),
"rrd_bytes": len(payload),
"timeline": "session_time",
"timeline_start_ns": 0,
"timeline_end_ns": 400_000_000,
}
inspector = RecordedMediaInspector(tmp_path / "prepared-media")
def prepare_media(
command: ReplayCommand,
_recording: MaterializedRecording,
) -> tuple[Any, ...]:
return tuple(
inspector.inspect(artifact, command)
for artifact in store.list_recorded_media(command.session_id)
)
manager = SessionRecordingPreparationManager(
SessionRecordingMaterializer(store.data_dir, exporter=exporter),
ready_preparer=prepare_media,
)
try:
manager.enqueue(store.prepare_replay(session.name))
deadline = time.monotonic() + 2
snapshot = manager.status(session.name)
while snapshot is not None and snapshot.state not in {"failed", "ready"}:
assert time.monotonic() < deadline
time.sleep(0.005)
snapshot = manager.status(session.name)
assert snapshot is not None
assert snapshot.state == "failed"
assert snapshot.recording is None
assert snapshot.recorded_media is None
finally:
manager.close()
def test_recorded_media_preparation_sidecar_reuses_and_rebuilds_generation(
tmp_path: Path,
monkeypatch: pytest.MonkeyPatch,
) -> None:
repository = tmp_path / "repo"
sessions = repository / "sessions"
session = make_legacy_session(sessions, "20260716T205632Z_viewer_live")
init, segment = make_recorded_h264_fixture()
writer = CameraArchiveWriter(session, "sensor.camera.private-left", 1)
writer.append("init", init)
writer.append(
"media",
segment,
host_epoch_ns=1_500_000_000,
host_monotonic_ns=2_500_000_000,
)
writer.close()
store = SessionStore(repository, data_dir=tmp_path / "data")
store.reconcile_archive(xgrids_k1_archive_source(sessions))
command = store.prepare_replay(session.name)
artifact = store.list_recorded_media(session.name)[0]
cache_root = tmp_path / "prepared-media"
initial = RecordedMediaInspector(cache_root).inspect(artifact, command)
sidecars = list(cache_root.glob("*.json"))
assert len(sidecars) == 1
sidecar = sidecars[0]
assert stat.S_IMODE(sidecar.stat().st_mode) == 0o600
assert str(session) not in sidecar.read_text(encoding="utf-8")
prepared = json.loads(sidecar.read_text(encoding="utf-8"))
assert prepared["schema_version"] == "missioncore.recorded-media-preparation/v3"
assert prepared["manifest"]["epochs"][0]["segments"][0]["random_access"] is True
assert prepared["manifest"]["epochs"][0]["segments"][0]["end_time_seconds"] == 0.5
def forbidden_reparse(*_args: object, **_kwargs: object) -> None:
raise AssertionError("restart must reuse the prepared media sidecar")
with monkeypatch.context() as context:
context.setattr(recorded_media_module, "_read_manifest", forbidden_reparse)
context.setattr(recorded_media_module, "_mp4_video_timing", forbidden_reparse)
restarted = RecordedMediaInspector(cache_root).inspect(artifact, command)
assert restarted.generation_sha256 == initial.generation_sha256
assert restarted.timeline_end_seconds == initial.timeline_end_seconds
assert [segment.end_time_seconds for segment in restarted.epochs[0].segments] == [0.5]
original_read_manifest = recorded_media_module._read_manifest
reparses = 0
def counted_reparse(*args: Any, **kwargs: Any) -> Any:
nonlocal reparses
reparses += 1
return original_read_manifest(*args, **kwargs)
old_body = json.loads(json.dumps(prepared))
old_body.pop("checksum_sha256")
old_body["schema_version"] = "missioncore.recorded-media-preparation/v2"
old_body["manifest"]["schema_version"] = "missioncore.observation-recorded-media/v3"
for epoch in old_body["manifest"]["epochs"]:
for old_segment in epoch["segments"]:
old_segment.pop("end_time_seconds")
old_checksum = hashlib.sha256(recorded_media_module._canonical_json(old_body)).hexdigest()
sidecar.write_bytes(
recorded_media_module._canonical_json({**old_body, "checksum_sha256": old_checksum})
)
with monkeypatch.context() as context:
context.setattr(recorded_media_module, "_read_manifest", counted_reparse)
migrated = RecordedMediaInspector(cache_root).inspect(artifact, command)
assert migrated.generation_sha256 == initial.generation_sha256
assert reparses == 1
migrated_sidecar = json.loads(sidecar.read_text(encoding="utf-8"))
assert migrated_sidecar["schema_version"] == "missioncore.recorded-media-preparation/v3"
assert migrated_sidecar["manifest"]["epochs"][0]["segments"][0]["random_access"] is True
assert migrated_sidecar["manifest"]["epochs"][0]["segments"][0]["end_time_seconds"] == 0.5
reparses = 0
sidecar.write_bytes(b"{corrupt")
with monkeypatch.context() as context:
context.setattr(recorded_media_module, "_read_manifest", counted_reparse)
repaired = RecordedMediaInspector(cache_root).inspect(artifact, command)
assert repaired.generation_sha256 == initial.generation_sha256
assert reparses == 1
assert json.loads(sidecar.read_text(encoding="utf-8"))["checksum_sha256"]
segment_path = writer.segments_dir / "1.m4s"
replacement = writer.segments_dir / "replacement.m4s"
replacement.write_bytes(segment_path.read_bytes())
replacement.replace(segment_path)
reparses = 0
with monkeypatch.context() as context:
context.setattr(recorded_media_module, "_read_manifest", counted_reparse)
replaced = RecordedMediaInspector(cache_root).inspect(artifact, command)
assert replaced.generation_sha256 == initial.generation_sha256
assert reparses == 1
def test_recorded_media_preparation_scavenges_crash_temp_files(tmp_path: Path) -> None:
cache_root = tmp_path / "prepared-media"
cache_root.mkdir()
crash_temp = cache_root / f".tmp-{'a' * 32}"
crash_temp.write_bytes(b"partial")
RecordedMediaInspector(cache_root)
assert crash_temp.exists() is False
def test_layout_routes_round_trip_versioned_document(tmp_path: Path) -> None:
store = SessionStore(tmp_path / "repo", data_dir=tmp_path / "data")
router = build_session_router(store)
put_route = endpoint(router, "/api/v1/workspace-layouts/{workspace_id}", "PUT")
get_route = endpoint(router, "/api/v1/workspace-layouts/{workspace_id}", "GET")
request = LayoutPutRequest(
version=2,
revision=0,
workspace_id="observation.spatial",
scene_settings={
"projection": "3d",
"point_size": 3.0,
"color_mode": "intensity",
"palette": "turbo",
"custom_color": "#35d7c1",
"accumulation_seconds": 12,
"show_points": True,
"show_trajectory": True,
"show_grid": True,
"show_labels": False,
"show_camera_frustums": True,
},
visible_source_ids=["sensor.lidar.primary"],
active_floating_source_id=None,
window_rects={"sensor.camera.left": {"x": 0.65, "y": 0.65, "width": 0.3, "height": 0.3}},
viewport_size={"width": 1440, "height": 900},
)
put_response = Response()
saved = put_route(
workspace_id="observation.spatial",
request=request,
response=put_response,
if_match='"0"',
)
get_response = Response()
loaded = get_route(workspace_id="observation.spatial", response=get_response)
assert saved == loaded
assert set(loaded) == {
"version",
"revision",
"workspace_id",
"scene_settings",
"visible_source_ids",
"active_floating_source_id",
"window_rects",
"viewport_size",
}
assert loaded["version"] == 2
assert loaded["workspace_id"] == "observation.spatial"
assert loaded["revision"] == 1
assert put_response.headers["etag"] == '"1"'
assert get_response.headers["etag"] == '"1"'
stale = request.model_copy(update={"revision": 0})
with pytest.raises(HTTPException) as conflict:
put_route(
workspace_id="observation.spatial",
request=stale,
response=Response(),
if_match='"0"',
)
assert conflict.value.status_code == 412
def test_layout_get_migrates_legacy_tool_windows_out_of_the_contract(tmp_path: Path) -> None:
store = SessionStore(tmp_path / "repo", data_dir=tmp_path / "data")
store.save_layout(
"observation.spatial",
schema_version=1,
expected_revision=0,
name="Пространственная сцена",
layout={
"scene_settings": {
"projection": "3d",
"point_size": 3.0,
"color_mode": "intensity",
"palette": "turbo",
"custom_color": "#35d7c1",
"accumulation_seconds": 12,
"show_points": True,
"show_trajectory": True,
"show_grid": True,
"show_labels": False,
"show_camera_frustums": True,
},
"tool_windows": {
"sources_open": True,
"display_open": True,
"layers_open": True,
"order": ["layers", "sources", "display"],
},
"visible_source_ids": [],
"active_floating_source_id": None,
"window_rects": {},
"viewport_size": {"width": 1440, "height": 900},
},
)
router = build_session_router(store)
get_route = endpoint(router, "/api/v1/workspace-layouts/{workspace_id}", "GET")
loaded = get_route(workspace_id="observation.spatial", response=Response())
assert loaded["version"] == 2
assert "tool_windows" not in loaded