From 1b6af397e5f0e71799ab089bfddceaefa4a1f68f Mon Sep 17 00:00:00 2001 From: raph666 Date: Sat, 18 Jul 2026 16:58:19 +0200 Subject: [PATCH] Add incremental REG3 frame extractor --- custom_components/arkteos/frame_extractor.py | 138 +++++++++++++++ tests/test_frame_extractor.py | 167 +++++++++++++++++++ 2 files changed, 305 insertions(+) create mode 100644 custom_components/arkteos/frame_extractor.py create mode 100644 tests/test_frame_extractor.py diff --git a/custom_components/arkteos/frame_extractor.py b/custom_components/arkteos/frame_extractor.py new file mode 100644 index 0000000..b8a4403 --- /dev/null +++ b/custom_components/arkteos/frame_extractor.py @@ -0,0 +1,138 @@ +"""Extracteur incrémental des trames REG3 observées sur le flux TCP. + +Ce module ne décode aucune valeur métier et ne dépend ni du réseau ni de Home +Assistant. Les signatures, tailles et cohérences d'en-tête sont limitées à ce +qui est observé dans ``docs/frame_protocol_analysis.md``. +""" + +from __future__ import annotations + +from typing import Final + + +FRAME_SIGNATURE: Final = b"\x55\x00" +HEADER_MIN_SIZE: Final = 12 +HEADER_LENGTH_OFFSET: Final = 10 +HEADER_SIZE_ADJUSTMENT: Final = 15 +FRAME_TYPE_OFFSET: Final = 8 + +FRAME_SIZES_BY_TYPE: Final[dict[int, int]] = { + 0x0C: 95, + 0x0A: 163, + 0x0B: 227, +} +KNOWN_FRAME_SIZES: Final[frozenset[int]] = frozenset(FRAME_SIZES_BY_TYPE.values()) + + +class FrameExtractor: + """Assemble des trames complètes depuis des fragments TCP arbitraires.""" + + def __init__(self, *, max_buffer_size: int = 4096) -> None: + if max_buffer_size < HEADER_MIN_SIZE: + raise ValueError( + f"max_buffer_size doit être supérieur ou égal à {HEADER_MIN_SIZE}" + ) + self._max_buffer_size = max_buffer_size + self._buffer = bytearray() + + @property + def buffered_bytes(self) -> bytes: + """Retourne une copie immuable des octets en attente.""" + + return bytes(self._buffer) + + @property + def buffered_size(self) -> int: + """Retourne la taille des octets en attente.""" + + return len(self._buffer) + + @property + def max_buffer_size(self) -> int: + """Retourne la limite configurable du tampon persistant.""" + + return self._max_buffer_size + + def feed(self, data: bytes) -> list[bytes]: + """Ajoute un fragment TCP et retourne les trames complètes extraites. + + Les données entrantes sont ajoutées par blocs bornés afin que le tampon + persistant ne dépasse jamais ``max_buffer_size``. + """ + + if not isinstance(data, bytes): + raise TypeError("data doit être de type bytes") + + frames = self._extract_available() + offset = 0 + while offset < len(data): + available = self._max_buffer_size - len(self._buffer) + if available == 0: + before = len(self._buffer) + frames.extend(self._extract_available()) + if len(self._buffer) == before: + # Cette branche ne doit être atteinte qu'avec un en-tête + # incomplet ou invalide ; avancer garantit la progression. + del self._buffer[0] + continue + + end = min(offset + available, len(data)) + self._buffer.extend(data[offset:end]) + offset = end + frames.extend(self._extract_available()) + + return frames + + def _extract_available(self) -> list[bytes]: + frames: list[bytes] = [] + + while True: + if len(self._buffer) < len(FRAME_SIGNATURE): + return frames + + signature_offset = self._buffer.find(FRAME_SIGNATURE) + if signature_offset < 0: + self._preserve_possible_split_signature() + return frames + + if signature_offset > 0: + del self._buffer[:signature_offset] + + if len(self._buffer) < HEADER_MIN_SIZE: + return frames + + frame_size = ( + int.from_bytes( + self._buffer[ + HEADER_LENGTH_OFFSET : HEADER_LENGTH_OFFSET + 2 + ], + "little", + ) + + HEADER_SIZE_ADJUSTMENT + ) + frame_type = self._buffer[FRAME_TYPE_OFFSET] + expected_size = FRAME_SIZES_BY_TYPE.get(frame_type) + + if ( + frame_size not in KNOWN_FRAME_SIZES + or frame_size > self._max_buffer_size + or expected_size != frame_size + ): + # Ne supprimer qu'un octet : une nouvelle signature peut + # commencer à l'octet suivant du flux reçu. + del self._buffer[0] + continue + + if len(self._buffer) < frame_size: + return frames + + frames.append(bytes(self._buffer[:frame_size])) + del self._buffer[:frame_size] + + def _preserve_possible_split_signature(self) -> None: + """Conserve seulement un dernier 0x55 en l'absence de signature.""" + + if self._buffer[-1] == FRAME_SIGNATURE[0]: + self._buffer[:] = FRAME_SIGNATURE[:1] + else: + self._buffer.clear() diff --git a/tests/test_frame_extractor.py b/tests/test_frame_extractor.py new file mode 100644 index 0000000..279dbcb --- /dev/null +++ b/tests/test_frame_extractor.py @@ -0,0 +1,167 @@ +"""Tests de l'extracteur incrémental REG3, indépendants de Home Assistant.""" + +from pathlib import Path + +from custom_components.arkteos.frame_extractor import FrameExtractor + + +FIXTURES_PATH = Path(__file__).with_name("fixtures") + + +def _load_frames() -> dict[str, bytes]: + return { + "metadata": (FIXTURES_PATH / "metadata_95.bin").read_bytes(), + "frigo": (FIXTURES_PATH / "frigo_163.bin").read_bytes(), + "regulation": (FIXTURES_PATH / "regulation_227.bin").read_bytes(), + } + + +def test_no_bytes() -> None: + extractor = FrameExtractor() + assert extractor.feed(b"") == [] + assert extractor.buffered_bytes == b"" + + +def test_complete_metadata_frame() -> None: + metadata = _load_frames()["metadata"] + assert FrameExtractor().feed(metadata) == [metadata] + + +def test_complete_frigo_frame() -> None: + frigo = _load_frames()["frigo"] + assert FrameExtractor().feed(frigo) == [frigo] + + +def test_complete_regulation_frame() -> None: + regulation = _load_frames()["regulation"] + assert FrameExtractor().feed(regulation) == [regulation] + + +def test_frame_split_in_two_calls() -> None: + frigo = _load_frames()["frigo"] + extractor = FrameExtractor() + assert extractor.feed(frigo[:80]) == [] + assert extractor.feed(frigo[80:]) == [frigo] + + +def test_frame_provided_byte_by_byte() -> None: + regulation = _load_frames()["regulation"] + extractor = FrameExtractor() + extracted = [frame for byte in regulation for frame in extractor.feed(bytes([byte]))] + assert extracted == [regulation] + + +def test_multiple_concatenated_frames() -> None: + frames = _load_frames() + sequence = frames["metadata"] + frames["frigo"] + frames["regulation"] + assert FrameExtractor().feed(sequence) == [ + frames["metadata"], + frames["frigo"], + frames["regulation"], + ] + + +def test_order_different_from_observed_cycle() -> None: + frames = _load_frames() + sequence = frames["frigo"] + frames["metadata"] + frames["regulation"] + assert FrameExtractor().feed(sequence) == [ + frames["frigo"], + frames["metadata"], + frames["regulation"], + ] + + +def test_noise_before_signature() -> None: + metadata = _load_frames()["metadata"] + assert FrameExtractor().feed(b"parasites" + metadata) == [metadata] + + +def test_signature_split_between_calls() -> None: + metadata = _load_frames()["metadata"] + extractor = FrameExtractor() + assert extractor.feed(b"parasites\x55") == [] + assert extractor.buffered_bytes == b"\x55" + assert extractor.feed(metadata[1:]) == [metadata] + + +def test_incoherent_type_header_resynchronizes() -> None: + metadata = _load_frames()["metadata"] + invalid = bytearray(metadata) + invalid[8] = 0x0A + assert FrameExtractor().feed(bytes(invalid) + metadata) == [metadata] + + +def test_unknown_length_header_resynchronizes() -> None: + metadata = _load_frames()["metadata"] + invalid = bytearray(metadata) + invalid[10:12] = (81).to_bytes(2, "little") + assert FrameExtractor().feed(bytes(invalid) + metadata) == [metadata] + + +def test_announced_size_above_limit_resynchronizes() -> None: + frames = _load_frames() + extractor = FrameExtractor(max_buffer_size=200) + assert extractor.feed(frames["regulation"] + frames["metadata"]) == [ + frames["metadata"] + ] + + +def test_false_signature_in_noise() -> None: + metadata = _load_frames()["metadata"] + false_header = b"\x55\x00" + b"\x00" * 10 + assert FrameExtractor().feed(b"\x10" + false_header + metadata) == [metadata] + + +def test_complete_frame_followed_by_incomplete_frame() -> None: + frames = _load_frames() + extractor = FrameExtractor() + assert extractor.feed(frames["metadata"] + frames["frigo"][:20]) == [ + frames["metadata"] + ] + assert extractor.buffered_bytes == frames["frigo"][:20] + assert extractor.feed(frames["frigo"][20:]) == [frames["frigo"]] + + +def test_variable_size_fragments() -> None: + frames = _load_frames() + sequence = frames["frigo"] + frames["metadata"] + frames["regulation"] + extractor = FrameExtractor() + extracted: list[bytes] = [] + fragment_sizes = (1, 7, 31, 2, 83, 19, 11, 131) + offset = 0 + for size in fragment_sizes: + extracted.extend(extractor.feed(sequence[offset : offset + size])) + offset += size + extracted.extend(extractor.feed(sequence[offset:])) + assert extracted == [frames["frigo"], frames["metadata"], frames["regulation"]] + + +def test_buffer_is_empty_after_complete_extraction() -> None: + metadata = _load_frames()["metadata"] + extractor = FrameExtractor() + extractor.feed(metadata) + assert extractor.buffered_size == 0 + + +def test_incomplete_buffer_is_preserved() -> None: + metadata = _load_frames()["metadata"] + extractor = FrameExtractor() + assert extractor.feed(metadata[:10]) == [] + assert extractor.buffered_bytes == metadata[:10] + + +def test_repeated_invalid_data_terminates() -> None: + invalid_header = b"\x55\x00" + b"\xff" * 10 + extractor = FrameExtractor() + assert extractor.feed(invalid_header * 100) == [] + assert extractor.buffered_size == 0 + + +def test_sequence_with_all_real_fixtures() -> None: + frames = _load_frames() + sequence = frames["regulation"] + frames["frigo"] + frames["metadata"] + assert FrameExtractor().feed(sequence) == [ + frames["regulation"], + frames["frigo"], + frames["metadata"], + ]