2289 lines
83 KiB
Python
2289 lines
83 KiB
Python
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, HTTPException, Request, Response
|
||
from fastapi.routing import APIRoute
|
||
|
||
import k1link.sessions.media as recorded_media_module
|
||
from k1link.compute import 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 (
|
||
MaterializedRecording,
|
||
RecordedMediaInspector,
|
||
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,
|
||
build_session_router,
|
||
)
|
||
|
||
|
||
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,
|
||
) -> 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" * 8,
|
||
)
|
||
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 = full_box(b"tfhd", track_id.to_bytes(4, "big"), flags=0x020000)
|
||
tfdt = full_box(b"tfdt", base_decode_time.to_bytes(4, "big"))
|
||
trun = full_box(b"trun", (1).to_bytes(4, "big"))
|
||
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},
|
||
)
|
||
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)
|
||
item = next(value for value in listing["items"] if value["id"] == binding.session_id)
|
||
detail = detail_route(session_id=binding.session_id)
|
||
|
||
assert item["lab"] == binding.as_dict()
|
||
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_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))
|
||
|
||
assert 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_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_catalog_read_never_prepares_a_cold_historical_session(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("catalog reads must not invoke the exporter")
|
||
|
||
materializer = SessionRecordingMaterializer(store.data_dir, exporter=exporter)
|
||
manager = SessionRecordingPreparationManager(materializer)
|
||
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"
|
||
|
||
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")
|
||
]
|
||
|
||
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_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/v3"
|
||
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),
|
||
"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]
|
||
|
||
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"
|
||
|
||
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_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")
|
||
|
||
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
|
||
|
||
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)
|
||
|
||
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
|