Capture the current unreleased K1 connection, recovery, lifecycle, viewer, and test work as a single known-bad baseline for subsequent fixes.
1051 lines
40 KiB
Python
1051 lines
40 KiB
Python
"""Bounded latest-wins primitives for derived live perception."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import json
|
|
import math
|
|
import statistics
|
|
import struct
|
|
import threading
|
|
import time
|
|
import zlib
|
|
from collections import deque
|
|
from collections.abc import Callable, Mapping, Sequence
|
|
from dataclasses import dataclass
|
|
from hashlib import sha256
|
|
from typing import Any, Final, Literal
|
|
|
|
import numpy as np
|
|
|
|
from k1link.data_plane import DecodedPointCloudView, DecodedPoseView
|
|
|
|
WORLD_STATE_SCHEMA = "missioncore.live-perception-world-state/v1"
|
|
TELEMETRY_SCHEMA = "missioncore.live-perception-telemetry/v1"
|
|
|
|
HealthState = Literal["healthy", "degraded", "stale", "unavailable"]
|
|
LiveIngressModality = Literal[
|
|
"control",
|
|
"camera-init",
|
|
"camera-frame",
|
|
"lidar",
|
|
"pose",
|
|
]
|
|
|
|
LIVE_INGRESS_SCHEMA: Final = "missioncore.live-perception-ingress/v1"
|
|
LIVE_INGRESS_WIRE_SCHEMA: Final = "missioncore.live-perception-wire/v1"
|
|
LIVE_RESULT_WIRE_SCHEMA: Final = "missioncore.live-perception-result-wire/v2"
|
|
LIVE_RESULT_MAGIC: Final = b"MCPR"
|
|
LIVE_RESULT_MAX_HEADER_BYTES: Final = 256 * 1024
|
|
LIVE_RESULT_MAX_PAYLOAD_BYTES: Final = 2 * 1024 * 1024
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class LivePerceptionResultFrame:
|
|
session_id: str
|
|
session_generation: int
|
|
frame_index: int
|
|
source_frame_index: int
|
|
session_seconds: float
|
|
captured_at_epoch_ns: int
|
|
image_jpeg: bytes
|
|
segmentation_mask: np.ndarray[Any, np.dtype[np.uint8]] | None
|
|
objects: tuple[dict[str, Any], ...]
|
|
delivery: dict[str, Any]
|
|
|
|
|
|
def encode_live_perception_result(
|
|
*,
|
|
session_id: str,
|
|
session_generation: int,
|
|
frame_index: int,
|
|
source_frame_index: int,
|
|
session_seconds: float,
|
|
captured_at_epoch_ns: int,
|
|
image_jpeg: bytes,
|
|
segmentation_mask: np.ndarray[Any, Any] | None,
|
|
objects: Sequence[Mapping[str, Any]],
|
|
delivery: Mapping[str, Any],
|
|
) -> bytes:
|
|
"""Encode one bounded, non-authoritative worker-to-viewer result frame."""
|
|
|
|
if (
|
|
not session_id
|
|
or len(session_id) > 160
|
|
or isinstance(session_generation, bool)
|
|
or session_generation < 1
|
|
or frame_index < 0
|
|
or source_frame_index < 0
|
|
or captured_at_epoch_ns < 0
|
|
or not math.isfinite(session_seconds)
|
|
or session_seconds < 0
|
|
or not 4 <= len(image_jpeg) <= 1024 * 1024
|
|
or not image_jpeg.startswith(b"\xff\xd8")
|
|
or not image_jpeg.endswith(b"\xff\xd9")
|
|
):
|
|
raise ValueError("live perception result identity or image is invalid")
|
|
normalized_objects = tuple(_normalize_live_result_object(value) for value in objects)
|
|
if len(normalized_objects) > 128:
|
|
raise ValueError("live perception result object count exceeds the bound")
|
|
mask_payload = b""
|
|
mask_shape: list[int] | None = None
|
|
if segmentation_mask is not None:
|
|
mask = np.asarray(segmentation_mask, dtype=np.uint8)
|
|
if mask.shape != (600, 800):
|
|
raise ValueError("live perception segmentation shape is invalid")
|
|
mask_payload = zlib.compress(mask.tobytes(order="C"), level=1)
|
|
mask_shape = [600, 800]
|
|
payload = bytes(image_jpeg) + mask_payload
|
|
if len(payload) > LIVE_RESULT_MAX_PAYLOAD_BYTES:
|
|
raise ValueError("live perception result payload exceeds the bound")
|
|
header = {
|
|
"schema_version": LIVE_RESULT_WIRE_SCHEMA,
|
|
"session_id": session_id,
|
|
"session_generation": session_generation,
|
|
"frame_index": frame_index,
|
|
"source_frame_index": source_frame_index,
|
|
"session_seconds": session_seconds,
|
|
"captured_at_epoch_ns": captured_at_epoch_ns,
|
|
"image": {"codec": "jpeg", "byte_length": len(image_jpeg)},
|
|
"segmentation": (
|
|
None
|
|
if mask_shape is None
|
|
else {
|
|
"codec": "zlib-uint8-c1",
|
|
"shape": mask_shape,
|
|
"byte_length": len(mask_payload),
|
|
}
|
|
),
|
|
"objects": normalized_objects,
|
|
"delivery": dict(delivery),
|
|
"payload_bytes": len(payload),
|
|
"payload_sha256": sha256(payload).hexdigest(),
|
|
"authority": "shadow-diagnostic-only",
|
|
"commands_enabled": False,
|
|
"navigation_or_safety_accepted": False,
|
|
}
|
|
encoded_header = json.dumps(
|
|
header,
|
|
sort_keys=True,
|
|
separators=(",", ":"),
|
|
allow_nan=False,
|
|
).encode("utf-8")
|
|
if len(encoded_header) > LIVE_RESULT_MAX_HEADER_BYTES:
|
|
raise ValueError("live perception result header exceeds the bound")
|
|
return LIVE_RESULT_MAGIC + struct.pack("!I", len(encoded_header)) + encoded_header + payload
|
|
|
|
|
|
def decode_live_perception_result(encoded: bytes) -> LivePerceptionResultFrame:
|
|
"""Validate and decode one result frame before it reaches the Rerun bridge."""
|
|
|
|
if len(encoded) < 10 or not encoded.startswith(LIVE_RESULT_MAGIC):
|
|
raise ValueError("live perception result frame is truncated")
|
|
header_length = struct.unpack("!I", encoded[4:8])[0]
|
|
if not 2 <= header_length <= LIVE_RESULT_MAX_HEADER_BYTES:
|
|
raise ValueError("live perception result header length is invalid")
|
|
boundary = 8 + header_length
|
|
if boundary > len(encoded):
|
|
raise ValueError("live perception result header is truncated")
|
|
try:
|
|
header = json.loads(encoded[8:boundary])
|
|
except (UnicodeDecodeError, json.JSONDecodeError) as exc:
|
|
raise ValueError("live perception result header is invalid") from exc
|
|
payload = encoded[boundary:]
|
|
if (
|
|
not isinstance(header, dict)
|
|
or header.get("schema_version") != LIVE_RESULT_WIRE_SCHEMA
|
|
or header.get("authority") != "shadow-diagnostic-only"
|
|
or header.get("commands_enabled") is not False
|
|
or header.get("navigation_or_safety_accepted") is not False
|
|
or header.get("payload_bytes") != len(payload)
|
|
or len(payload) > LIVE_RESULT_MAX_PAYLOAD_BYTES
|
|
or header.get("payload_sha256") != sha256(payload).hexdigest()
|
|
):
|
|
raise ValueError("live perception result contract is invalid")
|
|
image = header.get("image")
|
|
segmentation = header.get("segmentation")
|
|
objects = header.get("objects")
|
|
delivery = header.get("delivery")
|
|
if (
|
|
not isinstance(image, dict)
|
|
or image.get("codec") != "jpeg"
|
|
or not isinstance(image.get("byte_length"), int)
|
|
or not isinstance(objects, list)
|
|
or len(objects) > 128
|
|
or not isinstance(delivery, dict)
|
|
):
|
|
raise ValueError("live perception result content descriptor is invalid")
|
|
image_length = image["byte_length"]
|
|
if not 4 <= image_length <= min(len(payload), 1024 * 1024):
|
|
raise ValueError("live perception result image length is invalid")
|
|
image_jpeg = payload[:image_length]
|
|
if not image_jpeg.startswith(b"\xff\xd8") or not image_jpeg.endswith(b"\xff\xd9"):
|
|
raise ValueError("live perception result JPEG is invalid")
|
|
mask: np.ndarray[Any, np.dtype[np.uint8]] | None = None
|
|
if segmentation is None:
|
|
if len(payload) != image_length:
|
|
raise ValueError("live perception result has an undescribed payload tail")
|
|
else:
|
|
if (
|
|
not isinstance(segmentation, dict)
|
|
or segmentation.get("codec") != "zlib-uint8-c1"
|
|
or segmentation.get("shape") != [600, 800]
|
|
or not isinstance(segmentation.get("byte_length"), int)
|
|
or segmentation["byte_length"] != len(payload) - image_length
|
|
):
|
|
raise ValueError("live perception segmentation descriptor is invalid")
|
|
try:
|
|
raw_mask = zlib.decompress(payload[image_length:])
|
|
except zlib.error as exc:
|
|
raise ValueError("live perception segmentation payload is invalid") from exc
|
|
if len(raw_mask) != 600 * 800:
|
|
raise ValueError("live perception segmentation byte length is invalid")
|
|
mask = np.frombuffer(raw_mask, dtype=np.uint8).reshape((600, 800)).copy()
|
|
normalized_objects = tuple(_normalize_live_result_object(value) for value in objects)
|
|
frame_index = header.get("frame_index")
|
|
session_id = header.get("session_id")
|
|
session_generation = header.get("session_generation")
|
|
source_frame_index = header.get("source_frame_index")
|
|
session_seconds = header.get("session_seconds")
|
|
captured_at_epoch_ns = header.get("captured_at_epoch_ns")
|
|
if (
|
|
not isinstance(session_id, str)
|
|
or not session_id
|
|
or len(session_id) > 160
|
|
or not isinstance(session_generation, int)
|
|
or isinstance(session_generation, bool)
|
|
or session_generation < 1
|
|
or not isinstance(frame_index, int)
|
|
or isinstance(frame_index, bool)
|
|
or frame_index < 0
|
|
or not isinstance(source_frame_index, int)
|
|
or isinstance(source_frame_index, bool)
|
|
or source_frame_index < 0
|
|
or not isinstance(captured_at_epoch_ns, int)
|
|
or isinstance(captured_at_epoch_ns, bool)
|
|
or captured_at_epoch_ns < 0
|
|
or not isinstance(session_seconds, (int, float))
|
|
or isinstance(session_seconds, bool)
|
|
or not math.isfinite(float(session_seconds))
|
|
or float(session_seconds) < 0
|
|
):
|
|
raise ValueError("live perception result time identity is invalid")
|
|
return LivePerceptionResultFrame(
|
|
session_id=session_id,
|
|
session_generation=session_generation,
|
|
frame_index=frame_index,
|
|
source_frame_index=source_frame_index,
|
|
session_seconds=float(session_seconds),
|
|
captured_at_epoch_ns=captured_at_epoch_ns,
|
|
image_jpeg=image_jpeg,
|
|
segmentation_mask=mask,
|
|
objects=normalized_objects,
|
|
delivery=dict(delivery),
|
|
)
|
|
|
|
|
|
def _normalize_live_result_object(value: Mapping[str, Any]) -> dict[str, Any]:
|
|
if not isinstance(value, Mapping):
|
|
raise ValueError("live perception result object is invalid")
|
|
track_id = value.get("track_id")
|
|
label = value.get("label")
|
|
score = value.get("score")
|
|
bbox = value.get("bbox_xyxy")
|
|
if (
|
|
not isinstance(track_id, int)
|
|
or isinstance(track_id, bool)
|
|
or track_id < 0
|
|
or not isinstance(label, str)
|
|
or not 1 <= len(label) <= 64
|
|
or not isinstance(score, (int, float))
|
|
or isinstance(score, bool)
|
|
or not math.isfinite(float(score))
|
|
or not isinstance(bbox, Sequence)
|
|
or isinstance(bbox, (str, bytes))
|
|
or len(bbox) != 4
|
|
):
|
|
raise ValueError("live perception result object identity is invalid")
|
|
bbox_values = [float(item) for item in bbox]
|
|
if not all(math.isfinite(item) for item in bbox_values):
|
|
raise ValueError("live perception result 2D box is invalid")
|
|
normalized: dict[str, Any] = {
|
|
"track_id": track_id,
|
|
"label": label,
|
|
"score": float(score),
|
|
"bbox_xyxy": bbox_values,
|
|
}
|
|
distance = value.get("distance_smoothed_m", value.get("distance_median_m"))
|
|
if distance is not None:
|
|
if (
|
|
not isinstance(distance, (int, float))
|
|
or isinstance(distance, bool)
|
|
or not math.isfinite(float(distance))
|
|
or float(distance) < 0
|
|
):
|
|
raise ValueError("live perception result distance is invalid")
|
|
normalized["distance_m"] = float(distance)
|
|
else:
|
|
normalized["distance_m"] = None
|
|
cuboid_fields = (
|
|
("cuboid_center_map", 3),
|
|
("cuboid_half_size", 3),
|
|
("cuboid_quaternion_xyzw", 4),
|
|
)
|
|
present = [value.get(name) is not None for name, _ in cuboid_fields]
|
|
if any(present) and not all(present):
|
|
raise ValueError("live perception result cuboid is incomplete")
|
|
for name, length in cuboid_fields:
|
|
candidate = value.get(name)
|
|
if candidate is None:
|
|
normalized[name] = None
|
|
continue
|
|
if (
|
|
not isinstance(candidate, Sequence)
|
|
or isinstance(candidate, (str, bytes))
|
|
or len(candidate) != length
|
|
):
|
|
raise ValueError("live perception result cuboid geometry is invalid")
|
|
values = [float(item) for item in candidate]
|
|
if not all(math.isfinite(item) for item in values):
|
|
raise ValueError("live perception result cuboid contains non-finite values")
|
|
normalized[name] = values
|
|
return normalized
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class LiveIngressEvent:
|
|
"""One raw-first, derived-only event admitted to the shadow transport."""
|
|
|
|
ingress_sequence: int
|
|
session_id: str
|
|
session_generation: int
|
|
modality: LiveIngressModality
|
|
source_id: str
|
|
source_sequence: int
|
|
captured_at_epoch_ns: int
|
|
received_monotonic_ns: int
|
|
payload: bytes
|
|
|
|
def wire_bytes(self) -> bytes:
|
|
header = json.dumps(
|
|
{
|
|
"schema_version": LIVE_INGRESS_WIRE_SCHEMA,
|
|
"ingress_sequence": self.ingress_sequence,
|
|
"session_id": self.session_id,
|
|
"session_generation": self.session_generation,
|
|
"modality": self.modality,
|
|
"source_id": self.source_id,
|
|
"source_sequence": self.source_sequence,
|
|
"captured_at_epoch_ns": self.captured_at_epoch_ns,
|
|
"received_monotonic_ns": self.received_monotonic_ns,
|
|
"payload_bytes": len(self.payload),
|
|
"payload_sha256": sha256(self.payload).hexdigest(),
|
|
"authority": "shadow-diagnostic-only",
|
|
"commands_enabled": False,
|
|
"navigation_or_safety_accepted": False,
|
|
},
|
|
ensure_ascii=True,
|
|
separators=(",", ":"),
|
|
sort_keys=True,
|
|
).encode("utf-8")
|
|
return struct.pack("!I", len(header)) + header + self.payload
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class LiveIngressQueueSnapshot:
|
|
capacity: int
|
|
depth: int
|
|
maximum_depth: int
|
|
published: int
|
|
consumed: int
|
|
dropped_overflow: int
|
|
rejected_oversize: int
|
|
|
|
|
|
@dataclass(slots=True)
|
|
class _LiveIngressQueue:
|
|
capacity: int
|
|
items: deque[LiveIngressEvent]
|
|
maximum_depth: int = 0
|
|
published: int = 0
|
|
consumed: int = 0
|
|
dropped_overflow: int = 0
|
|
rejected_oversize: int = 0
|
|
|
|
|
|
class LivePerceptionIngress:
|
|
"""Exclusive, bounded fan-out from committed K1 evidence to one AI worker.
|
|
|
|
The ingress is deliberately not an acquisition source and has no command
|
|
surface. Camera and MQTT producers call it only after their raw evidence
|
|
commit has completed. Separate modality queues prevent camera bursts from
|
|
evicting pose or LiDAR observations.
|
|
"""
|
|
|
|
_CAPACITIES: Final[dict[LiveIngressModality, int]] = {
|
|
"control": 4,
|
|
"camera-init": 1,
|
|
"camera-frame": 2,
|
|
"lidar": 8,
|
|
"pose": 16,
|
|
}
|
|
_MAX_PAYLOAD_BYTES: Final[dict[LiveIngressModality, int]] = {
|
|
"control": 16 * 1024,
|
|
"camera-init": 1024 * 1024,
|
|
"camera-frame": 1024 * 1024,
|
|
"lidar": 2 * 1024 * 1024,
|
|
"pose": 2 * 1024 * 1024,
|
|
}
|
|
|
|
def __init__(self) -> None:
|
|
self._condition = threading.Condition()
|
|
self._queues = {
|
|
modality: _LiveIngressQueue(capacity, deque())
|
|
for modality, capacity in self._CAPACITIES.items()
|
|
}
|
|
self._ingress_sequence = 0
|
|
self._session_id: str | None = None
|
|
self._session_generation = 0
|
|
self._active = False
|
|
self._closed = False
|
|
self._consumer_id: str | None = None
|
|
self._results_accepted = 0
|
|
self._results_rejected_stale = 0
|
|
self._results_rejected_receiver = 0
|
|
|
|
def begin_session(self, session_id: str) -> None:
|
|
if not session_id or len(session_id) > 160:
|
|
raise ValueError("live perception session id is invalid")
|
|
with self._condition:
|
|
if self._closed:
|
|
raise RuntimeError("live perception ingress is closed")
|
|
if self._active:
|
|
if self._session_id == session_id:
|
|
return
|
|
raise RuntimeError("another live perception session is active")
|
|
# Queued items belong to one acquisition. A worker may be absent
|
|
# while a session ends, so clear bounded leftovers before the next
|
|
# session can become visible to that worker.
|
|
for queue in self._queues.values():
|
|
queue.items.clear()
|
|
self._session_id = session_id
|
|
self._session_generation += 1
|
|
self._active = True
|
|
self._publish_locked(
|
|
modality="control",
|
|
source_id="mission-core",
|
|
source_sequence=0,
|
|
captured_at_epoch_ns=time.time_ns(),
|
|
received_monotonic_ns=time.monotonic_ns(),
|
|
payload=b'{"event":"session-start"}',
|
|
)
|
|
|
|
def end_session(self, session_id: str) -> None:
|
|
with self._condition:
|
|
if not self._active or self._session_id != session_id:
|
|
return
|
|
self._publish_locked(
|
|
modality="control",
|
|
source_id="mission-core",
|
|
source_sequence=0,
|
|
captured_at_epoch_ns=time.time_ns(),
|
|
received_monotonic_ns=time.monotonic_ns(),
|
|
payload=b'{"event":"session-end"}',
|
|
)
|
|
self._active = False
|
|
|
|
def publish(
|
|
self,
|
|
*,
|
|
modality: LiveIngressModality,
|
|
source_id: str,
|
|
source_sequence: int,
|
|
captured_at_epoch_ns: int,
|
|
received_monotonic_ns: int,
|
|
payload: bytes,
|
|
) -> bool:
|
|
if modality == "control":
|
|
raise ValueError("control events are owned by the ingress lifecycle")
|
|
if not source_id or source_sequence < 0:
|
|
raise ValueError("live perception source identity is invalid")
|
|
if captured_at_epoch_ns < 0 or received_monotonic_ns < 0:
|
|
raise ValueError("live perception timestamps must be non-negative")
|
|
with self._condition:
|
|
if self._closed or not self._active:
|
|
return False
|
|
return self._publish_locked(
|
|
modality=modality,
|
|
source_id=source_id,
|
|
source_sequence=source_sequence,
|
|
captured_at_epoch_ns=captured_at_epoch_ns,
|
|
received_monotonic_ns=received_monotonic_ns,
|
|
payload=payload,
|
|
)
|
|
|
|
def open_consumer(self, consumer_id: str) -> None:
|
|
if not consumer_id or len(consumer_id) > 128:
|
|
raise ValueError("live perception consumer id is invalid")
|
|
with self._condition:
|
|
if self._closed:
|
|
raise RuntimeError("live perception ingress is closed")
|
|
if self._consumer_id is not None and self._consumer_id != consumer_id:
|
|
raise RuntimeError("live perception ingress already has a consumer")
|
|
self._consumer_id = consumer_id
|
|
|
|
def close_consumer(self, consumer_id: str) -> None:
|
|
with self._condition:
|
|
if self._consumer_id == consumer_id:
|
|
self._consumer_id = None
|
|
self._condition.notify_all()
|
|
|
|
def admit_result(
|
|
self,
|
|
*,
|
|
session_id: str,
|
|
session_generation: int,
|
|
receiver: Callable[[], bool],
|
|
) -> bool:
|
|
"""Atomically reject a late worker result before viewer admission."""
|
|
|
|
with self._condition:
|
|
if (
|
|
self._closed
|
|
or not self._active
|
|
or self._session_id != session_id
|
|
or self._session_generation != session_generation
|
|
):
|
|
self._results_rejected_stale += 1
|
|
return False
|
|
accepted = receiver()
|
|
if accepted:
|
|
self._results_accepted += 1
|
|
else:
|
|
self._results_rejected_receiver += 1
|
|
return accepted
|
|
|
|
def take_next(
|
|
self,
|
|
consumer_id: str,
|
|
*,
|
|
timeout: float | None = None,
|
|
) -> LiveIngressEvent | None:
|
|
with self._condition:
|
|
if self._consumer_id != consumer_id:
|
|
raise RuntimeError("live perception consumer lease is not active")
|
|
ready = self._condition.wait_for(
|
|
lambda: any(queue.items for queue in self._queues.values()) or self._closed,
|
|
timeout=timeout,
|
|
)
|
|
if not ready:
|
|
return None
|
|
candidates = [
|
|
(queue.items[0].ingress_sequence, modality, queue)
|
|
for modality, queue in self._queues.items()
|
|
if queue.items
|
|
]
|
|
if not candidates:
|
|
return None
|
|
_, _, selected = min(candidates, key=lambda item: item[0])
|
|
selected.consumed += 1
|
|
return selected.items.popleft()
|
|
|
|
def close(self) -> None:
|
|
with self._condition:
|
|
self._closed = True
|
|
self._active = False
|
|
self._condition.notify_all()
|
|
|
|
def snapshot(self) -> dict[str, Any]:
|
|
with self._condition:
|
|
return {
|
|
"schema_version": LIVE_INGRESS_SCHEMA,
|
|
"mode": "shadow-diagnostic-only",
|
|
"active": self._active,
|
|
"session_id": self._session_id,
|
|
"session_generation": self._session_generation,
|
|
"consumer_connected": self._consumer_id is not None,
|
|
"results_accepted": self._results_accepted,
|
|
"results_rejected_stale": self._results_rejected_stale,
|
|
"results_rejected_receiver": self._results_rejected_receiver,
|
|
"commands_enabled": False,
|
|
"navigation_or_safety_accepted": False,
|
|
"closed": self._closed,
|
|
"queues": {
|
|
modality: {
|
|
"capacity": queue.capacity,
|
|
"depth": len(queue.items),
|
|
"maximum_depth": queue.maximum_depth,
|
|
"published": queue.published,
|
|
"consumed": queue.consumed,
|
|
"dropped_overflow": queue.dropped_overflow,
|
|
"rejected_oversize": queue.rejected_oversize,
|
|
}
|
|
for modality, queue in self._queues.items()
|
|
},
|
|
}
|
|
|
|
def _publish_locked(
|
|
self,
|
|
*,
|
|
modality: LiveIngressModality,
|
|
source_id: str,
|
|
source_sequence: int,
|
|
captured_at_epoch_ns: int,
|
|
received_monotonic_ns: int,
|
|
payload: bytes,
|
|
) -> bool:
|
|
queue = self._queues[modality]
|
|
if len(payload) > self._MAX_PAYLOAD_BYTES[modality]:
|
|
queue.rejected_oversize += 1
|
|
return False
|
|
session_id = self._session_id
|
|
if session_id is None:
|
|
return False
|
|
self._ingress_sequence += 1
|
|
event = LiveIngressEvent(
|
|
ingress_sequence=self._ingress_sequence,
|
|
session_id=session_id,
|
|
session_generation=self._session_generation,
|
|
modality=modality,
|
|
source_id=source_id,
|
|
source_sequence=source_sequence,
|
|
captured_at_epoch_ns=captured_at_epoch_ns,
|
|
received_monotonic_ns=received_monotonic_ns,
|
|
payload=bytes(payload),
|
|
)
|
|
if len(queue.items) == queue.capacity:
|
|
queue.items.popleft()
|
|
queue.dropped_overflow += 1
|
|
queue.items.append(event)
|
|
queue.published += 1
|
|
queue.maximum_depth = max(queue.maximum_depth, len(queue.items))
|
|
self._condition.notify_all()
|
|
return True
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class QueueSnapshot:
|
|
capacity: int
|
|
depth: int
|
|
maximum_depth: int
|
|
published: int
|
|
consumed: int
|
|
dropped_overflow: int
|
|
dropped_superseded: int
|
|
closed: bool
|
|
|
|
@property
|
|
def dropped_total(self) -> int:
|
|
return self.dropped_overflow + self.dropped_superseded
|
|
|
|
|
|
class LatestWinsQueue[T]:
|
|
"""A bounded derived-data queue that never lets old preview work accumulate."""
|
|
|
|
def __init__(self, capacity: int) -> None:
|
|
if capacity < 1:
|
|
raise ValueError("latest-wins queue capacity must be positive")
|
|
self._capacity = capacity
|
|
self._items: deque[T] = deque()
|
|
self._condition = threading.Condition()
|
|
self._maximum_depth = 0
|
|
self._published = 0
|
|
self._consumed = 0
|
|
self._dropped_overflow = 0
|
|
self._dropped_superseded = 0
|
|
self._closed = False
|
|
|
|
def publish(self, item: T) -> T | None:
|
|
"""Publish fresh work and return the evicted oldest item, if any."""
|
|
|
|
with self._condition:
|
|
if self._closed:
|
|
raise RuntimeError("cannot publish to a closed latest-wins queue")
|
|
dropped: T | None = None
|
|
if len(self._items) == self._capacity:
|
|
dropped = self._items.popleft()
|
|
self._dropped_overflow += 1
|
|
self._items.append(item)
|
|
self._published += 1
|
|
self._maximum_depth = max(self._maximum_depth, len(self._items))
|
|
self._condition.notify()
|
|
return dropped
|
|
|
|
def take_next(self, timeout: float | None = None) -> T | None:
|
|
with self._condition:
|
|
ready = self._condition.wait_for(
|
|
lambda: bool(self._items) or self._closed,
|
|
timeout=timeout,
|
|
)
|
|
if not ready or not self._items:
|
|
return None
|
|
item = self._items.popleft()
|
|
self._consumed += 1
|
|
return item
|
|
|
|
def close(self) -> None:
|
|
with self._condition:
|
|
self._closed = True
|
|
self._condition.notify_all()
|
|
|
|
def snapshot(self) -> QueueSnapshot:
|
|
with self._condition:
|
|
return QueueSnapshot(
|
|
capacity=self._capacity,
|
|
depth=len(self._items),
|
|
maximum_depth=self._maximum_depth,
|
|
published=self._published,
|
|
consumed=self._consumed,
|
|
dropped_overflow=self._dropped_overflow,
|
|
dropped_superseded=self._dropped_superseded,
|
|
closed=self._closed,
|
|
)
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class LiveSensorBinding:
|
|
"""One bounded camera→LiDAR→pose match for live diagnostic fusion."""
|
|
|
|
state: Literal[
|
|
"fused-ready",
|
|
"lidar-unavailable",
|
|
"lidar-camera-delta-exceeded",
|
|
"pose-unavailable",
|
|
"pose-point-delta-exceeded",
|
|
]
|
|
point_cloud: DecodedPointCloudView | None
|
|
pose: DecodedPoseView | None
|
|
lidar_camera_delta_ms: float | None
|
|
pose_point_delta_ms: float | None
|
|
|
|
|
|
class LiveSensorSynchronizer:
|
|
"""Keep a small arrival-time window and bind sensors without back-pressure.
|
|
|
|
The synchronizer deliberately uses the already-recorded host arrival clock
|
|
carried by the shadow wire contract. It does not claim hardware-clock
|
|
synchronization. A short wait budget lets a LiDAR or pose event that is
|
|
already in flight reach the receiver while keeping the detector/world-state
|
|
latency bounded.
|
|
"""
|
|
|
|
def __init__(
|
|
self,
|
|
*,
|
|
maximum_lidar_camera_delta_ms: float,
|
|
maximum_pose_point_delta_ms: float,
|
|
capacity_per_modality: int = 32,
|
|
retention_seconds: float = 3.0,
|
|
) -> None:
|
|
if (
|
|
maximum_lidar_camera_delta_ms <= 0
|
|
or maximum_pose_point_delta_ms <= 0
|
|
or capacity_per_modality < 2
|
|
or retention_seconds <= 0
|
|
):
|
|
raise ValueError("live sensor synchronizer bounds are invalid")
|
|
self._maximum_lidar_camera_delta_ns = round(maximum_lidar_camera_delta_ms * 1_000_000)
|
|
self._maximum_pose_point_delta_ns = round(maximum_pose_point_delta_ms * 1_000_000)
|
|
self._capacity = capacity_per_modality
|
|
self._retention_ns = round(retention_seconds * 1_000_000_000)
|
|
self._condition = threading.Condition()
|
|
self._points: deque[DecodedPointCloudView] = deque()
|
|
self._poses: deque[DecodedPoseView] = deque()
|
|
self._published_points = 0
|
|
self._published_poses = 0
|
|
self._evicted_points = 0
|
|
self._evicted_poses = 0
|
|
self._maximum_point_depth = 0
|
|
self._maximum_pose_depth = 0
|
|
|
|
def publish_point_cloud(self, value: DecodedPointCloudView) -> None:
|
|
with self._condition:
|
|
self._points.append(value)
|
|
self._published_points += 1
|
|
self._evicted_points += self._prune(self._points)
|
|
self._maximum_point_depth = max(self._maximum_point_depth, len(self._points))
|
|
self._condition.notify_all()
|
|
|
|
def publish_pose(self, value: DecodedPoseView) -> None:
|
|
with self._condition:
|
|
self._poses.append(value)
|
|
self._published_poses += 1
|
|
self._evicted_poses += self._prune(self._poses)
|
|
self._maximum_pose_depth = max(self._maximum_pose_depth, len(self._poses))
|
|
self._condition.notify_all()
|
|
|
|
def bind_camera(
|
|
self,
|
|
captured_at_epoch_ns: int,
|
|
*,
|
|
wait_seconds: float = 0.0,
|
|
) -> LiveSensorBinding:
|
|
if captured_at_epoch_ns < 0 or wait_seconds < 0 or not math.isfinite(wait_seconds):
|
|
raise ValueError("camera synchronization input is invalid")
|
|
deadline = time.monotonic() + wait_seconds
|
|
with self._condition:
|
|
while True:
|
|
binding = self._binding(captured_at_epoch_ns)
|
|
if binding.state == "fused-ready" or wait_seconds == 0:
|
|
return binding
|
|
remaining = deadline - time.monotonic()
|
|
if remaining <= 0:
|
|
return binding
|
|
self._condition.wait(timeout=remaining)
|
|
|
|
def snapshot(self) -> dict[str, int | float]:
|
|
with self._condition:
|
|
return {
|
|
"capacity_per_modality": self._capacity,
|
|
"retention_seconds": self._retention_ns / 1_000_000_000,
|
|
"point_depth": len(self._points),
|
|
"pose_depth": len(self._poses),
|
|
"maximum_point_depth": self._maximum_point_depth,
|
|
"maximum_pose_depth": self._maximum_pose_depth,
|
|
"published_points": self._published_points,
|
|
"published_poses": self._published_poses,
|
|
"evicted_points": self._evicted_points,
|
|
"evicted_poses": self._evicted_poses,
|
|
}
|
|
|
|
def _binding(self, captured_at_epoch_ns: int) -> LiveSensorBinding:
|
|
if not self._points:
|
|
return LiveSensorBinding("lidar-unavailable", None, None, None, None)
|
|
point = min(
|
|
self._points,
|
|
key=lambda value: abs(value.context.captured_at_epoch_ns - captured_at_epoch_ns),
|
|
)
|
|
lidar_delta_ns = point.context.captured_at_epoch_ns - captured_at_epoch_ns
|
|
lidar_delta_ms = lidar_delta_ns / 1_000_000
|
|
if abs(lidar_delta_ns) > self._maximum_lidar_camera_delta_ns:
|
|
return LiveSensorBinding(
|
|
"lidar-camera-delta-exceeded",
|
|
point,
|
|
None,
|
|
lidar_delta_ms,
|
|
None,
|
|
)
|
|
if not self._poses:
|
|
return LiveSensorBinding(
|
|
"pose-unavailable",
|
|
point,
|
|
None,
|
|
lidar_delta_ms,
|
|
None,
|
|
)
|
|
pose = min(
|
|
self._poses,
|
|
key=lambda value: abs(
|
|
value.context.captured_at_epoch_ns - point.context.captured_at_epoch_ns
|
|
),
|
|
)
|
|
pose_delta_ns = pose.context.captured_at_epoch_ns - point.context.captured_at_epoch_ns
|
|
pose_delta_ms = pose_delta_ns / 1_000_000
|
|
if abs(pose_delta_ns) > self._maximum_pose_point_delta_ns:
|
|
return LiveSensorBinding(
|
|
"pose-point-delta-exceeded",
|
|
point,
|
|
pose,
|
|
lidar_delta_ms,
|
|
pose_delta_ms,
|
|
)
|
|
return LiveSensorBinding(
|
|
"fused-ready",
|
|
point,
|
|
pose,
|
|
lidar_delta_ms,
|
|
pose_delta_ms,
|
|
)
|
|
|
|
def _prune(self, values: deque[DecodedPointCloudView] | deque[DecodedPoseView]) -> int:
|
|
removed = 0
|
|
newest = values[-1].context.captured_at_epoch_ns
|
|
oldest_allowed = newest - self._retention_ns
|
|
while values and (
|
|
len(values) > self._capacity or values[0].context.captured_at_epoch_ns < oldest_allowed
|
|
):
|
|
values.popleft()
|
|
removed += 1
|
|
return removed
|
|
|
|
|
|
def classify_health(
|
|
*,
|
|
source_available: bool,
|
|
fusion_state: str,
|
|
result_age_ms: float,
|
|
stale_after_ms: float,
|
|
unavailable_after_ms: float,
|
|
) -> tuple[HealthState, tuple[str, ...]]:
|
|
"""Classify freshness separately from whether depth was available."""
|
|
|
|
if stale_after_ms <= 0 or unavailable_after_ms <= stale_after_ms:
|
|
raise ValueError("health thresholds are invalid")
|
|
if not source_available or result_age_ms >= unavailable_after_ms:
|
|
return "unavailable", ("source-unavailable",)
|
|
if result_age_ms >= stale_after_ms:
|
|
return "stale", ("result-age-exceeded",)
|
|
if fusion_state != "fused":
|
|
return "degraded", (fusion_state,)
|
|
return "healthy", ()
|
|
|
|
|
|
class WorldStateProjector:
|
|
"""Project accepted E6 observations into a control-facing, timestamped state."""
|
|
|
|
def __init__(self, *, velocity_history_limit_s: float = 1.0) -> None:
|
|
if velocity_history_limit_s <= 0:
|
|
raise ValueError("velocity history limit must be positive")
|
|
self._velocity_history_limit_s = velocity_history_limit_s
|
|
self._track_history: dict[int, deque[tuple[float, tuple[float, float, float]]]] = {}
|
|
|
|
def project(
|
|
self,
|
|
*,
|
|
frame: Mapping[str, Any],
|
|
lidar_positions: Mapping[int, Sequence[float]],
|
|
clearance: Mapping[str, Any],
|
|
) -> dict[str, Any]:
|
|
session_seconds = float(frame["session_seconds"])
|
|
objects: list[dict[str, Any]] = []
|
|
for raw in frame.get("objects", []):
|
|
if not str(raw.get("cuboid_status", "")).startswith("accepted-"):
|
|
continue
|
|
track_id = int(raw["track_id"])
|
|
center_map = _vector3(raw["cuboid_center_map"], "cuboid center")
|
|
half_size = _vector3(raw["cuboid_half_size"], "cuboid half-size")
|
|
quaternion = _vector4(raw["cuboid_quaternion_xyzw"], "cuboid quaternion")
|
|
velocity, velocity_status, velocity_residual = self._velocity(
|
|
track_id, session_seconds, center_map
|
|
)
|
|
speed = None
|
|
if velocity is not None:
|
|
speed = math.sqrt(sum(component * component for component in velocity))
|
|
lidar = lidar_positions.get(track_id)
|
|
objects.append(
|
|
{
|
|
"track_id": track_id,
|
|
"class": str(raw["association_group"]),
|
|
"detector_label": str(raw["label"]),
|
|
"confidence": float(raw["score"]),
|
|
"position_map_m": list(center_map),
|
|
"position_lidar_m": (
|
|
None if lidar is None else [float(value) for value in lidar]
|
|
),
|
|
"orientation_map_xyzw": list(quaternion),
|
|
"size_m": [2.0 * value for value in half_size],
|
|
"range_m": float(raw["distance_smoothed_m"]),
|
|
"velocity_map_mps": None if velocity is None else list(velocity),
|
|
"speed_mps": speed,
|
|
"velocity_status": velocity_status,
|
|
"velocity_residual_m": velocity_residual,
|
|
"support_points": int(raw["clustered_points"]),
|
|
"geometry": "point-supported-visible-surface-envelope",
|
|
}
|
|
)
|
|
self._prune(session_seconds)
|
|
return {
|
|
"schema_version": WORLD_STATE_SCHEMA,
|
|
"frame_index": int(frame["frame_index"]),
|
|
"source_frame_index": int(frame["source_frame_index"]),
|
|
"session_seconds": session_seconds,
|
|
"coordinate_frames": {
|
|
"world": "k1-map",
|
|
"sensor_relative": "k1-lidar",
|
|
"vehicle_body": "unavailable-no-rig-to-vehicle-transform",
|
|
},
|
|
"fusion_state": str(frame["state"]),
|
|
"objects": objects,
|
|
"object_count": len(objects),
|
|
"clearance": dict(clearance),
|
|
}
|
|
|
|
def _velocity(
|
|
self,
|
|
track_id: int,
|
|
session_seconds: float,
|
|
center_map: tuple[float, float, float],
|
|
) -> tuple[tuple[float, float, float] | None, str, float | None]:
|
|
history = self._track_history.setdefault(track_id, deque(maxlen=32))
|
|
if history and session_seconds <= history[-1][0]:
|
|
return None, "unavailable-nonmonotonic-time", None
|
|
history.append((session_seconds, center_map))
|
|
oldest = session_seconds - self._velocity_history_limit_s
|
|
while history and history[0][0] < oldest:
|
|
history.popleft()
|
|
if len(history) < 4 or history[-1][0] - history[0][0] < 0.4:
|
|
return None, "unavailable-insufficient-history", None
|
|
slopes: list[tuple[float, float, float]] = []
|
|
values = list(history)
|
|
for left, (left_time, left_center) in enumerate(values):
|
|
for right_time, right_center in values[left + 1 :]:
|
|
delta = right_time - left_time
|
|
if delta < 0.2:
|
|
continue
|
|
slopes.append(
|
|
(
|
|
(right_center[0] - left_center[0]) / delta,
|
|
(right_center[1] - left_center[1]) / delta,
|
|
(right_center[2] - left_center[2]) / delta,
|
|
)
|
|
)
|
|
if not slopes:
|
|
return None, "unavailable-insufficient-baseline", None
|
|
velocity = (
|
|
statistics.median(item[0] for item in slopes),
|
|
statistics.median(item[1] for item in slopes),
|
|
statistics.median(item[2] for item in slopes),
|
|
)
|
|
speed = math.sqrt(sum(component * component for component in velocity))
|
|
latest_time, latest_center = values[-1]
|
|
residuals = []
|
|
for observed_time, observed_center in values:
|
|
predicted = tuple(
|
|
latest - component * (latest_time - observed_time)
|
|
for latest, component in zip(latest_center, velocity, strict=True)
|
|
)
|
|
residuals.append(
|
|
math.sqrt(
|
|
sum(
|
|
(observed - expected) ** 2
|
|
for observed, expected in zip(observed_center, predicted, strict=True)
|
|
)
|
|
)
|
|
)
|
|
residual = statistics.median(residuals)
|
|
if speed > 20.0:
|
|
return None, "rejected-speed-bound", residual
|
|
if residual > 0.75:
|
|
return None, "rejected-position-residual", residual
|
|
return velocity, "diagnostic-robust-history", residual
|
|
|
|
def _prune(self, session_seconds: float) -> None:
|
|
oldest = session_seconds - self._velocity_history_limit_s
|
|
expired = [
|
|
track_id
|
|
for track_id, history in self._track_history.items()
|
|
if not history or history[-1][0] < oldest
|
|
]
|
|
for track_id in expired:
|
|
del self._track_history[track_id]
|
|
|
|
|
|
def wait_until(deadline: float) -> float:
|
|
"""Wait for a replay deadline and return non-negative scheduling lag seconds."""
|
|
|
|
remaining = deadline - time.perf_counter()
|
|
if remaining > 0:
|
|
time.sleep(remaining)
|
|
return max(0.0, time.perf_counter() - deadline)
|
|
|
|
|
|
def _vector3(value: Sequence[Any], label: str) -> tuple[float, float, float]:
|
|
if len(value) != 3:
|
|
raise ValueError(f"{label} must contain three values")
|
|
return float(value[0]), float(value[1]), float(value[2])
|
|
|
|
|
|
def _vector4(value: Sequence[Any], label: str) -> tuple[float, float, float, float]:
|
|
if len(value) != 4:
|
|
raise ValueError(f"{label} must contain four values")
|
|
return float(value[0]), float(value[1]), float(value[2]), float(value[3])
|