139 lines
4.4 KiB
Python
139 lines
4.4 KiB
Python
"""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()
|