refactor(lab): admit versioned graph ledgers

This commit is contained in:
DCCONSTRUCTIONS
2026-08-26 13:10:11 +03:00
parent 252cbc8495
commit b688347778
+60 -20
View File
@@ -38,13 +38,13 @@ from .threat_timeline import (
) )
FRAME_EVIDENCE_SCHEMA: Final = "missioncore.m48s-reference-graph-frame-evidence/v0" 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 EXPECTED_FRAME_COUNT: Final = 4_489
CAMERA_ACCUMULATION_WINDOW_SECONDS: Final = 2.0 CAMERA_ACCUMULATION_WINDOW_SECONDS: Final = 2.0
CAMERA_ACCUMULATION_POINT_LIMIT: Final = 20_000 CAMERA_ACCUMULATION_POINT_LIMIT: Final = 20_000
CAMERA_POINT_OVERLAY_SCHEMA: Final = "missioncore.m48s-camera-point-overlay/v1" 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":' _SOURCE_ENVELOPE_MARKER: Final = b'"source_envelope":'
_JSON_DECODER: Final = json.JSONDecoder() _JSON_DECODER: Final = json.JSONDecoder()
@@ -61,16 +61,37 @@ class _LedgerIndex:
class M48sReplayTimeline: class M48sReplayTimeline:
"""Read source-indexed chunks while preserving latest-wins world-state gaps.""" """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.repository_root = repository_root.resolve(strict=True)
self.result_root = result_root.resolve(strict=True) self.result_root = result_root.resolve(strict=True)
self.result_id = result_id self.result_id = result_id
self.frames_path = ( if (
self.result_root / "reference-graph-replay-frames.jsonl" Path(frames_name).name != frames_name
).resolve(strict=True) or Path(worker_result_name).name != worker_result_name
self.worker_path = ( or frame_evidence_schema not in {
self.result_root / "reference-graph-replay-worker-result.json" FRAME_EVIDENCE_SCHEMA,
).resolve(strict=True) 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 ( if (
self.frames_path.parent != self.result_root self.frames_path.parent != self.result_root
or self.worker_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") raise M48sReplayTimelineError("M4.8S source clock is not monotonic")
worker = _object(json.loads(self.worker_path.read_text("utf-8")), "worker result") worker = _object(json.loads(self.worker_path.read_text("utf-8")), "worker result")
self.outcomes = _terminal_outcomes(worker) 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._cache_lock = Lock()
self._chunk_json_cache: OrderedDict[tuple[int, int], bytes] = OrderedDict() self._chunk_json_cache: OrderedDict[tuple[int, int], bytes] = OrderedDict()
self._camera_point_json_cache: OrderedDict[int, bytes] = OrderedDict() self._camera_point_json_cache: OrderedDict[int, bytes] = OrderedDict()
@@ -432,8 +458,8 @@ class M48sReplayTimeline:
"camera_proposals": camera_proposals, "camera_proposals": camera_proposals,
"decision_counts": _decision_counts(assessments), "decision_counts": _decision_counts(assessments),
"camera_url": ( "camera_url": (
"/api/v1/laboratory/m48s/fixed-class-detector/" f"{self.camera_endpoint_root}/{self.result_id}"
f"{self.result_id}/timeline/frames/{sequence}/camera" f"/timeline/frames/{sequence}/camera"
), ),
"ground_truth": False, "ground_truth": False,
"authority": "replay-simulated", "authority": "replay-simulated",
@@ -447,7 +473,10 @@ class M48sReplayTimeline:
stream.seek(offset) stream.seek(offset)
line = stream.readline() line = stream.readline()
value = json.loads(line) 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") raise M48sReplayTimelineError("M4.8S frame row is invalid")
return value return value
@@ -456,6 +485,8 @@ def _index_ledger(
path: Path, path: Path,
source_times_ns: tuple[int, ...], source_times_ns: tuple[int, ...],
outcomes: dict[int, str], outcomes: dict[int, str],
*,
frame_evidence_schema: str,
) -> _LedgerIndex: ) -> _LedgerIndex:
offsets: dict[int, int] = {} offsets: dict[int, int] = {}
with path.open("rb") as stream: with path.open("rb") as stream:
@@ -464,7 +495,10 @@ def _index_ledger(
line = stream.readline() line = stream.readline()
if not line: if not line:
break break
envelope = _ledger_source_envelope(line) envelope = _ledger_source_envelope(
line,
frame_evidence_schema=frame_evidence_schema,
)
timestamps = _object(envelope.get("timestamps"), "source timestamps") timestamps = _object(envelope.get("timestamps"), "source timestamps")
sequence = envelope.get("sequence") sequence = envelope.get("sequence")
if ( if (
@@ -483,7 +517,11 @@ def _index_ledger(
return _LedgerIndex(offsets) 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. """Validate a ledger row while decoding only its small trailing envelope.
The full row can exceed 100 KiB because it contains the delivered world 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. the GIL for many seconds.
""" """
if ( schema_marker = (
line.count(_FRAME_EVIDENCE_SCHEMA_MARKER) != 1 f'"schema_version":"{frame_evidence_schema}"'.encode("ascii")
or line.count(_SOURCE_ENVELOPE_MARKER) != 1 )
): if line.count(schema_marker) != 1 or line.count(_SOURCE_ENVELOPE_MARKER) != 1:
raise M48sReplayTimelineError("M4.8S ledger schema changed") raise M48sReplayTimelineError("M4.8S ledger schema changed")
start = line.find(_SOURCE_ENVELOPE_MARKER) + len(_SOURCE_ENVELOPE_MARKER) start = line.find(_SOURCE_ENVELOPE_MARKER) + len(_SOURCE_ENVELOPE_MARKER)
try: try:
@@ -569,6 +607,8 @@ __all__ = [
"CAMERA_ACCUMULATION_POINT_LIMIT", "CAMERA_ACCUMULATION_POINT_LIMIT",
"CAMERA_ACCUMULATION_WINDOW_SECONDS", "CAMERA_ACCUMULATION_WINDOW_SECONDS",
"CAMERA_POINT_OVERLAY_SCHEMA", "CAMERA_POINT_OVERLAY_SCHEMA",
"FRAME_EVIDENCE_SCHEMA",
"M48R3_FRAME_EVIDENCE_SCHEMA",
"M48sReplayTimeline", "M48sReplayTimeline",
"M48sReplayTimelineError", "M48sReplayTimelineError",
] ]