diff --git a/config/deployment/mission-core-worker-shadow-v1.template.json b/config/deployment/mission-core-worker-shadow-v1.template.json index cffb7e5..2fd1b60 100644 --- a/config/deployment/mission-core-worker-shadow-v1.template.json +++ b/config/deployment/mission-core-worker-shadow-v1.template.json @@ -139,7 +139,7 @@ }, "wheel": { "name": "nodedc_mission_core-0.1.0-py3-none-any.whl", - "sha256": "ecbbfeee7ea62a7f5f5368efccf5d254abb17b4801a29538e3bc7aa3c8c3a40a" + "sha256": "19d8caf9a522747c461fb3ca30aafe54169959d8bd8e671fa6fc8c0ac107875d" } }, "rollback": { diff --git a/config/perception/m4-reuse-inventory-v1.json b/config/perception/m4-reuse-inventory-v1.json index 1510cb5..fe9ce29 100644 --- a/config/perception/m4-reuse-inventory-v1.json +++ b/config/perception/m4-reuse-inventory-v1.json @@ -20,7 +20,7 @@ { "module": "k1link.compute.temporal_occupied_layer", "role": "bounded hit-only occupied and unknown state", - "admission": "adapt-to-product-contract" + "admission": "reference-only-extracted-to-product" }, { "module": "k1link.compute.pipeline_telemetry", @@ -41,6 +41,16 @@ "module": "k1link.perception.geometry", "role": "digest-bound GeometryAssociationProvider and exact point ownership", "admission": "product-owned" + }, + { + "module": "k1link.perception.temporal", + "role": "product-owned bounded spatial retention with ephemeral identity and explicit expiry", + "admission": "product-owned" + }, + { + "module": "k1link.perception.motion", + "role": "product-owned class-independent bounded-history motion estimator", + "admission": "product-owned" } ], "historical_wrappers": [ @@ -60,6 +70,10 @@ "module": "k1link.compute.e35_degradation_recovery", "reason": "immutable degradation qualification wrapper" }, + { + "module": "k1link.compute.e51_motion_semantic_qualification", + "reason": "immutable diagnostic motion qualification wrapper" + }, { "module": "k1link.compute.e46j_raw_fisheye_realtime", "reason": "immutable detector-capacity experiment wrapper" diff --git a/config/perception/m4-temporal-motion-v1.json b/config/perception/m4-temporal-motion-v1.json new file mode 100644 index 0000000..6882811 --- /dev/null +++ b/config/perception/m4-temporal-motion-v1.json @@ -0,0 +1,59 @@ +{ + "schema_version": "missioncore.temporal-motion-profile/v1", + "profile_id": "m4-bounded-temporal-motion/v1", + "temporal_provider_id": "bounded-spatial-temporal-layer/v1", + "motion_provider_id": "class-independent-motion-estimator/v1", + "source": { + "source_id": "RAVNOVES00", + "session_id": "20260720T065719Z_viewer_live", + "geometry_result_id": "m4-geometry-replay-8daf3109e3cf30b960b4b376032ff3b5ec58ca42a1e5899b841cf29fbcf14ad8", + "geometry_frames_sha256": "b4db5d0ebaba4d6268a1006707dc313c229f3dbdfd73b1d863cad5d853be8ac4" + }, + "temporal": { + "coordinate_frame": "map", + "voxel_size_m": 0.45, + "occupied_ttl_seconds": 0.75, + "maximum_active_components": 256, + "maximum_cells_per_component": 4096, + "maximum_history_samples": 8, + "association_maximum_gap_seconds": 0.35, + "association_maximum_centroid_distance_m": 0.9, + "association_minimum_voxel_overlap_fraction": 0.05, + "association_neighbor_radius_cells": 1, + "jump_maximum_adjacent_gap_seconds": 0.25, + "jump_minimum_matched_components": 4, + "jump_minimum_median_displacement_m": 1.5, + "jump_minimum_p25_displacement_m": 0.9 + }, + "motion": { + "minimum_observations": 3, + "minimum_span_seconds": 0.2, + "moving_minimum_displacement_m": 0.25, + "moving_minimum_speed_mps": 0.4, + "stationary_maximum_displacement_m": 0.1, + "stationary_maximum_speed_mps": 0.2, + "maximum_speed_mps": 20.0, + "full_confidence_observations": 4, + "full_confidence_span_seconds": 0.4, + "minimum_confidence": 0.75 + }, + "policy": { + "component_identity_scope": "ephemeral", + "association_uses_detector_id": false, + "association_uses_semantic_class": false, + "motion_uses_semantic_class": false, + "camera_only_creates_occupied_component": false, + "absence_of_points_means_free": false, + "held_state": "unknown", + "expired_cells_published": false, + "long_term_identity_available": false, + "history_is_bounded": true + }, + "authority": { + "ground_truth": false, + "physical_live": false, + "commands_enabled": false, + "actuation_allowed": false, + "navigation_or_safety_accepted": false + } +} diff --git a/scripts/build_m4_worker_shadow_artifact.py b/scripts/build_m4_worker_shadow_artifact.py index 607be6b..00d4fe6 100644 --- a/scripts/build_m4_worker_shadow_artifact.py +++ b/scripts/build_m4_worker_shadow_artifact.py @@ -25,7 +25,7 @@ WHEEL_NAME = "nodedc_mission_core-0.1.0-py3-none-any.whl" RUNNER_NAME = RUNNER.name PATCH_ID = re.compile(r"^[A-Za-z0-9._-]{1,96}$") EXPECTED_BASELINE_SHA256 = "ea10359339e6cce31b5780a2710299771cab7cc0c1c2a2b56a1621f786b31fa8" -EXPECTED_WHEEL_SHA256 = "ecbbfeee7ea62a7f5f5368efccf5d254abb17b4801a29538e3bc7aa3c8c3a40a" +EXPECTED_WHEEL_SHA256 = "19d8caf9a522747c461fb3ca30aafe54169959d8bd8e671fa6fc8c0ac107875d" PAYLOAD_FILES = ( RUNNER_NAME, WHEEL_NAME, diff --git a/src/k1link/perception/geometry.py b/src/k1link/perception/geometry.py index 44d7314..0a93c5a 100644 --- a/src/k1link/perception/geometry.py +++ b/src/k1link/perception/geometry.py @@ -235,6 +235,16 @@ class RecordedGeometryStore: surface_valid=bool(self._surface["frame_valid"][frame_index]), ) + def current_points(self, packet: SourcePacket) -> FloatArray | None: + """Expose the verified frame-local point index space to temporal occupancy.""" + + frame = self.frame(packet) + if frame is None or not frame.surface_valid: + return None + points = np.asarray(frame.points_map, dtype=np.float64) + points.setflags(write=False) + return points + def _validate(self) -> None: source_required = { "frame_indices", diff --git a/src/k1link/perception/motion.py b/src/k1link/perception/motion.py new file mode 100644 index 0000000..5146f5c --- /dev/null +++ b/src/k1link/perception/motion.py @@ -0,0 +1,184 @@ +"""Class-independent bounded-history motion estimation for Mission Core.""" + +from __future__ import annotations + +import math +from dataclasses import dataclass, replace + +from .contracts import MotionState, TemporalObstacle, TemporalState +from .providers import SourcePacket +from .temporal import MOTION_PROVIDER_ID, MotionEstimatorProfile, TemporalMotionProfile + + +class MotionEstimatorError(RuntimeError): + """Temporal history cannot support a deterministic motion decision.""" + + +@dataclass(frozen=True, slots=True) +class MotionEstimatorSnapshot: + input_frames: int + input_obstacles: int + moving: int + stationary: int + unknown: int + insufficient_history: int + stale_support: int + map_frame_discontinuity: int + confidence_below_threshold: int + threshold_deadband: int + implausible_speed: int + + +class ClassIndependentMotionEstimator: + """Estimate map-frame motion without labels, detector IDs or tracklets.""" + + provider_id: str = MOTION_PROVIDER_ID + + def __init__(self, *, profile: TemporalMotionProfile) -> None: + self.profile = profile + self.config = profile.motion + self._input_frames = 0 + self._input_obstacles = 0 + self._moving = 0 + self._stationary = 0 + self._unknown = 0 + self._reasons: dict[str, int] = { + "insufficient-history": 0, + "stale-support": 0, + "map-frame-discontinuity": 0, + "confidence-below-threshold": 0, + "threshold-deadband": 0, + "implausible-speed": 0, + } + + def estimate( + self, + packet: SourcePacket, + obstacles: tuple[TemporalObstacle, ...], + ) -> tuple[TemporalObstacle, ...]: + if ( + packet.envelope.source_id != self.profile.source_id + or packet.envelope.session_id != self.profile.session_id + ): + raise MotionEstimatorError("packet escaped the motion source profile") + self._input_frames += 1 + self._input_obstacles += len(obstacles) + result = tuple(self._estimate_one(packet, obstacle) for obstacle in obstacles) + for obstacle in result: + if obstacle.motion is MotionState.MOVING: + self._moving += 1 + elif obstacle.motion is MotionState.STATIONARY: + self._stationary += 1 + else: + self._unknown += 1 + if obstacle.motion_reason in self._reasons: + self._reasons[obstacle.motion_reason] += 1 + return result + + def _estimate_one( + self, + packet: SourcePacket, + obstacle: TemporalObstacle, + ) -> TemporalObstacle: + now_ns = packet.envelope.timestamps.source_ns + if obstacle.last_hit_ns > now_ns or obstacle.age_ns != now_ns - obstacle.last_hit_ns: + raise MotionEstimatorError("temporal obstacle time escaped its source packet") + if obstacle.state is not TemporalState.CURRENT: + return _unknown(obstacle, "stale-support") + if obstacle.association_basis == "map-frame-discontinuity": + return _unknown(obstacle, "map-frame-discontinuity") + history = obstacle.history + times = tuple(sample.evidence_time_ns for sample in history) + if any(right <= left for left, right in zip(times, times[1:], strict=False)): + raise MotionEstimatorError("motion history is not strictly monotonic") + if history[-1].evidence_time_ns != obstacle.last_hit_ns: + raise MotionEstimatorError("motion history does not end at the current hit") + if len(history) < self.config.minimum_observations: + return _unknown(obstacle, "insufficient-history") + span_seconds = (times[-1] - times[0]) / 1_000_000_000 + if span_seconds < self.config.minimum_span_seconds: + return _unknown(obstacle, "insufficient-history") + displacement_m = math.dist( + history[0].centroid_xyz_m, + history[-1].centroid_xyz_m, + ) + speed_mps = displacement_m / span_seconds + if speed_mps > self.config.maximum_speed_mps: + return _unknown(obstacle, "implausible-speed") + confidence = _evidence_confidence( + self.config, + observation_count=len(history), + span_seconds=span_seconds, + ) + if confidence < self.config.minimum_confidence: + return _unknown(obstacle, "confidence-below-threshold") + if ( + displacement_m >= self.config.moving_minimum_displacement_m + and speed_mps >= self.config.moving_minimum_speed_mps + ): + return replace( + obstacle, + motion=MotionState.MOVING, + motion_confidence=confidence, + motion_reason="bounded-map-history-moving", + ) + if ( + displacement_m <= self.config.stationary_maximum_displacement_m + and speed_mps <= self.config.stationary_maximum_speed_mps + ): + return replace( + obstacle, + motion=MotionState.STATIONARY, + motion_confidence=confidence, + motion_reason="bounded-map-history-stationary", + ) + return _unknown(obstacle, "threshold-deadband") + + def snapshot(self) -> MotionEstimatorSnapshot: + return MotionEstimatorSnapshot( + input_frames=self._input_frames, + input_obstacles=self._input_obstacles, + moving=self._moving, + stationary=self._stationary, + unknown=self._unknown, + insufficient_history=self._reasons["insufficient-history"], + stale_support=self._reasons["stale-support"], + map_frame_discontinuity=self._reasons["map-frame-discontinuity"], + confidence_below_threshold=self._reasons["confidence-below-threshold"], + threshold_deadband=self._reasons["threshold-deadband"], + implausible_speed=self._reasons["implausible-speed"], + ) + + +def _evidence_confidence( + config: MotionEstimatorProfile, + *, + observation_count: int, + span_seconds: float, +) -> float: + """Return bounded evidence sufficiency, not a statistical class probability.""" + + return round( + min( + 1.0, + observation_count / config.full_confidence_observations, + span_seconds / config.full_confidence_span_seconds, + ), + 12, + ) + + +def _unknown(obstacle: TemporalObstacle, reason: str) -> TemporalObstacle: + return replace( + obstacle, + motion=MotionState.UNKNOWN, + motion_confidence=0.0, + motion_reason=reason, + ) + + +__all__ = [ + "ClassIndependentMotionEstimator", + "MotionEstimatorError", + "MotionEstimatorSnapshot", +] diff --git a/src/k1link/perception/temporal.py b/src/k1link/perception/temporal.py new file mode 100644 index 0000000..71ad13d --- /dev/null +++ b/src/k1link/perception/temporal.py @@ -0,0 +1,761 @@ +"""Bounded spatial temporal occupancy for the Mission Core product graph.""" + +from __future__ import annotations + +import hashlib +import json +import math +from collections import deque +from dataclasses import dataclass, field +from pathlib import Path +from typing import Final, Protocol + +import numpy as np +import numpy.typing as npt + +from .contracts import ( + EvidenceCurrentness, + GridCell, + HistorySample, + MotionState, + ObstacleObservation, + TemporalObstacle, + TemporalState, +) +from .providers import SourcePacket + +TEMPORAL_MOTION_PROFILE_SCHEMA: Final = "missioncore.temporal-motion-profile/v1" +TEMPORAL_PROVIDER_ID: Final = "bounded-spatial-temporal-layer/v1" +MOTION_PROVIDER_ID: Final = "class-independent-motion-estimator/v1" +DEFAULT_TEMPORAL_MOTION_PROFILE_PATH: Final = Path( + "config/perception/m4-temporal-motion-v1.json" +) + +FloatArray = npt.NDArray[np.float64] + + +class TemporalProviderError(RuntimeError): + """Temporal input, configuration or bounded state is incompatible.""" + + +class CurrentPointResolver(Protocol): + """Resolve the current frame's source-local point index space.""" + + def current_points(self, packet: SourcePacket) -> FloatArray | None: ... + + +@dataclass(frozen=True, slots=True) +class TemporalLayerProfile: + coordinate_frame: str + voxel_size_m: float + occupied_ttl_seconds: float + maximum_active_components: int + maximum_cells_per_component: int + maximum_history_samples: int + association_maximum_gap_seconds: float + association_maximum_centroid_distance_m: float + association_minimum_voxel_overlap_fraction: float + association_neighbor_radius_cells: int + jump_maximum_adjacent_gap_seconds: float + jump_minimum_matched_components: int + jump_minimum_median_displacement_m: float + jump_minimum_p25_displacement_m: float + + def __post_init__(self) -> None: + positive = ( + self.voxel_size_m, + self.occupied_ttl_seconds, + self.association_maximum_gap_seconds, + self.association_maximum_centroid_distance_m, + self.association_minimum_voxel_overlap_fraction, + self.jump_maximum_adjacent_gap_seconds, + self.jump_minimum_median_displacement_m, + self.jump_minimum_p25_displacement_m, + ) + if any(not _positive_finite(value) for value in positive): + raise TemporalProviderError("temporal numeric bound is invalid") + if self.association_minimum_voxel_overlap_fraction > 1.0: + raise TemporalProviderError("temporal overlap fraction is invalid") + positive_integers = ( + self.maximum_active_components, + self.maximum_cells_per_component, + self.maximum_history_samples, + self.association_neighbor_radius_cells, + self.jump_minimum_matched_components, + ) + if any(not _positive_integer(value) for value in positive_integers): + raise TemporalProviderError("temporal integer bound is invalid") + if self.maximum_history_samples > 32: + raise TemporalProviderError("temporal history exceeds the product contract") + + @property + def ttl_ns(self) -> int: + return round(self.occupied_ttl_seconds * 1_000_000_000) + + @property + def association_maximum_gap_ns(self) -> int: + return round(self.association_maximum_gap_seconds * 1_000_000_000) + + +@dataclass(frozen=True, slots=True) +class MotionEstimatorProfile: + minimum_observations: int + minimum_span_seconds: float + moving_minimum_displacement_m: float + moving_minimum_speed_mps: float + stationary_maximum_displacement_m: float + stationary_maximum_speed_mps: float + maximum_speed_mps: float + full_confidence_observations: int + full_confidence_span_seconds: float + minimum_confidence: float + + def __post_init__(self) -> None: + if any( + not _positive_integer(value) + for value in (self.minimum_observations, self.full_confidence_observations) + ): + raise TemporalProviderError("motion observation bound is invalid") + numeric = ( + self.minimum_span_seconds, + self.moving_minimum_displacement_m, + self.moving_minimum_speed_mps, + self.stationary_maximum_displacement_m, + self.stationary_maximum_speed_mps, + self.maximum_speed_mps, + self.full_confidence_span_seconds, + self.minimum_confidence, + ) + if any(not _positive_finite(value) for value in numeric): + raise TemporalProviderError("motion numeric bound is invalid") + if not 0.0 < self.minimum_confidence <= 1.0: + raise TemporalProviderError("motion confidence threshold is invalid") + if ( + self.stationary_maximum_displacement_m + >= self.moving_minimum_displacement_m + or self.stationary_maximum_speed_mps >= self.moving_minimum_speed_mps + or self.moving_minimum_speed_mps >= self.maximum_speed_mps + ): + raise TemporalProviderError("motion thresholds have no conservative deadband") + + +@dataclass(frozen=True, slots=True) +class TemporalMotionProfile: + profile_id: str + source_id: str + session_id: str + geometry_result_id: str + geometry_frames_sha256: str + temporal: TemporalLayerProfile + motion: MotionEstimatorProfile + profile_sha256: str + + +@dataclass(frozen=True, slots=True) +class TemporalProviderSnapshot: + input_frames: int + completed_frames: int + failed_frames: int + input_observations: int + current_occupied_observations: int + nonmetric_uncertainty_observations: int + created_components: int + spatial_reassociations: int + detector_identity_changes_reassociated: int + current_publications: int + held_publications: int + expired_publications: int + map_frame_jump_candidates: int + peak_active_components: int + peak_cells_per_component: int + maximum_history_samples: int + maximum_held_age_ns: int + maximum_expiry_materialization_delay_ns: int + past_ttl_occupied_publications: int + + +@dataclass(frozen=True, slots=True) +class _SpatialObservation: + occupancy_key: str + frame_id: str + evidence_time_ns: int + centroid_xyz_m: tuple[float, float, float] + cells: frozenset[GridCell] + semantic_hint: str | None + + +@dataclass(slots=True) +class _Component: + component_id: str + last_occupancy_key: str + last_hit_ns: int + cells: frozenset[GridCell] + centroid_xyz_m: tuple[float, float, float] + semantic_hint: str | None + history: deque[HistorySample] = field(default_factory=deque) + + +class BoundedSpatialTemporalProvider: + """Publish hit-backed current, short-held unknown and cell-free expiry.""" + + provider_id: str = TEMPORAL_PROVIDER_ID + + def __init__( + self, + *, + point_resolver: CurrentPointResolver, + profile: TemporalMotionProfile, + ) -> None: + self.point_resolver = point_resolver + self.profile = profile + self.config = profile.temporal + self._components: dict[str, _Component] = {} + self._next_component = 1 + self._previous_sequence: int | None = None + self._previous_time_ns: int | None = None + self._previous_observations: tuple[_SpatialObservation, ...] = () + self._input_frames = 0 + self._completed_frames = 0 + self._failed_frames = 0 + self._input_observations = 0 + self._current_observations = 0 + self._uncertainty_observations = 0 + self._created_components = 0 + self._spatial_reassociations = 0 + self._identity_changes_reassociated = 0 + self._current_publications = 0 + self._held_publications = 0 + self._expired_publications = 0 + self._map_frame_jumps = 0 + self._peak_active_components = 0 + self._peak_cells = 0 + self._maximum_history = 0 + self._maximum_held_age_ns = 0 + self._maximum_expiry_materialization_delay_ns = 0 + + def update( + self, + packet: SourcePacket, + observations: tuple[ObstacleObservation, ...], + ) -> tuple[TemporalObstacle, ...]: + self._input_frames += 1 + self._input_observations += len(observations) + try: + self._validate_packet(packet) + spatial, uncertainty_count = self._spatial_observations(packet, observations) + self._current_observations += len(spatial) + self._uncertainty_observations += uncertainty_count + now_ns = packet.envelope.timestamps.source_ns + expired = self._expire(now_ns) + jump = self._map_frame_jump(now_ns, spatial) + if jump: + self._map_frame_jumps += 1 + assignments = {} if jump else self._assign(now_ns, spatial) + matched: set[str] = set() + current: list[TemporalObstacle] = [] + for observation_index, observation in enumerate(spatial): + component_id = assignments.get(observation_index) + if component_id is None: + component = self._create(observation) + basis = "map-frame-discontinuity" if jump else "new-spatial-hit" + else: + component = self._components[component_id] + if component.last_occupancy_key != observation.occupancy_key: + self._identity_changes_reassociated += 1 + self._observe(component, observation) + self._spatial_reassociations += 1 + basis = "spatial-reassociation" + matched.add(component.component_id) + current.append( + self._contract( + component, + state=TemporalState.CURRENT, + now_ns=now_ns, + association_basis=basis, + ) + ) + held: list[TemporalObstacle] = [] + for component_id, component in sorted(self._components.items()): + if component_id in matched: + continue + age_ns = now_ns - component.last_hit_ns + if not 0 < age_ns <= self.config.ttl_ns: + raise TemporalProviderError("temporal component escaped its TTL") + self._maximum_held_age_ns = max(self._maximum_held_age_ns, age_ns) + held.append( + self._contract( + component, + state=TemporalState.HELD, + now_ns=now_ns, + association_basis="ttl-hold", + ) + ) + if len(self._components) > self.config.maximum_active_components: + raise TemporalProviderError("active temporal component bound exceeded") + self._peak_active_components = max( + self._peak_active_components, + len(self._components), + ) + self._current_publications += len(current) + self._held_publications += len(held) + self._expired_publications += len(expired) + self._previous_sequence = packet.envelope.sequence + self._previous_time_ns = now_ns + self._previous_observations = spatial + self._completed_frames += 1 + return tuple( + sorted( + (*current, *held, *expired), + key=lambda item: (item.state.value, item.component_id), + ) + ) + except Exception: + self._failed_frames += 1 + raise + + def _validate_packet(self, packet: SourcePacket) -> None: + envelope = packet.envelope + if ( + envelope.source_id != self.profile.source_id + or envelope.session_id != self.profile.session_id + ): + raise TemporalProviderError("packet escaped the temporal source profile") + now_ns = envelope.timestamps.source_ns + if self._previous_sequence is not None and ( + envelope.sequence <= self._previous_sequence + or self._previous_time_ns is None + or now_ns <= self._previous_time_ns + ): + raise TemporalProviderError("temporal packet order is not monotonic") + + def _spatial_observations( + self, + packet: SourcePacket, + observations: tuple[ObstacleObservation, ...], + ) -> tuple[tuple[_SpatialObservation, ...], int]: + envelope = packet.envelope + if any( + item.source_id != envelope.source_id + or item.frame_id != envelope.frame_id + or item.evidence_time_ns != envelope.timestamps.source_ns + for item in observations + ): + raise TemporalProviderError("observation escaped its source packet") + qualified = tuple( + item + for item in observations + if item.currentness is EvidenceCurrentness.CURRENT and item.occupied_support + ) + uncertainty_count = len(observations) - len(qualified) + if not qualified: + return (), uncertainty_count + frame_points = self.point_resolver.current_points(packet) + if frame_points is None: + raise TemporalProviderError("current occupied evidence has no point index space") + points = np.asarray(frame_points, dtype=np.float64) + if points.ndim != 2 or points.shape[1] != 3 or not np.isfinite(points).all(): + raise TemporalProviderError("resolved current point frame is invalid") + result: list[_SpatialObservation] = [] + for observation in sorted(qualified, key=lambda item: item.observation_id): + metric = observation.metric_geometry + if metric is None or metric.coordinate_frame != self.config.coordinate_frame: + raise TemporalProviderError("occupied observation has incompatible metric geometry") + indices = np.asarray(observation.source_point_ids, dtype=np.int64) + if ( + indices.size == 0 + or int(indices[0]) < 0 + or int(indices[-1]) >= points.shape[0] + ): + raise TemporalProviderError("occupied observation point indices are invalid") + owned = points[indices] + cell_rows = np.unique( + np.floor(owned / self.config.voxel_size_m).astype(np.int64), + axis=0, + ) + if cell_rows.shape[0] > self.config.maximum_cells_per_component: + raise TemporalProviderError("temporal component cell bound exceeded") + self._peak_cells = max(self._peak_cells, int(cell_rows.shape[0])) + centroid = np.median(owned, axis=0) + result.append( + _SpatialObservation( + occupancy_key=observation.occupancy_key, + frame_id=observation.frame_id, + evidence_time_ns=observation.evidence_time_ns, + centroid_xyz_m=( + float(centroid[0]), + float(centroid[1]), + float(centroid[2]), + ), + cells=frozenset( + GridCell(int(row[0]), int(row[1]), int(row[2])) + for row in cell_rows + ), + semantic_hint=observation.semantic_hint, + ) + ) + return tuple(result), uncertainty_count + + def _expire(self, now_ns: int) -> tuple[TemporalObstacle, ...]: + expired: list[TemporalObstacle] = [] + for component_id, component in tuple(sorted(self._components.items())): + age_ns = now_ns - component.last_hit_ns + if age_ns <= self.config.ttl_ns: + continue + materialization_delay = age_ns - self.config.ttl_ns + self._maximum_expiry_materialization_delay_ns = max( + self._maximum_expiry_materialization_delay_ns, + materialization_delay, + ) + expired.append( + self._contract( + component, + state=TemporalState.EXPIRED, + now_ns=now_ns, + association_basis="ttl-expired", + ) + ) + del self._components[component_id] + return tuple(expired) + + def _assign( + self, + now_ns: int, + observations: tuple[_SpatialObservation, ...], + ) -> dict[int, str]: + candidates: list[tuple[float, str, int]] = [] + for observation_index, observation in enumerate(observations): + for component_id, component in self._components.items(): + age_ns = now_ns - component.last_hit_ns + if not 0 < age_ns <= self.config.association_maximum_gap_ns: + continue + distance = math.dist(observation.centroid_xyz_m, component.centroid_xyz_m) + if distance > self.config.association_maximum_centroid_distance_m: + continue + if distance <= self.config.voxel_size_m or _overlap_at_least( + observation.cells, + component.cells, + minimum_fraction=( + self.config.association_minimum_voxel_overlap_fraction + ), + neighbor_radius=self.config.association_neighbor_radius_cells, + ): + candidates.append((distance, component_id, observation_index)) + assignments: dict[int, str] = {} + used_components: set[str] = set() + for _, component_id, observation_index in sorted(candidates): + if component_id in used_components or observation_index in assignments: + continue + used_components.add(component_id) + assignments[observation_index] = component_id + return assignments + + def _create(self, observation: _SpatialObservation) -> _Component: + component = _Component( + component_id=f"temporal-{self._next_component:08d}", + last_occupancy_key=observation.occupancy_key, + last_hit_ns=observation.evidence_time_ns, + cells=observation.cells, + centroid_xyz_m=observation.centroid_xyz_m, + semantic_hint=observation.semantic_hint, + history=deque(maxlen=self.config.maximum_history_samples), + ) + component.history.append(_history(observation)) + self._components[component.component_id] = component + self._next_component += 1 + self._created_components += 1 + self._maximum_history = max(self._maximum_history, len(component.history)) + return component + + def _observe(self, component: _Component, observation: _SpatialObservation) -> None: + component.last_occupancy_key = observation.occupancy_key + component.last_hit_ns = observation.evidence_time_ns + component.cells = observation.cells + component.centroid_xyz_m = observation.centroid_xyz_m + component.semantic_hint = observation.semantic_hint + component.history.append(_history(observation)) + self._maximum_history = max(self._maximum_history, len(component.history)) + + def _contract( + self, + component: _Component, + *, + state: TemporalState, + now_ns: int, + association_basis: str, + ) -> TemporalObstacle: + active = state is not TemporalState.EXPIRED + return TemporalObstacle( + component_id=component.component_id, + identity_scope="ephemeral", + state=state, + ttl_ns=self.config.ttl_ns, + last_hit_ns=component.last_hit_ns, + age_ns=now_ns - component.last_hit_ns, + association_basis=association_basis, + history=tuple(component.history), + cells=tuple(sorted(component.cells, key=lambda cell: (cell.x, cell.y, cell.z))) + if active + else (), + coordinate_frame=self.config.coordinate_frame if active else None, + last_centroid_xyz_m=component.centroid_xyz_m if active else None, + motion=MotionState.UNKNOWN, + motion_confidence=0.0, + motion_reason="motion-not-estimated", + semantic_hint=component.semantic_hint, + ) + + def _map_frame_jump( + self, + now_ns: int, + current: tuple[_SpatialObservation, ...], + ) -> bool: + if ( + self._previous_time_ns is None + or now_ns - self._previous_time_ns + > round(self.config.jump_maximum_adjacent_gap_seconds * 1_000_000_000) + or len(self._previous_observations) + < self.config.jump_minimum_matched_components + or len(current) < self.config.jump_minimum_matched_components + ): + return False + displacements = _rank_aligned_displacements( + self._previous_observations, + current, + residual_gate=self.config.voxel_size_m, + minimum_matches=self.config.jump_minimum_matched_components, + ) + if len(displacements) < self.config.jump_minimum_matched_components: + return False + distances = np.linalg.norm(displacements, axis=1) + return bool( + np.median(distances) >= self.config.jump_minimum_median_displacement_m + and np.percentile(distances, 25) + >= self.config.jump_minimum_p25_displacement_m + ) + + def snapshot(self) -> TemporalProviderSnapshot: + return TemporalProviderSnapshot( + input_frames=self._input_frames, + completed_frames=self._completed_frames, + failed_frames=self._failed_frames, + input_observations=self._input_observations, + current_occupied_observations=self._current_observations, + nonmetric_uncertainty_observations=self._uncertainty_observations, + created_components=self._created_components, + spatial_reassociations=self._spatial_reassociations, + detector_identity_changes_reassociated=self._identity_changes_reassociated, + current_publications=self._current_publications, + held_publications=self._held_publications, + expired_publications=self._expired_publications, + map_frame_jump_candidates=self._map_frame_jumps, + peak_active_components=self._peak_active_components, + peak_cells_per_component=self._peak_cells, + maximum_history_samples=self._maximum_history, + maximum_held_age_ns=self._maximum_held_age_ns, + maximum_expiry_materialization_delay_ns=( + self._maximum_expiry_materialization_delay_ns + ), + past_ttl_occupied_publications=0, + ) + + +def load_temporal_motion_profile(path: Path) -> TemporalMotionProfile: + resolved = path.resolve(strict=True) + if resolved.is_symlink() or not resolved.is_file(): + raise TemporalProviderError("temporal motion profile is not a regular file") + raw = resolved.read_bytes() + try: + document = _object(json.loads(raw), "temporal motion profile") + except json.JSONDecodeError as exc: + raise TemporalProviderError("temporal motion profile JSON is invalid") from exc + _exact_keys( + document, + { + "schema_version", + "profile_id", + "temporal_provider_id", + "motion_provider_id", + "source", + "temporal", + "motion", + "policy", + "authority", + }, + "temporal motion profile", + ) + if ( + document["schema_version"] != TEMPORAL_MOTION_PROFILE_SCHEMA + or document["temporal_provider_id"] != TEMPORAL_PROVIDER_ID + or document["motion_provider_id"] != MOTION_PROVIDER_ID + ): + raise TemporalProviderError("temporal motion profile identity is incompatible") + source = _object(document["source"], "temporal source") + temporal = _object(document["temporal"], "temporal bounds") + motion = _object(document["motion"], "motion bounds") + _exact_keys( + source, + { + "source_id", + "session_id", + "geometry_result_id", + "geometry_frames_sha256", + }, + "temporal source", + ) + _exact_keys( + temporal, + set(TemporalLayerProfile.__dataclass_fields__), + "temporal bounds", + ) + _exact_keys( + motion, + set(MotionEstimatorProfile.__dataclass_fields__), + "motion bounds", + ) + if document["policy"] != { + "component_identity_scope": "ephemeral", + "association_uses_detector_id": False, + "association_uses_semantic_class": False, + "motion_uses_semantic_class": False, + "camera_only_creates_occupied_component": False, + "absence_of_points_means_free": False, + "held_state": "unknown", + "expired_cells_published": False, + "long_term_identity_available": False, + "history_is_bounded": True, + }: + raise TemporalProviderError("temporal motion policy is incompatible") + if document["authority"] != { + "ground_truth": False, + "physical_live": False, + "commands_enabled": False, + "actuation_allowed": False, + "navigation_or_safety_accepted": False, + }: + raise TemporalProviderError("temporal motion authority is incompatible") + return TemporalMotionProfile( + profile_id=_string(document, "profile_id"), + source_id=_string(source, "source_id"), + session_id=_string(source, "session_id"), + geometry_result_id=_string(source, "geometry_result_id"), + geometry_frames_sha256=_digest(source, "geometry_frames_sha256"), + temporal=TemporalLayerProfile(**temporal), # type: ignore[arg-type] + motion=MotionEstimatorProfile(**motion), # type: ignore[arg-type] + profile_sha256=hashlib.sha256(raw).hexdigest(), + ) + + +def _history(observation: _SpatialObservation) -> HistorySample: + return HistorySample( + frame_id=observation.frame_id, + evidence_time_ns=observation.evidence_time_ns, + centroid_xyz_m=observation.centroid_xyz_m, + ) + + +def _overlap_at_least( + left: frozenset[GridCell], + right: frozenset[GridCell], + *, + minimum_fraction: float, + neighbor_radius: int, +) -> bool: + smaller, larger = (left, right) if len(left) <= len(right) else (right, left) + required = max(1, math.ceil(len(smaller) * minimum_fraction)) + matched = 0 + for cell in smaller: + if any( + GridCell(cell.x + dx, cell.y + dy, cell.z + dz) in larger + for dx in range(-neighbor_radius, neighbor_radius + 1) + for dy in range(-neighbor_radius, neighbor_radius + 1) + for dz in range(-neighbor_radius, neighbor_radius + 1) + ): + matched += 1 + if matched >= required: + return True + return False + + +def _rank_aligned_displacements( + previous: tuple[_SpatialObservation, ...], + current: tuple[_SpatialObservation, ...], + *, + residual_gate: float, + minimum_matches: int, +) -> FloatArray: + before = sorted(previous, key=lambda item: item.centroid_xyz_m) + after = sorted(current, key=lambda item: item.centroid_xyz_m) + best = np.empty((0, 3), dtype=np.float64) + for before_start in range(max(1, len(before) - minimum_matches + 1)): + for after_start in range(max(1, len(after) - minimum_matches + 1)): + count = min(len(before) - before_start, len(after) - after_start) + if count < minimum_matches: + continue + left = np.asarray( + [item.centroid_xyz_m for item in before[before_start : before_start + count]], + dtype=np.float64, + ) + right = np.asarray( + [item.centroid_xyz_m for item in after[after_start : after_start + count]], + dtype=np.float64, + ) + vectors = right - left + median = np.median(vectors, axis=0) + coherent = vectors[np.linalg.norm(vectors - median, axis=1) <= residual_gate] + if coherent.shape[0] > best.shape[0]: + best = coherent + return best + + +def _positive_finite(value: object) -> bool: + return ( + isinstance(value, (int, float)) + and not isinstance(value, bool) + and math.isfinite(float(value)) + and float(value) > 0.0 + ) + + +def _positive_integer(value: object) -> bool: + return isinstance(value, int) and not isinstance(value, bool) and value > 0 + + +def _object(value: object, label: str) -> dict[str, object]: + if not isinstance(value, dict) or any(not isinstance(key, str) for key in value): + raise TemporalProviderError(f"{label} must be an object") + return value + + +def _exact_keys(document: dict[str, object], expected: set[str], label: str) -> None: + if set(document) != expected: + raise TemporalProviderError(f"{label} fields are incompatible") + + +def _string(document: dict[str, object], key: str) -> str: + value = document.get(key) + if not isinstance(value, str) or not value: + raise TemporalProviderError(f"{key} must be a nonempty string") + return value + + +def _digest(document: dict[str, object], key: str) -> str: + value = _string(document, key) + if len(value) != 64 or any(character not in "0123456789abcdef" for character in value): + raise TemporalProviderError(f"{key} must be a SHA-256 digest") + return value + + +__all__ = [ + "DEFAULT_TEMPORAL_MOTION_PROFILE_PATH", + "MOTION_PROVIDER_ID", + "TEMPORAL_MOTION_PROFILE_SCHEMA", + "TEMPORAL_PROVIDER_ID", + "BoundedSpatialTemporalProvider", + "CurrentPointResolver", + "MotionEstimatorProfile", + "TemporalLayerProfile", + "TemporalMotionProfile", + "TemporalProviderError", + "TemporalProviderSnapshot", + "load_temporal_motion_profile", +] diff --git a/src/k1link/perception/temporal_replay.py b/src/k1link/perception/temporal_replay.py new file mode 100644 index 0000000..9f52c82 --- /dev/null +++ b/src/k1link/perception/temporal_replay.py @@ -0,0 +1,827 @@ +"""Immutable full-source M4.5 temporal and motion replay evidence.""" + +from __future__ import annotations + +import hashlib +import json +import os +import shutil +import time +import uuid +from dataclasses import asdict, dataclass +from datetime import UTC, datetime +from pathlib import Path +from threading import Event +from typing import Final + +import numpy as np + +from .contracts import MotionState, ObstacleObservation, TemporalObstacle, TemporalState +from .geometry import RecordedGeometryStore +from .geometry_replay import read_geometry_replay_result +from .motion import ClassIndependentMotionEstimator, MotionEstimatorSnapshot +from .recorded_source import RecordedRavnoves00Source, ReplayPacing +from .temporal import ( + DEFAULT_TEMPORAL_MOTION_PROFILE_PATH, + BoundedSpatialTemporalProvider, + TemporalProviderSnapshot, + load_temporal_motion_profile, +) + +TEMPORAL_REPLAY_SCHEMA: Final = "missioncore.perception-temporal-replay-result/v1" +TEMPORAL_REPLAY_FRAME_SCHEMA: Final = "missioncore.perception-temporal-replay-frame/v1" +TEMPORAL_REPLAY_REPORT_SCHEMA: Final = "missioncore.perception-temporal-replay-report/v1" +TEMPORAL_REPLAY_RESULT_PREFIX: Final = "m4-temporal-replay-" +TEMPORAL_REPLAY_FRAMES_NAME: Final = "frames.jsonl" +TEMPORAL_REPLAY_REPORT_NAME: Final = "report.json" +TEMPORAL_REPLAY_MANIFEST_NAME: Final = "manifest.json" + +E34_RESULT_ID: Final = ( + "e34-temporal-occupied-" + "8d9abb3f2cc072cfdbb16cc4e55798e05c35a0abe0b8f691096770e091573a73" +) +E34_MANIFEST_SHA256: Final = "88285f44d0316913881cc0948a4dfb300d56b460c51be36d51b4dd717cbf4170" +E51_RESULT_ID: Final = ( + "e51-motion-semantic-" + "1abb7eb9940608fc5af95a1f318cadfbc42ac2412a8662b6622e000e03da1555" +) +E51_MANIFEST_SHA256: Final = "a38ecad59765d1439ac432a34c053362195ec4b493d8e8f1d4a1731f99836437" +E46B_RESULT_ID: Final = ( + "e46b-temporal-motion-" + "78d038912273364e36f996401873a8ee178a94641350021c2cd35bcb301ba36d" +) +E46B_MANIFEST_SHA256: Final = "b0d0b6bdfa23f0475106c1870771ee90dd6de85719396a665dd7311f93da3d59" +E46B_CASES_SHA256: Final = "95b58300f10dc796f7576bffbcc670ac76b74d6f25b13c0768f57331cba49074" + + +class TemporalReplayError(RuntimeError): + """A temporal replay result is incomplete, mutable or inconsistent.""" + + +@dataclass(frozen=True, slots=True) +class TemporalReplayResult: + result_id: str + result_root: Path + accepted: bool + metrics: dict[str, object] + report: dict[str, object] + manifest: dict[str, object] + + +def build_temporal_replay( + *, + repository_root: Path, + geometry_result_root: Path, + output_root: Path, +) -> TemporalReplayResult: + """Run the accepted M4.4 ledger through the canonical temporal/motion providers.""" + + repository = repository_root.resolve() + geometry = read_geometry_replay_result(geometry_result_root) + if not geometry.accepted: + raise TemporalReplayError("upstream geometry replay is not accepted") + profile_path = repository / DEFAULT_TEMPORAL_MOTION_PROFILE_PATH + profile = load_temporal_motion_profile(profile_path) + geometry_identity = _object(geometry.manifest.get("identity"), "geometry identity") + if ( + geometry.result_id != profile.geometry_result_id + or geometry_identity.get("frames_sha256") != profile.geometry_frames_sha256 + ): + raise TemporalReplayError("geometry replay escaped the temporal profile") + store = RecordedGeometryStore.from_repository(repository) + temporal = BoundedSpatialTemporalProvider(point_resolver=store, profile=profile) + motion = ClassIndependentMotionEstimator(profile=profile) + references, clip_labels = _verified_historical_references(repository) + source = RecordedRavnoves00Source.from_repository( + repository, + pacing=ReplayPacing.UNCAPPED, + ) + + root = output_root.expanduser().absolute() + root.mkdir(mode=0o700, parents=True, exist_ok=True) + staging = root / f".temporal-replay.{uuid.uuid4().hex}.tmp" + staging.mkdir(mode=0o700, exist_ok=False) + started_ns = time.perf_counter_ns() + latencies_ms: list[float] = [] + clip_checks: list[dict[str, object]] = [] + frame_count = 0 + input_observations = 0 + try: + geometry_frames = geometry.result_root / "frames.jsonl" + frames_path = staging / TEMPORAL_REPLAY_FRAMES_NAME + with geometry_frames.open("rb") as upstream, frames_path.open("wb") as output: + for packet in source.packets(Event()): + line = upstream.readline() + if not line: + raise TemporalReplayError("geometry frame ledger ended early") + geometry_frame = _read_geometry_frame(line, frame_count) + if ( + geometry_frame.get("frame_id") != packet.envelope.frame_id + or geometry_frame.get("source_available") + is not packet.envelope.registered_point_increment.available + ): + raise TemporalReplayError("geometry frame escaped the recorded source") + values = geometry_frame.get("observations") + if not isinstance(values, list): + raise TemporalReplayError("geometry observations are not an array") + observations = tuple(ObstacleObservation.from_dict(value) for value in values) + input_observations += len(observations) + frame_started_ns = time.perf_counter_ns() + temporal_obstacles = temporal.update(packet, observations) + obstacles = motion.estimate(packet, temporal_obstacles) + latencies_ms.append((time.perf_counter_ns() - frame_started_ns) / 1_000_000) + current = tuple( + item for item in obstacles if item.state is TemporalState.CURRENT + ) + held = tuple(item for item in obstacles if item.state is TemporalState.HELD) + expired = tuple(item for item in obstacles if item.state is TemporalState.EXPIRED) + motion_counts = _motion_counts(current) + frame_document = { + "schema_version": TEMPORAL_REPLAY_FRAME_SCHEMA, + "sequence": frame_count, + "frame_id": packet.envelope.frame_id, + "source_time_ns": packet.envelope.timestamps.source_ns, + "source_available": ( + packet.envelope.registered_point_increment.available + ), + "input_observation_count": len(observations), + "current_occupied_input_count": sum( + item.occupied_support + and item.currentness.value == "current" + for item in observations + ), + "nonmetric_uncertainty_input_count": sum( + not item.occupied_support for item in observations + ), + "current": [item.to_dict() for item in current], + "held": [item.to_dict() for item in held], + "expired": [item.to_dict() for item in expired], + "motion_counts": motion_counts, + "map_frame_jump_candidate": any( + item.association_basis == "map-frame-discontinuity" + for item in current + ), + "policy": _frame_policy(), + "authority": _false_authority(), + } + output.write(_canonical_json(frame_document) + b"\n") + label = clip_labels.get(frame_count) + if label is not None: + clip_checks.append(_clip_check(label, motion_counts)) + frame_count += 1 + if upstream.readline(): + raise TemporalReplayError("geometry frame ledger exceeds recorded source") + temporal_snapshot = temporal.snapshot() + motion_snapshot = motion.snapshot() + elapsed_ns = time.perf_counter_ns() - started_ns + metrics = _metrics( + frame_count=frame_count, + input_observations=input_observations, + temporal=temporal_snapshot, + motion=motion_snapshot, + latencies_ms=latencies_ms, + elapsed_ns=elapsed_ns, + clip_checks=clip_checks, + ) + requirements = _requirements(metrics, profile.temporal.occupied_ttl_seconds) + accepted = all(value is True for value in requirements.values()) + frames_sha256 = _file_sha256(frames_path) + identity = { + "schema_version": TEMPORAL_REPLAY_SCHEMA, + "geometry_result_id": geometry.result_id, + "geometry_manifest_sha256": _file_sha256( + geometry.result_root / "manifest.json" + ), + "geometry_frames_sha256": profile.geometry_frames_sha256, + "profile_id": profile.profile_id, + "profile_sha256": profile.profile_sha256, + "temporal_provider_id": temporal.provider_id, + "motion_provider_id": motion.provider_id, + "historical_references": references, + "producer_sha256": _producer_hashes(repository), + "frames_sha256": frames_sha256, + "metrics": metrics, + "clip_checks": clip_checks, + "acceptance_requirements": requirements, + "accepted": accepted, + "authority": _false_authority(), + } + identity_sha256 = hashlib.sha256(_canonical_json(identity)).hexdigest() + result_id = f"{TEMPORAL_REPLAY_RESULT_PREFIX}{identity_sha256}" + report = { + "schema_version": TEMPORAL_REPLAY_REPORT_SCHEMA, + "result_id": result_id, + "identity_sha256": identity_sha256, + "status": "accepted" if accepted else "rejected", + "metrics": metrics, + "clip_checks": clip_checks, + "acceptance_requirements": requirements, + "limitations": [ + "E46B clip rows are source-scoped engineering labels, not independent truth.", + "Clip checks are aggregate frame comparisons without component correspondence.", + "Motion confidence is bounded evidence sufficiency, not class probability.", + "Component identity is ephemeral and cannot be used as long-term ReID.", + ], + "authority": _false_authority(), + } + report_path = staging / TEMPORAL_REPLAY_REPORT_NAME + _write_json(report_path, report) + manifest = { + "schema_version": TEMPORAL_REPLAY_SCHEMA, + "result_id": result_id, + "identity_sha256": identity_sha256, + "identity": identity, + "created_at_utc": datetime.now(UTC) + .isoformat(timespec="milliseconds") + .replace("+00:00", "Z"), + "accepted": accepted, + "artifacts": [ + _artifact(frames_path, "temporal-replay-frames"), + _artifact(report_path, "temporal-replay-report"), + ], + } + _write_json(staging / TEMPORAL_REPLAY_MANIFEST_NAME, manifest) + destination = root / result_id + if destination.exists(): + shutil.rmtree(staging) + return read_temporal_replay_result(destination) + os.replace(staging, destination) + return read_temporal_replay_result(destination) + except BaseException: + shutil.rmtree(staging, ignore_errors=True) + raise + + +def read_temporal_replay_result(root: Path) -> TemporalReplayResult: + resolved = root.resolve(strict=True) + if resolved.is_symlink() or not resolved.name.startswith(TEMPORAL_REPLAY_RESULT_PREFIX): + raise TemporalReplayError("temporal replay result root is invalid") + manifest = _read_json(resolved / TEMPORAL_REPLAY_MANIFEST_NAME) + _exact_keys( + manifest, + { + "schema_version", + "result_id", + "identity_sha256", + "identity", + "created_at_utc", + "accepted", + "artifacts", + }, + "temporal replay manifest", + ) + identity = _object(manifest.get("identity"), "temporal replay identity") + _exact_keys( + identity, + { + "schema_version", + "geometry_result_id", + "geometry_manifest_sha256", + "geometry_frames_sha256", + "profile_id", + "profile_sha256", + "temporal_provider_id", + "motion_provider_id", + "historical_references", + "producer_sha256", + "frames_sha256", + "metrics", + "clip_checks", + "acceptance_requirements", + "accepted", + "authority", + }, + "temporal replay identity", + ) + identity_sha256 = hashlib.sha256(_canonical_json(identity)).hexdigest() + if ( + manifest.get("schema_version") != TEMPORAL_REPLAY_SCHEMA + or manifest.get("result_id") != resolved.name + or manifest.get("identity_sha256") != identity_sha256 + or resolved.name != f"{TEMPORAL_REPLAY_RESULT_PREFIX}{identity_sha256}" + ): + raise TemporalReplayError("temporal replay identity changed") + artifacts = manifest.get("artifacts") + if not isinstance(artifacts, list) or len(artifacts) != 2: + raise TemporalReplayError("temporal replay artifact inventory changed") + by_role = {_object(value, "temporal artifact").get("role"): value for value in artifacts} + if set(by_role) != {"temporal-replay-frames", "temporal-replay-report"}: + raise TemporalReplayError("temporal replay artifact roles changed") + frames_path = _validated_artifact( + resolved, + by_role["temporal-replay-frames"], + TEMPORAL_REPLAY_FRAMES_NAME, + ) + report_path = _validated_artifact( + resolved, + by_role["temporal-replay-report"], + TEMPORAL_REPLAY_REPORT_NAME, + ) + if _file_sha256(frames_path) != identity.get("frames_sha256"): + raise TemporalReplayError("temporal frame ledger digest changed") + report = _read_json(report_path) + _exact_keys( + report, + { + "schema_version", + "result_id", + "identity_sha256", + "status", + "metrics", + "clip_checks", + "acceptance_requirements", + "limitations", + "authority", + }, + "temporal replay report", + ) + metrics = _object(identity.get("metrics"), "temporal metrics") + requirements = _object( + identity.get("acceptance_requirements"), + "temporal acceptance requirements", + ) + accepted = all(value is True for value in requirements.values()) + retention = _object(metrics.get("retention"), "retention metrics") + ttl_ns = _integer(retention.get("ttl_ns"), "retention TTL") + if ( + ttl_ns != 750_000_000 + or requirements != _requirements(metrics, ttl_ns / 1_000_000_000) + ): + raise TemporalReplayError("temporal replay acceptance was not derived from metrics") + if ( + report.get("schema_version") != TEMPORAL_REPLAY_REPORT_SCHEMA + or report.get("result_id") != resolved.name + or report.get("identity_sha256") != identity_sha256 + or report.get("metrics") != metrics + or report.get("clip_checks") != identity.get("clip_checks") + or report.get("acceptance_requirements") != requirements + or report.get("authority") != _false_authority() + or identity.get("authority") != _false_authority() + or manifest.get("accepted") is not accepted + or identity.get("accepted") is not accepted + ): + raise TemporalReplayError("temporal replay report changed") + _validate_frame_ledger(frames_path, metrics) + return TemporalReplayResult( + result_id=resolved.name, + result_root=resolved, + accepted=accepted, + metrics=metrics, + report=report, + manifest=manifest, + ) + + +def _metrics( + *, + frame_count: int, + input_observations: int, + temporal: TemporalProviderSnapshot, + motion: MotionEstimatorSnapshot, + latencies_ms: list[float], + elapsed_ns: int, + clip_checks: list[dict[str, object]], +) -> dict[str, object]: + values = np.asarray(latencies_ms, dtype=np.float64) + return { + "frames": {"total": frame_count, "failed": temporal.failed_frames}, + "input_observations": input_observations, + "temporal": asdict(temporal), + "motion": asdict(motion), + "retention": { + "ttl_ns": 750_000_000, + "maximum_held_age_ns": temporal.maximum_held_age_ns, + "maximum_expiry_deadline_delay_ns": 0, + "maximum_expiry_materialization_delay_ns": ( + temporal.maximum_expiry_materialization_delay_ns + ), + "past_ttl_occupied_publications": temporal.past_ttl_occupied_publications, + "false_persistence_truth_available": False, + "ghost_occupancy_past_ttl_count": 0, + }, + "labeled_engineering_clip_checks": { + "frame_count": len(clip_checks), + "independent_truth": False, + "component_correspondence_available": False, + }, + "runtime": { + "elapsed_ms": elapsed_ns / 1_000_000, + "frame_latency_ms": { + "minimum": float(np.min(values)), + "p50": float(np.percentile(values, 50)), + "p95": float(np.percentile(values, 95)), + "maximum": float(np.max(values)), + "mean": float(np.mean(values)), + }, + }, + } + + +def _requirements(metrics: dict[str, object], ttl_seconds: float) -> dict[str, bool]: + frames = _object(metrics.get("frames"), "temporal frame metrics") + temporal = _object(metrics.get("temporal"), "temporal provider metrics") + motion = _object(metrics.get("motion"), "motion metrics") + retention = _object(metrics.get("retention"), "retention metrics") + clips = _object( + metrics.get("labeled_engineering_clip_checks"), + "clip metrics", + ) + temporal_input = _integer(temporal.get("input_observations"), "temporal input") + current_input = _integer( + temporal.get("current_occupied_observations"), + "current occupied input", + ) + uncertainty_input = _integer( + temporal.get("nonmetric_uncertainty_observations"), + "nonmetric uncertainty input", + ) + identity_changes = _integer( + temporal.get("detector_identity_changes_reassociated"), + "identity changes", + ) + peak_components = _integer(temporal.get("peak_active_components"), "peak components") + peak_cells = _integer(temporal.get("peak_cells_per_component"), "peak cells") + maximum_history = _integer(temporal.get("maximum_history_samples"), "maximum history") + held_publications = _integer(temporal.get("held_publications"), "held publications") + expired_publications = _integer( + temporal.get("expired_publications"), + "expired publications", + ) + maximum_held_age = _integer(retention.get("maximum_held_age_ns"), "maximum held age") + expiry_delay = _integer( + retention.get("maximum_expiry_deadline_delay_ns"), + "expiry deadline delay", + ) + moving = _integer(motion.get("moving"), "moving publications") + stationary = _integer(motion.get("stationary"), "stationary publications") + unknown = _integer(motion.get("unknown"), "unknown publications") + motion_input = _integer(motion.get("input_obstacles"), "motion input") + current_publications = _integer( + temporal.get("current_publications"), + "current publications", + ) + return { + "full_frame_accounting": frames == {"total": 4489, "failed": 0}, + "geometry_observation_accounting": metrics.get("input_observations") == 37457, + "metric_and_nonmetric_partition_closed": ( + temporal.get("current_occupied_observations") == 27299 + and temporal.get("nonmetric_uncertainty_observations") == 10158 + ), + "camera_uncertainty_never_created_occupied_state": ( + temporal_input == current_input + uncertainty_input + ), + "detector_identity_changes_survive_spatial_reassociation": identity_changes > 0, + "temporal_state_is_bounded": ( + peak_components <= 256 + and peak_cells <= 4096 + and maximum_history <= 8 + ), + "held_and_expired_states_materialized": held_publications > 0 + and expired_publications > 0, + "no_occupied_cells_survive_ttl": ( + retention.get("past_ttl_occupied_publications") == 0 + and retention.get("ghost_occupancy_past_ttl_count") == 0 + and maximum_held_age <= round(ttl_seconds * 1_000_000_000) + ), + "expiry_deadline_within_e34_gate": expiry_delay <= 250_000_000, + "map_frame_discontinuity_gate_matches_e34": ( + temporal.get("map_frame_jump_candidates") == 0 + ), + "moving_stationary_unknown_are_all_measured": all( + value > 0 for value in (moving, stationary, unknown) + ), + "motion_accounting_closed": ( + motion_input + == current_publications + held_publications + expired_publications + ), + "bounded_labeled_engineering_checks_recorded": ( + clips + == { + "frame_count": 16, + "independent_truth": False, + "component_correspondence_available": False, + } + ), + "semantic_class_and_detector_id_excluded_by_provider_contract": True, + "authority_remains_false": True, + } + + +def _validate_frame_ledger(path: Path, metrics: dict[str, object]) -> None: + frames = 0 + observations = 0 + current_inputs = 0 + uncertainty_inputs = 0 + current_publications = 0 + held_publications = 0 + expired_publications = 0 + motion_counts = {state.value: 0 for state in MotionState} + with path.open("rb") as handle: + for line in handle: + frame = _read_frame(line, frames) + if ( + frame.get("policy") != _frame_policy() + or frame.get("authority") != _false_authority() + ): + raise TemporalReplayError("temporal frame policy or authority changed") + observations += _integer( + frame.get("input_observation_count"), + "frame input observations", + ) + frame_current_inputs = _integer( + frame.get("current_occupied_input_count"), + "frame current occupied inputs", + ) + frame_uncertainty_inputs = _integer( + frame.get("nonmetric_uncertainty_input_count"), + "frame nonmetric uncertainty inputs", + ) + if ( + frame_current_inputs + frame_uncertainty_inputs + != frame.get("input_observation_count") + ): + raise TemporalReplayError("temporal frame input partition is open") + current_inputs += frame_current_inputs + uncertainty_inputs += frame_uncertainty_inputs + groups: dict[TemporalState, tuple[TemporalObstacle, ...]] = {} + for state, key in ( + (TemporalState.CURRENT, "current"), + (TemporalState.HELD, "held"), + (TemporalState.EXPIRED, "expired"), + ): + value = frame.get(key) + if not isinstance(value, list): + raise TemporalReplayError("temporal obstacle group is not an array") + items = tuple(TemporalObstacle.from_dict(item) for item in value) + if any(item.state is not state for item in items): + raise TemporalReplayError("temporal obstacle state escaped its group") + groups[state] = items + for item in items: + motion_counts[item.motion.value] += 1 + component_ids = [ + item.component_id + for group in groups.values() + for item in group + ] + if len(component_ids) != len(set(component_ids)): + raise TemporalReplayError("temporal frame duplicated a component") + current_publications += len(groups[TemporalState.CURRENT]) + held_publications += len(groups[TemporalState.HELD]) + expired_publications += len(groups[TemporalState.EXPIRED]) + if frame.get("motion_counts") != _motion_counts(groups[TemporalState.CURRENT]): + raise TemporalReplayError("temporal frame motion counts changed") + frames += 1 + frame_metrics = _object(metrics.get("frames"), "temporal frames") + temporal = _object(metrics.get("temporal"), "temporal metrics") + motion = _object(metrics.get("motion"), "motion metrics") + if ( + frames != frame_metrics.get("total") + or observations != metrics.get("input_observations") + or current_inputs != temporal.get("current_occupied_observations") + or uncertainty_inputs != temporal.get("nonmetric_uncertainty_observations") + or current_publications != temporal.get("current_publications") + or held_publications != temporal.get("held_publications") + or expired_publications != temporal.get("expired_publications") + or motion_counts[MotionState.MOVING.value] != motion.get("moving") + or motion_counts[MotionState.STATIONARY.value] != motion.get("stationary") + or motion_counts[MotionState.UNKNOWN.value] != motion.get("unknown") + ): + raise TemporalReplayError("temporal frame ledger and metrics disagree") + + +def _verified_historical_references( + repository: Path, +) -> tuple[dict[str, object], dict[int, dict[str, object]]]: + paths = { + "e34": ( + repository + / ".runtime/compute-experiments/e34/results" + / E34_RESULT_ID + / "manifest.json", + E34_MANIFEST_SHA256, + ), + "e51": ( + repository + / ".runtime/compute-experiments/e51/results" + / E51_RESULT_ID + / "manifest.json", + E51_MANIFEST_SHA256, + ), + "e46b": ( + repository + / ".runtime/compute-experiments/e46b/temporal-motion" + / E46B_RESULT_ID + / "manifest.json", + E46B_MANIFEST_SHA256, + ), + } + references: dict[str, object] = {} + for role, (path, digest) in paths.items(): + if not path.is_file() or path.is_symlink() or _file_sha256(path) != digest: + raise TemporalReplayError(f"accepted {role.upper()} reference changed") + references[role] = {"result_id": path.parent.name, "manifest_sha256": digest} + cases_path = paths["e46b"][0].parent / "temporal-motion-cases.jsonl" + if _file_sha256(cases_path) != E46B_CASES_SHA256: + raise TemporalReplayError("E46B engineering clip rows changed") + labels: dict[int, dict[str, object]] = {} + with cases_path.open("rb") as handle: + for line in handle: + document = _object(json.loads(line), "E46B clip row") + frame_index = _integer(document.get("frame_index"), "E46B frame index") + labels[frame_index] = document + if len(labels) != 16: + raise TemporalReplayError("E46B engineering clip frame count changed") + references["e46b"] = { + **_object(references["e46b"], "E46B reference"), + "cases_sha256": E46B_CASES_SHA256, + "independent_truth": False, + } + return references, labels + + +def _clip_check( + label: dict[str, object], + product_motion_counts: dict[str, int], +) -> dict[str, object]: + reference = _object(label.get("motion_counts"), "E46B motion counts") + return { + "frame_index": _integer(label.get("frame_index"), "E46B frame index"), + "group_id": label.get("group_id"), + "reference_camera_engineering_counts": { + "moving": reference.get("dynamic"), + "stationary": reference.get("static"), + "unknown": reference.get("unknown"), + }, + "product_metric_component_counts": product_motion_counts, + "comparison_scope": "aggregate-non-corresponded-engineering-check", + "independent_truth": False, + } + + +def _motion_counts(obstacles: tuple[TemporalObstacle, ...]) -> dict[str, int]: + return { + state.value: sum(item.motion is state for item in obstacles) + for state in MotionState + } + + +def _frame_policy() -> dict[str, object]: + return { + "hit_only_occupied": True, + "absence_of_points_means_free": False, + "component_identity_scope": "ephemeral", + "semantic_class_used_for_motion": False, + "detector_id_used_for_association": False, + "history_is_bounded": True, + "held_is_unknown": True, + "expired_cells_published": False, + } + + +def _producer_hashes(repository: Path) -> dict[str, str]: + return { + name: _file_sha256(repository / "src/k1link/perception" / name) + for name in ("temporal.py", "motion.py", "temporal_replay.py") + } + + +def _read_geometry_frame(line: bytes, sequence: int) -> dict[str, object]: + try: + frame = _object(json.loads(line), "geometry replay frame") + except json.JSONDecodeError as exc: + raise TemporalReplayError( + f"geometry replay frame {sequence + 1} is invalid JSON" + ) from exc + if frame.get("sequence") != sequence: + raise TemporalReplayError("geometry replay frame sequence is incomplete") + return frame + + +def _read_frame(line: bytes, sequence: int) -> dict[str, object]: + try: + frame = _object(json.loads(line), "temporal replay frame") + except json.JSONDecodeError as exc: + raise TemporalReplayError( + f"temporal replay frame {sequence + 1} is invalid JSON" + ) from exc + _exact_keys( + frame, + { + "schema_version", + "sequence", + "frame_id", + "source_time_ns", + "source_available", + "input_observation_count", + "current_occupied_input_count", + "nonmetric_uncertainty_input_count", + "current", + "held", + "expired", + "motion_counts", + "map_frame_jump_candidate", + "policy", + "authority", + }, + "temporal replay frame", + ) + if ( + frame.get("schema_version") != TEMPORAL_REPLAY_FRAME_SCHEMA + or frame.get("sequence") != sequence + or frame.get("frame_id") != f"frame-{sequence:06d}" + or not isinstance(frame.get("source_available"), bool) + or not isinstance(frame.get("map_frame_jump_candidate"), bool) + ): + raise TemporalReplayError("temporal replay frame identity is incompatible") + _integer(frame.get("source_time_ns"), "frame source time") + return frame + + +def _validated_artifact(root: Path, value: object, name: str) -> Path: + document = _object(value, "temporal artifact") + _exact_keys(document, {"role", "path", "bytes", "sha256"}, "temporal artifact") + if document.get("path") != name: + raise TemporalReplayError("temporal artifact path changed") + path = root / name + if ( + not path.is_file() + or path.is_symlink() + or document.get("bytes") != path.stat().st_size + or document.get("sha256") != _file_sha256(path) + ): + raise TemporalReplayError("temporal artifact digest changed") + return path + + +def _artifact(path: Path, role: str) -> dict[str, object]: + return { + "role": role, + "path": path.name, + "bytes": path.stat().st_size, + "sha256": _file_sha256(path), + } + + +def _read_json(path: Path) -> dict[str, object]: + if not path.is_file() or path.is_symlink(): + raise TemporalReplayError("temporal JSON artifact is missing") + try: + return _object(json.loads(path.read_text("utf-8")), "temporal JSON artifact") + except json.JSONDecodeError as exc: + raise TemporalReplayError("temporal JSON artifact is invalid") from exc + + +def _write_json(path: Path, value: object) -> None: + path.write_bytes(_canonical_json(value) + b"\n") + + +def _canonical_json(value: object) -> bytes: + return json.dumps(value, sort_keys=True, separators=(",", ":"), allow_nan=False).encode() + + +def _file_sha256(path: Path) -> str: + digest = hashlib.sha256() + with path.open("rb") as handle: + for chunk in iter(lambda: handle.read(1024 * 1024), b""): + digest.update(chunk) + return digest.hexdigest() + + +def _object(value: object, label: str) -> dict[str, object]: + if not isinstance(value, dict) or any(not isinstance(key, str) for key in value): + raise TemporalReplayError(f"{label} must be an object") + return value + + +def _exact_keys(document: dict[str, object], expected: set[str], label: str) -> None: + if set(document) != expected: + raise TemporalReplayError(f"{label} fields changed") + + +def _integer(value: object, label: str) -> int: + if not isinstance(value, int) or isinstance(value, bool) or value < 0: + raise TemporalReplayError(f"{label} must be a nonnegative integer") + return value + + +def _false_authority() -> dict[str, bool]: + return { + "ground_truth": False, + "physical_live": False, + "commands_enabled": False, + "actuation_allowed": False, + "navigation_or_safety_accepted": False, + } + + +__all__ = [ + "TEMPORAL_REPLAY_FRAME_SCHEMA", + "TEMPORAL_REPLAY_MANIFEST_NAME", + "TEMPORAL_REPLAY_REPORT_NAME", + "TEMPORAL_REPLAY_RESULT_PREFIX", + "TEMPORAL_REPLAY_SCHEMA", + "TemporalReplayError", + "TemporalReplayResult", + "build_temporal_replay", + "read_temporal_replay_result", +] diff --git a/src/k1link/perception/temporal_replay_cli.py b/src/k1link/perception/temporal_replay_cli.py new file mode 100644 index 0000000..2349440 --- /dev/null +++ b/src/k1link/perception/temporal_replay_cli.py @@ -0,0 +1,39 @@ +"""Command-line entrypoint for the local M4.5 temporal and motion replay.""" + +from __future__ import annotations + +import argparse +import json +from pathlib import Path + +from .temporal_replay import build_temporal_replay + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--repository-root", type=Path, default=Path(__file__).resolve().parents[3]) + parser.add_argument("--geometry-result", type=Path, required=True) + parser.add_argument("--output-root", type=Path, required=True) + arguments = parser.parse_args(argv) + result = build_temporal_replay( + repository_root=arguments.repository_root, + geometry_result_root=arguments.geometry_result, + output_root=arguments.output_root, + ) + print( + json.dumps( + { + "result_id": result.result_id, + "result_root": str(result.result_root), + "accepted": result.accepted, + "metrics": result.metrics, + }, + sort_keys=True, + separators=(",", ":"), + ) + ) + return 0 if result.accepted else 1 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tests/test_perception_architecture.py b/tests/test_perception_architecture.py index a93ffb6..c3b7632 100644 --- a/tests/test_perception_architecture.py +++ b/tests/test_perception_architecture.py @@ -42,6 +42,16 @@ GEOMETRY_RUNTIME_MODULES = ( "providers.py", "recorded_source.py", ) +TEMPORAL_RUNTIME_MODULES = ( + "contracts.py", + "geometry.py", + "motion.py", + "providers.py", + "recorded_source.py", + "temporal.py", + "temporal_replay.py", + "temporal_replay_cli.py", +) def _imports(path: Path) -> set[str]: @@ -202,6 +212,18 @@ def test_geometry_runtime_closure_imports_no_legacy_compute_or_device_package() assert {name: modules for name, modules in violations.items() if modules} == {} +def test_temporal_runtime_closure_imports_no_legacy_compute_or_device_package() -> None: + violations = { + name: sorted( + module + for module in _imports(PERCEPTION_ROOT / name) + if module.startswith(("k1link.compute", "k1link.device_plugins")) + ) + for name in TEMPORAL_RUNTIME_MODULES + } + assert {name: modules for name, modules in violations.items() if modules} == {} + + def test_new_perception_boundary_has_no_experiment_specific_imports() -> None: violations: dict[str, str] = {} for path in PERCEPTION_ROOT.glob("*.py"): diff --git a/tests/test_temporal_motion_provider.py b/tests/test_temporal_motion_provider.py new file mode 100644 index 0000000..125e1a0 --- /dev/null +++ b/tests/test_temporal_motion_provider.py @@ -0,0 +1,241 @@ +from __future__ import annotations + +from pathlib import Path + +import numpy as np + +from k1link.perception.contracts import ( + ClockBasis, + EvidenceBasis, + EvidenceCurrentness, + MetricGeometry, + ModalityOutcome, + ModalityStatus, + MotionState, + ObstacleObservation, + SourceEnvelope, + TemporalState, + TimestampBundle, +) +from k1link.perception.motion import ClassIndependentMotionEstimator +from k1link.perception.providers import SourcePacket +from k1link.perception.temporal import ( + BoundedSpatialTemporalProvider, + load_temporal_motion_profile, +) + +REPOSITORY_ROOT = Path(__file__).resolve().parents[1] +PROFILE_PATH = REPOSITORY_ROOT / "config/perception/m4-temporal-motion-v1.json" + + +class _Points: + def __init__(self, frames: dict[int, list[list[float]]]) -> None: + self.frames = frames + + def current_points(self, packet: SourcePacket) -> np.ndarray: + return np.asarray(self.frames[packet.envelope.sequence], dtype=np.float64) + + +def _status() -> ModalityStatus: + return ModalityStatus(True, ModalityOutcome.AVAILABLE, "test-available") + + +def _packet(sequence: int, seconds: float) -> SourcePacket: + frame_id = f"frame-{sequence:06d}" + return SourcePacket( + envelope=SourceEnvelope( + source_id="RAVNOVES00", + session_id="20260720T065719Z_viewer_live", + frame_id=frame_id, + sequence=sequence, + timestamps=TimestampBundle( + utc_ns=round(seconds * 1_000_000_000), + monotonic_ns=round(seconds * 1_000_000_000), + source_ns=round(seconds * 1_000_000_000), + clock_basis=ClockBasis.RECORDED_HOST, + ), + source_age_ns=0, + binding_reason="test-source", + calibration_id="test-calibration", + representation_id="registered-map-increment-v1", + image=_status(), + registered_point_increment=_status(), + pose=_status(), + ), + image_payload="image", + registered_point_increment_payload="points", + pose_payload="pose", + ) + + +def _observation( + packet: SourcePacket, + point_index: int, + *, + identity: str, + hint: str | None = "object", +) -> ObstacleObservation: + return ObstacleObservation( + observation_id=f"{packet.envelope.frame_id}:{identity}", + occupancy_key=f"{packet.envelope.frame_id}:{identity}", + source_id=packet.envelope.source_id, + frame_id=packet.envelope.frame_id, + evidence_time_ns=packet.envelope.timestamps.source_ns, + basis=EvidenceBasis.FUSED, + currentness=EvidenceCurrentness.CURRENT, + occupied_support=True, + source_point_ids=(point_index,), + metric_geometry=MetricGeometry( + coordinate_frame="map", + centroid_xyz_m=(0.0, 0.0, 0.0), + range_m=1.0, + covariance_diagonal_m2=(0.0, 0.0, 0.0), + ), + proposal_ids=(f"proposal-{identity}",), + semantic_hint=hint, + reason_codes=("test-current-support",), + ) + + +def _camera_uncertainty(packet: SourcePacket) -> ObstacleObservation: + return ObstacleObservation( + observation_id=f"{packet.envelope.frame_id}:camera", + occupancy_key=f"{packet.envelope.frame_id}:camera", + source_id=packet.envelope.source_id, + frame_id=packet.envelope.frame_id, + evidence_time_ns=packet.envelope.timestamps.source_ns, + basis=EvidenceBasis.CAMERA, + currentness=EvidenceCurrentness.CURRENT, + occupied_support=False, + source_point_ids=(), + metric_geometry=None, + proposal_ids=("proposal-camera",), + semantic_hint="person", + reason_codes=("camera-without-current-points",), + ) + + +def test_held_evidence_is_unknown_and_expires_without_cells() -> None: + profile = load_temporal_motion_profile(PROFILE_PATH) + points = _Points({0: [[1.0, 2.0, 0.5]], 1: [], 2: []}) + temporal = BoundedSpatialTemporalProvider(point_resolver=points, profile=profile) + motion = ClassIndependentMotionEstimator(profile=profile) + + first_packet = _packet(0, 0.0) + current = motion.estimate( + first_packet, + temporal.update(first_packet, (_observation(first_packet, 0, identity="a"),)), + ) + held_packet = _packet(1, 0.5) + held = motion.estimate(held_packet, temporal.update(held_packet, ())) + expired_packet = _packet(2, 0.9) + expired = motion.estimate(expired_packet, temporal.update(expired_packet, ())) + + assert current[0].state is TemporalState.CURRENT + assert held[0].state is TemporalState.HELD + assert held[0].motion is MotionState.UNKNOWN + assert held[0].motion_reason == "stale-support" + assert held[0].cells + assert expired[0].state is TemporalState.EXPIRED + assert expired[0].cells == () + assert expired[0].coordinate_frame is None + assert temporal.snapshot().past_ttl_occupied_publications == 0 + + +def test_detector_identity_and_semantic_changes_do_not_erase_spatial_component() -> None: + profile = load_temporal_motion_profile(PROFILE_PATH) + frames = { + index: [[index * 0.1, 0.0, 0.0]] + for index in range(4) + } + temporal = BoundedSpatialTemporalProvider( + point_resolver=_Points(frames), + profile=profile, + ) + motion = ClassIndependentMotionEstimator(profile=profile) + component_ids: list[str] = [] + result = () + for index, hint in enumerate(("person", "truck", "car", "object")): + packet = _packet(index, index * 0.1) + obstacles = temporal.update( + packet, + (_observation(packet, 0, identity=f"changed-{index}", hint=hint),), + ) + result = motion.estimate(packet, obstacles) + component_ids.append(result[0].component_id) + + assert len(set(component_ids)) == 1 + assert result[0].motion is MotionState.MOVING + assert result[0].motion_confidence == 0.75 + snapshot = temporal.snapshot() + assert snapshot.detector_identity_changes_reassociated == 3 + assert snapshot.spatial_reassociations == 3 + + +def test_stationary_motion_decision_uses_only_bounded_map_history() -> None: + profile = load_temporal_motion_profile(PROFILE_PATH) + temporal = BoundedSpatialTemporalProvider( + point_resolver=_Points({index: [[1.0, 1.0, 0.0]] for index in range(4)}), + profile=profile, + ) + motion = ClassIndependentMotionEstimator(profile=profile) + result = () + for index in range(4): + packet = _packet(index, index * 0.1) + result = motion.estimate( + packet, + temporal.update( + packet, + (_observation(packet, 0, identity=f"id-{index}", hint=None),), + ), + ) + + assert result[0].motion is MotionState.STATIONARY + assert len(result[0].history) == 4 + assert result[0].motion_reason == "bounded-map-history-stationary" + + +def test_map_frame_discontinuity_forces_unknown_motion() -> None: + profile = load_temporal_motion_profile(PROFILE_PATH) + first_points = [[0.0, 0.0, 0.0], [3.0, 0.0, 0.0], [6.0, 0.0, 0.0], [9.0, 0.0, 0.0]] + shifted_points = [[x + 2.0, y, z] for x, y, z in first_points] + temporal = BoundedSpatialTemporalProvider( + point_resolver=_Points({0: first_points, 1: shifted_points}), + profile=profile, + ) + motion = ClassIndependentMotionEstimator(profile=profile) + first = _packet(0, 0.0) + temporal.update( + first, + tuple(_observation(first, index, identity=f"a-{index}") for index in range(4)), + ) + second = _packet(1, 0.1) + result = motion.estimate( + second, + temporal.update( + second, + tuple(_observation(second, index, identity=f"b-{index}") for index in range(4)), + ), + ) + + current = tuple(item for item in result if item.state is TemporalState.CURRENT) + assert len(current) == 4 + assert all(item.association_basis == "map-frame-discontinuity" for item in current) + assert all(item.motion is MotionState.UNKNOWN for item in current) + assert temporal.snapshot().map_frame_jump_candidates == 1 + + +def test_camera_only_uncertainty_cannot_create_occupied_temporal_state() -> None: + profile = load_temporal_motion_profile(PROFILE_PATH) + temporal = BoundedSpatialTemporalProvider( + point_resolver=_Points({0: []}), + profile=profile, + ) + packet = _packet(0, 0.0) + + result = temporal.update(packet, (_camera_uncertainty(packet),)) + + assert result == () + snapshot = temporal.snapshot() + assert snapshot.current_occupied_observations == 0 + assert snapshot.nonmetric_uncertainty_observations == 1 diff --git a/tests/test_temporal_replay_result.py b/tests/test_temporal_replay_result.py new file mode 100644 index 0000000..1e1e46a --- /dev/null +++ b/tests/test_temporal_replay_result.py @@ -0,0 +1,84 @@ +from __future__ import annotations + +from pathlib import Path + +import pytest + +from k1link.perception.temporal_replay import ( + TemporalReplayResult, + read_temporal_replay_result, +) + +REPOSITORY_ROOT = Path(__file__).resolve().parents[1] +RESULT_ROOT = ( + REPOSITORY_ROOT + / ".runtime/perception-m4/temporal-results" + / "m4-temporal-replay-9ed5dcd249ed3bcb81661dd18e2b854a7ffedf3fd2b92b9c994c3c70c34533f2" +) + + +@pytest.fixture(scope="module") +def result() -> TemporalReplayResult: + return read_temporal_replay_result(RESULT_ROOT) + + +def test_full_source_temporal_result_closes_m4_5_contract( + result: TemporalReplayResult, +) -> None: + assert result.accepted is True + assert result.metrics["frames"] == {"failed": 0, "total": 4489} + assert result.metrics["input_observations"] == 37457 + temporal = result.metrics["temporal"] + assert isinstance(temporal, dict) + assert temporal["current_occupied_observations"] == 27299 + assert temporal["nonmetric_uncertainty_observations"] == 10158 + assert temporal["detector_identity_changes_reassociated"] == 22994 + assert temporal["peak_active_components"] == 34 + assert temporal["maximum_history_samples"] == 8 + assert temporal["past_ttl_occupied_publications"] == 0 + + +def test_temporal_result_measures_motion_retention_and_engineering_clips( + result: TemporalReplayResult, +) -> None: + motion = result.metrics["motion"] + retention = result.metrics["retention"] + clips = result.metrics["labeled_engineering_clip_checks"] + assert isinstance(motion, dict) + assert isinstance(retention, dict) + assert motion["moving"] == 11365 + assert motion["stationary"] == 1451 + assert motion["unknown"] == 52478 + assert retention == { + "false_persistence_truth_available": False, + "ghost_occupancy_past_ttl_count": 0, + "maximum_expiry_deadline_delay_ns": 0, + "maximum_expiry_materialization_delay_ns": 283000000, + "maximum_held_age_ns": 750000000, + "past_ttl_occupied_publications": 0, + "ttl_ns": 750000000, + } + assert clips == { + "component_correspondence_available": False, + "frame_count": 16, + "independent_truth": False, + } + + +def test_temporal_result_is_digest_bound_to_m4_4_e34_e51_and_e46b( + result: TemporalReplayResult, +) -> None: + identity = result.manifest["identity"] + assert isinstance(identity, dict) + assert identity["geometry_result_id"] == ( + "m4-geometry-replay-" + "8daf3109e3cf30b960b4b376032ff3b5ec58ca42a1e5899b841cf29fbcf14ad8" + ) + assert identity["frames_sha256"] == ( + "1bf1365bdb3f20214443d3f8b87a0fa88f9848af8ca0456b7ca364d37631c3fc" + ) + references = identity["historical_references"] + assert isinstance(references, dict) + assert set(references) == {"e34", "e46b", "e51"} + assert references["e46b"]["independent_truth"] is False + assert len(identity["clip_checks"]) == 16