refactor(python): centralize SETProtocol v2 codec
This commit is contained in:
3
python/setprotocol/__init__.py
Normal file
3
python/setprotocol/__init__.py
Normal file
@@ -0,0 +1,3 @@
|
||||
"""Cross-platform Python facade for the canonical SETProtocol core."""
|
||||
|
||||
from .core import * # noqa: F401,F403
|
||||
671
python/setprotocol/core.py
Normal file
671
python/setprotocol/core.py
Normal file
@@ -0,0 +1,671 @@
|
||||
"""Эталонная Python-реализация SET protocol v2.
|
||||
|
||||
Источник wire-контракта находится в ``templates/c/set-protocol/PROTOCOL.md``.
|
||||
Модуль не импортирует Qt и используется SETGUI, host-тестами и утилитами.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import binascii
|
||||
import hashlib
|
||||
import struct
|
||||
from dataclasses import dataclass
|
||||
from enum import IntEnum, IntFlag
|
||||
from typing import Iterable
|
||||
|
||||
SOF = b"\xA5\x5A"
|
||||
VERSION = 0x02
|
||||
HEADER_SIZE = 14
|
||||
CRC_SIZE = 4
|
||||
MAX_PAYLOAD_SIZE = 512
|
||||
FRAME_MAX = HEADER_SIZE + MAX_PAYLOAD_SIZE + CRC_SIZE
|
||||
|
||||
NODE_LOCAL = 0x0000
|
||||
NODE_BROADCAST = 0xFFFF
|
||||
|
||||
|
||||
class SetProtocolError(ValueError):
|
||||
"""Кадр или payload нарушает обязательный контракт SETP v2."""
|
||||
|
||||
|
||||
class FrameFlag(IntFlag):
|
||||
RESPONSE = 0x01
|
||||
EVENT = 0x02
|
||||
ERROR = 0x04
|
||||
ACK_REQUIRED = 0x08
|
||||
MORE = 0x10
|
||||
PRIORITY = 0x20
|
||||
|
||||
|
||||
KNOWN_FLAGS = int(
|
||||
FrameFlag.RESPONSE
|
||||
| FrameFlag.EVENT
|
||||
| FrameFlag.ERROR
|
||||
| FrameFlag.ACK_REQUIRED
|
||||
| FrameFlag.MORE
|
||||
| FrameFlag.PRIORITY
|
||||
)
|
||||
|
||||
|
||||
class MessageType(IntEnum):
|
||||
PING = 0x0001
|
||||
DEVICE_INFO = 0x0002
|
||||
CAPABILITIES = 0x0003
|
||||
DIAGNOSTICS = 0x0008
|
||||
READ = 0x0009
|
||||
WRITE = 0x000A
|
||||
LOG_READ = 0x0010
|
||||
CATALOG = 0x0011
|
||||
SUBSCRIBE = 0x0012
|
||||
PUBLISH = 0x0013
|
||||
UNSUBSCRIBE = 0x0014
|
||||
|
||||
FW_BEGIN = 0x0100
|
||||
FW_DATA = 0x0101
|
||||
FW_END = 0x0102
|
||||
FW_ABORT = 0x0103
|
||||
FW_STATUS = 0x0104
|
||||
FW_ACTIVATE = 0x0105
|
||||
|
||||
|
||||
class Status(IntEnum):
|
||||
OK = 0
|
||||
INVALID_ARGUMENT = 1
|
||||
INVALID_LENGTH = 2
|
||||
NOT_FOUND = 3
|
||||
ACCESS_DENIED = 4
|
||||
BUSY = 5
|
||||
NO_PROVIDER = 6
|
||||
INTERNAL = 7
|
||||
CRC = 8
|
||||
SEQUENCE = 9
|
||||
NO_SPACE = 10
|
||||
UNSUPPORTED = 11
|
||||
VERIFY_FAILED = 12
|
||||
AUTH_FAILED = 13
|
||||
TIMEOUT = 14
|
||||
WRONG_STATE = 15
|
||||
|
||||
|
||||
class Interface(IntEnum):
|
||||
NONE = 0
|
||||
RS232 = 1
|
||||
RS485 = 2
|
||||
CAN = 3
|
||||
USB_CDC = 4
|
||||
ETHERNET_TCP = 5
|
||||
ETHERNET_UDP = 6
|
||||
|
||||
|
||||
class Feature(IntFlag):
|
||||
READ = 1 << 0
|
||||
WRITE = 1 << 1
|
||||
CATALOG = 1 << 2
|
||||
SUBSCRIBE = 1 << 3
|
||||
LOG_READ = 1 << 4
|
||||
FIRMWARE = 1 << 5
|
||||
DIAGNOSTICS = 1 << 6
|
||||
|
||||
|
||||
class ValueEncoding(IntEnum):
|
||||
U16 = 1
|
||||
I16 = 2
|
||||
U32 = 3
|
||||
I32 = 4
|
||||
F32 = 5
|
||||
BYTES = 6
|
||||
|
||||
|
||||
class FirmwareFlag(IntFlag):
|
||||
RESUME = 0x01
|
||||
SIGNED = 0x02
|
||||
ERASE_SLOT = 0x04
|
||||
ACTIVATE = 0x08
|
||||
|
||||
|
||||
def crc32_ieee(data: bytes | bytearray | memoryview) -> int:
|
||||
return binascii.crc32(data) & 0xFFFFFFFF
|
||||
|
||||
|
||||
def _u16(value: int, name: str) -> int:
|
||||
if not 0 <= value <= 0xFFFF:
|
||||
raise SetProtocolError(f"{name} вне диапазона u16")
|
||||
return value
|
||||
|
||||
|
||||
def _u32(value: int, name: str) -> int:
|
||||
if not 0 <= value <= 0xFFFFFFFF:
|
||||
raise SetProtocolError(f"{name} вне диапазона u32")
|
||||
return value
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class Frame:
|
||||
message_type: int
|
||||
sequence: int
|
||||
payload: bytes = b""
|
||||
flags: FrameFlag = FrameFlag(0)
|
||||
source: int = NODE_LOCAL
|
||||
destination: int = NODE_LOCAL
|
||||
|
||||
def __post_init__(self) -> None:
|
||||
_u16(int(self.message_type), "message_type")
|
||||
_u16(self.sequence, "sequence")
|
||||
_u16(self.source, "source")
|
||||
_u16(self.destination, "destination")
|
||||
if int(self.flags) & ~KNOWN_FLAGS:
|
||||
raise SetProtocolError("зарезервированные флаги должны быть нулевыми")
|
||||
if len(self.payload) > MAX_PAYLOAD_SIZE:
|
||||
raise SetProtocolError(f"payload превышает {MAX_PAYLOAD_SIZE} байт")
|
||||
|
||||
|
||||
def build_frame(frame: Frame) -> bytes:
|
||||
protected = struct.pack(
|
||||
"<BBHHHHH",
|
||||
VERSION,
|
||||
int(frame.flags),
|
||||
int(frame.message_type),
|
||||
frame.source,
|
||||
frame.destination,
|
||||
frame.sequence,
|
||||
len(frame.payload),
|
||||
) + frame.payload
|
||||
return SOF + protected + struct.pack("<I", crc32_ieee(protected))
|
||||
|
||||
|
||||
def decode_datagram(data: bytes | bytearray | memoryview) -> Frame:
|
||||
"""Разбирает ровно один SETP-кадр из UDP или другого datagram-носителя."""
|
||||
packet = bytes(data)
|
||||
if len(packet) < HEADER_SIZE + CRC_SIZE or packet[:2] != SOF:
|
||||
raise SetProtocolError("датаграмма не содержит полный SETP-кадр")
|
||||
if packet[2] != VERSION:
|
||||
raise SetProtocolError(f"неподдерживаемая версия {packet[2]}")
|
||||
if packet[3] & ~KNOWN_FLAGS:
|
||||
raise SetProtocolError("зарезервированные флаги должны быть нулевыми")
|
||||
payload_length = int.from_bytes(packet[12:14], "little")
|
||||
if payload_length > MAX_PAYLOAD_SIZE:
|
||||
raise SetProtocolError("payload датаграммы превышает лимит")
|
||||
expected_length = HEADER_SIZE + payload_length + CRC_SIZE
|
||||
if len(packet) != expected_length:
|
||||
raise SetProtocolError("длина UDP-датаграммы не совпадает с SETP header")
|
||||
if crc32_ieee(packet[2:-CRC_SIZE]) != int.from_bytes(packet[-CRC_SIZE:], "little"):
|
||||
raise SetProtocolError("неверный CRC32 датаграммы")
|
||||
_, flags, message_type, source, destination, sequence, _ = struct.unpack_from(
|
||||
"<BBHHHHH", packet, 2
|
||||
)
|
||||
return Frame(
|
||||
message_type=message_type,
|
||||
sequence=sequence,
|
||||
payload=packet[HEADER_SIZE:-CRC_SIZE],
|
||||
flags=FrameFlag(flags),
|
||||
source=source,
|
||||
destination=destination,
|
||||
)
|
||||
|
||||
|
||||
@dataclass
|
||||
class ParserStats:
|
||||
frames: int = 0
|
||||
crc_errors: int = 0
|
||||
version_errors: int = 0
|
||||
length_errors: int = 0
|
||||
flag_errors: int = 0
|
||||
stray_bytes: int = 0
|
||||
overflows: int = 0
|
||||
|
||||
|
||||
class FrameParser:
|
||||
"""Потоковый parser для UART, USB CDC и Ethernet TCP."""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self._buffer = bytearray()
|
||||
self.stats = ParserStats()
|
||||
|
||||
def reset(self) -> None:
|
||||
self._buffer.clear()
|
||||
|
||||
@property
|
||||
def buffered_bytes(self) -> int:
|
||||
return len(self._buffer)
|
||||
|
||||
def feed(self, data: bytes | bytearray | memoryview) -> list[Frame]:
|
||||
if not data:
|
||||
return []
|
||||
self._buffer.extend(data)
|
||||
|
||||
frames: list[Frame] = []
|
||||
while True:
|
||||
sof_index = self._buffer.find(SOF)
|
||||
if sof_index < 0:
|
||||
keep = 1 if self._buffer[-1:] == SOF[:1] else 0
|
||||
self.stats.stray_bytes += len(self._buffer) - keep
|
||||
self._buffer[:] = self._buffer[-1:] if keep else b""
|
||||
break
|
||||
if sof_index:
|
||||
self.stats.stray_bytes += sof_index
|
||||
del self._buffer[:sof_index]
|
||||
if len(self._buffer) < HEADER_SIZE:
|
||||
break
|
||||
if self._buffer[2] != VERSION:
|
||||
self.stats.version_errors += 1
|
||||
del self._buffer[0]
|
||||
continue
|
||||
flags = self._buffer[3]
|
||||
if flags & ~KNOWN_FLAGS:
|
||||
self.stats.flag_errors += 1
|
||||
del self._buffer[0]
|
||||
continue
|
||||
payload_length = int.from_bytes(self._buffer[12:14], "little")
|
||||
if payload_length > MAX_PAYLOAD_SIZE:
|
||||
self.stats.length_errors += 1
|
||||
del self._buffer[0]
|
||||
continue
|
||||
total = HEADER_SIZE + payload_length + CRC_SIZE
|
||||
if len(self._buffer) < total:
|
||||
break
|
||||
packet = bytes(self._buffer[:total])
|
||||
protected = packet[2:-CRC_SIZE]
|
||||
received_crc = int.from_bytes(packet[-CRC_SIZE:], "little")
|
||||
if crc32_ieee(protected) != received_crc:
|
||||
self.stats.crc_errors += 1
|
||||
del self._buffer[0]
|
||||
continue
|
||||
(
|
||||
_version,
|
||||
_flags,
|
||||
message_type,
|
||||
source,
|
||||
destination,
|
||||
sequence,
|
||||
_payload_length,
|
||||
) = struct.unpack_from("<BBHHHHH", packet, 2)
|
||||
frames.append(
|
||||
Frame(
|
||||
message_type=message_type,
|
||||
sequence=sequence,
|
||||
payload=packet[HEADER_SIZE : HEADER_SIZE + payload_length],
|
||||
flags=FrameFlag(flags),
|
||||
source=source,
|
||||
destination=destination,
|
||||
)
|
||||
)
|
||||
self.stats.frames += 1
|
||||
del self._buffer[:total]
|
||||
return frames
|
||||
|
||||
|
||||
def encode_response(status: int | Status, body: bytes = b"") -> bytes:
|
||||
return struct.pack("<H", _u16(int(status), "status")) + body
|
||||
|
||||
|
||||
def decode_response(frame: Frame) -> tuple[int, bytes]:
|
||||
if not frame.flags & FrameFlag.RESPONSE:
|
||||
raise SetProtocolError("кадр не является ответом")
|
||||
if len(frame.payload) < 2:
|
||||
raise SetProtocolError("ответ не содержит status")
|
||||
return int.from_bytes(frame.payload[:2], "little"), frame.payload[2:]
|
||||
|
||||
|
||||
DEVICE_INFO_SCHEMA_VERSION = 1
|
||||
DEVICE_INFO_FIXED_SIZE = 25
|
||||
DEVICE_MODEL_MAX = 63
|
||||
CAPABILITIES_SCHEMA_VERSION = 1
|
||||
CAPABILITIES_SIZE = 20
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class DeviceInfo:
|
||||
"""Канонический body ответа DEVICE_INFO, без двухбайтового status."""
|
||||
|
||||
device_class: int
|
||||
hardware_version: int
|
||||
firmware_version: int
|
||||
dictionary_version: int
|
||||
serial_number: int
|
||||
model: str
|
||||
schema_version: int = DEVICE_INFO_SCHEMA_VERSION
|
||||
|
||||
def encode(self) -> bytes:
|
||||
_u16(self.schema_version, "schema_version")
|
||||
if self.schema_version != DEVICE_INFO_SCHEMA_VERSION:
|
||||
raise SetProtocolError("неподдерживаемая схема DEVICE_INFO")
|
||||
_u16(self.device_class, "device_class")
|
||||
_u32(self.hardware_version, "hardware_version")
|
||||
_u32(self.firmware_version, "firmware_version")
|
||||
_u32(self.dictionary_version, "dictionary_version")
|
||||
if not 0 <= self.serial_number <= 0xFFFFFFFFFFFFFFFF:
|
||||
raise SetProtocolError("serial_number вне диапазона u64")
|
||||
model = self.model.encode("utf-8")
|
||||
if len(model) > DEVICE_MODEL_MAX:
|
||||
raise SetProtocolError("model длиннее 63 байт UTF-8")
|
||||
return struct.pack(
|
||||
"<HHIIIQB",
|
||||
self.schema_version,
|
||||
self.device_class,
|
||||
self.hardware_version,
|
||||
self.firmware_version,
|
||||
self.dictionary_version,
|
||||
self.serial_number,
|
||||
len(model),
|
||||
) + model
|
||||
|
||||
@classmethod
|
||||
def decode(cls, payload: bytes) -> "DeviceInfo":
|
||||
if len(payload) < DEVICE_INFO_FIXED_SIZE:
|
||||
raise SetProtocolError("DEVICE_INFO короче фиксированного заголовка")
|
||||
values = struct.unpack_from("<HHIIIQB", payload)
|
||||
model_length = values[-1]
|
||||
if values[0] != DEVICE_INFO_SCHEMA_VERSION:
|
||||
raise SetProtocolError(f"неподдерживаемая схема DEVICE_INFO {values[0]}")
|
||||
if model_length > DEVICE_MODEL_MAX or len(payload) != DEVICE_INFO_FIXED_SIZE + model_length:
|
||||
raise SetProtocolError("DEVICE_INFO содержит неверную длину model")
|
||||
try:
|
||||
model = payload[DEVICE_INFO_FIXED_SIZE:].decode("utf-8")
|
||||
except UnicodeDecodeError as error:
|
||||
raise SetProtocolError("DEVICE_INFO model не является UTF-8") from error
|
||||
return cls(
|
||||
device_class=values[1],
|
||||
hardware_version=values[2],
|
||||
firmware_version=values[3],
|
||||
dictionary_version=values[4],
|
||||
serial_number=values[5],
|
||||
model=model,
|
||||
schema_version=values[0],
|
||||
)
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class Capabilities:
|
||||
"""Канонический body ответа CAPABILITIES, без status."""
|
||||
|
||||
max_payload: int
|
||||
interface_mask: int
|
||||
feature_flags: int
|
||||
max_read_items: int
|
||||
max_write_items: int
|
||||
max_subscriptions: int
|
||||
max_publish_items: int
|
||||
schema_version: int = CAPABILITIES_SCHEMA_VERSION
|
||||
|
||||
def encode(self) -> bytes:
|
||||
if self.schema_version != CAPABILITIES_SCHEMA_VERSION:
|
||||
raise SetProtocolError("неподдерживаемая схема CAPABILITIES")
|
||||
if not 1 <= self.max_payload <= MAX_PAYLOAD_SIZE:
|
||||
raise SetProtocolError("max_payload вне базового профиля")
|
||||
_u32(self.interface_mask, "interface_mask")
|
||||
_u32(self.feature_flags, "feature_flags")
|
||||
for name in (
|
||||
"max_read_items", "max_write_items", "max_subscriptions",
|
||||
"max_publish_items",
|
||||
):
|
||||
_u16(getattr(self, name), name)
|
||||
return struct.pack(
|
||||
"<HHIIHHHH",
|
||||
self.schema_version,
|
||||
self.max_payload,
|
||||
self.interface_mask,
|
||||
self.feature_flags,
|
||||
self.max_read_items,
|
||||
self.max_write_items,
|
||||
self.max_subscriptions,
|
||||
self.max_publish_items,
|
||||
)
|
||||
|
||||
@classmethod
|
||||
def decode(cls, payload: bytes) -> "Capabilities":
|
||||
if len(payload) != CAPABILITIES_SIZE:
|
||||
raise SetProtocolError("CAPABILITIES должен содержать ровно 20 байт")
|
||||
values = struct.unpack("<HHIIHHHH", payload)
|
||||
result = cls(
|
||||
max_payload=values[1],
|
||||
interface_mask=values[2],
|
||||
feature_flags=values[3],
|
||||
max_read_items=values[4],
|
||||
max_write_items=values[5],
|
||||
max_subscriptions=values[6],
|
||||
max_publish_items=values[7],
|
||||
schema_version=values[0],
|
||||
)
|
||||
# Каноническая повторная кодировка выполняет все проверки диапазонов.
|
||||
result.encode()
|
||||
return result
|
||||
|
||||
@property
|
||||
def interfaces(self) -> tuple[Interface, ...]:
|
||||
return tuple(
|
||||
interface for interface in Interface
|
||||
if interface is not Interface.NONE and self.interface_mask & (1 << interface)
|
||||
)
|
||||
|
||||
@property
|
||||
def features(self) -> Feature:
|
||||
return Feature(self.feature_flags)
|
||||
|
||||
|
||||
FW_BEGIN_FIXED_SIZE = 58
|
||||
FW_DATA_HEADER_SIZE = 12
|
||||
FW_END_SIZE = 40
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class FirmwareBegin:
|
||||
image_size: int
|
||||
image_crc32: int
|
||||
image_version: int
|
||||
base_address: int
|
||||
slot: int
|
||||
block_size: int
|
||||
sha256: bytes
|
||||
flags: FirmwareFlag = FirmwareFlag(0)
|
||||
signing_key_id: int = 0
|
||||
signature: bytes = b""
|
||||
|
||||
def encode(self) -> bytes:
|
||||
_u32(self.image_size, "image_size")
|
||||
_u32(self.image_crc32, "image_crc32")
|
||||
_u32(self.image_version, "image_version")
|
||||
_u32(self.base_address, "base_address")
|
||||
_u32(self.signing_key_id, "signing_key_id")
|
||||
if not 0 <= self.slot <= 0xFF:
|
||||
raise SetProtocolError("slot вне диапазона u8")
|
||||
if not 1 <= self.block_size <= MAX_PAYLOAD_SIZE - FW_DATA_HEADER_SIZE:
|
||||
raise SetProtocolError("block_size не помещается в FW_DATA")
|
||||
if len(self.sha256) != 32:
|
||||
raise SetProtocolError("SHA-256 должен содержать 32 байта")
|
||||
if self.flags & FirmwareFlag.SIGNED and not self.signature:
|
||||
raise SetProtocolError("SIGNED требует подпись")
|
||||
_u16(len(self.signature), "signature_length")
|
||||
payload = struct.pack(
|
||||
"<IIIIBBH32sIH",
|
||||
self.image_size,
|
||||
self.image_crc32,
|
||||
self.image_version,
|
||||
self.base_address,
|
||||
self.slot,
|
||||
int(self.flags),
|
||||
self.block_size,
|
||||
self.sha256,
|
||||
self.signing_key_id,
|
||||
len(self.signature),
|
||||
) + self.signature
|
||||
if len(payload) > MAX_PAYLOAD_SIZE:
|
||||
raise SetProtocolError("FW_BEGIN не помещается в payload")
|
||||
return payload
|
||||
|
||||
@classmethod
|
||||
def decode(cls, payload: bytes) -> "FirmwareBegin":
|
||||
if len(payload) < FW_BEGIN_FIXED_SIZE:
|
||||
raise SetProtocolError("FW_BEGIN короче фиксированного заголовка")
|
||||
values = struct.unpack_from("<IIIIBBH32sIH", payload)
|
||||
signature_length = values[-1]
|
||||
if len(payload) != FW_BEGIN_FIXED_SIZE + signature_length:
|
||||
raise SetProtocolError("FW_BEGIN содержит неверную длину подписи")
|
||||
result = cls(
|
||||
image_size=values[0],
|
||||
image_crc32=values[1],
|
||||
image_version=values[2],
|
||||
base_address=values[3],
|
||||
slot=values[4],
|
||||
flags=FirmwareFlag(values[5]),
|
||||
block_size=values[6],
|
||||
sha256=values[7],
|
||||
signing_key_id=values[8],
|
||||
signature=payload[FW_BEGIN_FIXED_SIZE:],
|
||||
)
|
||||
# Повторная кодировка выполняет все семантические проверки.
|
||||
result.encode()
|
||||
return result
|
||||
|
||||
|
||||
def firmware_manifest(begin: FirmwareBegin) -> bytes:
|
||||
"""Канонические байты, которые подписывает поставщик прошивки."""
|
||||
return b"SETPFW2\0" + struct.pack(
|
||||
"<IIIIB32s",
|
||||
begin.image_size,
|
||||
begin.image_crc32,
|
||||
begin.image_version,
|
||||
begin.base_address,
|
||||
begin.slot,
|
||||
begin.sha256,
|
||||
)
|
||||
|
||||
|
||||
def firmware_begin_for_image(
|
||||
image: bytes,
|
||||
*,
|
||||
image_version: int,
|
||||
base_address: int,
|
||||
slot: int,
|
||||
block_size: int = 480,
|
||||
) -> FirmwareBegin:
|
||||
return FirmwareBegin(
|
||||
image_size=len(image),
|
||||
image_crc32=crc32_ieee(image),
|
||||
image_version=image_version,
|
||||
base_address=base_address,
|
||||
slot=slot,
|
||||
block_size=block_size,
|
||||
sha256=hashlib.sha256(image).digest(),
|
||||
flags=FirmwareFlag.RESUME,
|
||||
)
|
||||
|
||||
|
||||
def encode_firmware_data(offset: int, data: bytes, flags: int = 0) -> bytes:
|
||||
_u32(offset, "offset")
|
||||
_u16(flags, "flags")
|
||||
if not data or len(data) > MAX_PAYLOAD_SIZE - FW_DATA_HEADER_SIZE:
|
||||
raise SetProtocolError("неверный размер блока FW_DATA")
|
||||
return struct.pack("<IHHI", offset, len(data), flags, crc32_ieee(data)) + data
|
||||
|
||||
|
||||
def decode_firmware_data(payload: bytes) -> tuple[int, int, bytes]:
|
||||
if len(payload) < FW_DATA_HEADER_SIZE:
|
||||
raise SetProtocolError("FW_DATA короче заголовка")
|
||||
offset, length, flags, expected_crc = struct.unpack_from("<IHHI", payload)
|
||||
data = payload[FW_DATA_HEADER_SIZE:]
|
||||
if not data or len(data) != length:
|
||||
raise SetProtocolError("FW_DATA содержит неверную длину блока")
|
||||
if crc32_ieee(data) != expected_crc:
|
||||
raise SetProtocolError("FW_DATA содержит неверный CRC блока")
|
||||
return offset, flags, data
|
||||
|
||||
|
||||
def encode_subscribe(subscription_id: int, period_ms: int, addresses: Iterable[int]) -> bytes:
|
||||
items = tuple(addresses)
|
||||
_u16(subscription_id, "subscription_id")
|
||||
_u32(period_ms, "period_ms")
|
||||
_u16(len(items), "address_count")
|
||||
for address in items:
|
||||
_u32(address, "address")
|
||||
payload = struct.pack("<HIH", subscription_id, period_ms, len(items))
|
||||
payload += b"".join(struct.pack("<I", address) for address in items)
|
||||
if len(payload) > MAX_PAYLOAD_SIZE:
|
||||
raise SetProtocolError("SUBSCRIBE не помещается в payload")
|
||||
return payload
|
||||
|
||||
|
||||
def decode_subscribe(payload: bytes) -> tuple[int, int, tuple[int, ...]]:
|
||||
if len(payload) < 8:
|
||||
raise SetProtocolError("SUBSCRIBE короче заголовка")
|
||||
subscription_id, period_ms, count = struct.unpack_from("<HIH", payload)
|
||||
if len(payload) != 8 + count * 4:
|
||||
raise SetProtocolError("SUBSCRIBE содержит неверное число адресов")
|
||||
addresses = struct.unpack_from(f"<{count}I", payload, 8) if count else ()
|
||||
return subscription_id, period_ms, tuple(addresses)
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class PublishItem:
|
||||
address: int
|
||||
encoding: ValueEncoding
|
||||
element_count: int
|
||||
data: bytes
|
||||
|
||||
|
||||
def encode_publish(
|
||||
subscription_id: int,
|
||||
sample_sequence: int,
|
||||
timestamp_ms: int,
|
||||
items: Iterable[PublishItem],
|
||||
) -> bytes:
|
||||
values = tuple(items)
|
||||
_u16(subscription_id, "subscription_id")
|
||||
_u16(sample_sequence, "sample_sequence")
|
||||
_u32(timestamp_ms, "timestamp_ms")
|
||||
_u16(len(values), "item_count")
|
||||
chunks = [struct.pack("<HHIH", subscription_id, sample_sequence, timestamp_ms, len(values))]
|
||||
for item in values:
|
||||
_u32(item.address, "address")
|
||||
if not 1 <= item.element_count <= 0xFF:
|
||||
raise SetProtocolError("element_count вне диапазона 1..255")
|
||||
if not item.data or len(item.data) > 0xFFFF:
|
||||
raise SetProtocolError("элемент PUBLISH не содержит данные")
|
||||
chunks.append(
|
||||
struct.pack(
|
||||
"<IBBH",
|
||||
item.address,
|
||||
int(item.encoding),
|
||||
item.element_count,
|
||||
len(item.data),
|
||||
)
|
||||
+ item.data
|
||||
)
|
||||
payload = b"".join(chunks)
|
||||
if len(payload) > MAX_PAYLOAD_SIZE:
|
||||
raise SetProtocolError("PUBLISH не помещается в один payload")
|
||||
return payload
|
||||
|
||||
|
||||
def decode_publish(payload: bytes) -> tuple[int, int, int, tuple[PublishItem, ...]]:
|
||||
if len(payload) < 10:
|
||||
raise SetProtocolError("PUBLISH короче заголовка")
|
||||
subscription_id, sample_sequence, timestamp_ms, count = struct.unpack_from(
|
||||
"<HHIH", payload
|
||||
)
|
||||
offset = 10
|
||||
items: list[PublishItem] = []
|
||||
for _ in range(count):
|
||||
if offset + 8 > len(payload):
|
||||
raise SetProtocolError("PUBLISH обрывается в заголовке элемента")
|
||||
address, encoding, element_count, data_length = struct.unpack_from(
|
||||
"<IBBH", payload, offset
|
||||
)
|
||||
offset += 8
|
||||
if not element_count or not data_length or offset + data_length > len(payload):
|
||||
raise SetProtocolError("PUBLISH содержит повреждённый элемент")
|
||||
try:
|
||||
value_encoding = ValueEncoding(encoding)
|
||||
except ValueError as error:
|
||||
raise SetProtocolError(f"неизвестный encoding {encoding}") from error
|
||||
items.append(
|
||||
PublishItem(
|
||||
address=address,
|
||||
encoding=value_encoding,
|
||||
element_count=element_count,
|
||||
data=payload[offset : offset + data_length],
|
||||
)
|
||||
)
|
||||
offset += data_length
|
||||
if offset != len(payload):
|
||||
raise SetProtocolError("PUBLISH содержит хвост после элементов")
|
||||
return subscription_id, sample_sequence, timestamp_ms, tuple(items)
|
||||
Reference in New Issue
Block a user