Add incremental REG3 frame extractor
This commit is contained in:
@@ -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()
|
||||
Reference in New Issue
Block a user