From b9b067213f576a6676b2b165bf278bc3a4955e83 Mon Sep 17 00:00:00 2001 From: QiuSW <105186638@qq.com> Date: Fri, 28 Aug 2026 22:27:04 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E5=BB=BA=E7=AB=8B=20Brain=20=E5=8F=AF?= =?UTF-8?q?=E6=9B=BF=E6=8D=A2=E8=A7=86=E9=A2=91=E8=A7=A3=E7=A0=81=E6=B5=81?= =?UTF-8?q?=E6=B0=B4=E7=BA=BF=20(#13)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- Brain/src/yovision_brain/decode/__init__.py | 16 ++ Brain/src/yovision_brain/decode/models.py | 35 ++++ Brain/src/yovision_brain/decode/pipeline.py | 47 ++++++ Brain/src/yovision_brain/decode/raw_rgb.py | 45 ++++++ Brain/src/yovision_brain/decode/y4m.py | 167 ++++++++++++++++++++ Brain/tests/decode/test_pipeline.py | 104 ++++++++++++ Brain/tests/fixtures/decode/README.md | 5 + 7 files changed, 419 insertions(+) create mode 100644 Brain/src/yovision_brain/decode/__init__.py create mode 100644 Brain/src/yovision_brain/decode/models.py create mode 100644 Brain/src/yovision_brain/decode/pipeline.py create mode 100644 Brain/src/yovision_brain/decode/raw_rgb.py create mode 100644 Brain/src/yovision_brain/decode/y4m.py create mode 100644 Brain/tests/decode/test_pipeline.py create mode 100644 Brain/tests/fixtures/decode/README.md diff --git a/Brain/src/yovision_brain/decode/__init__.py b/Brain/src/yovision_brain/decode/__init__.py new file mode 100644 index 0000000..1958331 --- /dev/null +++ b/Brain/src/yovision_brain/decode/__init__.py @@ -0,0 +1,16 @@ +"""Replaceable Brain-internal video decode pipeline.""" + +from .models import DecodedFrame, DecoderBackend, DecoderError +from .pipeline import DecoderPipeline, decode_packets +from .raw_rgb import RawRGBDecoder +from .y4m import Y4MDecoder + +__all__ = [ + "DecodedFrame", + "DecoderBackend", + "DecoderError", + "DecoderPipeline", + "RawRGBDecoder", + "Y4MDecoder", + "decode_packets", +] diff --git a/Brain/src/yovision_brain/decode/models.py b/Brain/src/yovision_brain/decode/models.py new file mode 100644 index 0000000..609fdf8 --- /dev/null +++ b/Brain/src/yovision_brain/decode/models.py @@ -0,0 +1,35 @@ +"""Decode-layer ports and frame model.""" + +from __future__ import annotations + +from dataclasses import dataclass +from typing import Iterable, Iterator, Protocol + +from yovision_brain.input import CancellationToken, InputPacket + + +class DecoderError(RuntimeError): + """A safe and actionable decode failure.""" + + +@dataclass(frozen=True, slots=True) +class DecodedFrame: + sequence: int + timestamp_ns: int + logical_device_id: str + profile_id: str + width: int + height: int + pixel_format: str + payload: bytes + dimensions_changed: bool = False + + +class DecoderBackend(Protocol): + media_formats: frozenset[str] + + def decode( + self, + packets: Iterable[InputPacket], + cancellation: CancellationToken | None = None, + ) -> Iterator[DecodedFrame]: ... diff --git a/Brain/src/yovision_brain/decode/pipeline.py b/Brain/src/yovision_brain/decode/pipeline.py new file mode 100644 index 0000000..fa837cd --- /dev/null +++ b/Brain/src/yovision_brain/decode/pipeline.py @@ -0,0 +1,47 @@ +"""Decoder selection independent of concrete codec libraries.""" + +from __future__ import annotations + +from collections.abc import Iterable, Iterator +from itertools import chain + +from yovision_brain.input import CancellationToken, InputPacket + +from .models import DecodedFrame, DecoderBackend, DecoderError +from .raw_rgb import RawRGBDecoder +from .y4m import Y4MDecoder + + +class DecoderPipeline: + def __init__(self, backends: Iterable[DecoderBackend] | None = None) -> None: + selected = tuple(backends) if backends is not None else (RawRGBDecoder(), Y4MDecoder()) + self._backends: dict[str, DecoderBackend] = {} + for backend in selected: + for media_format in backend.media_formats: + if media_format in self._backends: + raise ValueError(f"duplicate decoder for media format {media_format!r}") + self._backends[media_format] = backend + + def decode( + self, + packets: Iterable[InputPacket], + cancellation: CancellationToken | None = None, + ) -> Iterator[DecodedFrame]: + iterator = iter(packets) + if cancellation is not None and cancellation.cancelled: + return + try: + first = next(iterator) + except StopIteration: + return + backend = self._backends.get(first.media_format) + if backend is None: + raise DecoderError(f"no decoder registered for media format {first.media_format!r}") + yield from backend.decode(chain((first,), iterator), cancellation) + + +def decode_packets( + packets: Iterable[InputPacket], + cancellation: CancellationToken | None = None, +) -> Iterator[DecodedFrame]: + return DecoderPipeline().decode(packets, cancellation) diff --git a/Brain/src/yovision_brain/decode/raw_rgb.py b/Brain/src/yovision_brain/decode/raw_rgb.py new file mode 100644 index 0000000..2dcf98f --- /dev/null +++ b/Brain/src/yovision_brain/decode/raw_rgb.py @@ -0,0 +1,45 @@ +"""Pass-through decoder for deterministic RGB24 synthetic frames.""" + +from __future__ import annotations + +from collections.abc import Iterable, Iterator + +from yovision_brain.input import CancellationToken, InputPacket + +from .models import DecodedFrame, DecoderError + + +class RawRGBDecoder: + media_formats = frozenset({"rgb24"}) + + def decode( + self, + packets: Iterable[InputPacket], + cancellation: CancellationToken | None = None, + ) -> Iterator[DecodedFrame]: + previous_dimensions: tuple[int, int] | None = None + for packet in packets: + if cancellation is not None and cancellation.cancelled: + return + if packet.media_format != "rgb24": + raise DecoderError(f"raw RGB decoder does not support {packet.media_format!r}") + expected = packet.width * packet.height * 3 + if len(packet.payload) != expected: + raise DecoderError( + f"RGB24 frame {packet.sequence} has {len(packet.payload)} bytes; expected {expected}" + ) + if packet.timestamp_ns is None: + raise DecoderError(f"RGB24 frame {packet.sequence} has no source timestamp") + dimensions = (packet.width, packet.height) + yield DecodedFrame( + sequence=packet.sequence, + timestamp_ns=packet.timestamp_ns, + logical_device_id=packet.logical_device_id, + profile_id=packet.profile_id, + width=packet.width, + height=packet.height, + pixel_format="rgb24", + payload=packet.payload, + dimensions_changed=previous_dimensions is not None and dimensions != previous_dimensions, + ) + previous_dimensions = dimensions diff --git a/Brain/src/yovision_brain/decode/y4m.py b/Brain/src/yovision_brain/decode/y4m.py new file mode 100644 index 0000000..4215547 --- /dev/null +++ b/Brain/src/yovision_brain/decode/y4m.py @@ -0,0 +1,167 @@ +"""Minimal streaming YUV4MPEG2 decoder for anonymous local fixtures. + +The backend intentionally supports only uncompressed C444 streams. Production +codecs and RTSP belong behind the same decoder port in later tasks. +""" + +from __future__ import annotations + +from collections.abc import Iterable, Iterator +from dataclasses import dataclass + +from yovision_brain.input import CancellationToken, InputPacket + +from .models import DecodedFrame, DecoderError + +_MAX_HEADER_BYTES = 4096 +_MAX_FRAME_BYTES = 256 * 1024 * 1024 + + +class _Cancelled(Exception): + pass + + +class _PacketReader: + def __init__( + self, + packets: Iterable[InputPacket], + cancellation: CancellationToken | None, + ) -> None: + self._packets = iter(packets) + self._cancellation = cancellation + self._buffer = bytearray() + self._ended = False + self.first_packet: InputPacket | None = None + + def _fill(self) -> bool: + if self._cancellation is not None and self._cancellation.cancelled: + raise _Cancelled + if self._ended: + return False + try: + packet = next(self._packets) + except StopIteration: + self._ended = True + return False + if packet.media_format != "container-bytes": + raise DecoderError(f"Y4M decoder does not support {packet.media_format!r}") + if self.first_packet is None: + self.first_packet = packet + else: + first = self.first_packet + if (packet.logical_device_id, packet.profile_id) != ( + first.logical_device_id, + first.profile_id, + ): + raise DecoderError("input identity changed inside one local video stream") + self._buffer.extend(packet.payload) + return True + + def line(self, *, allow_clean_eof: bool = False) -> bytes | None: + while True: + newline = self._buffer.find(b"\n") + if newline >= 0: + result = bytes(self._buffer[:newline]) + del self._buffer[: newline + 1] + return result + if len(self._buffer) > _MAX_HEADER_BYTES: + raise DecoderError("Y4M header exceeds the safe size limit") + if not self._fill(): + if not self._buffer and allow_clean_eof: + return None + raise DecoderError("truncated Y4M header") + + def exact(self, size: int) -> bytes: + while len(self._buffer) < size: + if not self._fill(): + raise DecoderError("truncated Y4M frame payload") + result = bytes(self._buffer[:size]) + del self._buffer[:size] + return result + + +@dataclass(frozen=True, slots=True) +class _Header: + width: int + height: int + fps_numerator: int + fps_denominator: int + + +def _positive_int(value: bytes, field: str) -> int: + try: + result = int(value) + except ValueError as exc: + raise DecoderError(f"invalid Y4M {field}") from exc + if result <= 0: + raise DecoderError(f"invalid Y4M {field}") + return result + + +def _parse_header(line: bytes) -> _Header: + parts = line.split() + if not parts or parts[0] != b"YUV4MPEG2": + raise DecoderError("unsupported local video format; expected YUV4MPEG2") + fields = {part[:1]: part[1:] for part in parts[1:] if len(part) > 1} + if fields.get(b"C", b"444") not in {b"444", b"444jpeg"}: + raise DecoderError("unsupported Y4M chroma; only C444 is supported") + width = _positive_int(fields.get(b"W", b""), "width") + height = _positive_int(fields.get(b"H", b""), "height") + fps_parts = fields.get(b"F", b"").split(b":", 1) + if len(fps_parts) != 2: + raise DecoderError("invalid Y4M frame rate") + header = _Header( + width=width, + height=height, + fps_numerator=_positive_int(fps_parts[0], "frame rate numerator"), + fps_denominator=_positive_int(fps_parts[1], "frame rate denominator"), + ) + if header.width * header.height * 3 > _MAX_FRAME_BYTES: + raise DecoderError("Y4M frame exceeds the safe size limit") + return header + + +class Y4MDecoder: + media_formats = frozenset({"container-bytes"}) + + def decode( + self, + packets: Iterable[InputPacket], + cancellation: CancellationToken | None = None, + ) -> Iterator[DecodedFrame]: + reader = _PacketReader(packets, cancellation) + try: + header_line = reader.line() + assert header_line is not None + header = _parse_header(header_line) + first = reader.first_packet + if first is None: + raise DecoderError("local video input is empty") + if (first.width, first.height) != (header.width, header.height): + raise DecoderError( + "Y4M dimensions do not match the configured input profile " + f"({header.width}x{header.height} != {first.width}x{first.height})" + ) + interval_ns = round(1_000_000_000 * header.fps_denominator / header.fps_numerator) + frame_size = header.width * header.height * 3 + sequence = 0 + while True: + frame_header = reader.line(allow_clean_eof=True) + if frame_header is None: + return + if frame_header != b"FRAME": + raise DecoderError(f"invalid Y4M frame header at frame {sequence}") + payload = reader.exact(frame_size) + yield DecodedFrame( + sequence=sequence, + timestamp_ns=sequence * interval_ns, + logical_device_id=first.logical_device_id, + profile_id=first.profile_id, + width=header.width, + height=header.height, + pixel_format="yuv444p", + payload=payload, + ) + sequence += 1 + except _Cancelled: + return diff --git a/Brain/tests/decode/test_pipeline.py b/Brain/tests/decode/test_pipeline.py new file mode 100644 index 0000000..beb2ae3 --- /dev/null +++ b/Brain/tests/decode/test_pipeline.py @@ -0,0 +1,104 @@ +from __future__ import annotations + +from pathlib import Path + +import pytest + +from yovision_brain.config import parse_input_config +from yovision_brain.decode import DecoderError, DecoderPipeline, decode_packets +from yovision_brain.input import CancellationToken, InputPacket, LocalFileInput, SyntheticInput + + +def synthetic_packets(): + config = parse_input_config( + { + "schema": "brain.internal.input/v1", + "logical_device_id": "synthetic-01", + "profile": {"id": "main", "width": 2, "height": 1, "fps": 5}, + "source": {"kind": "synthetic", "seed": 3, "frame_count": 2}, + } + ) + return SyntheticInput(config).packets() + + +def local_packets(path: Path, *, width: int = 2, height: int = 1, chunk_size: int = 5): + config = parse_input_config( + { + "schema": "brain.internal.input/v1", + "logical_device_id": "local-01", + "profile": {"id": "archive", "width": width, "height": height, "fps": 25}, + "source": {"kind": "local_file", "path": str(path), "chunk_size": chunk_size}, + } + ) + return LocalFileInput(config).packets() + + +def test_rgb24_pipeline_preserves_order_timestamps_and_metadata() -> None: + frames = list(decode_packets(synthetic_packets())) + assert [frame.sequence for frame in frames] == [0, 1] + assert [frame.timestamp_ns for frame in frames] == [0, 200_000_000] + assert all(frame.logical_device_id == "synthetic-01" for frame in frames) + assert all(frame.profile_id == "main" for frame in frames) + assert all((frame.width, frame.height, frame.pixel_format) == (2, 1, "rgb24") for frame in frames) + + +def test_rgb24_dimension_change_is_explicit() -> None: + packets = [ + InputPacket(0, 0, "camera", "main", 1, 1, "rgb24", b"abc"), + InputPacket(1, 1, "camera", "main", 2, 1, "rgb24", b"abcdef"), + ] + frames = list(decode_packets(packets)) + assert [frame.dimensions_changed for frame in frames] == [False, True] + + +def test_invalid_rgb_payload_and_unsupported_format_are_clear() -> None: + bad = [InputPacket(0, 0, "camera", "main", 2, 2, "rgb24", b"short")] + with pytest.raises(DecoderError, match="expected 12"): + list(decode_packets(bad)) + unknown = [InputPacket(0, 0, "camera", "main", 1, 1, "opaque", b"data")] + with pytest.raises(DecoderError, match="no decoder registered"): + list(DecoderPipeline().decode(unknown)) + + +def test_y4m_local_video_decodes_across_input_chunks(tmp_path: Path) -> None: + video = tmp_path / "anonymous.y4m" + first, second = b"abcdef", b"ghijkl" + video.write_bytes(b"YUV4MPEG2 W2 H1 F25:1 C444\nFRAME\n" + first + b"FRAME\n" + second) + frames = list(decode_packets(local_packets(video))) + assert [frame.payload for frame in frames] == [first, second] + assert [frame.timestamp_ns for frame in frames] == [0, 40_000_000] + assert all(frame.pixel_format == "yuv444p" for frame in frames) + assert all((frame.width, frame.height) == (2, 1) for frame in frames) + assert all(frame.profile_id == "archive" for frame in frames) + + +def test_y4m_clean_eof_and_cancellation_are_normal(tmp_path: Path) -> None: + video = tmp_path / "empty.y4m" + video.write_bytes(b"YUV4MPEG2 W2 H1 F25:1 C444\n") + assert list(decode_packets(local_packets(video))) == [] + + token = CancellationToken() + token.cancel() + assert list(decode_packets(local_packets(video), token)) == [] + + +@pytest.mark.parametrize( + ("payload", "message"), + [ + (b"not-video\n", "expected YUV4MPEG2"), + (b"YUV4MPEG2 W2 H1 F25:1 C420\n", "only C444"), + (b"YUV4MPEG2 W2 H1 F25:1 C444\nFRAME\nabc", "truncated Y4M frame"), + ], +) +def test_y4m_damage_and_unsupported_content_are_clear(tmp_path: Path, payload: bytes, message: str) -> None: + video = tmp_path / "broken.y4m" + video.write_bytes(payload) + with pytest.raises(DecoderError, match=message): + list(decode_packets(local_packets(video))) + + +def test_y4m_profile_dimension_mismatch_is_rejected(tmp_path: Path) -> None: + video = tmp_path / "mismatch.y4m" + video.write_bytes(b"YUV4MPEG2 W2 H1 F25:1 C444\n") + with pytest.raises(DecoderError, match="do not match"): + list(decode_packets(local_packets(video, width=3))) diff --git a/Brain/tests/fixtures/decode/README.md b/Brain/tests/fixtures/decode/README.md new file mode 100644 index 0000000..7ed954b --- /dev/null +++ b/Brain/tests/fixtures/decode/README.md @@ -0,0 +1,5 @@ +# Brain decode fixtures + +Decode tests generate tiny anonymous YUV4MPEG2 streams at runtime. Do not add +customer recordings, camera credentials, machine-specific codec paths, or +large model/media artifacts to this directory. -- 2.34.1