diff --git a/apps/control-station/src/components/RerunViewport.tsx b/apps/control-station/src/components/RerunViewport.tsx index 66df3ba..5900a42 100644 --- a/apps/control-station/src/components/RerunViewport.tsx +++ b/apps/control-station/src/components/RerunViewport.tsx @@ -152,6 +152,7 @@ interface RerunNativeReceiver { interface LoadedNativePerceptionSource { receiver: RerunNativeReceiver; + descriptorUrl: string; sourceUrl: string; byteLength: number; } @@ -596,7 +597,7 @@ export function resolveRecordedPerceptionViewerSourceUrl( return endpoint.href; } -/** Verify the immutable URL with a four-byte range request before Rerun opens it. */ +/** Resolve the immutable overlay generation before Rerun opens its native URL. */ export async function probeRecordedPerceptionViewerSource( sourceUrl: string, { @@ -608,7 +609,7 @@ export async function probeRecordedPerceptionViewerSource( signal?: AbortSignal; fetcher?: typeof globalThis.fetch; }, -): Promise { +): Promise<{ sourceUrl: string; byteLength: number }> { const base = new URL(origin); const endpoint = new URL(sourceUrl, base.origin); const allowedParameters = ["application_id", "generation", "recording_id"]; @@ -626,57 +627,29 @@ export async function probeRecordedPerceptionViewerSource( throw new Error("Unsafe recorded perception viewer source"); } const response = await fetcher(endpoint.href, { - method: "GET", + method: "HEAD", credentials: "same-origin", headers: { Accept: "application/vnd.rerun.rrd", - Range: "bytes=0-3", }, signal, }); const contentType = response.headers.get("Content-Type")?.split(";", 1)[0].trim(); - const contentRange = response.headers.get("Content-Range"); - const rangedLength = contentRange?.match(/^bytes 0-3\/(\d+)$/)?.[1]; - const declaredLength = Number( - rangedLength ?? response.headers.get("Content-Length"), - ); + const declaredLength = Number(response.headers.get("Content-Length")); + const etag = response.headers.get("ETag")?.match(/^"([a-f0-9]{64})"$/)?.[1]; if ( - ![200, 206].includes(response.status) || + response.status !== 200 || contentType !== "application/vnd.rerun.rrd" || + response.headers.get("X-Rerun-Format") !== "RRF2" || + !etag || !Number.isSafeInteger(declaredLength) || declaredLength < 4 || declaredLength > MAX_PERCEPTION_BYTES ) { throw new Error("Invalid recorded perception viewer response"); } - const prefix = new Uint8Array(4); - let receivedBytes = 0; - const reader = response.body?.getReader(); - if (reader) { - while (receivedBytes < prefix.byteLength) { - const { done, value } = await reader.read(); - if (done) break; - const count = Math.min(value.byteLength, prefix.byteLength - receivedBytes); - prefix.set(value.subarray(0, count), receivedBytes); - receivedBytes += count; - } - await reader.cancel(); - } else { - const payload = new Uint8Array(await response.arrayBuffer()); - const count = Math.min(payload.byteLength, prefix.byteLength); - prefix.set(payload.subarray(0, count)); - receivedBytes = count; - } - if ( - receivedBytes !== prefix.byteLength || - prefix[0] !== 0x52 || - prefix[1] !== 0x52 || - prefix[2] !== 0x46 || - prefix[3] !== 0x32 - ) { - throw new Error("Invalid recorded perception viewer RRD"); - } - return declaredLength; + endpoint.searchParams.set("overlay_generation", etag); + return { sourceUrl: endpoint.href, byteLength: declaredLength }; } export function RerunViewport({ @@ -1636,7 +1609,7 @@ export function RerunViewport({ return; } const loaded = loadedNativePerceptionSourceRef.current; - if (loaded?.receiver === receiver && loaded.sourceUrl === sourceUrl) { + if (loaded?.receiver === receiver && loaded.descriptorUrl === sourceUrl) { onPerceptionLoadChange?.({ phase: "ready", receivedBytes: loaded.byteLength, @@ -1657,17 +1630,18 @@ export function RerunViewport({ void probeRecordedPerceptionViewerSource(sourceUrl, { origin: window.location.origin, signal: abort.signal, - }).then((byteLength) => { + }).then(({ sourceUrl: immutableSourceUrl, byteLength }) => { if ( abort.signal.aborted || perceptionReceiverRef.current !== receiver || recordedIdentityRef.current !== identity || !receiver.ready() ) return; - receiver.open(sourceUrl); + receiver.open(immutableSourceUrl); loadedNativePerceptionSourceRef.current = { receiver, - sourceUrl, + descriptorUrl: sourceUrl, + sourceUrl: immutableSourceUrl, byteLength, }; onPerceptionLoadChange?.({ diff --git a/apps/control-station/test/observationSources.test.mjs b/apps/control-station/test/observationSources.test.mjs index 733bfb5..6ac6a78 100644 --- a/apps/control-station/test/observationSources.test.mjs +++ b/apps/control-station/test/observationSources.test.mjs @@ -686,23 +686,56 @@ test("LAB perception sidecar is streamed by native Rerun from one generation-bou `&recording_id=recording-001&generation=${"b".repeat(64)}`, ); - const byteLength = await probeRecordedPerceptionViewerSource(sourceUrl, { + const overlayGeneration = "c".repeat(64); + const probe = await probeRecordedPerceptionViewerSource(sourceUrl, { origin: "http://127.0.0.1:5174", fetcher: async (input, init) => { assert.equal(String(input), sourceUrl); - assert.equal(init.method, "GET"); - assert.equal(new Headers(init.headers).get("Range"), "bytes=0-3"); - return new Response(Uint8Array.from([0x52, 0x52, 0x46, 0x32]), { - status: 206, + assert.equal(init.method, "HEAD"); + assert.equal(new Headers(init.headers).get("Range"), null); + return new Response(null, { + status: 200, headers: { "Content-Type": "application/vnd.rerun.rrd", - "Content-Length": "4", - "Content-Range": "bytes 0-3/393203594", + "Content-Length": "186058411", + "ETag": `"${overlayGeneration}"`, + "X-Rerun-Format": "RRF2", }, }); }, }); - assert.equal(byteLength, 393_203_594); + assert.deepEqual(probe, { + sourceUrl: `${sourceUrl}&overlay_generation=${overlayGeneration}`, + byteLength: 186_058_411, + }); +}); + +test("LAB native source rejects an unsealed overlay descriptor", async () => { + const endpoint = + `http://127.0.0.1:5174/api/v1/laboratory/vegetation-shadow/` + + `lab-v1-vegetation-shadow-${"a".repeat(64)}/canonical-overlay.rrd`; + const sourceUrl = resolveRecordedPerceptionViewerSourceUrl( + endpoint, + { applicationId: "nodedc_mission_core_recorded", recordingId: "recording-001" }, + "b".repeat(64), + "http://127.0.0.1:5174", + ); + + await assert.rejects( + probeRecordedPerceptionViewerSource(sourceUrl, { + origin: "http://127.0.0.1:5174", + fetcher: async () => new Response(null, { + status: 200, + headers: { + "Content-Type": "application/vnd.rerun.rrd", + "Content-Length": "186058411", + "ETag": `"${"c".repeat(64)}"`, + "X-Rerun-Format": "RRF1", + }, + }), + }), + /Invalid recorded perception viewer response/, + ); }); test("recorded replay creates an isolated source catalog without live device bindings", () => { diff --git a/src/k1link/laboratory/canonical_rerun_overlay.py b/src/k1link/laboratory/canonical_rerun_overlay.py index ddfda57..abf8992 100644 --- a/src/k1link/laboratory/canonical_rerun_overlay.py +++ b/src/k1link/laboratory/canonical_rerun_overlay.py @@ -14,6 +14,7 @@ import io import json import os import subprocess +import sys import tempfile import threading import zipfile @@ -30,9 +31,9 @@ from PIL import Image APPLICATION_ID: Final = "nodedc_mission_core_recorded" SESSION_TIMELINE: Final = "session_time" -RENDERER_VERSION: Final = "upstream-rerun-0.36.3-v2" +RENDERER_VERSION: Final = "upstream-rerun-0.36.3-encoded-optimized-v4" MAX_SOURCE_BYTES: Final = 768 * 1024 * 1024 -MAX_OVERLAY_BYTES: Final = 512 * 1024 * 1024 +MAX_OVERLAY_BYTES: Final = 256 * 1024 * 1024 class CanonicalLabOverlayError(RuntimeError): @@ -136,6 +137,7 @@ def canonical_lab_overlay( source = _verified_camera_source(root, route, jobs) proxy = _camera_proxy(source, int(route["frame_count"]), ffmpeg, cache) _render_overlay(temporary, root, route, recording_id, proxy) + _optimize_overlay(temporary) stat = temporary.stat() if stat.st_size < 4 or stat.st_size > MAX_OVERLAY_BYTES: raise CanonicalLabOverlayError("canonical LAB overlay size is invalid") @@ -324,6 +326,7 @@ def _render_overlay( frame_times = _frame_times(root, route) layers = route["layers"] archives: dict[str, zipfile.ZipFile] = {} + palettes: dict[str, tuple[int, ...]] = {} recording = rr.RecordingStream(APPLICATION_ID, recording_id=recording_id) try: recording.set_sinks(rr.FileSink(output, write_footer=True)) @@ -346,6 +349,7 @@ def _render_overlay( classes = taxonomy.get("classes") if isinstance(taxonomy, dict) else None if not isinstance(classes, list): raise CanonicalLabOverlayError("semantic taxonomy is invalid") + palettes[layer_id] = _semantic_palette(classes) context = rr.AnnotationContext( [ rr.ClassDescription( @@ -386,7 +390,15 @@ def _render_overlay( for layer_id, mask in masks.items(): recording.log( f"/perception/camera/segmentation/{layer_id}", - rr.SegmentationImage(mask, opacity=0.72, draw_order=1.0), + rr.EncodedImage( + contents=_encoded_semantic_png( + mask, + palettes[layer_id], + ), + media_type="image/png", + opacity=0.72, + draw_order=1.0, + ), ) boxes, labels = semantic_component_boxes(masks["city"], index) if boxes: @@ -414,6 +426,55 @@ def _render_overlay( recording.disconnect() +def _optimize_overlay(path: Path) -> None: + """Compact thousands of per-frame chunks before the browser admits the store.""" + + optimized = path.with_name(f".{path.stem}.{uuid4().hex}.optimized.rrd") + try: + completed = subprocess.run( + [ + sys.executable, + "-m", + "rerun", + "rrd", + "optimize", + "--profile", + "object-store", + "--max-size", + "4MiB", + "--max-rows", + "512", + "--num-pass", + "20", + str(path), + "-o", + str(optimized), + ], + check=False, + capture_output=True, + timeout=120, + ) + if ( + completed.returncode != 0 + or not optimized.is_file() + or optimized.is_symlink() + or optimized.stat().st_size < 4 + or optimized.stat().st_size > MAX_OVERLAY_BYTES + ): + raise CanonicalLabOverlayError( + f"canonical LAB overlay optimization failed: {completed.stderr[-1000:]!r}" + ) + with optimized.open("rb") as stream: + if stream.read(4) != b"RRF2": + raise CanonicalLabOverlayError("canonical LAB overlay optimization is invalid") + os.chmod(optimized, 0o600) + os.replace(optimized, path) + except subprocess.TimeoutExpired as exc: + raise CanonicalLabOverlayError("canonical LAB overlay optimization timed out") from exc + finally: + optimized.unlink(missing_ok=True) + + def _frame_times(root: Path, route: dict[str, Any]) -> np.ndarray: descriptor = route.get("timeline") relative = descriptor.get("path") if isinstance(descriptor, dict) else None @@ -465,6 +526,54 @@ def _read_mask(archive: zipfile.ZipFile, sequence: int) -> np.ndarray: return mask +def _semantic_palette(classes: list[object]) -> tuple[int, ...]: + """Return one complete indexed-PNG palette from a sealed taxonomy.""" + + palette = [0] * (256 * 3) + seen: set[int] = set() + for item in classes: + if not isinstance(item, dict): + raise CanonicalLabOverlayError("semantic taxonomy class is invalid") + class_id = item.get("class_id") + color = item.get("color_rgb") + if ( + not isinstance(class_id, int) + or isinstance(class_id, bool) + or not 0 <= class_id <= 255 + or class_id in seen + or not isinstance(color, list) + or len(color) != 3 + or any( + not isinstance(channel, int) + or isinstance(channel, bool) + or not 0 <= channel <= 255 + for channel in color + ) + ): + raise CanonicalLabOverlayError("semantic taxonomy palette is invalid") + seen.add(class_id) + offset = class_id * 3 + palette[offset : offset + 3] = color + if not seen: + raise CanonicalLabOverlayError("semantic taxonomy palette is empty") + return tuple(palette) + + +def _encoded_semantic_png(mask: np.ndarray, palette: tuple[int, ...]) -> bytes: + """Keep semantic frames compressed in the Rerun store instead of expanding 7 GiB.""" + + if mask.shape != (600, 800) or mask.dtype != np.uint8 or len(palette) != 256 * 3: + raise CanonicalLabOverlayError("semantic frame encoding input is invalid") + image = Image.fromarray(mask) + image.putpalette(palette) + output = io.BytesIO() + image.save(output, format="PNG", compress_level=6) + payload = output.getvalue() + if not payload.startswith(b"\x89PNG\r\n\x1a\n"): + raise CanonicalLabOverlayError("semantic frame encoding failed") + return payload + + def semantic_component_boxes( mask: np.ndarray, _sequence: int, @@ -589,8 +698,11 @@ def _restore_cached( def _artifact_is_regular(artifact: CanonicalLabOverlayArtifact) -> bool: try: + with artifact.path.open("rb") as stream: + magic = stream.read(4) return ( not artifact.path.is_symlink() + and magic == b"RRF2" and artifact.path.stat().st_size == artifact.byte_length and _sha256(artifact.path) == artifact.sha256 ) diff --git a/src/k1link/web/vegetation_shadow_lab_api.py b/src/k1link/web/vegetation_shadow_lab_api.py index 3aa8d64..f3e75d0 100644 --- a/src/k1link/web/vegetation_shadow_lab_api.py +++ b/src/k1link/web/vegetation_shadow_lab_api.py @@ -22,6 +22,7 @@ from PIL import Image from pydantic import BaseModel, ConfigDict, Field from k1link.laboratory.canonical_rerun_overlay import ( + CanonicalLabOverlayArtifact, CanonicalLabOverlayError, _mask_component_boxes, canonical_lab_overlay, @@ -298,12 +299,12 @@ def _build_vegetation_lab_router( }, ) - async def canonical_rerun_overlay_response( + async def canonical_rerun_overlay_artifact( result_id: str, request: CanonicalLabRerunRequest, *, expected_base_generation_sha256: str | None = None, - ) -> FileResponse: + ) -> CanonicalLabOverlayArtifact: """Project LAB-only evidence into the base recording's native clock.""" if ( @@ -349,6 +350,11 @@ def _build_vegetation_lab_router( status_code=503, detail="Canonical LAB Rerun overlay failed verification", ) from exc + return artifact + + def canonical_rerun_overlay_file_response( + artifact: CanonicalLabOverlayArtifact, + ) -> FileResponse: return FileResponse( artifact.path, media_type="application/vnd.rerun.rrd", @@ -366,7 +372,47 @@ def _build_vegetation_lab_router( ) -> FileResponse: """Resolve the sealed sidecar for bounded non-viewer consumers.""" - return await canonical_rerun_overlay_response(result_id, request) + artifact = await canonical_rerun_overlay_artifact(result_id, request) + return canonical_rerun_overlay_file_response(artifact) + + @router.head("/{result_id}/canonical-overlay.rrd") + async def describe_canonical_rerun_overlay( + result_id: str, + application_id: Annotated[ + Literal["nodedc_mission_core_recorded"], + Query(), + ], + recording_id: Annotated[ + str, + Query( + min_length=1, + max_length=128, + pattern=r"^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$", + ), + ], + generation: Annotated[str, Query(pattern=r"^[a-f0-9]{64}$")], + ) -> Response: + """Describe the exact AI generation before Rerun opens its immutable URL.""" + + artifact = await canonical_rerun_overlay_artifact( + result_id, + CanonicalLabRerunRequest( + application_id=application_id, + recording_id=recording_id, + ), + expected_base_generation_sha256=generation, + ) + return Response( + status_code=200, + media_type="application/vnd.rerun.rrd", + headers={ + "Cache-Control": "private, no-store", + "Content-Length": str(artifact.byte_length), + "ETag": f'"{artifact.sha256}"', + "X-Content-Type-Options": "nosniff", + "X-Rerun-Format": "RRF2", + }, + ) @router.get("/{result_id}/canonical-overlay.rrd") async def stream_canonical_rerun_overlay( @@ -384,10 +430,11 @@ def _build_vegetation_lab_router( ), ], generation: Annotated[str, Query(pattern=r"^[a-f0-9]{64}$")], + overlay_generation: Annotated[str, Query(pattern=r"^[a-f0-9]{64}$")], ) -> FileResponse: """Stream one immutable LAB sidecar through Rerun's native HTTP receiver.""" - return await canonical_rerun_overlay_response( + artifact = await canonical_rerun_overlay_artifact( result_id, CanonicalLabRerunRequest( application_id=application_id, @@ -395,6 +442,9 @@ def _build_vegetation_lab_router( ), expected_base_generation_sha256=generation, ) + if overlay_generation != artifact.sha256: + raise HTTPException(status_code=412, detail="Canonical overlay generation changed") + return canonical_rerun_overlay_file_response(artifact) @router.get("/{result_id}/timeline") def get_canonical_route_timeline(result_id: str) -> dict[str, object]: diff --git a/tests/test_vegetation_shadow_lab.py b/tests/test_vegetation_shadow_lab.py index 54bd945..b64b17d 100644 --- a/tests/test_vegetation_shadow_lab.py +++ b/tests/test_vegetation_shadow_lab.py @@ -14,6 +14,7 @@ from fastapi import FastAPI from fastapi.testclient import TestClient from PIL import Image +import k1link.laboratory.canonical_rerun_overlay as canonical_overlay_module import k1link.laboratory.vegetation_policy_review as policy_review_module import k1link.laboratory.vegetation_policy_video as policy_video_module import k1link.laboratory.vegetation_shadow_lab as vegetation_lab_module @@ -22,6 +23,9 @@ from k1link.laboratory import LaboratoryEvidenceRegistry from k1link.laboratory.canonical_rerun_overlay import ( CanonicalLabOverlayArtifact, _artifact_is_regular, + _encoded_semantic_png, + _optimize_overlay, + _semantic_palette, _video_reference_timestamps, ) from k1link.laboratory.evidence_report import verify_laboratory_evidence_result @@ -90,6 +94,64 @@ def test_canonical_overlay_memory_cache_rejects_same_size_tampering( assert not _artifact_is_regular(artifact) +def test_canonical_overlay_keeps_semantics_as_palette_encoded_png() -> None: + mask = np.zeros((600, 800), dtype=np.uint8) + mask[120:420, 200:600] = 7 + palette = _semantic_palette( + [ + {"class_id": 0, "color_rgb": [0, 0, 0]}, + {"class_id": 7, "color_rgb": [255, 47, 128]}, + ] + ) + + encoded = _encoded_semantic_png(mask, palette) + + assert len(encoded) < mask.nbytes // 20 + with Image.open(io.BytesIO(encoded)) as image: + assert image.mode == "P" + assert image.getpixel((0, 0)) == 0 + assert image.getpixel((300, 300)) == 7 + assert image.getpalette()[7 * 3 : 7 * 3 + 3] == [255, 47, 128] + + +def test_canonical_overlay_compacts_chunks_before_cache_publication( + tmp_path: Path, + monkeypatch, +) -> None: + source = tmp_path / "source.rrd" + source.write_bytes(b"RRF2-source") + + def optimize(command: list[str], **options: object) -> SimpleNamespace: + assert command[:4] == [ + canonical_overlay_module.sys.executable, + "-m", + "rerun", + "rrd", + ] + assert command[4:13] == [ + "optimize", + "--profile", + "object-store", + "--max-size", + "4MiB", + "--max-rows", + "512", + "--num-pass", + "20", + ] + assert command[13] == str(source) + assert command[14] == "-o" + Path(command[15]).write_bytes(b"RRF2-optimized") + assert options == {"check": False, "capture_output": True, "timeout": 120} + return SimpleNamespace(returncode=0, stderr=b"") + + monkeypatch.setattr(canonical_overlay_module.subprocess, "run", optimize) + + _optimize_overlay(source) + + assert source.read_bytes() == b"RRF2-optimized" + + def test_canonical_overlay_get_is_generation_bound_and_range_streamable( tmp_path: Path, monkeypatch, @@ -148,12 +210,28 @@ def test_canonical_overlay_get_is_generation_bound_and_range_streamable( ) client = TestClient(app) endpoint = f"/api/v1/laboratory/vegetation-shadow/{result_id}/canonical-overlay.rrd" + descriptor = client.head( + endpoint, + params={ + "application_id": "nodedc_mission_core_recorded", + "recording_id": recording_id, + "generation": generation, + }, + ) + assert descriptor.status_code == 200 + assert descriptor.content == b"" + assert descriptor.headers["content-length"] == str(artifact.byte_length) + assert descriptor.headers["etag"] == f'"{artifact.sha256}"' + assert descriptor.headers["x-rerun-format"] == "RRF2" + assert descriptor.headers["cache-control"] == "private, no-store" + response = client.get( endpoint, params={ "application_id": "nodedc_mission_core_recorded", "recording_id": recording_id, "generation": generation, + "overlay_generation": artifact.sha256, }, headers={"Range": "bytes=0-3"}, ) @@ -163,7 +241,7 @@ def test_canonical_overlay_get_is_generation_bound_and_range_streamable( assert response.headers["etag"] == f'"{artifact.sha256}"' assert response.headers["cache-control"].endswith("immutable") - stale = client.get( + stale = client.head( endpoint, params={ "application_id": "nodedc_mission_core_recorded", @@ -173,6 +251,17 @@ def test_canonical_overlay_get_is_generation_bound_and_range_streamable( ) assert stale.status_code == 412 + stale_overlay = client.get( + endpoint, + params={ + "application_id": "nodedc_mission_core_recorded", + "recording_id": recording_id, + "generation": generation, + "overlay_generation": "d" * 64, + }, + ) + assert stale_overlay.status_code == 412 + def test_route_playback_chunk_descriptor_seals_only_requested_binary_window() -> None: points = np.arange(18, dtype="