fix(replay): stream sealed perception overlays

This commit is contained in:
DCCONSTRUCTIONS
2026-07-30 00:51:14 +03:00
parent 091ab2671e
commit 2ab9e45548
10 changed files with 380 additions and 44 deletions
+73 -20
View File
@@ -26,7 +26,10 @@ from k1link.artifacts import write_json_atomic
from k1link.sessions import SessionIntegrityError
from .jobs import CameraComputeJob, validate_camera_compute_job
from .results import RecordedPerceptionOverlayError
from .results import (
RecordedPerceptionOverlayArtifact,
RecordedPerceptionOverlayError,
)
RESULT_SCHEMA = "missioncore.e10-integrated-perception-result/v1"
IDENTITY_SCHEMA = "missioncore.e10-integrated-perception-identity/v1"
@@ -397,6 +400,39 @@ class IntegratedPerceptionOverlayStore:
application_id: str,
recording_id: str,
) -> bytes | None:
"""Compatibility adapter for in-process consumers that require bytes."""
artifact = self.materialize(
session_id,
application_id=application_id,
recording_id=recording_id,
)
if artifact is None:
return None
try:
payload = artifact.path.read_bytes()
except OSError as exc:
raise RecordedPerceptionOverlayError(
"integrated perception cache became unavailable"
) from exc
if (
len(payload) != artifact.byte_length
or hashlib.sha256(payload).hexdigest() != artifact.sha256
):
raise RecordedPerceptionOverlayError(
"integrated perception cache changed after validation"
)
return payload
def materialize(
self,
session_id: str,
*,
application_id: str,
recording_id: str,
) -> RecordedPerceptionOverlayArtifact | None:
"""Return a verified file-backed overlay without copying it into heap."""
if application_id != "nodedc_mission_core_recorded":
raise ValueError("integrated perception application id is invalid")
if (
@@ -414,7 +450,7 @@ class IntegratedPerceptionOverlayStore:
key,
state="ready",
phase="ready",
byte_length=len(cached),
byte_length=cached.byte_length,
)
return cached
self._set_status(key, state="preparing", phase="queued")
@@ -431,7 +467,7 @@ class IntegratedPerceptionOverlayStore:
)
output = cache / f"{recording_id}.rrd"
sidecar = output.with_suffix(".rrd.cache.json")
cached = _read_cache(
cached = _read_cache_artifact(
output,
sidecar,
result_id=result.result_id,
@@ -442,7 +478,7 @@ class IntegratedPerceptionOverlayStore:
key,
state="ready",
phase="ready",
byte_length=len(cached),
byte_length=cached.byte_length,
)
return cached
self._set_status(key, state="preparing", phase="rendering")
@@ -483,7 +519,17 @@ class IntegratedPerceptionOverlayStore:
phase="ready",
byte_length=len(payload),
)
return payload
published = _read_cache_artifact(
output,
sidecar,
result_id=result.result_id,
recording_id=recording_id,
)
if published is None:
raise RecordedPerceptionOverlayError(
"integrated perception cache publication failed"
)
return published
except BaseException:
self._set_status(key, state="error", phase="error")
raise
@@ -556,7 +602,7 @@ class IntegratedPerceptionOverlayStore:
self,
session_id: str,
recording_id: str,
) -> bytes | None:
) -> RecordedPerceptionOverlayArtifact | None:
cache_root = _private_directory(self.cache_root)
session_cache = _private_child(cache_root, session_id)
admission = self._load_admission(session_cache, session_id)
@@ -570,15 +616,15 @@ class IntegratedPerceptionOverlayStore:
recording_id,
)
if recovered is not None:
_, payload = recovered
return payload
_, artifact = recovered
return artifact
admission = None
if admission is None or admission.result_id != latest.result_id:
return None
result_cache = _private_child(session_cache, admission.result_id)
output = result_cache / f"{recording_id}.rrd"
sidecar = output.with_suffix(".rrd.cache.json")
return _read_cache(
return _read_cache_artifact(
output,
sidecar,
result_id=admission.result_id,
@@ -631,17 +677,17 @@ class IntegratedPerceptionOverlayStore:
session_cache: Path,
descriptor: _ResultDescriptor,
recording_id: str,
) -> tuple[_OverlayAdmission, bytes] | None:
) -> tuple[_OverlayAdmission, RecordedPerceptionOverlayArtifact] | None:
cache = _private_child(session_cache, descriptor.result_id)
output = cache / f"{recording_id}.rrd"
sidecar = output.with_suffix(".rrd.cache.json")
payload = _read_cache(
artifact = _read_cache_artifact(
output,
sidecar,
result_id=descriptor.result_id,
recording_id=recording_id,
)
if payload is None:
if artifact is None:
return None
admission = _OverlayAdmission(
session_id=descriptor.session_id,
@@ -650,7 +696,7 @@ class IntegratedPerceptionOverlayStore:
result_created_at_utc=descriptor.created_at_utc,
)
self._write_admission_document(session_cache, admission)
return admission, payload
return admission, artifact
def _write_admission(self, result: IntegratedPerceptionResult) -> None:
descriptor = _read_result_descriptor(result.result_root)
@@ -1183,16 +1229,19 @@ def _read_result_descriptor(result_root: Path) -> _ResultDescriptor:
)
def _read_cache(
def _read_cache_artifact(
output: Path,
sidecar: Path,
*,
result_id: str,
recording_id: str,
) -> bytes | None:
) -> RecordedPerceptionOverlayArtifact | None:
try:
value = _read_object(sidecar, sidecar.parent)
payload = output.read_bytes()
metadata = _confined_file(output, sidecar.parent)
with output.open("rb") as stream:
magic = stream.read(4)
digest = _sha256(output)
except (OSError, SessionIntegrityError):
return None
if (
@@ -1200,12 +1249,16 @@ def _read_cache(
or value.get("renderer_version") != OVERLAY_RENDERER_VERSION
or value.get("result_id") != result_id
or value.get("recording_id") != recording_id
or value.get("byte_length") != len(payload)
or value.get("sha256") != hashlib.sha256(payload).hexdigest()
or not payload.startswith(b"RRF2")
or value.get("byte_length") != metadata.st_size
or value.get("sha256") != digest
or magic != b"RRF2"
):
return None
return payload
return RecordedPerceptionOverlayArtifact(
path=output.resolve(strict=True),
byte_length=metadata.st_size,
sha256=digest,
)
def _read_object(path: Path, root: Path) -> dict[str, Any]: