"""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()