feat(perception): qualify E35 degradation recovery

This commit is contained in:
DCCONSTRUCTIONS
2026-07-27 17:53:59 +03:00
parent 621084fcd6
commit b894945344
21 changed files with 3985 additions and 17 deletions
+623
View File
@@ -0,0 +1,623 @@
"""Deterministic E35 source degradation transforms over TrackGeometry v1.
The transforms are pure and frame-local. They never alter the accepted E32
source, infer free space, retain evidence whose required channel is absent, or
grant runtime authority.
"""
from __future__ import annotations
import hashlib
import json
import math
from dataclasses import dataclass
from enum import StrEnum
from typing import Any, Final
import numpy as np
from .track_geometry import (
PointSlab,
TrackGeometry,
TrackGeometryCurrentness,
TrackGeometryEvidenceState,
TrackGeometryFrame,
TrackGeometryMetricBasis,
TrackGeometryOwnerKind,
)
E35_SCENARIO_SCHEMA: Final = "missioncore.e35-degradation-scenario/v1"
E35_TRANSFORMATION_SCHEMA: Final = "missioncore.e35-frame-transformation/v1"
class DegradationRecoveryError(ValueError):
"""An E35 scenario or transformed frame violates the frozen contract."""
class DegradationKind(StrEnum):
CAMERA_LOSS = "camera-loss"
LIDAR_LOSS = "lidar-loss"
POSE_STALENESS = "pose-staleness"
DELAYED_FRAMES = "delayed-frames"
BOUNDED_DROP = "bounded-drop"
TIMING_OFFSET = "timing-offset"
@dataclass(frozen=True, slots=True)
class DegradationScenario:
"""One bounded deterministic source transformation."""
scenario_id: str
kind: DegradationKind
frame_start: int
frame_end: int
parameters: dict[str, Any]
def __post_init__(self) -> None:
if (
not self.scenario_id
or len(self.scenario_id) > 80
or self.scenario_id != self.kind.value
or not isinstance(self.frame_start, int)
or isinstance(self.frame_start, bool)
or not isinstance(self.frame_end, int)
or isinstance(self.frame_end, bool)
or self.frame_start < 1
or self.frame_end < self.frame_start
):
raise DegradationRecoveryError("degradation scenario bounds are invalid")
self._validate_parameters()
def _validate_parameters(self) -> None:
expected: dict[DegradationKind, set[str]] = {
DegradationKind.CAMERA_LOSS: {"drop_camera_observations"},
DegradationKind.LIDAR_LOSS: {"drop_lidar_points"},
DegradationKind.POSE_STALENESS: {"pose_age_seconds"},
DegradationKind.DELAYED_FRAMES: {
"delay_seconds",
"late_result_policy",
},
DegradationKind.BOUNDED_DROP: {"drop_every_nth_frame"},
DegradationKind.TIMING_OFFSET: {"camera_lidar_offset_ms"},
}
if set(self.parameters) != expected[self.kind]:
raise DegradationRecoveryError(
"degradation scenario parameters are incompatible"
)
if self.kind is DegradationKind.CAMERA_LOSS:
_require_true(self.parameters["drop_camera_observations"])
elif self.kind is DegradationKind.LIDAR_LOSS:
_require_true(self.parameters["drop_lidar_points"])
elif self.kind is DegradationKind.POSE_STALENESS:
_positive_number(self.parameters["pose_age_seconds"], "pose age")
elif self.kind is DegradationKind.DELAYED_FRAMES:
_positive_number(self.parameters["delay_seconds"], "delivery delay")
if self.parameters["late_result_policy"] != "discard":
raise DegradationRecoveryError(
"late E35 results must be discarded"
)
elif self.kind is DegradationKind.BOUNDED_DROP:
value = self.parameters["drop_every_nth_frame"]
if (
not isinstance(value, int)
or isinstance(value, bool)
or value < 2
):
raise DegradationRecoveryError(
"bounded drop cadence is invalid"
)
else:
value = self.parameters["camera_lidar_offset_ms"]
if (
not isinstance(value, int)
or isinstance(value, bool)
or abs(value) <= 100
or abs(value) > 1_000
):
raise DegradationRecoveryError(
"timing offset must exceed the admitted binding"
)
@property
def frame_count(self) -> int:
return self.frame_end - self.frame_start + 1
def active(self, frame_index: int) -> bool:
return self.frame_start <= frame_index <= self.frame_end
def to_dict(self) -> dict[str, Any]:
return {
"schema_version": E35_SCENARIO_SCHEMA,
"scenario_id": self.scenario_id,
"kind": self.kind.value,
"frame_start": self.frame_start,
"frame_end": self.frame_end,
"parameters": self.parameters,
}
@classmethod
def from_dict(cls, value: object) -> DegradationScenario:
document = _object(value, "degradation scenario")
if set(document) != {
"scenario_id",
"kind",
"frame_start",
"frame_end",
"parameters",
}:
raise DegradationRecoveryError(
"degradation scenario fields are incompatible"
)
try:
kind = DegradationKind(
_string(document.get("kind"), "degradation kind")
)
except ValueError as exc:
raise DegradationRecoveryError(
"degradation scenario kind is invalid"
) from exc
parameters = _object(
document.get("parameters"),
"degradation scenario parameters",
)
return cls(
scenario_id=_string(document.get("scenario_id"), "scenario id"),
kind=kind,
frame_start=_integer(document.get("frame_start"), "frame start"),
frame_end=_integer(document.get("frame_end"), "frame end"),
parameters=parameters,
)
@dataclass(frozen=True, slots=True)
class TransformedTrackGeometryFrame:
"""One E35 frame and its explicit transformation document."""
frame: TrackGeometryFrame
transformation: dict[str, Any]
def transform_track_geometry_frame(
frame: TrackGeometryFrame,
scenario: DegradationScenario,
) -> TransformedTrackGeometryFrame:
"""Apply one scenario without mutating the accepted source frame."""
if not scenario.active(frame.frame_index):
return TransformedTrackGeometryFrame(
frame=frame,
transformation=_transformation(
frame,
scenario,
active=False,
action="none",
channels=_nominal_channels(frame),
transformed_frame=frame,
),
)
if scenario.kind is DegradationKind.CAMERA_LOSS:
transformed = _camera_loss(frame)
action = "camera-observations-removed"
channels = {
"camera": "unavailable",
"lidar": _lidar_state(frame),
"pose": "available",
"delivery": "on-time",
}
elif scenario.kind is DegradationKind.LIDAR_LOSS:
transformed = _metric_unavailable(
frame,
reason="e35-lidar-unavailable",
source_available=False,
)
action = "lidar-points-withheld"
channels = {
"camera": "available",
"lidar": "unavailable",
"pose": "available",
"delivery": "on-time",
}
elif scenario.kind is DegradationKind.POSE_STALENESS:
transformed = _metric_unavailable(
frame,
reason="e35-pose-stale",
source_available=frame.source_available,
)
action = "map-points-withheld-for-stale-pose"
channels = {
"camera": "available",
"lidar": _lidar_state(frame),
"pose": "stale",
"delivery": "on-time",
}
elif scenario.kind is DegradationKind.DELAYED_FRAMES:
transformed = _empty_frame(frame)
action = "late-frame-discarded"
channels = {
"camera": "late-discarded",
"lidar": "late-discarded",
"pose": "late-discarded",
"delivery": "late-discarded",
}
elif scenario.kind is DegradationKind.BOUNDED_DROP:
cadence = int(scenario.parameters["drop_every_nth_frame"])
dropped = (frame.frame_index - scenario.frame_start) % cadence == 0
transformed = _empty_frame(frame) if dropped else frame
action = "input-frame-dropped" if dropped else "bounded-drop-pass"
channels = (
{
"camera": "dropped",
"lidar": "dropped",
"pose": "dropped",
"delivery": "dropped",
}
if dropped
else _nominal_channels(frame)
)
else:
transformed = _timing_offset(frame)
action = "camera-lidar-evidence-split"
channels = {
"camera": "offset",
"lidar": _lidar_state(frame),
"pose": "available",
"delivery": "on-time",
}
return TransformedTrackGeometryFrame(
frame=transformed,
transformation=_transformation(
frame,
scenario,
active=True,
action=action,
channels=channels,
transformed_frame=transformed,
),
)
def frame_digest(frame: TrackGeometryFrame) -> str:
"""Digest one complete frame without serializing large point arrays."""
digest = hashlib.sha256()
compact = {
"binding": frame.binding.to_dict(),
"frame_index": frame.frame_index,
"source_frame_index": frame.source_frame_index,
"session_seconds": frame.session_seconds,
"source_available": frame.source_available,
"source_point_count": frame.point_slab.source_point_count,
"coordinate_frame": frame.point_slab.coordinate_frame,
"owner_keys": list(frame.point_slab.owner_keys),
"geometries": [geometry.to_dict() for geometry in frame.geometries],
}
digest.update(_canonical_json(compact))
digest.update(frame.point_slab.source_indices.astype("<i8", copy=False).tobytes())
digest.update(frame.point_slab.points_xyz_m.astype("<f4", copy=False).tobytes())
digest.update(frame.point_slab.owner_indices.astype("<u4", copy=False).tobytes())
return digest.hexdigest()
def _camera_loss(frame: TrackGeometryFrame) -> TrackGeometryFrame:
geometries: list[TrackGeometry] = []
owner_map: dict[str, str | None] = {}
for geometry in frame.geometries:
if geometry.owner_kind is TrackGeometryOwnerKind.GEOMETRY_CLUSTER:
geometries.append(geometry)
if geometry.metric_basis is TrackGeometryMetricBasis.CURRENT_POINTS:
owner_map[geometry.owner_key] = geometry.owner_key
continue
if geometry.metric_basis is not TrackGeometryMetricBasis.CURRENT_POINTS:
owner_map[geometry.owner_key] = None
continue
owner_key = f"e35-camera-loss:{geometry.owner_key}"
geometries.append(
TrackGeometry(
owner_key=owner_key,
owner_kind=TrackGeometryOwnerKind.GEOMETRY_CLUSTER,
evidence_state=TrackGeometryEvidenceState.GEOMETRY_ONLY,
currentness=TrackGeometryCurrentness.CURRENT,
metric_basis=TrackGeometryMetricBasis.CURRENT_POINTS,
reason_codes=(
"e35-camera-unavailable",
"geometry-retained-without-semantics",
),
range_m=geometry.range_m,
)
)
owner_map[geometry.owner_key] = owner_key
return _frame_with(
frame,
geometries=geometries,
point_slab=_remap_slab(frame, geometries, owner_map),
source_available=frame.source_available,
)
def _metric_unavailable(
frame: TrackGeometryFrame,
*,
reason: str,
source_available: bool,
) -> TrackGeometryFrame:
geometries: list[TrackGeometry] = []
for geometry in frame.geometries:
if geometry.owner_kind is TrackGeometryOwnerKind.GEOMETRY_CLUSTER:
continue
if geometry.metric_basis is TrackGeometryMetricBasis.CURRENT_POINTS:
geometries.append(
TrackGeometry(
owner_key=geometry.owner_key,
owner_kind=TrackGeometryOwnerKind.CAMERA_TRACK,
evidence_state=TrackGeometryEvidenceState.CAMERA_ONLY,
currentness=TrackGeometryCurrentness.CURRENT,
metric_basis=TrackGeometryMetricBasis.UNAVAILABLE,
reason_codes=(reason, "camera-remains-nonmetric"),
semantic_track_id=geometry.semantic_track_id,
semantic_label=geometry.semantic_label,
bbox_xyxy=geometry.bbox_xyxy,
)
)
else:
geometries.append(geometry)
return _frame_with(
frame,
geometries=geometries,
point_slab=_empty_slab(frame),
source_available=source_available,
)
def _timing_offset(frame: TrackGeometryFrame) -> TrackGeometryFrame:
geometries: list[TrackGeometry] = []
owner_map: dict[str, str | None] = {}
for geometry in frame.geometries:
if (
geometry.owner_kind is TrackGeometryOwnerKind.CAMERA_TRACK
and geometry.metric_basis is TrackGeometryMetricBasis.CURRENT_POINTS
):
geometries.append(
TrackGeometry(
owner_key=geometry.owner_key,
owner_kind=TrackGeometryOwnerKind.CAMERA_TRACK,
evidence_state=TrackGeometryEvidenceState.CAMERA_ONLY,
currentness=TrackGeometryCurrentness.CURRENT,
metric_basis=TrackGeometryMetricBasis.UNAVAILABLE,
reason_codes=(
"e35-camera-lidar-offset-unqualified",
"camera-remains-nonmetric",
),
semantic_track_id=geometry.semantic_track_id,
semantic_label=geometry.semantic_label,
bbox_xyxy=geometry.bbox_xyxy,
)
)
geometry_owner = f"e35-timing-offset:{geometry.owner_key}"
geometries.append(
TrackGeometry(
owner_key=geometry_owner,
owner_kind=TrackGeometryOwnerKind.GEOMETRY_CLUSTER,
evidence_state=TrackGeometryEvidenceState.GEOMETRY_ONLY,
currentness=TrackGeometryCurrentness.CURRENT,
metric_basis=TrackGeometryMetricBasis.CURRENT_POINTS,
reason_codes=(
"e35-camera-lidar-offset-unqualified",
"geometry-retained-without-semantics",
),
range_m=geometry.range_m,
)
)
owner_map[geometry.owner_key] = geometry_owner
else:
geometries.append(geometry)
if geometry.metric_basis is TrackGeometryMetricBasis.CURRENT_POINTS:
owner_map[geometry.owner_key] = geometry.owner_key
return _frame_with(
frame,
geometries=geometries,
point_slab=_remap_slab(frame, geometries, owner_map),
source_available=frame.source_available,
)
def _empty_frame(frame: TrackGeometryFrame) -> TrackGeometryFrame:
return _frame_with(
frame,
geometries=[],
point_slab=_empty_slab(frame),
source_available=False,
)
def _frame_with(
frame: TrackGeometryFrame,
*,
geometries: list[TrackGeometry],
point_slab: PointSlab,
source_available: bool,
) -> TrackGeometryFrame:
return TrackGeometryFrame(
binding=frame.binding,
frame_index=frame.frame_index,
source_frame_index=frame.source_frame_index,
session_seconds=frame.session_seconds,
source_available=source_available,
point_slab=point_slab,
geometries=tuple(geometries),
)
def _empty_slab(frame: TrackGeometryFrame) -> PointSlab:
return PointSlab(
frame_index=frame.frame_index,
source_frame_index=frame.source_frame_index,
source_point_count=frame.point_slab.source_point_count,
coordinate_frame=frame.point_slab.coordinate_frame,
owner_keys=(),
source_indices=np.empty(0, dtype="<i8"),
points_xyz_m=np.empty((0, 3), dtype="<f4"),
owner_indices=np.empty(0, dtype="<u4"),
)
def _remap_slab(
frame: TrackGeometryFrame,
geometries: list[TrackGeometry],
owner_map: dict[str, str | None],
) -> PointSlab:
owner_keys = tuple(
geometry.owner_key
for geometry in geometries
if geometry.metric_basis is TrackGeometryMetricBasis.CURRENT_POINTS
)
owner_indices_by_key = {
owner_key: owner_index
for owner_index, owner_key in enumerate(owner_keys)
}
source_indices: list[int] = []
points: list[list[float]] = []
owner_indices: list[int] = []
for row_index, source_index in enumerate(frame.point_slab.source_indices):
old_owner = frame.point_slab.owner_keys[
int(frame.point_slab.owner_indices[row_index])
]
new_owner = owner_map.get(old_owner)
if new_owner is None:
continue
source_indices.append(int(source_index))
points.append(frame.point_slab.points_xyz_m[row_index].tolist())
owner_indices.append(owner_indices_by_key[new_owner])
return PointSlab(
frame_index=frame.frame_index,
source_frame_index=frame.source_frame_index,
source_point_count=frame.point_slab.source_point_count,
coordinate_frame=frame.point_slab.coordinate_frame,
owner_keys=owner_keys,
source_indices=np.asarray(source_indices, dtype="<i8"),
points_xyz_m=(
np.asarray(points, dtype="<f4")
if points
else np.empty((0, 3), dtype="<f4")
),
owner_indices=np.asarray(owner_indices, dtype="<u4"),
)
def _transformation(
original_frame: TrackGeometryFrame,
scenario: DegradationScenario,
*,
active: bool,
action: str,
channels: dict[str, str],
transformed_frame: TrackGeometryFrame,
) -> dict[str, Any]:
return {
"schema_version": E35_TRANSFORMATION_SCHEMA,
"scenario_id": scenario.scenario_id,
"kind": scenario.kind.value,
"frame_index": original_frame.frame_index,
"active": active,
"action": action,
"parameters": scenario.parameters if active else {},
"channels": channels,
"original_frame_sha256": frame_digest(original_frame),
"transformed_frame_sha256": frame_digest(transformed_frame),
"original": _frame_counts(original_frame),
"transformed": _frame_counts(transformed_frame),
"policy": {
"absence_of_points_means_free": False,
"late_results_may_reenter": False,
"persistent_reconstruction_modified": False,
},
"authority": _authority(),
}
def _frame_counts(frame: TrackGeometryFrame) -> dict[str, int | bool]:
return {
"source_available": frame.source_available,
"point_rows": frame.point_slab.row_count,
"camera_tracks": sum(
geometry.owner_kind is TrackGeometryOwnerKind.CAMERA_TRACK
for geometry in frame.geometries
),
"geometry_clusters": sum(
geometry.owner_kind is TrackGeometryOwnerKind.GEOMETRY_CLUSTER
for geometry in frame.geometries
),
"agree": sum(
geometry.evidence_state is TrackGeometryEvidenceState.AGREE
for geometry in frame.geometries
),
}
def _nominal_channels(frame: TrackGeometryFrame) -> dict[str, str]:
return {
"camera": "available",
"lidar": _lidar_state(frame),
"pose": "available",
"delivery": "on-time",
}
def _lidar_state(frame: TrackGeometryFrame) -> str:
return "available" if frame.source_available else "source-unavailable"
def _canonical_json(value: object) -> bytes:
return json.dumps(
value,
ensure_ascii=False,
sort_keys=True,
separators=(",", ":"),
allow_nan=False,
).encode("utf-8")
def _authority() -> dict[str, bool]:
return {
"commands_enabled": False,
"navigation_or_safety_accepted": False,
}
def _require_true(value: object) -> None:
if value is not True:
raise DegradationRecoveryError(
"degradation source removal must be enabled"
)
def _positive_number(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 DegradationRecoveryError(f"{label} is invalid")
return float(value)
def _integer(value: object, label: str) -> int:
if not isinstance(value, int) or isinstance(value, bool):
raise DegradationRecoveryError(f"{label} is invalid")
return value
def _string(value: object, label: str) -> str:
if not isinstance(value, str) or not value:
raise DegradationRecoveryError(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 DegradationRecoveryError(f"{label} must be an object")
return value
File diff suppressed because it is too large Load Diff