from __future__ import annotations import json import struct from pathlib import Path from typing import Any, cast import lz4.block import pytest from fastapi import APIRouter from fastapi.routing import APIRoute from k1link.compute import ( K1_LIDAR_PACK_V2_PROFILE, K1LocalSurfaceV1, LidarPipelineStage, LidarReadiness, LidarReplayError, LidarReplayPackV2, assess_lidar_profile, build_k1_local_surface, build_lidar_replay_pack_v2, verify_lidar_replay_equivalence, ) from k1link.device_plugins.xgrids_k1.mqtt.capture import FRAME_HEADER, RAW_MAGIC from k1link.web.lidar_api import build_lidar_router def _varint(value: int) -> bytes: encoded = bytearray() while value > 0x7F: encoded.append((value & 0x7F) | 0x80) value >>= 7 encoded.append(value) return bytes(encoded) def _key(number: int, wire_type: int) -> bytes: return _varint((number << 3) | wire_type) def _uint(number: int, value: int) -> bytes: return _key(number, 0) + _varint(value) def _sint(number: int, value: int) -> bytes: zigzag = (value << 1) ^ (value >> 63) return _uint(number, zigzag & 0xFFFFFFFFFFFFFFFF) def _bytes(number: int, value: bytes) -> bytes: return _key(number, 2) + _varint(len(value)) + value def _fixed32(number: int, value: float) -> bytes: return _key(number, 5) + struct.pack(" bytes: return _key(number, 1) + struct.pack(" bytes: return b"".join( ( _uint(1, sequence), _sint(2, stamp), _sint(3, scaler), _bytes(4, b"device-redacted"), _bytes(5, b"session-redacted"), ) ) def _pcl_payload(sequence: int, stamp: int, first_rgbi: int) -> bytes: point_1 = ( _sint(1, 1000 * sequence) + _sint(2, -2000) + _sint(3, 500) + _uint(4, first_rgbi) ) point_2 = ( _sint(1, -250) + _sint(2, sequence) + _sint(3, 4000) + _uint(4, 0xAABBCC09) ) report = ( _bytes(1, _header(sequence, stamp)) + _bytes(2, point_1) + _bytes(2, point_2) ) compressed = lz4.block.compress(report, store_size=False) return _uint(3, len(report)) + _bytes(4, compressed) def _pose_payload(sequence: int, stamp: int) -> bytes: position = ( _fixed64(1, float(sequence)) + _fixed64(2, -2.5) + _fixed64(3, 3.75) ) orientation = ( _fixed64(1, 0.0) + _fixed64(2, 0.0) + _fixed64(3, 0.0) + _fixed64(4, 1.0) ) pose = _bytes(1, position) + _bytes(2, orientation) stamped = _sint(1, stamp + 1) + _bytes(2, pose) return ( _bytes(1, _header(sequence, stamp)) + _bytes(2, stamped) + _fixed32(3, 12.5) + _fixed32(4, 0.001) ) def _capture(tmp_path: Path) -> Path: session = tmp_path / "session-001" root = session / "captures" / "mqtt_live" root.mkdir(parents=True) messages = [ ("lixel/application/report/lio_pcl", _pcl_payload(10, 1_000, 0x11223344)), ("lixel/application/report/lio_pose", _pose_payload(20, 1_010)), ("lixel/application/report/lio_pcl", _pcl_payload(11, 2_000, 0x55667788)), ("lixel/application/report/lio_pose", _pose_payload(21, 2_010)), ] raw = bytearray(RAW_MAGIC) metadata: list[str] = [] for sequence, (topic, payload) in enumerate(messages, start=1): topic_bytes = topic.encode() raw.extend(FRAME_HEADER.pack(len(topic_bytes), len(payload))) raw.extend(topic_bytes) raw.extend(payload) metadata.append( json.dumps( { "schema_version": 1, "record_type": "message", "sequence": sequence, "received_at_epoch_ns": 1_000_000_000 + sequence * 10_000_000, "received_monotonic_ns": 5_000_000_000 + sequence * 10_000_000, }, separators=(",", ":"), ) ) capture = root / "mqtt.raw.k1mqtt" capture.write_bytes(raw) (root / "mqtt.metadata.jsonl").write_text( "\n".join(metadata) + "\n", encoding="utf-8", ) (root / "mqtt.timeline.origin.json").write_text( json.dumps( { "schema_version": 1, "started_at_epoch_ns": 1_000_000_000, "started_monotonic_ns": 5_000_000_000, } ), encoding="utf-8", ) return capture def _endpoint(router: APIRouter, path: str) -> object: for route in router.routes: if ( isinstance(route, APIRoute) and route.path == path and "GET" in (route.methods or set()) ): return route.endpoint raise AssertionError(f"GET {path} route is missing") def test_lidar_replay_v2_retains_fields_and_passes_equivalence(tmp_path: Path) -> None: capture = _capture(tmp_path) output = build_lidar_replay_pack_v2(capture, tmp_path / "packs") pack = LidarReplayPackV2(output) try: assert pack.identity["session_id"] == "session-001" assert pack.identity["lidar_evidence_profile"] == K1_LIDAR_PACK_V2_PROFILE.to_dict() assert pack.point_frame_count == 2 assert pack.pose_frame_count == 2 assert pack.point_count == 4 point = pack.point_frame(0) assert point.capture_sequence == 1 assert point.received_at_epoch_ns == 1_010_000_000 assert point.received_monotonic_ns == 5_010_000_000 assert point.header_seq == 10 assert point.header_stamp == 1_000 assert point.scaler == 1_000 assert point.raw_xyz.tolist() == [[10_000, -2_000, 500], [-250, 10, 4_000]] assert point.xyz_map.tolist() == [[10.0, -2.0, 0.5], [-0.25, 0.01, 4.0]] assert point.rgbi.tolist() == [0x11223344, 0xAABBCC09] assert point.intensity.tolist() == [0x44, 0x09] assert point.decoded_view().intensities == bytes((0x44, 0x09)) pose = pack.pose_frame(0) assert pose.capture_sequence == 2 assert pose.header_seq == 20 assert pose.pose_stamp == 1_011 assert pose.position_map == (20.0, -2.5, 3.75) assert pose.decoded_view().frame_id == "map" assert pack.quality["field_retention"]["raw_rgbi"] is True assert pack.quality["field_retention"]["host_monotonic_ns"] is True assert pack.quality["sensor_range_m"]["sample_count"] == 0 assert pack.quality["sensor_range_status"].startswith("unavailable-") assert pack.quality["pose_binding"]["coverage_fraction"] == pytest.approx(1.0) assert pack.equivalence["status"] == "passed" assert pack.equivalence["array_mismatches"] == 0 rerun = verify_lidar_replay_equivalence(capture, pack) assert rerun["status"] == "passed" surface_output = build_k1_local_surface( pack, tmp_path / "local-surfaces", display_name="synthetic replay v2 local surface", ) surface = K1LocalSurfaceV1(surface_output) try: assert surface.identity["source_pack_id"] == pack.pack_id assert ( surface.identity["source_pack_identity_sha256"] == pack.manifest["identity_sha256"] ) assert ( surface.identity["source_schema_version"] == "missioncore.lidar-replay-pack/v2" ) assert ( surface.identity["source_representation"] == "lossless-lidar-replay-v2" ) assert surface.identity["frame_count"] == pack.point_frame_count assert surface.identity["point_count"] == pack.point_count surface_source = cast(dict[str, Any], surface.report["source"]) assert surface_source["intensity_available"] is True detail = surface.frame_detail(pack, 0) assert detail["source_pack_id"] == pack.pack_id assert detail["source_frame_index"] == 1 detail_pose = cast(dict[str, Any], detail["pose"]) assert detail_pose["binding_age_ms"] == pytest.approx(10.0) finally: surface.close() finally: pack.close() assert build_lidar_replay_pack_v2(capture, tmp_path / "packs") == output def test_lidar_replay_v2_unblocks_detector_input_but_not_mapping() -> None: assessments = { item.stage: item for item in assess_lidar_profile(K1_LIDAR_PACK_V2_PROFILE) } assert ( assessments[LidarPipelineStage.LIDAR_3D_DETECTION].readiness is LidarReadiness.DEGRADED ) assert ( assessments[LidarPipelineStage.LIDAR_INERTIAL_SLAM].readiness is LidarReadiness.BLOCKED ) assert ( assessments[LidarPipelineStage.NVIDIA_NVBLOX].readiness is LidarReadiness.BLOCKED ) def test_lidar_replay_v2_fails_closed_without_exact_host_timing( tmp_path: Path, ) -> None: capture = _capture(tmp_path) capture.with_name("mqtt.metadata.jsonl").unlink() with pytest.raises(LidarReplayError, match="exact host timing"): build_lidar_replay_pack_v2(capture, tmp_path / "packs") def test_lidar_api_exposes_path_free_quality_and_equivalence(tmp_path: Path) -> None: capture = _capture(tmp_path) root = tmp_path / "packs" output = build_lidar_replay_pack_v2(capture, root) router = build_lidar_router(root_provider=lambda: root) catalog_route = _endpoint(router, "/api/v1/lidar/replay-packs") detail_route = _endpoint(router, "/api/v1/lidar/replay-packs/{pack_id}") catalog = catalog_route(limit=20) # type: ignore[operator] detail = detail_route(pack_id=output.name) # type: ignore[operator] assert catalog["schema_version"] == "missioncore.lidar-replay-pack-catalog/v1" assert catalog["valid_total"] == 1 assert catalog["invalid_total"] == 0 assert catalog["items"][0]["equivalence_status"] == "passed" assert detail["pack"]["pack_id"] == output.name assert detail["quality"]["field_retention"]["raw_rgbi"] is True assert detail["equivalence"]["status"] == "passed" assert detail["access"] == "read-only" serialized = repr({"catalog": catalog, "detail": detail}) assert str(tmp_path) not in serialized