672 lines
22 KiB
Python
672 lines
22 KiB
Python
"""Эталонная 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)
|