From b688347778f23e2c14fb11fb391055acc0704b55 Mon Sep 17 00:00:00 2001 From: DCCONSTRUCTIONS Date: Wed, 26 Aug 2026 13:10:11 +0300 Subject: [PATCH] refactor(lab): admit versioned graph ledgers --- src/k1link/perception/m48s_replay_timeline.py | 80 ++++++++++++++----- 1 file changed, 60 insertions(+), 20 deletions(-) diff --git a/src/k1link/perception/m48s_replay_timeline.py b/src/k1link/perception/m48s_replay_timeline.py index d2d66e1..2840a77 100644 --- a/src/k1link/perception/m48s_replay_timeline.py +++ b/src/k1link/perception/m48s_replay_timeline.py @@ -38,13 +38,13 @@ from .threat_timeline import ( ) FRAME_EVIDENCE_SCHEMA: Final = "missioncore.m48s-reference-graph-frame-evidence/v0" +M48R3_FRAME_EVIDENCE_SCHEMA: Final = ( + "missioncore.m48s-reference-graph-frame-evidence/v1" +) EXPECTED_FRAME_COUNT: Final = 4_489 CAMERA_ACCUMULATION_WINDOW_SECONDS: Final = 2.0 CAMERA_ACCUMULATION_POINT_LIMIT: Final = 20_000 CAMERA_POINT_OVERLAY_SCHEMA: Final = "missioncore.m48s-camera-point-overlay/v1" -_FRAME_EVIDENCE_SCHEMA_MARKER: Final = ( - b'"schema_version":"missioncore.m48s-reference-graph-frame-evidence/v0"' -) _SOURCE_ENVELOPE_MARKER: Final = b'"source_envelope":' _JSON_DECODER: Final = json.JSONDecoder() @@ -61,16 +61,37 @@ class _LedgerIndex: class M48sReplayTimeline: """Read source-indexed chunks while preserving latest-wins world-state gaps.""" - def __init__(self, *, repository_root: Path, result_root: Path, result_id: str) -> None: + def __init__( + self, + *, + repository_root: Path, + result_root: Path, + result_id: str, + frames_name: str = "reference-graph-replay-frames.jsonl", + worker_result_name: str = "reference-graph-replay-worker-result.json", + frame_evidence_schema: str = FRAME_EVIDENCE_SCHEMA, + camera_endpoint_root: str = ( + "/api/v1/laboratory/m48s/fixed-class-detector" + ), + ) -> None: self.repository_root = repository_root.resolve(strict=True) self.result_root = result_root.resolve(strict=True) self.result_id = result_id - self.frames_path = ( - self.result_root / "reference-graph-replay-frames.jsonl" - ).resolve(strict=True) - self.worker_path = ( - self.result_root / "reference-graph-replay-worker-result.json" - ).resolve(strict=True) + if ( + Path(frames_name).name != frames_name + or Path(worker_result_name).name != worker_result_name + or frame_evidence_schema not in { + FRAME_EVIDENCE_SCHEMA, + M48R3_FRAME_EVIDENCE_SCHEMA, + } + or not camera_endpoint_root.startswith("/api/v1/laboratory/") + or camera_endpoint_root.endswith("/") + ): + raise M48sReplayTimelineError("M4.8S timeline binding is invalid") + self.frame_evidence_schema = frame_evidence_schema + self.camera_endpoint_root = camera_endpoint_root + self.frames_path = (self.result_root / frames_name).resolve(strict=True) + self.worker_path = (self.result_root / worker_result_name).resolve(strict=True) if ( self.frames_path.parent != self.result_root or self.worker_path.parent != self.result_root @@ -96,7 +117,12 @@ class M48sReplayTimeline: raise M48sReplayTimelineError("M4.8S source clock is not monotonic") worker = _object(json.loads(self.worker_path.read_text("utf-8")), "worker result") self.outcomes = _terminal_outcomes(worker) - self.index = _index_ledger(self.frames_path, self.source_times_ns, self.outcomes) + self.index = _index_ledger( + self.frames_path, + self.source_times_ns, + self.outcomes, + frame_evidence_schema=self.frame_evidence_schema, + ) self._cache_lock = Lock() self._chunk_json_cache: OrderedDict[tuple[int, int], bytes] = OrderedDict() self._camera_point_json_cache: OrderedDict[int, bytes] = OrderedDict() @@ -432,8 +458,8 @@ class M48sReplayTimeline: "camera_proposals": camera_proposals, "decision_counts": _decision_counts(assessments), "camera_url": ( - "/api/v1/laboratory/m48s/fixed-class-detector/" - f"{self.result_id}/timeline/frames/{sequence}/camera" + f"{self.camera_endpoint_root}/{self.result_id}" + f"/timeline/frames/{sequence}/camera" ), "ground_truth": False, "authority": "replay-simulated", @@ -447,7 +473,10 @@ class M48sReplayTimeline: stream.seek(offset) line = stream.readline() value = json.loads(line) - if not isinstance(value, dict) or value.get("schema_version") != FRAME_EVIDENCE_SCHEMA: + if ( + not isinstance(value, dict) + or value.get("schema_version") != self.frame_evidence_schema + ): raise M48sReplayTimelineError("M4.8S frame row is invalid") return value @@ -456,6 +485,8 @@ def _index_ledger( path: Path, source_times_ns: tuple[int, ...], outcomes: dict[int, str], + *, + frame_evidence_schema: str, ) -> _LedgerIndex: offsets: dict[int, int] = {} with path.open("rb") as stream: @@ -464,7 +495,10 @@ def _index_ledger( line = stream.readline() if not line: break - envelope = _ledger_source_envelope(line) + envelope = _ledger_source_envelope( + line, + frame_evidence_schema=frame_evidence_schema, + ) timestamps = _object(envelope.get("timestamps"), "source timestamps") sequence = envelope.get("sequence") if ( @@ -483,7 +517,11 @@ def _index_ledger( return _LedgerIndex(offsets) -def _ledger_source_envelope(line: bytes) -> dict[str, object]: +def _ledger_source_envelope( + line: bytes, + *, + frame_evidence_schema: str, +) -> dict[str, object]: """Validate a ledger row while decoding only its small trailing envelope. The full row can exceed 100 KiB because it contains the delivered world @@ -492,10 +530,10 @@ def _ledger_source_envelope(line: bytes) -> dict[str, object]: the GIL for many seconds. """ - if ( - line.count(_FRAME_EVIDENCE_SCHEMA_MARKER) != 1 - or line.count(_SOURCE_ENVELOPE_MARKER) != 1 - ): + schema_marker = ( + f'"schema_version":"{frame_evidence_schema}"'.encode("ascii") + ) + if line.count(schema_marker) != 1 or line.count(_SOURCE_ENVELOPE_MARKER) != 1: raise M48sReplayTimelineError("M4.8S ledger schema changed") start = line.find(_SOURCE_ENVELOPE_MARKER) + len(_SOURCE_ENVELOPE_MARKER) try: @@ -569,6 +607,8 @@ __all__ = [ "CAMERA_ACCUMULATION_POINT_LIMIT", "CAMERA_ACCUMULATION_WINDOW_SECONDS", "CAMERA_POINT_OVERLAY_SCHEMA", + "FRAME_EVIDENCE_SCHEMA", + "M48R3_FRAME_EVIDENCE_SCHEMA", "M48sReplayTimeline", "M48sReplayTimelineError", ]