417 lines
13 KiB
Python
417 lines
13 KiB
Python
"""
|
|
Transport packets for synchronous BASE + ROI JPEG composite frames.
|
|
|
|
The wire header has a fixed 32-byte, padding-free network-byte-order layout::
|
|
|
|
!4sBBHIHHHHIII
|
|
|
|
Offset Size Field
|
|
0 4 magic (b"SRV1")
|
|
4 1 version
|
|
5 1 object_type (1=BASE, 2=ROI)
|
|
6 2 header_size
|
|
8 4 composite_frame_id
|
|
12 2 fragment_index
|
|
14 2 fragment_count
|
|
16 2 payload_length
|
|
18 2 flags (reserved, must be zero)
|
|
20 4 jpeg_size
|
|
24 4 object_crc32
|
|
28 4 packet_crc32
|
|
|
|
packet_crc32 is calculated over the complete header with packet_crc32 set to
|
|
zero, followed by the packet payload. object_crc32 is calculated over the
|
|
complete JPEG before fragmentation and checked after reassembly.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
from dataclasses import dataclass, field
|
|
from enum import IntEnum
|
|
import struct
|
|
import zlib
|
|
|
|
|
|
MAGIC = b"SRV1"
|
|
VERSION = 1
|
|
HEADER_FORMAT = "!4sBBHIHHHHIII"
|
|
HEADER_SIZE = struct.calcsize(HEADER_FORMAT)
|
|
MAX_PAYLOAD_LENGTH = 0xFFFF
|
|
MAX_FRAGMENT_COUNT = 0xFFFF
|
|
MAX_JPEG_SIZE = 0xFFFFFFFF
|
|
|
|
|
|
class ObjectType(IntEnum):
|
|
"""JPEG object carried by a packet."""
|
|
|
|
BASE = 1
|
|
ROI = 2
|
|
|
|
|
|
class VideoPacketError(ValueError):
|
|
"""Base class for video transport validation errors."""
|
|
|
|
|
|
class PacketCRCError(VideoPacketError):
|
|
"""The packet header or payload failed its CRC32 check."""
|
|
|
|
|
|
class ObjectCRCError(VideoPacketError):
|
|
"""A completely reassembled JPEG failed its CRC32 check."""
|
|
|
|
|
|
class ObjectConsistencyError(VideoPacketError):
|
|
"""Fragments for one object contain inconsistent metadata or data."""
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class VideoPacket:
|
|
"""Decoded and CRC-checked transport packet."""
|
|
|
|
composite_frame_id: int
|
|
object_type: ObjectType
|
|
fragment_index: int
|
|
fragment_count: int
|
|
jpeg_size: int
|
|
object_crc32: int
|
|
payload: bytes
|
|
packet_crc32: int = 0
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class CompositeFrame:
|
|
"""Atomically reassembled BASE and ROI JPEGs for one source update."""
|
|
|
|
composite_frame_id: int
|
|
base_jpeg: bytes
|
|
roi_jpeg: bytes
|
|
|
|
|
|
def crc32(data: bytes) -> int:
|
|
"""Return an unsigned IEEE CRC32."""
|
|
|
|
return zlib.crc32(data) & 0xFFFFFFFF
|
|
|
|
|
|
def _validate_packet_fields(packet: VideoPacket) -> None:
|
|
if not 0 <= packet.composite_frame_id <= 0xFFFFFFFF:
|
|
raise VideoPacketError("composite_frame_id is outside uint32")
|
|
if packet.object_type not in (ObjectType.BASE, ObjectType.ROI):
|
|
raise VideoPacketError("unsupported object_type")
|
|
if not 1 <= packet.fragment_count <= MAX_FRAGMENT_COUNT:
|
|
raise VideoPacketError("fragment_count is outside uint16")
|
|
if not 0 <= packet.fragment_index < packet.fragment_count:
|
|
raise VideoPacketError("fragment_index is outside fragment_count")
|
|
if not 1 <= len(packet.payload) <= MAX_PAYLOAD_LENGTH:
|
|
raise VideoPacketError("payload length is outside uint16")
|
|
if not 1 <= packet.jpeg_size <= MAX_JPEG_SIZE:
|
|
raise VideoPacketError("jpeg_size is outside uint32")
|
|
if len(packet.payload) > packet.jpeg_size:
|
|
raise VideoPacketError("payload is larger than the JPEG object")
|
|
if not 0 <= packet.object_crc32 <= 0xFFFFFFFF:
|
|
raise VideoPacketError("object_crc32 is outside uint32")
|
|
|
|
|
|
def _pack_header(packet: VideoPacket, packet_crc32: int) -> bytes:
|
|
return struct.pack(
|
|
HEADER_FORMAT,
|
|
MAGIC,
|
|
VERSION,
|
|
int(packet.object_type),
|
|
HEADER_SIZE,
|
|
packet.composite_frame_id,
|
|
packet.fragment_index,
|
|
packet.fragment_count,
|
|
len(packet.payload),
|
|
0,
|
|
packet.jpeg_size,
|
|
packet.object_crc32,
|
|
packet_crc32,
|
|
)
|
|
|
|
|
|
def encode_packet(packet: VideoPacket) -> bytes:
|
|
"""Serialize a packet and calculate its packet CRC32."""
|
|
|
|
if not isinstance(packet, VideoPacket):
|
|
raise TypeError("packet must be VideoPacket")
|
|
_validate_packet_fields(packet)
|
|
header_with_zero_crc = _pack_header(packet, 0)
|
|
packet_crc32 = crc32(header_with_zero_crc + packet.payload)
|
|
return _pack_header(packet, packet_crc32) + packet.payload
|
|
|
|
|
|
def decode_packet(wire_packet: bytes) -> VideoPacket:
|
|
"""Deserialize a packet and validate its fixed header and packet CRC32."""
|
|
|
|
if not isinstance(wire_packet, (bytes, bytearray)):
|
|
raise TypeError("wire_packet must be bytes or bytearray")
|
|
wire_packet = bytes(wire_packet)
|
|
if len(wire_packet) < HEADER_SIZE + 1:
|
|
raise VideoPacketError("packet is shorter than header plus payload")
|
|
|
|
(
|
|
magic,
|
|
version,
|
|
object_type_value,
|
|
header_size,
|
|
composite_frame_id,
|
|
fragment_index,
|
|
fragment_count,
|
|
payload_length,
|
|
flags,
|
|
jpeg_size,
|
|
object_crc32,
|
|
received_packet_crc32,
|
|
) = struct.unpack(HEADER_FORMAT, wire_packet[:HEADER_SIZE])
|
|
|
|
if magic != MAGIC:
|
|
raise VideoPacketError("invalid packet magic")
|
|
if version != VERSION:
|
|
raise VideoPacketError("unsupported packet version")
|
|
if header_size != HEADER_SIZE:
|
|
raise VideoPacketError("invalid fixed header size")
|
|
if flags != 0:
|
|
raise VideoPacketError("reserved flags must be zero")
|
|
try:
|
|
object_type = ObjectType(object_type_value)
|
|
except ValueError as error:
|
|
raise VideoPacketError("unsupported object_type") from error
|
|
|
|
if len(wire_packet) != HEADER_SIZE + payload_length:
|
|
raise VideoPacketError("packet length does not match payload_length")
|
|
|
|
packet = VideoPacket(
|
|
composite_frame_id=composite_frame_id,
|
|
object_type=object_type,
|
|
fragment_index=fragment_index,
|
|
fragment_count=fragment_count,
|
|
jpeg_size=jpeg_size,
|
|
object_crc32=object_crc32,
|
|
payload=wire_packet[HEADER_SIZE:],
|
|
packet_crc32=received_packet_crc32,
|
|
)
|
|
_validate_packet_fields(packet)
|
|
calculated_packet_crc32 = crc32(
|
|
_pack_header(packet, 0) + packet.payload
|
|
)
|
|
if calculated_packet_crc32 != received_packet_crc32:
|
|
raise PacketCRCError(
|
|
"packet CRC mismatch: "
|
|
f"received 0x{received_packet_crc32:08X}, "
|
|
f"calculated 0x{calculated_packet_crc32:08X}"
|
|
)
|
|
return packet
|
|
|
|
|
|
def packetize_jpeg(
|
|
jpeg: bytes,
|
|
composite_frame_id: int,
|
|
object_type: ObjectType,
|
|
max_payload_length: int,
|
|
) -> list[bytes]:
|
|
"""Split one non-empty JPEG object into serialized packets."""
|
|
|
|
if not isinstance(jpeg, (bytes, bytearray)):
|
|
raise TypeError("jpeg must be bytes or bytearray")
|
|
jpeg = bytes(jpeg)
|
|
if not jpeg:
|
|
raise ValueError("JPEG object must not be empty")
|
|
if len(jpeg) > MAX_JPEG_SIZE:
|
|
raise ValueError("JPEG object is larger than uint32")
|
|
if not 1 <= max_payload_length <= MAX_PAYLOAD_LENGTH:
|
|
raise ValueError("max_payload_length is outside uint16")
|
|
if not 0 <= composite_frame_id <= 0xFFFFFFFF:
|
|
raise ValueError("composite_frame_id is outside uint32")
|
|
try:
|
|
object_type = ObjectType(object_type)
|
|
except ValueError as error:
|
|
raise ValueError("unsupported object_type") from error
|
|
|
|
fragment_count = (
|
|
len(jpeg) + max_payload_length - 1
|
|
) // max_payload_length
|
|
if fragment_count > MAX_FRAGMENT_COUNT:
|
|
raise ValueError("JPEG needs more than 65535 fragments")
|
|
|
|
object_crc32 = crc32(jpeg)
|
|
packets = []
|
|
for fragment_index in range(fragment_count):
|
|
start = fragment_index * max_payload_length
|
|
payload = jpeg[start:start + max_payload_length]
|
|
packets.append(
|
|
encode_packet(
|
|
VideoPacket(
|
|
composite_frame_id=composite_frame_id,
|
|
object_type=object_type,
|
|
fragment_index=fragment_index,
|
|
fragment_count=fragment_count,
|
|
jpeg_size=len(jpeg),
|
|
object_crc32=object_crc32,
|
|
payload=payload,
|
|
)
|
|
)
|
|
)
|
|
return packets
|
|
|
|
|
|
@dataclass
|
|
class _ObjectAssembly:
|
|
composite_frame_id: int
|
|
object_type: ObjectType
|
|
fragment_count: int
|
|
jpeg_size: int
|
|
object_crc32: int
|
|
fragments: dict[int, bytes] = field(default_factory=dict)
|
|
|
|
@classmethod
|
|
def from_packet(cls, packet: VideoPacket) -> "_ObjectAssembly":
|
|
return cls(
|
|
composite_frame_id=packet.composite_frame_id,
|
|
object_type=packet.object_type,
|
|
fragment_count=packet.fragment_count,
|
|
jpeg_size=packet.jpeg_size,
|
|
object_crc32=packet.object_crc32,
|
|
)
|
|
|
|
def add(self, packet: VideoPacket) -> bool:
|
|
metadata = (
|
|
packet.composite_frame_id,
|
|
packet.object_type,
|
|
packet.fragment_count,
|
|
packet.jpeg_size,
|
|
packet.object_crc32,
|
|
)
|
|
expected = (
|
|
self.composite_frame_id,
|
|
self.object_type,
|
|
self.fragment_count,
|
|
self.jpeg_size,
|
|
self.object_crc32,
|
|
)
|
|
if metadata != expected:
|
|
raise ObjectConsistencyError(
|
|
"fragment metadata conflicts with object assembly"
|
|
)
|
|
|
|
existing = self.fragments.get(packet.fragment_index)
|
|
if existing is not None:
|
|
if existing != packet.payload:
|
|
raise ObjectConsistencyError(
|
|
"different payload for an existing fragment index"
|
|
)
|
|
return False
|
|
self.fragments[packet.fragment_index] = packet.payload
|
|
return True
|
|
|
|
def missing_fragments(self) -> tuple[int, ...]:
|
|
return tuple(
|
|
index
|
|
for index in range(self.fragment_count)
|
|
if index not in self.fragments
|
|
)
|
|
|
|
def assemble_if_complete(self) -> bytes | None:
|
|
if self.missing_fragments():
|
|
return None
|
|
jpeg = b"".join(
|
|
self.fragments[index]
|
|
for index in range(self.fragment_count)
|
|
)
|
|
if len(jpeg) != self.jpeg_size:
|
|
raise ObjectConsistencyError(
|
|
"reassembled JPEG length does not match jpeg_size"
|
|
)
|
|
calculated_object_crc32 = crc32(jpeg)
|
|
if calculated_object_crc32 != self.object_crc32:
|
|
raise ObjectCRCError(
|
|
"object CRC mismatch: "
|
|
f"received 0x{self.object_crc32:08X}, "
|
|
f"calculated 0x{calculated_object_crc32:08X}"
|
|
)
|
|
return jpeg
|
|
|
|
|
|
class CompositeReassembler:
|
|
"""Reassemble interleaved BASE/ROI packets and publish atomic frames."""
|
|
|
|
def __init__(self) -> None:
|
|
self._objects: dict[
|
|
tuple[int, ObjectType], _ObjectAssembly
|
|
] = {}
|
|
self._completed: dict[
|
|
int, dict[ObjectType, bytes]
|
|
] = {}
|
|
self._finalized_frame_ids: set[int] = set()
|
|
self.duplicate_packets = 0
|
|
|
|
def ingest(self, wire_packet: bytes) -> CompositeFrame | None:
|
|
"""Accept one packet and return a frame only when BASE and ROI exist."""
|
|
|
|
packet = decode_packet(wire_packet)
|
|
if packet.composite_frame_id in self._finalized_frame_ids:
|
|
self.duplicate_packets += 1
|
|
return None
|
|
|
|
key = (packet.composite_frame_id, packet.object_type)
|
|
assembly = self._objects.get(key)
|
|
if assembly is None:
|
|
assembly = _ObjectAssembly.from_packet(packet)
|
|
self._objects[key] = assembly
|
|
|
|
if not assembly.add(packet):
|
|
self.duplicate_packets += 1
|
|
return None
|
|
|
|
jpeg = assembly.assemble_if_complete()
|
|
if jpeg is None:
|
|
return None
|
|
|
|
frame_parts = self._completed.setdefault(
|
|
packet.composite_frame_id, {}
|
|
)
|
|
frame_parts[packet.object_type] = jpeg
|
|
if not all(
|
|
object_type in frame_parts
|
|
for object_type in (ObjectType.BASE, ObjectType.ROI)
|
|
):
|
|
return None
|
|
|
|
frame = CompositeFrame(
|
|
composite_frame_id=packet.composite_frame_id,
|
|
base_jpeg=frame_parts[ObjectType.BASE],
|
|
roi_jpeg=frame_parts[ObjectType.ROI],
|
|
)
|
|
self._finalized_frame_ids.add(packet.composite_frame_id)
|
|
self._completed.pop(packet.composite_frame_id, None)
|
|
for object_type in (ObjectType.BASE, ObjectType.ROI):
|
|
self._objects.pop(
|
|
(packet.composite_frame_id, object_type), None
|
|
)
|
|
return frame
|
|
|
|
def missing_fragments(
|
|
self,
|
|
composite_frame_id: int,
|
|
object_type: ObjectType,
|
|
) -> tuple[int, ...] | None:
|
|
"""Return missing indexes, or None if this object has not started."""
|
|
|
|
assembly = self._objects.get(
|
|
(composite_frame_id, ObjectType(object_type))
|
|
)
|
|
if assembly is None:
|
|
return None
|
|
return assembly.missing_fragments()
|
|
|
|
def object_is_complete(
|
|
self,
|
|
composite_frame_id: int,
|
|
object_type: ObjectType,
|
|
) -> bool:
|
|
"""Report whether one CRC-validated object awaits its counterpart."""
|
|
|
|
return ObjectType(object_type) in self._completed.get(
|
|
composite_frame_id, {}
|
|
)
|