diff --git a/docs/13_LIDAR_WORKER_PRODUCT_AND_ROADMAP.md b/docs/13_LIDAR_WORKER_PRODUCT_AND_ROADMAP.md index 2cc79fc..f8af7fa 100644 --- a/docs/13_LIDAR_WORKER_PRODUCT_AND_ROADMAP.md +++ b/docs/13_LIDAR_WORKER_PRODUCT_AND_ROADMAP.md @@ -645,19 +645,26 @@ with a complete `missioncore.laboratory-method/v1` manifest. The current E30 AI-assisted engineering generation is `e30-engineering-generation-62a4fea10dea9b77f69ceac1af5bf0e4928d9c7716083c22258a03670fe5bd4f`. It covers all `486` selected items: `403` confirmed, `81` corrected and `2` -retained as `insufficient-evidence` human exceptions. The earlier five-case -generation and its human draft remain immutable history. Frames 1213, 162 and -1823 are now automatic engineering outcomes; only geometry frames 2622 and -4147 require a human decision. It found no systematic camera↔LiDAR -registration failure in the review set. The dominant actionable signals are -detector errors (`110`, including barrier false positives and 21 missed -class-bearing objects), sparse occupied support (`98`), held-track time -freshness (`92`) and self points (`10`). This is diagnostic engineering -evidence, not human ground truth, detection accuracy or safety acceptance. The -exception budget is 1% only for diagnostic residual work and never overrides a -repeated-cause or high-impact blocker. The current queue is `2 / 486` (`0.41%`); -E31 remains blocked until those two decisions freeze the E30 minimum correction -set. +initially routed human exceptions. The immutable human generation +`e30-review-generation-7982a882558d0be690b4c7092e328c080bfcbf52478a220452be7e887a588250` +closes both: geometry `2622:4` is background/noise and `4147:3` is a real +occupied object. The earlier five-case generation and drafts remain immutable +history. E30 found no systematic camera↔LiDAR registration failure in the +review set. The dominant actionable signals remain detector errors (`110`, +including barrier false positives and 21 missed class-bearing objects), sparse +occupied support (`98`), held-track time freshness (`92`) and self points +(`10`). This is diagnostic engineering evidence, not human ground truth, +detection accuracy or safety acceptance. + +E31 accepted immutable source-scoped profile +`e31-source-qualification-b2460a5eb143688c7eea6821b2277e13aea79868abe81d83f7e78548c119159a`. +E32 then produced +`e32-track-geometry-a14ca0e7fb3850ca0dfa3c41634e1b490a2d58ab74d101afc6d6921fbdb0e6fd`: +all 4,489 E29 frames reproduce exactly, all source/object/point claims close, +the self and exact correction set is applied without a generic geometry mask, +and exclusive PointSlab ownership is enforced. Conflict count remains `38`; +the result is accepted only as the source-scoped diagnostic/shadow input for +E33 and is not promoted as a detector-accuracy or staleness improvement. Execution is strictly sequential through E33: E30 determines what E31 is allowed to change; E31 determines the E32 profile; E32 determines the E33 diff --git a/docs/16_ARCHITECTURE_AUDIT_EXECUTION_ROADMAP.md b/docs/16_ARCHITECTURE_AUDIT_EXECUTION_ROADMAP.md index 013fb44..d888667 100644 --- a/docs/16_ARCHITECTURE_AUDIT_EXECUTION_ROADMAP.md +++ b/docs/16_ARCHITECTURE_AUDIT_EXECUTION_ROADMAP.md @@ -171,9 +171,31 @@ RAVNOVES00 diagnostic binding and does not claim that host arrival is hardware firing time. The factory KB4 identity is exact, but a measured calibration-target residual and physical body/mount dimensions are unavailable. The accepted profile is therefore source-session scoped and -cannot transfer to another mount. A6/E32 full replay is the next critical-path -implementation and must bind every published frame to both `TrackGeometry v1` -and the accepted E31 profile. +cannot transfer to another mount. + +A6/E32 is complete in immutable result +`e32-track-geometry-a14ca0e7fb3850ca0dfa3c41634e1b490a2d58ab74d101afc6d6921fbdb0e6fd`. +It reproduces all 4,489 E29 frames exactly before correction and binds every +published frame to `TrackGeometry v1` plus the accepted E31 profile. All +20,513 semantic observations, 21,321 geometry clusters and 2,125,813 qualified +point claims close without hidden loss. The replay removes 385 source-scoped +self observations and three evidenced geometry clusters, withholds 1,558 +unqualified ranges and resolves 4,461 overlapping point claims. Of 709 +arbitrated semantic observations, 563 retain `agree` and 146 conservatively +become `unknown`. The 38 E29 conflicts remain 38; A6 is accepted as a +diagnostic/shadow contract, not as detector-accuracy improvement. A7/E33 +recorded-source-paced worker execution is now the critical path. + +- [x] Reproduce all 4,489 immutable E29 frames with the exact frozen profile + before applying E31/E30 changes. +- [x] Apply only the admitted semantic self-mask, two exact geometry + corrections and the complete A3 human exception set. +- [x] Enforce exclusive PointSlab ownership and journal every arbitration, + excluded claim and conservative state transition. +- [x] Compare E29/E32 by status, class, range, 60-second scene interval and + cause while preserving source availability and unknown/free-space policy. +- [x] Add a digest-bound compact binary PointSlab encoding and strict + TrackGeometryFrame reconstruction/validation. ### A3 residual and human-exception policy diff --git a/docs/adr/0025-e32-full-replay-and-exclusive-point-ownership.md b/docs/adr/0025-e32-full-replay-and-exclusive-point-ownership.md new file mode 100644 index 0000000..bbcc77c --- /dev/null +++ b/docs/adr/0025-e32-full-replay-and-exclusive-point-ownership.md @@ -0,0 +1,92 @@ +# ADR 0025 — E32 full replay and exclusive point ownership + +Date: 2026-07-27 + +Status: accepted for source-scoped diagnostic/shadow use + +## Context + +A5 defined `TrackGeometry v1` and `PointSlab`, but E29 did not retain the exact +frame-local point indices that produced each semantic and geometry result. +E32 therefore had to recover those indices without changing the frozen E29 +thresholds, apply only the accepted E31/E30 corrections and prove that the +translated result had no hidden source loss. + +The accepted E31 profile is limited to the immutable RAVNOVES00 source session. +It admits one normalized person self-mask, rejects a generic geometry mask and +binds two exact geometry self-corrections. The A3 human generation additionally +rejects geometry cluster `2622:4` as background/noise and retains `4147:3` as a +real occupied object. + +## Decision + +1. E32 recomputes E29 support indices using the exact E29 profile and immutable + camera, LiDAR, pose and local-surface inputs. +2. Every recomputed frame must equal the stored E29 frame before any E31/E30 + correction is applied. A difference in any observation, cluster, metric, + policy or frame binding fails the replay closed. +3. E32 applies only: + - the admitted E31 person self-mask by normalized bbox centre; + - the two exact E31 geometry correction item locators; + - the complete immutable A3 human geometry dispositions. +4. The rejected generic E31 geometry mask is never applied. E29 thresholds are + never retuned and excluded points are never reclassified as free space. +5. `PointSlab` grants one owner to each published source point. When two E29 + camera observations claim the same point, E32 uses a deterministic + non-threshold policy: smallest bbox, then higher detector score, lower track + id and lower observation ordinal. A track that retains points stays `agree`; + a track that loses all exclusive points becomes explicit `unknown`. +6. Camera-only and conflict observations cannot publish an E29 range derived + from support that failed the E29 qualification threshold. Their semantic + observation remains, but `metric_basis=unavailable` and `range_m=null`. +7. A held observation is published only when an earlier current observation + for the same track exists in the replay. A held observation without such + provenance is explicitly journalled and omitted instead of inventing + `held_from_frame_index`. +8. Compact storage uses deterministic `.npy` arrays for frame offsets, + frame-local source indices, map-frame float32 points and local owner + indices. The JSONL frame record retains the complete TrackGeometry table and + PointSlab reference. The public reader reconstructs and validates the exact + `TrackGeometryFrame`. +9. Every excluded or altered E29 product is written to + `e29-e32-changes.jsonl` with its source locator, before/after state, reason, + decision identity when applicable and affected claim count. +10. E32 remains diagnostic. It grants no command, navigation, traversability, + free-space or safety authority and does not modify the persistent map. + +## Accepted result + +The immutable result is +`e32-track-geometry-a14ca0e7fb3850ca0dfa3c41634e1b490a2d58ab74d101afc6d6921fbdb0e6fd`. + +- 4,489 / 4,489 E29 frames reproduce exactly. +- Source availability remains 3,928 available and 561 unavailable frames. +- 20,513 semantic observations close as 20,119 published, 385 source-scoped + self-mask exclusions and 9 held observations without prior-current + provenance. +- 21,321 geometry clusters close as 21,318 published, two exact E31 + self-corrections and one A3 human background/noise exclusion. The A3 + object-present cluster remains published. +- 2,125,813 E29 qualified point claims close as 2,119,302 published rows, + 2,050 explicitly excluded claims and 4,461 overlapping claims removed by + ownership arbitration. +- 709 semantic observations required point-ownership arbitration: 563 retained + `agree`; 146 became `unknown`. +- 1,558 ranges backed only by unqualified semantic support are withheld. +- E29/E32 conflict count remains 38 and no conflict observation changes. E32 + is not presented as a detector-accuracy improvement. +- `agree` changes from 6,341 to 6,195, `single-source-camera` from 13,246 to + 12,861, `unknown` from 888 to 1,025 and geometry-only from 21,321 to 21,318. + +## Consequences + +- A6 is complete as a reproducible, source-scoped TrackGeometry input for the + A7/E33 recorded-source-paced worker shadow. +- E33 receives a closed ownership and source-accounting contract instead of + ambiguous E29 point claims. +- The unchanged 38 conflicts and retained held observations remain perception + limitations. Runtime qualification must not describe them as corrected. +- The 146 conservative `agree → unknown` transitions are an intentional + consequence of exclusive ownership, not threshold degradation. +- A second source or changed mount must pass a new qualification; this E32 + result cannot be transferred by assumption. diff --git a/src/k1link/compute/__init__.py b/src/k1link/compute/__init__.py index 73d9ac1..3c14017 100644 --- a/src/k1link/compute/__init__.py +++ b/src/k1link/compute/__init__.py @@ -17,6 +17,16 @@ from .e31_source_qualification import ( E31SourceQualificationProfile, build_e31_source_qualification, ) +from .e32_track_geometry_replay import ( + E32_TRACK_GEOMETRY_RECORD_SCHEMA, + E32_TRACK_GEOMETRY_REPLAY_SCHEMA, + E32_TRACK_GEOMETRY_REPORT_SCHEMA, + E32TrackGeometryReplay, + E32TrackGeometryReplayError, + build_e32_track_geometry_replay, + e32_track_geometry_frame, + read_e32_track_geometry_replay, +) from .evaluation_pack import ( ANNOTATION_CONTRACT_SCHEMA, EVALUATION_PACK_SCHEMA, @@ -354,6 +364,14 @@ __all__ = [ "assess_lidar_profile", "build_lidar_replay_pack_v2", "build_e31_source_qualification", + "build_e32_track_geometry_replay", + "read_e32_track_geometry_replay", + "e32_track_geometry_frame", + "E32_TRACK_GEOMETRY_REPLAY_SCHEMA", + "E32_TRACK_GEOMETRY_REPORT_SCHEMA", + "E32_TRACK_GEOMETRY_RECORD_SCHEMA", + "E32TrackGeometryReplay", + "E32TrackGeometryReplayError", "build_lidar_ground_annotation_template", "build_lidar_ground_benchmark", "build_k1_local_surface", diff --git a/src/k1link/compute/e32_track_geometry_replay.py b/src/k1link/compute/e32_track_geometry_replay.py new file mode 100644 index 0000000..9efcc72 --- /dev/null +++ b/src/k1link/compute/e32_track_geometry_replay.py @@ -0,0 +1,2006 @@ +"""Full E32 replay from immutable E29 evidence into TrackGeometry v1. + +E32 does not retune E29. It reproduces every E29 frame with the exact frozen +profile, applies only the admitted E31 source-scoped mask and the immutable E30 +decision chain, then publishes contract-validated TrackGeometry frames. Point +ownership is stored in deterministic ``.npy`` slabs so the frame stream stays +compact without weakening the PointSlab semantics. +""" + +from __future__ import annotations + +import hashlib +import json +import math +import os +import re +import shutil +import time +from collections import Counter +from collections.abc import Iterable, Mapping +from dataclasses import dataclass, fields +from pathlib import Path +from typing import Any, Final, cast + +import numpy as np +import numpy.typing as npt + +from k1link.artifacts import utc_now_iso +from k1link.device_plugins.xgrids_k1.analyze.calibrated_projection import ( + project_map_points_kb4, +) + +from .e31_source_qualification import ( + E31_BINDING_PROFILE_SCHEMA, + _load_e30_chain, + _safe_result_root, + _validate_content_identity, + _verified_artifacts, +) +from .e32_track_geometry_storage import ( + E32_FRAME_OFFSETS_NAME, + E32_FRAMES_NAME, + E32_OWNER_INDICES_NAME, + E32_POINT_SLAB_REFERENCE_SCHEMA, + E32_POINTS_NAME, + E32_SOURCE_INDICES_NAME, + E32_TRACK_GEOMETRY_RECORD_SCHEMA, + E32TrackGeometryStorageError, + frame_from_record, + load_point_storage, + record, + validate_storage, + write_point_storage, +) +from .lidar_field_review import E10LidarFieldSource +from .lidar_local_surface import K1LocalSurfaceV1 +from .semantic_geometry_fusion import ( + CAMERA_GEOMETRY_FRAME_SCHEMA, + CAMERA_GEOMETRY_FUSION_SCHEMA, + CameraGeometryFusionProfile, + _claimed_indices, + _fusion_frame, + _geometry_cluster_supports, + _GeometryClusterSupport, + _projection_profile, + _semantic_support, + _SemanticSupport, +) +from .sensor_representation import K1_LIO_PCL_CAPABILITIES +from .track_geometry import ( + POINT_SLAB_SCHEMA, + TRACK_GEOMETRY_FRAME_SCHEMA, + TRACK_GEOMETRY_SCHEMA, + PointSlab, + TrackGeometry, + TrackGeometryCurrentness, + TrackGeometryEvidenceState, + TrackGeometryFrame, + TrackGeometryMetricBasis, + TrackGeometryOwnerKind, + TrackGeometrySourceBinding, +) + +E32_TRACK_GEOMETRY_REPLAY_SCHEMA: Final = "missioncore.e32-track-geometry-replay/v1" +E32_TRACK_GEOMETRY_REPORT_SCHEMA: Final = "missioncore.e32-track-geometry-report/v1" +E32_CHANGE_SCHEMA: Final = "missioncore.e32-e29-change/v1" + +E32_MANIFEST_NAME: Final = "manifest.json" +E32_REPORT_NAME: Final = "track-geometry-report.json" +E32_CHANGES_NAME: Final = "e29-e32-changes.jsonl" + +_E29_RESULT_ID = re.compile(r"^e29-camera-geometry-[a-f0-9]{64}$") +_E31_RESULT_ID = re.compile(r"^e31-source-qualification-[a-f0-9]{64}$") +_E32_RESULT_ID = re.compile(r"^e32-track-geometry-[a-f0-9]{64}$") +_E10_RESULT_ID = re.compile(r"^e10-integrated-perception-[a-f0-9]{64}$") +_SHA256 = re.compile(r"^[a-f0-9]{64}$") + +Float64Array = npt.NDArray[np.float64] +Int64Array = npt.NDArray[np.int64] +UInt32Array = npt.NDArray[np.uint32] + + +class E32TrackGeometryReplayError(RuntimeError): + """The E32 input chain, replay or immutable result is invalid.""" + + +@dataclass(frozen=True, slots=True) +class E32TrackGeometryReplay: + """One immutable E32 result.""" + + result_root: Path + result_id: str + manifest: dict[str, Any] + report: dict[str, Any] + + +def read_e32_track_geometry_replay(root: Path) -> E32TrackGeometryReplay: + """Open and fully validate one immutable E32 result.""" + + resolved = _safe_result_root(root, _E32_RESULT_ID) + manifest = _read_json(resolved / E32_MANIFEST_NAME) + identity = _object(manifest.get("identity"), "E32 identity") + return _read_existing_result(resolved, identity) + + +def e32_track_geometry_frame( + result: E32TrackGeometryReplay, + frame_index: int, +) -> TrackGeometryFrame: + """Reconstruct one exact TrackGeometryFrame from compact E32 storage.""" + + _required_nonnegative_int(frame_index, "E32 frame index") + result_identity = _object(result.manifest.get("identity"), "E32 identity") + frame_count = _required_positive_int( + result_identity.get("frame_count"), + "E32 frame count", + ) + if frame_index >= frame_count: + raise E32TrackGeometryReplayError("E32 frame index is out of range") + artifacts = _verified_artifacts(result.result_root, result.manifest) + try: + arrays = load_point_storage(artifacts, frame_count=frame_count) + except E32TrackGeometryStorageError as exc: + raise E32TrackGeometryReplayError(str(exc)) from exc + frames_path = artifacts["track-geometry-frames"] + record_value: dict[str, Any] | None = None + with frames_path.open("r", encoding="utf-8") as stream: + for current_index, line in enumerate(stream): + if current_index == frame_index: + record_value = record(line, expected_frame_index=frame_index) + break + if record_value is None: + raise E32TrackGeometryReplayError("E32 frame record is missing") + binding = TrackGeometrySourceBinding.from_dict( + result_identity.get("track_geometry_binding") + ) + try: + return frame_from_record( + record_value=record_value, + binding=binding, + frame_offsets=arrays[0], + source_indices=arrays[1], + points=arrays[2], + owner_indices=arrays[3], + ) + except E32TrackGeometryStorageError as exc: + raise E32TrackGeometryReplayError(str(exc)) from exc + + +@dataclass(frozen=True, slots=True) +class _E29Source: + root: Path + manifest: dict[str, Any] + frames_path: Path + report_path: Path + profile: CameraGeometryFusionProfile + + +@dataclass(frozen=True, slots=True) +class _CameraSource: + root: Path + result: dict[str, Any] + fusion_frames_path: Path + + +@dataclass(frozen=True, slots=True) +class _E31Source: + root: Path + manifest: dict[str, Any] + binding_profile: dict[str, Any] + binding_profile_sha256: str + + +@dataclass(frozen=True, slots=True) +class _CorrectionPlan: + semantic_rectangle_normalized_xyxy: tuple[float, float, float, float] + semantic_class_allowlist: frozenset[str] + image_width: int + image_height: int + exact_geometry_corrections: dict[tuple[int, int], str] + human_geometry_dispositions: dict[tuple[int, int], tuple[str, str]] + + +@dataclass(frozen=True, slots=True) +class _FrameTranslation: + frame: TrackGeometryFrame + record: dict[str, object] + point_source_indices: Int64Array + point_coordinates: npt.NDArray[np.float32] + point_owner_indices: UInt32Array + changes: tuple[dict[str, object], ...] + semantic_published: int + semantic_masked: int + semantic_unpublishable_held: int + geometry_published: int + geometry_exact_excluded: int + geometry_human_excluded: int + baseline_qualified_points: int + published_qualified_points: int + excluded_qualified_points: int + ownership_overlap_claims: int + unqualified_semantic_support_points: int + unqualified_ranges_withheld: int + + +def build_e32_track_geometry_replay( + *, + camera_result_root: Path, + e29_result_root: Path, + source_pack_root: Path, + local_surface_root: Path, + e31_result_root: Path, + materialization_root: Path, + engineering_generation_root: Path, + human_generation_root: Path, + output_root: Path, +) -> E32TrackGeometryReplay: + """Build or verify the full 4,489-frame E32 diagnostic replay.""" + + source = E10LidarFieldSource(source_pack_root) + surface = K1LocalSurfaceV1(local_surface_root) + try: + camera = _load_camera_source(camera_result_root) + e29 = _load_e29_source(e29_result_root) + e31 = _load_e31_source(e31_result_root) + chain = _load_e30_chain( + materialization_root=materialization_root, + engineering_generation_root=engineering_generation_root, + human_generation_root=human_generation_root, + ) + _validate_input_chain( + source=source, + surface=surface, + camera=camera, + e29=e29, + e31=e31, + chain=chain, + ) + corrections = _correction_plan(e31.binding_profile, chain) + binding = _track_geometry_binding(source=source, e31=e31) + identity = _result_identity( + source=source, + surface=surface, + camera=camera, + e29=e29, + e31=e31, + chain=chain, + binding=binding, + corrections=corrections, + ) + identity_sha256 = hashlib.sha256(_canonical_json(identity)).hexdigest() + result_id = f"e32-track-geometry-{identity_sha256}" + destination = output_root.expanduser().absolute() + destination.mkdir(mode=0o700, parents=True, exist_ok=True) + result_root = destination / result_id + if result_root.exists(): + return _read_existing_result(result_root, identity) + + staging = destination / f".{result_id}.{os.getpid()}.incomplete" + staging.mkdir(mode=0o700, exist_ok=False) + try: + report = _run_replay( + staging=staging, + result_id=result_id, + identity=identity, + source=source, + surface=surface, + camera=camera, + e29=e29, + binding=binding, + corrections=corrections, + ) + _write_json(staging / E32_REPORT_NAME, report) + artifacts = [ + _artifact("track-geometry-frames", staging / E32_FRAMES_NAME), + _artifact("e29-e32-changes", staging / E32_CHANGES_NAME), + _artifact("frame-point-offsets", staging / E32_FRAME_OFFSETS_NAME), + _artifact("point-source-indices", staging / E32_SOURCE_INDICES_NAME), + _artifact("point-coordinates-map-f32", staging / E32_POINTS_NAME), + _artifact("point-owner-indices", staging / E32_OWNER_INDICES_NAME), + _artifact("track-geometry-report", staging / E32_REPORT_NAME), + ] + manifest = { + "schema_version": E32_TRACK_GEOMETRY_REPLAY_SCHEMA, + "result_id": result_id, + "identity_sha256": identity_sha256, + "identity": identity, + "created_at_utc": report["created_at_utc"], + "classification": "private-derived-perception-diagnostic", + "ground_truth": False, + "lab_published": False, + "artifacts": artifacts, + "authority": _authority(), + } + _write_json(staging / E32_MANIFEST_NAME, manifest) + os.replace(staging, result_root) + except BaseException: + shutil.rmtree(staging, ignore_errors=True) + raise + return _read_existing_result(result_root, identity) + finally: + source.close() + surface.close() + + +def _run_replay( + *, + staging: Path, + result_id: str, + identity: dict[str, object], + source: E10LidarFieldSource, + surface: K1LocalSurfaceV1, + camera: _CameraSource, + e29: _E29Source, + binding: TrackGeometrySourceBinding, + corrections: _CorrectionPlan, +) -> dict[str, Any]: + started = time.perf_counter() + arrays = source.arrays + offsets = arrays["cloud_offsets"] + all_points = arrays["cloud_points_map"] + all_classes = surface.arrays["point_class"] + all_heights = surface.arrays["point_height_m"] + projection = _projection_profile(source) + last_current_frame: dict[int, int] = {} + + frame_offsets = [0] + point_source_parts: list[Int64Array] = [] + point_coordinate_parts: list[npt.NDArray[np.float32]] = [] + point_owner_parts: list[UInt32Array] = [] + baseline_counts: Counter[tuple[str, str, str, str, str]] = Counter() + e32_counts: Counter[tuple[str, str, str, str, str]] = Counter() + change_counts: Counter[str] = Counter() + totals: Counter[str] = Counter() + source_availability: Counter[str] = Counter() + frame_processing_ms: list[float] = [] + + frames_path = staging / E32_FRAMES_NAME + changes_path = staging / E32_CHANGES_NAME + frame_count = 0 + with ( + camera.fusion_frames_path.open("r", encoding="utf-8") as camera_stream, + e29.frames_path.open("r", encoding="utf-8") as e29_stream, + frames_path.open("x", encoding="utf-8") as frame_output, + changes_path.open("x", encoding="utf-8") as change_output, + ): + for frame_index, lines in enumerate( + zip(camera_stream, e29_stream, strict=True) + ): + frame_started = time.perf_counter() + fusion_line, e29_line = lines + fusion_frame = _fusion_frame( + fusion_line, + expected_frame_index=frame_index, + source=source, + ) + e29_frame = _e29_frame(e29_line, expected_frame_index=frame_index) + start = int(offsets[frame_index]) + end = int(offsets[frame_index + 1]) + frame_points = np.asarray(all_points[start:end], dtype=np.float64) + frame_classes = all_classes[start:end] + frame_heights = all_heights[start:end] + source_available = bool(arrays["sample_available"][frame_index]) + surface_valid = bool(surface.arrays["frame_valid"][frame_index]) + projected = None + if source_available and surface_valid: + position = arrays["pose_positions_map"][frame_index] + orientation = arrays["pose_quaternions_map_from_lidar"][frame_index] + projected = project_map_points_kb4( + frame_points, + position_map_xyz=cast( + tuple[float, float, float], + tuple(float(value) for value in position), + ), + orientation_map_from_lidar_xyzw=cast( + tuple[float, float, float, float], + tuple(float(value) for value in orientation), + ), + profile=projection, + ) + semantic_supports = tuple( + _semantic_support( + item, + projected=projected, + frame_points_map=frame_points, + point_class=frame_classes, + point_height_m=frame_heights, + source_available=source_available, + surface_valid=surface_valid, + profile=e29.profile, + ) + for item in fusion_frame["objects"] + ) + geometry_supports = tuple( + _geometry_cluster_supports( + points_map=frame_points, + point_class=frame_classes, + point_height_m=frame_heights, + sensor_position_map=np.asarray( + arrays["pose_positions_map"][frame_index], + dtype=np.float64, + ), + claimed_source_indices=_claimed_indices(semantic_supports), + profile=e29.profile, + ) + ) + _require_exact_e29_reproduction( + e29_frame=e29_frame, + frame_index=frame_index, + fusion_frame=fusion_frame, + source_available=source_available, + surface_valid=surface_valid, + semantic_supports=semantic_supports, + geometry_supports=geometry_supports, + ) + translation = _translate_frame( + frame_index=frame_index, + source_frame_index=int(fusion_frame["source_frame_index"]), + session_seconds=float(fusion_frame["session_seconds"]), + source_available=source_available, + frame_points=frame_points, + semantic_supports=semantic_supports, + geometry_supports=geometry_supports, + binding=binding, + corrections=corrections, + last_current_frame=last_current_frame, + ) + for semantic_support in semantic_supports: + _count_baseline( + baseline_counts, + semantic_support.document, + session_seconds=float(fusion_frame["session_seconds"]), + geometry_only=False, + ) + for geometry_support in geometry_supports: + _count_baseline( + baseline_counts, + geometry_support.document, + session_seconds=float(fusion_frame["session_seconds"]), + geometry_only=True, + ) + for geometry in translation.frame.geometries: + _count_e32( + e32_counts, + geometry, + session_seconds=float(fusion_frame["session_seconds"]), + ) + _write_jsonl(frame_output, translation.record) + for change in translation.changes: + _write_jsonl(change_output, change) + reason = str(change["reason"]) + change_counts[reason] += 1 + before = change.get("before") + if isinstance(before, dict) and before.get("geometry_status") == "conflict": + totals["changed_conflict_observations"] += 1 + if reason == "point-ownership-arbitration": + after = change.get("after") + if not isinstance(after, dict): + raise E32TrackGeometryReplayError( + "point ownership change outcome is missing" + ) + if after.get("evidence_state") == "unknown": + totals["ownership_downgraded_to_unknown"] += 1 + else: + totals["ownership_retained_agree"] += 1 + point_source_parts.append(translation.point_source_indices) + point_coordinate_parts.append(translation.point_coordinates) + point_owner_parts.append(translation.point_owner_indices) + frame_offsets.append(frame_offsets[-1] + translation.point_source_indices.size) + totals.update( + { + "semantic_baseline": len(semantic_supports), + "semantic_published": translation.semantic_published, + "semantic_masked": translation.semantic_masked, + "semantic_unpublishable_held": ( + translation.semantic_unpublishable_held + ), + "geometry_baseline": len(geometry_supports), + "geometry_published": translation.geometry_published, + "geometry_exact_excluded": translation.geometry_exact_excluded, + "geometry_human_excluded": translation.geometry_human_excluded, + "baseline_qualified_points": ( + translation.baseline_qualified_points + ), + "published_qualified_points": ( + translation.published_qualified_points + ), + "excluded_qualified_points": ( + translation.excluded_qualified_points + ), + "ownership_overlap_claims": ( + translation.ownership_overlap_claims + ), + "unqualified_semantic_support_points": ( + translation.unqualified_semantic_support_points + ), + "unqualified_ranges_withheld": ( + translation.unqualified_ranges_withheld + ), + } + ) + source_availability[ + "available" if source_available else "unavailable" + ] += 1 + frame_processing_ms.append( + (time.perf_counter() - frame_started) * 1_000.0 + ) + frame_count += 1 + if frame_count != source.frame_count: + raise E32TrackGeometryReplayError("E32 input streams are incomplete") + + point_source_indices = _concatenate( + point_source_parts, + dtype=np.dtype(" _FrameTranslation: + geometries: list[TrackGeometry] = [] + owner_keys: list[str] = [] + point_sources: list[Int64Array] = [] + point_coordinates: list[npt.NDArray[np.float32]] = [] + point_owners: list[UInt32Array] = [] + changes: list[dict[str, object]] = [] + semantic_published = 0 + semantic_masked = 0 + semantic_unpublishable_held = 0 + geometry_published = 0 + geometry_exact_excluded = 0 + geometry_human_excluded = 0 + baseline_qualified_points = 0 + published_qualified_points = 0 + excluded_qualified_points = 0 + ownership_overlap_claims = 0 + unqualified_semantic_support_points = 0 + unqualified_ranges_withheld = 0 + semantic_allocations = _allocate_semantic_point_ownership( + semantic_supports, + corrections=corrections, + ) + + for observation_index, semantic_support in enumerate(semantic_supports): + document = semantic_support.document + is_qualified = document["geometry_status"] == "agree" + if is_qualified: + baseline_qualified_points += int( + semantic_support.occupied_source_indices.size + ) + else: + unqualified_semantic_support_points += int( + semantic_support.occupied_source_indices.size + ) + if _semantic_masked(document, corrections): + semantic_masked += 1 + excluded_qualified_points += ( + int(semantic_support.occupied_source_indices.size) + if is_qualified + else 0 + ) + changes.append( + _change( + frame_index=frame_index, + source_frame_index=source_frame_index, + session_seconds=session_seconds, + kind="semantic-observation", + locator_index=observation_index, + reason="e31-semantic-self-mask", + before=document, + excluded_point_count=( + int(semantic_support.occupied_source_indices.size) + if is_qualified + else 0 + ), + ) + ) + continue + track_id = _required_nonnegative_int(document.get("track_id"), "semantic track id") + current = document.get("semantic_current") is True + held_from = None + if not current: + held_from = last_current_frame.get(track_id) + if held_from is None: + semantic_unpublishable_held += 1 + changes.append( + _change( + frame_index=frame_index, + source_frame_index=source_frame_index, + session_seconds=session_seconds, + kind="semantic-observation", + locator_index=observation_index, + reason="held-without-prior-current-provenance", + before=document, + excluded_point_count=0, + ) + ) + continue + owned_indices = semantic_allocations.get( + observation_index, + semantic_support.occupied_source_indices, + ) + ownership_removed = int( + semantic_support.occupied_source_indices.size - owned_indices.size + ) + if ownership_removed: + ownership_overlap_claims += ownership_removed + changes.append( + _change( + frame_index=frame_index, + source_frame_index=source_frame_index, + session_seconds=session_seconds, + kind="semantic-observation", + locator_index=observation_index, + reason="point-ownership-arbitration", + before=document, + after={ + "evidence_state": ( + "agree" if owned_indices.size else "unknown" + ), + "original_claim_point_count": int( + semantic_support.occupied_source_indices.size + ), + "owned_point_count": int(owned_indices.size), + "winner_policy": ( + "smallest-bbox-then-score-track-and-observation" + ), + }, + excluded_point_count=ownership_removed, + ) + ) + unqualified_range_withheld = ( + document.get("geometry_status") != "agree" + and document.get("range_m") is not None + ) + if unqualified_range_withheld: + unqualified_ranges_withheld += 1 + changes.append( + _change( + frame_index=frame_index, + source_frame_index=source_frame_index, + session_seconds=session_seconds, + kind="semantic-observation", + locator_index=observation_index, + reason="unqualified-range-withheld", + before=document, + after={ + "evidence_state": document["geometry_status"], + "range_m": None, + "metric_basis": "unavailable", + }, + excluded_point_count=0, + ) + ) + geometry = _semantic_track_geometry( + document, + held_from_frame_index=held_from, + owned_point_count=int(owned_indices.size), + ownership_reduced=ownership_removed > 0, + unqualified_range_withheld=unqualified_range_withheld, + ) + geometries.append(geometry) + semantic_published += 1 + if current: + last_current_frame[track_id] = frame_index + if geometry.metric_basis is TrackGeometryMetricBasis.CURRENT_POINTS: + _append_owner_points( + owner_key=geometry.owner_key, + source_indices=owned_indices, + frame_points=frame_points, + owner_keys=owner_keys, + point_sources=point_sources, + point_coordinates=point_coordinates, + point_owners=point_owners, + ) + published_qualified_points += int( + owned_indices.size + ) + + for cluster_index, geometry_support in enumerate(geometry_supports): + baseline_qualified_points += int( + geometry_support.occupied_source_indices.size + ) + locator = (frame_index, cluster_index) + if locator in corrections.exact_geometry_corrections: + geometry_exact_excluded += 1 + excluded_qualified_points += int( + geometry_support.occupied_source_indices.size + ) + changes.append( + _change( + frame_index=frame_index, + source_frame_index=source_frame_index, + session_seconds=session_seconds, + kind="geometry-only-cluster", + locator_index=cluster_index, + reason="e31-exact-geometry-correction", + before=geometry_support.document, + excluded_point_count=int( + geometry_support.occupied_source_indices.size + ), + decision_item_id=corrections.exact_geometry_corrections[locator], + ) + ) + continue + human = corrections.human_geometry_dispositions.get(locator) + if human is not None and human[0] == "background-or-noise": + geometry_human_excluded += 1 + excluded_qualified_points += int( + geometry_support.occupied_source_indices.size + ) + changes.append( + _change( + frame_index=frame_index, + source_frame_index=source_frame_index, + session_seconds=session_seconds, + kind="geometry-only-cluster", + locator_index=cluster_index, + reason="a3-human-background-or-noise", + before=geometry_support.document, + excluded_point_count=int( + geometry_support.occupied_source_indices.size + ), + decision_item_id=human[1], + ) + ) + continue + reason_codes = ["e29-unassociated-occupied-component"] + if human is not None: + if human[0] != "object-present": + raise E32TrackGeometryReplayError( + "unresolved human geometry disposition reached E32" + ) + reason_codes.append("a3-human-object-present") + geometry = TrackGeometry( + owner_key=f"geometry:{cluster_index}", + owner_kind=TrackGeometryOwnerKind.GEOMETRY_CLUSTER, + evidence_state=TrackGeometryEvidenceState.GEOMETRY_ONLY, + currentness=TrackGeometryCurrentness.CURRENT, + metric_basis=TrackGeometryMetricBasis.CURRENT_POINTS, + reason_codes=tuple(reason_codes), + range_m=_required_positive_float( + geometry_support.document.get("nearest_range_m"), + "geometry nearest range", + ), + ) + geometries.append(geometry) + geometry_published += 1 + _append_owner_points( + owner_key=geometry.owner_key, + source_indices=geometry_support.occupied_source_indices, + frame_points=frame_points, + owner_keys=owner_keys, + point_sources=point_sources, + point_coordinates=point_coordinates, + point_owners=point_owners, + ) + published_qualified_points += int( + geometry_support.occupied_source_indices.size + ) + + source_indices = _concatenate( + point_sources, + dtype=np.dtype(" TrackGeometry: + track_id = _required_nonnegative_int(document.get("track_id"), "semantic track id") + label = _required_string(document.get("label"), "semantic label") + bbox = _bbox(document.get("bbox_xyxy")) + reason = _required_string(document.get("geometry_reason"), "geometry reason") + if held_from_frame_index is not None: + return TrackGeometry( + owner_key=f"track:{track_id}", + owner_kind=TrackGeometryOwnerKind.CAMERA_TRACK, + evidence_state=TrackGeometryEvidenceState.UNKNOWN, + currentness=TrackGeometryCurrentness.HELD, + metric_basis=TrackGeometryMetricBasis.UNAVAILABLE, + reason_codes=(reason, "e29-held-camera-observation"), + semantic_track_id=track_id, + semantic_label=label, + bbox_xyxy=bbox, + held_from_frame_index=held_from_frame_index, + ) + status = _required_string(document.get("geometry_status"), "geometry status") + evidence_state = TrackGeometryEvidenceState(status) + if evidence_state is TrackGeometryEvidenceState.AGREE and owned_point_count == 0: + return TrackGeometry( + owner_key=f"track:{track_id}", + owner_kind=TrackGeometryOwnerKind.CAMERA_TRACK, + evidence_state=TrackGeometryEvidenceState.UNKNOWN, + currentness=TrackGeometryCurrentness.CURRENT, + metric_basis=TrackGeometryMetricBasis.UNAVAILABLE, + reason_codes=(reason, "e32-point-ownership-collision"), + semantic_track_id=track_id, + semantic_label=label, + bbox_xyxy=bbox, + ) + metric_basis = ( + TrackGeometryMetricBasis.CURRENT_POINTS + if evidence_state is TrackGeometryEvidenceState.AGREE + else TrackGeometryMetricBasis.UNAVAILABLE + ) + range_m = ( + _required_positive_float(document.get("range_m"), "semantic range") + if evidence_state is TrackGeometryEvidenceState.AGREE + else None + ) + return TrackGeometry( + owner_key=f"track:{track_id}", + owner_kind=TrackGeometryOwnerKind.CAMERA_TRACK, + evidence_state=evidence_state, + currentness=TrackGeometryCurrentness.CURRENT, + metric_basis=metric_basis, + reason_codes=tuple( + [ + reason, + *(["e32-exclusive-point-ownership"] if ownership_reduced else []), + *( + ["e32-unqualified-range-withheld"] + if unqualified_range_withheld + else [] + ), + ] + ), + semantic_track_id=track_id, + semantic_label=label, + bbox_xyxy=bbox, + range_m=range_m, + ) + + +def _allocate_semantic_point_ownership( + semantic_supports: tuple[_SemanticSupport, ...], + *, + corrections: _CorrectionPlan, +) -> dict[int, Int64Array]: + eligible = { + observation_index: support + for observation_index, support in enumerate(semantic_supports) + if support.document.get("geometry_status") == "agree" + and not _semantic_masked(support.document, corrections) + } + claims: dict[int, list[int]] = {} + for observation_index, support in eligible.items(): + for source_index in support.occupied_source_indices: + claims.setdefault(int(source_index), []).append(observation_index) + winners: dict[int, int] = {} + for source_index, candidates in claims.items(): + winners[source_index] = min( + candidates, + key=lambda observation_index: _semantic_owner_priority( + eligible[observation_index].document, + observation_index, + ), + ) + allocations: dict[int, Int64Array] = {} + for observation_index, support in eligible.items(): + allocations[observation_index] = np.asarray( + [ + int(source_index) + for source_index in support.occupied_source_indices + if winners[int(source_index)] == observation_index + ], + dtype=" tuple[float, float, int, int]: + bbox = _bbox(document.get("bbox_xyxy")) + area = (bbox[2] - bbox[0]) * (bbox[3] - bbox[1]) + score = _optional_float(document.get("score")) or 0.0 + track_id = _required_nonnegative_int(document.get("track_id"), "semantic track id") + return (area, -score, track_id, observation_index) + + +def _append_owner_points( + *, + owner_key: str, + source_indices: Int64Array, + frame_points: Float64Array, + owner_keys: list[str], + point_sources: list[Int64Array], + point_coordinates: list[npt.NDArray[np.float32]], + point_owners: list[UInt32Array], +) -> None: + if source_indices.size == 0: + raise E32TrackGeometryReplayError("current metric owner has no source points") + owner_index = len(owner_keys) + owner_keys.append(owner_key) + point_sources.append(np.asarray(source_indices, dtype=" bool: + label = document.get("label") + if label not in corrections.semantic_class_allowlist: + return False + bbox = _bbox(document.get("bbox_xyxy")) + left, top, right, bottom = corrections.semantic_rectangle_normalized_xyxy + center_x = ((bbox[0] + bbox[2]) * 0.5) / corrections.image_width + center_y = ((bbox[1] + bbox[3]) * 0.5) / corrections.image_height + return left <= center_x <= right and top <= center_y <= bottom + + +def _require_exact_e29_reproduction( + *, + e29_frame: dict[str, Any], + frame_index: int, + fusion_frame: dict[str, Any], + source_available: bool, + surface_valid: bool, + semantic_supports: tuple[_SemanticSupport, ...], + geometry_supports: tuple[_GeometryClusterSupport, ...], +) -> None: + reproduced = { + "schema_version": CAMERA_GEOMETRY_FRAME_SCHEMA, + "frame_index": frame_index, + "source_frame_index": fusion_frame["source_frame_index"], + "session_seconds": fusion_frame["session_seconds"], + "source_available": source_available, + "local_surface_valid": surface_valid, + "semantic_observations": [support.document for support in semantic_supports], + "geometry_only_occupied": [support.document for support in geometry_supports], + "policy": { + "camera_owns_semantics": True, + "lidar_owns_metric_geometry": True, + "absence_of_points_means_free": False, + "unknown_is_occupied": True, + }, + "authority": _authority(), + } + if reproduced != e29_frame: + raise E32TrackGeometryReplayError( + f"current replay does not exactly reproduce immutable E29 frame {frame_index}" + ) + + +def _load_camera_source(root: Path) -> _CameraSource: + candidate = root.expanduser().absolute() + if candidate.is_symlink(): + raise E32TrackGeometryReplayError("camera result cannot be a symlink") + resolved = candidate.resolve(strict=True) + if not resolved.is_dir() or _E10_RESULT_ID.fullmatch(resolved.name) is None: + raise E32TrackGeometryReplayError("camera result id is invalid") + result = _read_json(resolved / "result.json") + identity = _object(result.get("identity"), "camera result identity") + identity_sha256 = result.get("identity_sha256") + if ( + result.get("schema_version") + != "missioncore.e10-integrated-perception-result/v1" + or result.get("result_id") != resolved.name + or not isinstance(identity_sha256, str) + or _SHA256.fullmatch(identity_sha256) is None + or hashlib.sha256(_canonical_json(identity)).hexdigest() != identity_sha256 + or resolved.name != f"e10-integrated-perception-{identity_sha256}" + ): + raise E32TrackGeometryReplayError("camera result identity is invalid") + artifacts = _verified_camera_artifacts(resolved, result) + fusion = artifacts.get("e10-fusion-frames") + if fusion is None: + raise E32TrackGeometryReplayError("camera fusion frame artifact is missing") + return _CameraSource(root=resolved, result=result, fusion_frames_path=fusion) + + +def _load_e29_source(root: Path) -> _E29Source: + resolved = _safe_result_root(root, _E29_RESULT_ID) + manifest = _read_json(resolved / "manifest.json") + _validate_e29_identity( + root=resolved, + manifest=manifest, + ) + artifacts = _verified_artifacts(resolved, manifest) + frames_path = artifacts.get("camera-geometry-frames") + report_path = artifacts.get("camera-geometry-report") + if frames_path is None or report_path is None: + raise E32TrackGeometryReplayError("E29 artifacts are incomplete") + identity = _object(manifest.get("identity"), "E29 identity") + profile_document = _object(identity.get("profile"), "E29 profile") + profile_fields = {item.name for item in fields(CameraGeometryFusionProfile)} + profile = CameraGeometryFusionProfile( + **{name: profile_document[name] for name in profile_fields} + ) + if profile.to_dict() != profile_document: + raise E32TrackGeometryReplayError("E29 profile cannot be reproduced exactly") + return _E29Source( + root=resolved, + manifest=manifest, + frames_path=frames_path, + report_path=report_path, + profile=profile, + ) + + +def _validate_e29_identity( + *, + root: Path, + manifest: Mapping[str, object], +) -> None: + identity = manifest.get("identity") + identity_sha256 = manifest.get("identity_sha256") + if ( + manifest.get("schema_version") != CAMERA_GEOMETRY_FUSION_SCHEMA + or manifest.get("result_id") != root.name + or not isinstance(identity, dict) + or not isinstance(identity_sha256, str) + or _SHA256.fullmatch(identity_sha256) is None + or hashlib.sha256(_canonical_json(identity)).hexdigest() != identity_sha256 + or root.name != f"e29-camera-geometry-{identity_sha256}" + or identity.get("authority") != _authority() + ): + raise E32TrackGeometryReplayError("E29 result identity is invalid") + + +def _load_e31_source(root: Path) -> _E31Source: + resolved = _safe_result_root(root, _E31_RESULT_ID) + manifest = _read_json(resolved / "manifest.json") + _validate_content_identity( + root=resolved, + manifest=manifest, + schema="missioncore.e31-source-qualification/v1", + prefix="e31-source-qualification-", + ) + artifacts = _verified_artifacts(resolved, manifest) + binding_path = artifacts.get("binding-profile") + if binding_path is None: + raise E32TrackGeometryReplayError("E31 binding profile is missing") + binding_profile = _read_json(binding_path) + if ( + binding_profile.get("schema_version") != E31_BINDING_PROFILE_SCHEMA + or binding_profile.get("status") != "accepted" + or binding_profile.get("eligible_for_e32") is not True + or binding_profile.get("authority") != _authority() + ): + raise E32TrackGeometryReplayError("E31 binding profile is not eligible") + return _E31Source( + root=resolved, + manifest=manifest, + binding_profile=binding_profile, + binding_profile_sha256=_sha256(binding_path), + ) + + +def _validate_input_chain( + *, + source: E10LidarFieldSource, + surface: K1LocalSurfaceV1, + camera: _CameraSource, + e29: _E29Source, + e31: _E31Source, + chain: Any, +) -> None: + e29_identity = _object(e29.manifest.get("identity"), "E29 identity") + e31_identity = _object(e31.manifest.get("identity"), "E31 identity") + e31_source = _object(e31_identity.get("source"), "E31 source") + camera_artifact = _artifact_by_kind(camera.result, "e10-fusion-frames") + materialization_identity = _object( + chain.materialization_manifest.get("identity"), + "materialization identity", + ) + materialization_source = _object( + materialization_identity.get("source"), + "materialization source", + ) + e31_chain = _object( + e31.binding_profile.get("a3_decision_chain"), + "E31 A3 decision chain", + ) + if ( + e29_identity.get("frame_count") != source.frame_count + or e29_identity.get("source_pack_id") != source.pack_id + or e29_identity.get("local_surface_model_id") != surface.model_id + or e29_identity.get("source_result_id") != camera.root.name + or e29_identity.get("source_fusion_frames_sha256") + != camera_artifact.get("sha256") + or surface.identity.get("source_pack_id") != source.pack_id + or surface.identity.get("frame_count") != source.frame_count + or surface.identity.get("point_count") != source.point_count + or e31_source.get("source_pack_id") != source.pack_id + or e31_source.get("local_surface_model_id") != surface.model_id + or e31_source.get("materialization_id") + != chain.materialization_manifest.get("result_id") + or e31_source.get("engineering_generation_id") + != chain.engineering_manifest.get("result_id") + or e31_source.get("human_generation_id") + != chain.human_manifest.get("result_id") + or e31_chain.get("engineering_generation_id") + != chain.engineering_manifest.get("result_id") + or e31_chain.get("human_generation_id") + != chain.human_manifest.get("result_id") + or materialization_source.get("e29_result_id") != e29.root.name + or materialization_source.get("e29_manifest_sha256") + != _sha256(e29.root / "manifest.json") + ): + raise E32TrackGeometryReplayError("E32 immutable source chain is inconsistent") + + +def _correction_plan(binding_profile: dict[str, Any], chain: Any) -> _CorrectionPlan: + mount = _object( + binding_profile.get("mount_and_self_mask"), + "E31 mount and self mask", + ) + if mount.get("generic_geometry_point_mask") is not None: + raise E32TrackGeometryReplayError("generic E31 geometry mask is not allowed") + semantic = _object(mount.get("semantic_mask"), "E31 semantic self mask") + calibration = _object(binding_profile.get("calibration"), "E31 calibration") + rectangle = semantic.get("rectangle_normalized_xyxy") + allowlist = semantic.get("class_allowlist") + resolution = calibration.get("resolution") + if ( + semantic.get("status") != "admitted" + or semantic.get("application_rule") != "bbox-center-inside-rectangle" + or not isinstance(rectangle, list) + or len(rectangle) != 4 + or not isinstance(allowlist, list) + or not allowlist + or not isinstance(resolution, list) + or len(resolution) != 2 + ): + raise E32TrackGeometryReplayError("E31 semantic self mask is incompatible") + rectangle_values = tuple(float(value) for value in rectangle) + if ( + not np.isfinite(rectangle_values).all() + or not 0.0 <= rectangle_values[0] < rectangle_values[2] <= 1.0 + or not 0.0 <= rectangle_values[1] < rectangle_values[3] <= 1.0 + ): + raise E32TrackGeometryReplayError("E31 semantic self mask bounds are invalid") + items = {str(item["item_id"]): item for item in chain.items} + exact: dict[tuple[int, int], str] = {} + exact_ids = mount.get("exact_geometry_correction_item_ids") + if not isinstance(exact_ids, list): + raise E32TrackGeometryReplayError("E31 exact geometry corrections changed") + for value in exact_ids: + item_id = _required_string(value, "exact geometry correction item id") + item = items.get(item_id) + if item is None: + raise E32TrackGeometryReplayError("E31 exact correction item is missing") + locator = _geometry_locator(item) + if locator in exact: + raise E32TrackGeometryReplayError("duplicate exact geometry correction") + exact[locator] = item_id + + human: dict[tuple[int, int], tuple[str, str]] = {} + for decision in chain.human_decisions: + item_id = _required_string(decision.get("item_id"), "human decision item id") + item = items.get(item_id) + if item is None: + raise E32TrackGeometryReplayError("human decision item is missing") + disposition = _required_string( + decision.get("disposition"), + "human geometry disposition", + ) + if disposition == "insufficient-evidence": + raise E32TrackGeometryReplayError( + "unresolved human exception cannot enter E32" + ) + locator = _geometry_locator(item) + if locator in human or locator in exact: + raise E32TrackGeometryReplayError("geometry correction ownership overlaps") + human[locator] = (disposition, item_id) + return _CorrectionPlan( + semantic_rectangle_normalized_xyxy=cast( + tuple[float, float, float, float], + rectangle_values, + ), + semantic_class_allowlist=frozenset( + _required_string(value, "semantic mask class") for value in allowlist + ), + image_width=_required_positive_int(resolution[0], "calibration width"), + image_height=_required_positive_int(resolution[1], "calibration height"), + exact_geometry_corrections=exact, + human_geometry_dispositions=human, + ) + + +def _track_geometry_binding( + *, + source: E10LidarFieldSource, + e31: _E31Source, +) -> TrackGeometrySourceBinding: + scope = _object(e31.binding_profile.get("scope"), "E31 scope") + time_binding = _object( + e31.binding_profile.get("time_binding"), + "E31 time binding", + ) + calibration = _object( + e31.binding_profile.get("calibration"), + "E31 calibration", + ) + return TrackGeometrySourceBinding( + source_pack_id=source.pack_id, + source_session_id=_required_string(scope.get("session_id"), "source session id"), + representation_profile_id=K1_LIO_PCL_CAPABILITIES.profile_id, + e31_qualification_id=e31.root.name, + calibration_sha256=_required_sha256( + calibration.get("content_identity_sha256"), + "calibration identity", + ), + coordinate_frame="map", + time_basis=_required_string(time_binding.get("basis"), "time basis"), + selected_offset_ms=_required_int( + time_binding.get("selected_host_arrival_offset_ms"), + "selected host arrival offset", + ), + ) + + +def _result_identity( + *, + source: E10LidarFieldSource, + surface: K1LocalSurfaceV1, + camera: _CameraSource, + e29: _E29Source, + e31: _E31Source, + chain: Any, + binding: TrackGeometrySourceBinding, + corrections: _CorrectionPlan, +) -> dict[str, object]: + e29_identity = _object(e29.manifest.get("identity"), "E29 identity") + camera_artifact = _artifact_by_kind(camera.result, "e10-fusion-frames") + return { + "schema_version": E32_TRACK_GEOMETRY_REPLAY_SCHEMA, + "source": { + "camera_result_id": camera.root.name, + "fusion_frames_sha256": camera_artifact["sha256"], + "e29_result_id": e29.root.name, + "e29_frames_sha256": _sha256(e29.frames_path), + "e29_report_sha256": _sha256(e29.report_path), + "e29_original_producer_sha256": e29_identity["producer_sha256"], + "source_pack_id": source.pack_id, + "source_pack_artifact_sha256": _source_artifact_sha256(source), + "local_surface_model_id": surface.model_id, + "local_surface_artifact_sha256": _surface_artifact_sha256(surface), + "e31_result_id": e31.root.name, + "e31_binding_profile_sha256": e31.binding_profile_sha256, + "materialization_id": chain.materialization_manifest["result_id"], + "engineering_generation_id": chain.engineering_manifest["result_id"], + "human_generation_id": chain.human_manifest["result_id"], + }, + "frame_count": source.frame_count, + "timeline_start_seconds": float(source.arrays["session_seconds"][0]), + "timeline_end_seconds": float(source.arrays["session_seconds"][-1]), + "track_geometry_binding": binding.to_dict(), + "method": { + "profile_id": "e32-exact-e29-to-track-geometry/v1", + "e29_profile": e29.profile.to_dict(), + "exact_e29_frame_reproduction_required": True, + "semantic_mask_application": "bbox-center-inside-normalized-rectangle", + "semantic_mask_rectangle_normalized_xyxy": list( + corrections.semantic_rectangle_normalized_xyxy + ), + "semantic_mask_class_allowlist": sorted( + corrections.semantic_class_allowlist + ), + "semantic_mask_resolution": [ + corrections.image_width, + corrections.image_height, + ], + "exact_geometry_correction_count": len( + corrections.exact_geometry_corrections + ), + "human_geometry_disposition_count": len( + corrections.human_geometry_dispositions + ), + "generic_geometry_mask_applied": False, + "threshold_tuning_allowed": False, + "point_storage": "deterministic-npy-frame-local-point-slab/v1", + "track_geometry_schema": TRACK_GEOMETRY_SCHEMA, + "track_geometry_frame_schema": TRACK_GEOMETRY_FRAME_SCHEMA, + "point_slab_schema": POINT_SLAB_SCHEMA, + "absence_of_points_means_free": False, + }, + "replay_dependencies": { + "e29_replay_implementation_sha256": _sha256( + Path(__file__).with_name("semantic_geometry_fusion.py") + ), + "track_geometry_contract_sha256": _sha256( + Path(__file__).with_name("track_geometry.py") + ), + "e32_storage_adapter_sha256": _sha256( + Path(__file__).with_name("e32_track_geometry_storage.py") + ), + }, + "producer_sha256": _sha256(Path(__file__).resolve(strict=True)), + "authority": _authority(), + } + + +def _count_baseline( + counter: Counter[tuple[str, str, str, str, str]], + document: Mapping[str, object], + *, + session_seconds: float, + geometry_only: bool, +) -> None: + if geometry_only: + status = "single-source-geometry" + label = "__geometry__" + range_m = _optional_float(document.get("nearest_range_m")) + cause = "e29-unassociated-occupied-component" + else: + status = _required_string(document.get("geometry_status"), "E29 status") + label = _required_string(document.get("label"), "E29 label") + range_m = _optional_float(document.get("range_m")) + cause = _required_string(document.get("geometry_reason"), "E29 reason") + counter[ + ( + status, + label, + _range_bucket(range_m), + _scene_bucket(session_seconds), + cause, + ) + ] += 1 + + +def _count_e32( + counter: Counter[tuple[str, str, str, str, str]], + geometry: TrackGeometry, + *, + session_seconds: float, +) -> None: + counter[ + ( + geometry.evidence_state.value, + geometry.semantic_label or "__geometry__", + _range_bucket(geometry.range_m), + _scene_bucket(session_seconds), + geometry.reason_codes[0], + ) + ] += 1 + + +def _comparison_document( + baseline: Counter[tuple[str, str, str, str, str]], + current: Counter[tuple[str, str, str, str, str]], +) -> dict[str, Any]: + dimensions = ("status", "class", "range", "scene", "cause") + result: dict[str, Any] = {} + for dimension_index, dimension in enumerate(dimensions): + baseline_grouped: Counter[str] = Counter() + current_grouped: Counter[str] = Counter() + for key, count in baseline.items(): + baseline_grouped[key[dimension_index]] += count + for key, count in current.items(): + current_grouped[key[dimension_index]] += count + values = sorted(set(baseline_grouped) | set(current_grouped)) + result[f"by_{dimension}"] = { + value: { + "e29": baseline_grouped[value], + "e32": current_grouped[value], + "delta": current_grouped[value] - baseline_grouped[value], + } + for value in values + } + return result + + +def _change( + *, + frame_index: int, + source_frame_index: int, + session_seconds: float, + kind: str, + locator_index: int, + reason: str, + before: Mapping[str, object], + excluded_point_count: int, + after: Mapping[str, object] | None = None, + decision_item_id: str | None = None, +) -> dict[str, object]: + return { + "schema_version": E32_CHANGE_SCHEMA, + "frame_index": frame_index, + "source_frame_index": source_frame_index, + "session_seconds": session_seconds, + "kind": kind, + "locator_index": locator_index, + "reason": reason, + "decision_item_id": decision_item_id, + "before": dict(before), + "after": None if after is None else dict(after), + "excluded_point_count": excluded_point_count, + "absence_of_points_means_free": False, + "authority": _authority(), + } + + +def _geometry_locator(item: Mapping[str, object]) -> tuple[int, int]: + binding = _object(item.get("evidence_binding"), "E30 evidence binding") + locator = _object(item.get("e29_locator"), "E30 E29 locator") + if locator.get("kind") != "geometry-only-cluster": + raise E32TrackGeometryReplayError("E32 geometry correction locator changed") + return ( + _required_nonnegative_int(binding.get("frame_index"), "correction frame index"), + _required_nonnegative_int(locator.get("cluster_index"), "correction cluster index"), + ) + + +def _e29_frame(line: str, *, expected_frame_index: int) -> dict[str, Any]: + try: + value = json.loads(line) + except json.JSONDecodeError as exc: + raise E32TrackGeometryReplayError("E29 frame JSON is invalid") from exc + if ( + not isinstance(value, dict) + or value.get("schema_version") != CAMERA_GEOMETRY_FRAME_SCHEMA + or value.get("frame_index") != expected_frame_index + ): + raise E32TrackGeometryReplayError("E29 frame order or schema changed") + return value + + +def _read_existing_result( + root: Path, + expected_identity: dict[str, object], +) -> E32TrackGeometryReplay: + resolved = _safe_result_root(root, _E32_RESULT_ID) + manifest = _read_json(resolved / E32_MANIFEST_NAME) + _validate_content_identity( + root=resolved, + manifest=manifest, + schema=E32_TRACK_GEOMETRY_REPLAY_SCHEMA, + prefix="e32-track-geometry-", + ) + artifacts = _verified_artifacts(resolved, manifest) + report_path = artifacts.get("track-geometry-report") + required = { + "track-geometry-frames", + "e29-e32-changes", + "frame-point-offsets", + "point-source-indices", + "point-coordinates-map-f32", + "point-owner-indices", + "track-geometry-report", + } + if ( + set(artifacts) != required + or manifest.get("identity") != expected_identity + or report_path is None + ): + raise E32TrackGeometryReplayError("existing E32 result is incompatible") + report = _read_json(report_path) + if ( + report.get("schema_version") != E32_TRACK_GEOMETRY_REPORT_SCHEMA + or report.get("result_id") != resolved.name + or report.get("identity") != expected_identity + or report.get("status") != "accepted-diagnostic-track-geometry-replay" + or report.get("authority") != _authority() + ): + raise E32TrackGeometryReplayError("existing E32 report is incompatible") + try: + validate_storage( + artifacts=artifacts, + identity=expected_identity, + ) + except E32TrackGeometryStorageError as exc: + raise E32TrackGeometryReplayError(str(exc)) from exc + return E32TrackGeometryReplay( + result_root=resolved, + result_id=resolved.name, + manifest=manifest, + report=report, + ) + + +def _verified_camera_artifacts( + root: Path, + result: Mapping[str, object], +) -> dict[str, Path]: + artifacts = result.get("artifacts") + if not isinstance(artifacts, list): + raise E32TrackGeometryReplayError("camera result artifacts are invalid") + output: dict[str, Path] = {} + for value in artifacts: + artifact = _object(value, "camera result artifact") + kind = _required_string(artifact.get("kind"), "camera artifact kind") + path_value = _required_string(artifact.get("path"), "camera artifact path") + relative = Path(path_value) + if relative.is_absolute() or ".." in relative.parts or kind in output: + raise E32TrackGeometryReplayError("camera artifact path is unsafe") + path = root / relative + expected_size = _required_nonnegative_int( + artifact.get("byte_length"), + "camera artifact size", + ) + expected_sha = _required_sha256( + artifact.get("sha256"), + "camera artifact digest", + ) + if ( + path.is_symlink() + or not path.is_file() + or path.stat().st_size != expected_size + or _sha256(path) != expected_sha + ): + raise E32TrackGeometryReplayError("camera result artifact changed") + output[kind] = path + return output + + +def _artifact_by_kind( + result: Mapping[str, object], + kind: str, +) -> dict[str, object]: + artifacts = result.get("artifacts") + if not isinstance(artifacts, list): + raise E32TrackGeometryReplayError("camera result artifacts are invalid") + matches = [ + _object(value, "camera result artifact") + for value in artifacts + if isinstance(value, dict) and value.get("kind") == kind + ] + if len(matches) != 1: + raise E32TrackGeometryReplayError(f"camera artifact {kind} is not unique") + return matches[0] + + +def _source_artifact_sha256(source: E10LidarFieldSource) -> str: + artifact = _object(source.manifest.get("artifact"), "source pack artifact") + return _required_sha256(artifact.get("sha256"), "source pack artifact digest") + + +def _surface_artifact_sha256(surface: K1LocalSurfaceV1) -> str: + artifacts = surface.manifest.get("artifacts") + if not isinstance(artifacts, list): + raise E32TrackGeometryReplayError("local surface artifacts are invalid") + matches = [ + _object(value, "local surface artifact") + for value in artifacts + if isinstance(value, dict) and value.get("role") == "local-surface" + ] + if len(matches) != 1: + raise E32TrackGeometryReplayError("local surface artifact is not unique") + return _required_sha256( + matches[0].get("sha256"), + "local surface artifact digest", + ) + + +def _artifact(role: str, path: Path) -> dict[str, object]: + return { + "role": role, + "path": path.name, + "byte_length": path.stat().st_size, + "sha256": _sha256(path), + } + + +def _concatenate( + values: Iterable[npt.NDArray[Any]], + *, + dtype: np.dtype[Any], + trailing_shape: tuple[int, ...], +) -> npt.NDArray[Any]: + parts = tuple(np.asarray(value, dtype=dtype) for value in values if value.shape[0]) + if not parts: + return np.empty((0, *trailing_shape), dtype=dtype) + return np.concatenate(parts, axis=0).astype(dtype, copy=False) + + +def _range_bucket(value: float | None) -> str: + if value is None: + return "unavailable" + if value < 3.0: + return "near" + if value < 7.0: + return "middle" + return "far" + + +def _scene_bucket(session_seconds: float) -> str: + start = math.floor(session_seconds / 60.0) * 60 + return f"{start:03d}-{start + 60:03d}s" + + +def _distribution(values: Iterable[float]) -> dict[str, float | int | None]: + array = np.asarray(tuple(values), dtype=np.float64) + if array.size == 0: + return { + "sample_count": 0, + "minimum": None, + "p50": None, + "p95": None, + "maximum": None, + "mean": None, + } + return { + "sample_count": int(array.size), + "minimum": float(np.min(array)), + "p50": float(np.percentile(array, 50)), + "p95": float(np.percentile(array, 95)), + "maximum": float(np.max(array)), + "mean": float(np.mean(array)), + } + + +def _bbox(value: object) -> tuple[float, float, float, float]: + if not isinstance(value, list) or len(value) != 4: + raise E32TrackGeometryReplayError("semantic bbox is invalid") + result = tuple(float(item) for item in value) + if ( + not np.isfinite(result).all() + or result[2] <= result[0] + or result[3] <= result[1] + ): + raise E32TrackGeometryReplayError("semantic bbox is invalid") + return cast(tuple[float, float, float, float], result) + + +def _optional_float(value: object) -> float | None: + if value is None: + return None + if not isinstance(value, (int, float)) or isinstance(value, bool): + raise E32TrackGeometryReplayError("optional numeric value is invalid") + result = float(value) + if not math.isfinite(result): + raise E32TrackGeometryReplayError("optional numeric value is not finite") + return result + + +def _required_positive_float(value: object, label: str) -> float: + result = _optional_float(value) + if result is None or result <= 0.0: + raise E32TrackGeometryReplayError(f"{label} is invalid") + return result + + +def _required_nonnegative_float(value: object, label: str) -> float: + result = _optional_float(value) + if result is None or result < 0.0: + raise E32TrackGeometryReplayError(f"{label} is invalid") + return result + + +def _required_bool(value: object, label: str) -> bool: + if not isinstance(value, bool): + raise E32TrackGeometryReplayError(f"{label} is invalid") + return value + + +def _required_nonnegative_int(value: object, label: str) -> int: + result = _required_int(value, label) + if result < 0: + raise E32TrackGeometryReplayError(f"{label} is invalid") + return result + + +def _required_positive_int(value: object, label: str) -> int: + result = _required_int(value, label) + if result <= 0: + raise E32TrackGeometryReplayError(f"{label} is invalid") + return result + + +def _required_int(value: object, label: str) -> int: + if not isinstance(value, int) or isinstance(value, bool): + raise E32TrackGeometryReplayError(f"{label} is invalid") + return value + + +def _required_string(value: object, label: str) -> str: + if not isinstance(value, str) or not value or len(value) > 256: + raise E32TrackGeometryReplayError(f"{label} is invalid") + return value + + +def _required_sha256(value: object, label: str) -> str: + result = _required_string(value, label) + if _SHA256.fullmatch(result) is None: + raise E32TrackGeometryReplayError(f"{label} is invalid") + return result + + +def _object(value: object, label: str) -> dict[str, Any]: + if not isinstance(value, dict) or any(not isinstance(key, str) for key in value): + raise E32TrackGeometryReplayError(f"{label} must be an object") + return value + + +def _authority() -> dict[str, bool]: + return { + "commands_enabled": False, + "navigation_or_safety_accepted": False, + } + + +def _read_json(path: Path) -> dict[str, Any]: + if path.is_symlink() or not path.is_file(): + raise E32TrackGeometryReplayError(f"required JSON file is missing: {path.name}") + try: + value = json.loads(path.read_text(encoding="utf-8")) + except (OSError, json.JSONDecodeError) as exc: + raise E32TrackGeometryReplayError(f"invalid JSON file: {path.name}") from exc + return _object(value, path.name) + + +def _write_json(path: Path, value: Mapping[str, object]) -> None: + with path.open("x", encoding="utf-8") as stream: + json.dump( + value, + stream, + ensure_ascii=False, + sort_keys=True, + separators=(",", ":"), + allow_nan=False, + ) + stream.write("\n") + + +def _write_jsonl(stream: Any, value: Mapping[str, object]) -> None: + stream.write( + json.dumps( + value, + ensure_ascii=False, + sort_keys=True, + separators=(",", ":"), + allow_nan=False, + ) + + "\n" + ) + + +def _canonical_json(value: object) -> bytes: + return json.dumps( + value, + ensure_ascii=False, + sort_keys=True, + separators=(",", ":"), + allow_nan=False, + ).encode("utf-8") + + +def _sha256(path: Path) -> str: + digest = hashlib.sha256() + with path.open("rb") as stream: + for chunk in iter(lambda: stream.read(1024 * 1024), b""): + digest.update(chunk) + return digest.hexdigest() diff --git a/src/k1link/compute/e32_track_geometry_storage.py b/src/k1link/compute/e32_track_geometry_storage.py new file mode 100644 index 0000000..94476ad --- /dev/null +++ b/src/k1link/compute/e32_track_geometry_storage.py @@ -0,0 +1,347 @@ +"""Compact, strict storage adapter for E32 TrackGeometry frames.""" + +from __future__ import annotations + +import json +import math +from collections.abc import Mapping +from pathlib import Path +from typing import Any, Final, cast + +import numpy as np +import numpy.typing as npt + +from .track_geometry import ( + POINT_SLAB_SCHEMA, + TRACK_GEOMETRY_FRAME_SCHEMA, + PointSlab, + TrackGeometry, + TrackGeometryFrame, + TrackGeometrySourceBinding, +) + +E32_TRACK_GEOMETRY_RECORD_SCHEMA: Final = "missioncore.e32-track-geometry-record/v1" +E32_POINT_SLAB_REFERENCE_SCHEMA: Final = "missioncore.e32-point-slab-reference/v1" + +E32_FRAMES_NAME: Final = "track-geometry-frames.jsonl" +E32_FRAME_OFFSETS_NAME: Final = "frame-point-offsets.npy" +E32_SOURCE_INDICES_NAME: Final = "point-source-indices.npy" +E32_POINTS_NAME: Final = "point-coordinates-map-f32.npy" +E32_OWNER_INDICES_NAME: Final = "point-owner-indices.npy" + +Int64Array = npt.NDArray[np.int64] +UInt32Array = npt.NDArray[np.uint32] + + +class E32TrackGeometryStorageError(ValueError): + """Compact E32 storage no longer satisfies the TrackGeometry contract.""" + + +def write_point_storage( + *, + staging: Path, + frame_offsets: Int64Array, + source_indices: Int64Array, + points: npt.NDArray[np.float32], + owner_indices: UInt32Array, +) -> None: + """Write deterministic non-pickle arrays for the compact frame stream.""" + + _save_npy(staging / E32_FRAME_OFFSETS_NAME, frame_offsets) + _save_npy(staging / E32_SOURCE_INDICES_NAME, source_indices) + _save_npy(staging / E32_POINTS_NAME, points) + _save_npy(staging / E32_OWNER_INDICES_NAME, owner_indices) + + +def validate_storage( + *, + artifacts: Mapping[str, Path], + identity: Mapping[str, object], +) -> None: + """Reconstruct and validate every persisted TrackGeometry frame.""" + + frame_count = _positive_int(identity.get("frame_count"), "E32 frame count") + frame_offsets, source_indices, points, owner_indices = load_point_storage( + artifacts, + frame_count=frame_count, + ) + binding = TrackGeometrySourceBinding.from_dict( + identity.get("track_geometry_binding") + ) + observed_frames = 0 + with artifacts["track-geometry-frames"].open("r", encoding="utf-8") as stream: + for expected_frame_index, line in enumerate(stream): + if expected_frame_index >= frame_count: + raise E32TrackGeometryStorageError( + "E32 frame stream has extra rows" + ) + value = record(line, expected_frame_index=expected_frame_index) + frame_from_record( + record_value=value, + binding=binding, + frame_offsets=frame_offsets, + source_indices=source_indices, + points=points, + owner_indices=owner_indices, + ) + observed_frames += 1 + if observed_frames != frame_count: + raise E32TrackGeometryStorageError("E32 frame stream is incomplete") + + +def load_point_storage( + artifacts: Mapping[str, Path], + *, + frame_count: int, +) -> tuple[ + Int64Array, + Int64Array, + npt.NDArray[np.float32], + UInt32Array, +]: + """Open the four digest-verified E32 arrays as read-only memory maps.""" + + try: + frame_offsets = np.load( + artifacts["frame-point-offsets"], + allow_pickle=False, + mmap_mode="r", + ) + source_indices = np.load( + artifacts["point-source-indices"], + allow_pickle=False, + mmap_mode="r", + ) + points = np.load( + artifacts["point-coordinates-map-f32"], + allow_pickle=False, + mmap_mode="r", + ) + owner_indices = np.load( + artifacts["point-owner-indices"], + allow_pickle=False, + mmap_mode="r", + ) + except (KeyError, OSError, ValueError) as exc: + raise E32TrackGeometryStorageError( + "E32 point storage is unreadable" + ) from exc + if ( + frame_offsets.dtype != np.dtype(" TrackGeometryFrame: + """Reconstruct one TrackGeometryFrame from its JSON row and slab slices.""" + + expected_keys = { + "schema_version", + "track_geometry_frame_schema", + "frame_index", + "source_frame_index", + "session_seconds", + "source_available", + "point_slab", + "geometries", + "policy", + "authority", + } + if ( + set(record_value) != expected_keys + or record_value.get("schema_version") != E32_TRACK_GEOMETRY_RECORD_SCHEMA + or record_value.get("track_geometry_frame_schema") + != TRACK_GEOMETRY_FRAME_SCHEMA + or record_value.get("authority") != _authority() + or record_value.get("policy") + != { + "camera_owns_semantics": True, + "one_owner_per_source_point": True, + "current_held_persistent_are_separate": True, + "absence_of_points_means_free": False, + "unknown_remains_unknown": True, + } + ): + raise E32TrackGeometryStorageError( + "E32 frame record contract changed" + ) + frame_index = _nonnegative_int(record_value.get("frame_index"), "frame index") + if frame_index + 1 >= frame_offsets.size: + raise E32TrackGeometryStorageError( + "E32 frame point offset is missing" + ) + row_start = int(frame_offsets[frame_index]) + row_end = int(frame_offsets[frame_index + 1]) + slab_reference = _object( + record_value.get("point_slab"), + "E32 PointSlab reference", + ) + if ( + set(slab_reference) + != { + "schema_version", + "contract_schema", + "source_point_count", + "coordinate_frame", + "owner_keys", + "row_count", + } + or slab_reference.get("schema_version") + != E32_POINT_SLAB_REFERENCE_SCHEMA + or slab_reference.get("contract_schema") != POINT_SLAB_SCHEMA + or slab_reference.get("row_count") != row_end - row_start + ): + raise E32TrackGeometryStorageError( + "E32 PointSlab reference changed" + ) + owner_key_values = slab_reference.get("owner_keys") + geometry_values = record_value.get("geometries") + if not isinstance(owner_key_values, list) or not isinstance( + geometry_values, + list, + ): + raise E32TrackGeometryStorageError( + "E32 frame owner or geometry table changed" + ) + slab = PointSlab( + frame_index=frame_index, + source_frame_index=_nonnegative_int( + record_value.get("source_frame_index"), + "source frame index", + ), + source_point_count=_nonnegative_int( + slab_reference.get("source_point_count"), + "source point count", + ), + coordinate_frame=_string( + slab_reference.get("coordinate_frame"), + "point coordinate frame", + ), + owner_keys=tuple( + _string(value, "point owner key") for value in owner_key_values + ), + source_indices=np.asarray(source_indices[row_start:row_end], dtype=" dict[str, Any]: + """Parse one ordered compact frame record.""" + + try: + value = json.loads(line) + except json.JSONDecodeError as exc: + raise E32TrackGeometryStorageError( + "E32 frame record JSON is invalid" + ) from exc + result = _object(value, "E32 frame record") + if result.get("frame_index") != expected_frame_index: + raise E32TrackGeometryStorageError( + "E32 frame record order changed" + ) + return result + + +def _save_npy(path: Path, value: npt.NDArray[Any]) -> None: + if path.exists(): + raise E32TrackGeometryStorageError( + "E32 point artifact already exists" + ) + with path.open("xb") as stream: + np.save(stream, value, allow_pickle=False) + + +def _nonnegative_float(value: object, label: str) -> float: + if ( + not isinstance(value, (int, float)) + or isinstance(value, bool) + or not math.isfinite(float(value)) + or float(value) < 0.0 + ): + raise E32TrackGeometryStorageError(f"{label} is invalid") + return float(value) + + +def _nonnegative_int(value: object, label: str) -> int: + if not isinstance(value, int) or isinstance(value, bool) or value < 0: + raise E32TrackGeometryStorageError(f"{label} is invalid") + return value + + +def _positive_int(value: object, label: str) -> int: + result = _nonnegative_int(value, label) + if result == 0: + raise E32TrackGeometryStorageError(f"{label} is invalid") + return result + + +def _boolean(value: object, label: str) -> bool: + if not isinstance(value, bool): + raise E32TrackGeometryStorageError(f"{label} is invalid") + return value + + +def _string(value: object, label: str) -> str: + if not isinstance(value, str) or not value or len(value) > 256: + raise E32TrackGeometryStorageError(f"{label} is invalid") + return value + + +def _object(value: object, label: str) -> dict[str, Any]: + if not isinstance(value, dict) or any( + not isinstance(key, str) for key in value + ): + raise E32TrackGeometryStorageError(f"{label} must be an object") + return value + + +def _authority() -> dict[str, bool]: + return { + "commands_enabled": False, + "navigation_or_safety_accepted": False, + } diff --git a/tests/test_e32_track_geometry_replay.py b/tests/test_e32_track_geometry_replay.py new file mode 100644 index 0000000..6be7772 --- /dev/null +++ b/tests/test_e32_track_geometry_replay.py @@ -0,0 +1,432 @@ +from __future__ import annotations + +from collections import Counter + +import numpy as np +import pytest + +from k1link.compute.e32_track_geometry_replay import ( + E32TrackGeometryReplayError, + _comparison_document, + _CorrectionPlan, + _require_exact_e29_reproduction, + _translate_frame, +) +from k1link.compute.e32_track_geometry_storage import frame_from_record +from k1link.compute.semantic_geometry_fusion import ( + CAMERA_GEOMETRY_FRAME_SCHEMA, + _GeometryClusterSupport, + _SemanticSupport, +) +from k1link.compute.sensor_representation import K1_LIO_PCL_CAPABILITIES +from k1link.compute.track_geometry import ( + TrackGeometryCurrentness, + TrackGeometryEvidenceState, + TrackGeometryMetricBasis, + TrackGeometrySourceBinding, +) + + +def _binding() -> TrackGeometrySourceBinding: + return TrackGeometrySourceBinding( + source_pack_id="e10-lidar-pack-" + "a" * 64, + source_session_id="source-session", + representation_profile_id=K1_LIO_PCL_CAPABILITIES.profile_id, + e31_qualification_id="e31-source-qualification-" + "b" * 64, + calibration_sha256="c" * 64, + coordinate_frame="map", + time_basis="nearest-host-arrival-best-effort", + selected_offset_ms=0, + ) + + +def _semantic( + *, + track_id: int, + label: str, + status: str, + indices: list[int], + bbox: list[float], + current: bool = True, +) -> _SemanticSupport: + return _SemanticSupport( + document={ + "source_track_id": track_id, + "track_id": track_id, + "label": label, + "association_group": label, + "score": 0.9, + "bbox_xyxy": bbox, + "semantic_current": current, + "camera_motion_state": "unknown", + "camera_motion_confidence": None, + "motion_state": "unknown", + "motion_status": "unknown", + "unknown_is_occupied": True, + "navigation_or_safety_accepted": False, + "geometry_status": status, + "geometry_reason": ( + "camera-semantic-with-connected-occupied-lidar-support" + if status == "agree" + else ( + "semantic-observation-not-current" + if not current + else "camera-semantic-without-qualified-occupied-lidar-support" + ) + ), + "range_m": 4.0 if status == "agree" else None, + "occupied_centroid_map_xyz_m": None, + "occupied_height_range_m": None, + "support": { + "projected_points_in_bbox": len(indices), + "classified_points_in_bbox": len(indices), + "surface_points_in_bbox": 0, + "occupied_points_in_bbox": len(indices), + "below_surface_points_in_bbox": 0, + "connected_occupied_points": len(indices), + "connected_occupied_voxels": int(bool(indices)), + }, + }, + occupied_source_indices=np.asarray(indices, dtype=np.int64), + ) + + +def _geometry(indices: list[int], range_m: float) -> _GeometryClusterSupport: + return _GeometryClusterSupport( + document={ + "geometry_status": "single-source-geometry", + "semantic_class": None, + "point_count": len(indices), + "voxel_count": 1, + "centroid_map_xyz_m": [0.0, 0.0, 0.0], + "bounds_map_xyz_m": [[0.0, 0.0, 0.0], [1.0, 1.0, 1.0]], + "height_range_m": [0.2, 1.0], + "nearest_range_m": range_m, + "unknown_is_occupied": True, + "navigation_or_safety_accepted": False, + }, + occupied_source_indices=np.asarray(indices, dtype=np.int64), + ) + + +def test_e32_translation_applies_only_bound_corrections_and_closes_point_ownership() -> None: + points = np.asarray( + [[float(index), 0.0, 1.0] for index in range(8)], + dtype=np.float64, + ) + semantic_supports = ( + _semantic( + track_id=7, + label="car", + status="agree", + indices=[0, 1], + bbox=[10.0, 10.0, 100.0, 100.0], + ), + _semantic( + track_id=8, + label="person", + status="single-source-camera", + indices=[2], + bbox=[300.0, 500.0, 340.0, 580.0], + ), + _semantic( + track_id=9, + label="car", + status="unknown", + indices=[], + bbox=[200.0, 100.0, 250.0, 150.0], + current=False, + ), + ) + geometry_supports = ( + _geometry([3, 4], 3.0), + _geometry([5], 4.0), + _geometry([6, 7], 5.0), + ) + corrections = _CorrectionPlan( + semantic_rectangle_normalized_xyxy=(0.30, 0.75, 0.50, 1.0), + semantic_class_allowlist=frozenset({"person"}), + image_width=800, + image_height=600, + exact_geometry_corrections={(12, 0): "e30-review-item-" + "d" * 64}, + human_geometry_dispositions={ + (12, 1): ("background-or-noise", "e30-review-item-" + "e" * 64), + (12, 2): ("object-present", "e30-review-item-" + "f" * 64), + }, + ) + + translated = _translate_frame( + frame_index=12, + source_frame_index=120, + session_seconds=42.0, + source_available=True, + frame_points=points, + semantic_supports=semantic_supports, + geometry_supports=geometry_supports, + binding=_binding(), + corrections=corrections, + last_current_frame={}, + ) + + frame = translated.frame + assert [geometry.owner_key for geometry in frame.geometries] == [ + "track:7", + "geometry:2", + ] + assert frame.point_slab.owner_keys == ("track:7", "geometry:2") + assert frame.point_slab.source_indices.tolist() == [0, 1, 6, 7] + assert frame.point_slab.owner_indices.tolist() == [0, 0, 1, 1] + assert frame.geometries[1].reason_codes == ( + "e29-unassociated-occupied-component", + "a3-human-object-present", + ) + assert translated.semantic_published == 1 + assert translated.semantic_masked == 1 + assert translated.semantic_unpublishable_held == 1 + assert translated.geometry_published == 1 + assert translated.geometry_exact_excluded == 1 + assert translated.geometry_human_excluded == 1 + assert translated.baseline_qualified_points == 7 + assert translated.published_qualified_points == 4 + assert translated.excluded_qualified_points == 3 + assert translated.ownership_overlap_claims == 0 + assert translated.unqualified_semantic_support_points == 1 + assert [change["reason"] for change in translated.changes] == [ + "e31-semantic-self-mask", + "held-without-prior-current-provenance", + "e31-exact-geometry-correction", + "a3-human-background-or-noise", + ] + restored = frame_from_record( + record_value=translated.record, + binding=_binding(), + frame_offsets=np.asarray([0] * 13 + [4], dtype=" None: + held = _semantic( + track_id=11, + label="truck", + status="unknown", + indices=[], + bbox=[20.0, 20.0, 60.0, 60.0], + current=False, + ) + translated = _translate_frame( + frame_index=15, + source_frame_index=150, + session_seconds=45.0, + source_available=False, + frame_points=np.empty((0, 3), dtype=np.float64), + semantic_supports=(held,), + geometry_supports=(), + binding=_binding(), + corrections=_CorrectionPlan( + semantic_rectangle_normalized_xyxy=(0.30, 0.75, 0.50, 1.0), + semantic_class_allowlist=frozenset({"person"}), + image_width=800, + image_height=600, + exact_geometry_corrections={}, + human_geometry_dispositions={}, + ), + last_current_frame={11: 13}, + ) + + geometry = translated.frame.geometries[0] + assert geometry.currentness is TrackGeometryCurrentness.HELD + assert geometry.evidence_state is TrackGeometryEvidenceState.UNKNOWN + assert geometry.metric_basis is TrackGeometryMetricBasis.UNAVAILABLE + assert geometry.held_from_frame_index == 13 + assert translated.frame.point_slab.row_count == 0 + + +def test_e32_arbitrates_overlapping_camera_claims_without_duplicate_points() -> None: + larger = _semantic( + track_id=20, + label="car", + status="agree", + indices=[0, 1], + bbox=[10.0, 10.0, 100.0, 100.0], + ) + smaller = _semantic( + track_id=21, + label="person", + status="agree", + indices=[0, 1], + bbox=[20.0, 20.0, 40.0, 70.0], + ) + translated = _translate_frame( + frame_index=20, + source_frame_index=200, + session_seconds=50.0, + source_available=True, + frame_points=np.asarray( + [[1.0, 0.0, 1.0], [2.0, 0.0, 1.0]], + dtype=np.float64, + ), + semantic_supports=(larger, smaller), + geometry_supports=(), + binding=_binding(), + corrections=_CorrectionPlan( + semantic_rectangle_normalized_xyxy=(0.30, 0.75, 0.50, 1.0), + semantic_class_allowlist=frozenset({"person"}), + image_width=800, + image_height=600, + exact_geometry_corrections={}, + human_geometry_dispositions={}, + ), + last_current_frame={}, + ) + + assert translated.frame.point_slab.source_indices.tolist() == [0, 1] + assert translated.frame.point_slab.owner_keys == ("track:21",) + assert translated.frame.geometries[0].evidence_state is TrackGeometryEvidenceState.UNKNOWN + assert translated.frame.geometries[0].reason_codes[-1] == ( + "e32-point-ownership-collision" + ) + assert translated.frame.geometries[1].evidence_state is TrackGeometryEvidenceState.AGREE + assert translated.baseline_qualified_points == 4 + assert translated.published_qualified_points == 2 + assert translated.ownership_overlap_claims == 2 + assert translated.excluded_qualified_points == 0 + assert translated.changes[0]["reason"] == "point-ownership-arbitration" + + +def test_e32_withholds_unqualified_e29_range_and_retains_camera_state() -> None: + camera_only = _semantic( + track_id=30, + label="car", + status="single-source-camera", + indices=[0], + bbox=[10.0, 10.0, 100.0, 100.0], + ) + camera_only.document["range_m"] = 6.0 + translated = _translate_frame( + frame_index=30, + source_frame_index=300, + session_seconds=60.0, + source_available=True, + frame_points=np.asarray([[1.0, 0.0, 1.0]], dtype=np.float64), + semantic_supports=(camera_only,), + geometry_supports=(), + binding=_binding(), + corrections=_CorrectionPlan( + semantic_rectangle_normalized_xyxy=(0.30, 0.75, 0.50, 1.0), + semantic_class_allowlist=frozenset({"person"}), + image_width=800, + image_height=600, + exact_geometry_corrections={}, + human_geometry_dispositions={}, + ), + last_current_frame={}, + ) + + geometry = translated.frame.geometries[0] + assert geometry.evidence_state is TrackGeometryEvidenceState.CAMERA_ONLY + assert geometry.metric_basis is TrackGeometryMetricBasis.UNAVAILABLE + assert geometry.range_m is None + assert geometry.reason_codes[-1] == "e32-unqualified-range-withheld" + assert translated.unqualified_ranges_withheld == 1 + assert translated.changes[0]["reason"] == "unqualified-range-withheld" + + +def test_e32_replay_rejects_any_e29_reproduction_drift() -> None: + semantic = _semantic( + track_id=1, + label="car", + status="single-source-camera", + indices=[], + bbox=[10.0, 10.0, 20.0, 20.0], + ) + fusion_frame = { + "source_frame_index": 10, + "session_seconds": 3.0, + } + baseline = { + "schema_version": CAMERA_GEOMETRY_FRAME_SCHEMA, + "frame_index": 0, + "source_frame_index": 10, + "session_seconds": 3.0, + "source_available": False, + "local_surface_valid": False, + "semantic_observations": [semantic.document], + "geometry_only_occupied": [], + "policy": { + "camera_owns_semantics": True, + "lidar_owns_metric_geometry": True, + "absence_of_points_means_free": False, + "unknown_is_occupied": True, + }, + "authority": { + "commands_enabled": False, + "navigation_or_safety_accepted": False, + }, + } + _require_exact_e29_reproduction( + e29_frame=baseline, + frame_index=0, + fusion_frame=fusion_frame, + source_available=False, + surface_valid=False, + semantic_supports=(semantic,), + geometry_supports=(), + ) + baseline["semantic_observations"] = [] + with pytest.raises(E32TrackGeometryReplayError, match="exactly reproduce"): + _require_exact_e29_reproduction( + e29_frame=baseline, + frame_index=0, + fusion_frame=fusion_frame, + source_available=False, + surface_valid=False, + semantic_supports=(semantic,), + geometry_supports=(), + ) + + +def test_e32_comparison_exposes_status_class_range_scene_and_cause_deltas() -> None: + baseline = { + ( + "agree", + "car", + "middle", + "000-060s", + "connected-support", + ): 2, + ( + "single-source-geometry", + "__geometry__", + "near", + "000-060s", + "unassociated", + ): 1, + } + current = { + ( + "agree", + "car", + "middle", + "000-060s", + "connected-support", + ): 1, + } + + comparison = _comparison_document( + Counter(baseline), + Counter(current), + ) + + assert comparison["by_status"]["agree"] == { + "e29": 2, + "e32": 1, + "delta": -1, + } + assert comparison["by_class"]["__geometry__"]["delta"] == -1 + assert comparison["by_range"]["near"]["delta"] == -1 + assert comparison["by_scene"]["000-060s"]["delta"] == -2 + assert comparison["by_cause"]["unassociated"]["delta"] == -1