Skip to content

Wire Protocol

Packet

Каждый пакет на проводе выглядит так:

block-beta
  columns 5
  a["Data Section<br/>(variable)"]
  b["PacketType<br/>1 byte"]
  c["Timestamp<br/>8 bytes"]
  d["Nonce<br/>8 bytes"]
  e["Seq<br/>8 bytes"]

Все целые числа — little-endian.

Поле Размер Описание
Data variable Сериализованные данные (см. ниже)
PacketType 1 byte 0=Standard, 1=Durable, 2=Invulnerable
Timestamp 8 bytes Наносекунды с UNIX_EPOCH. 0 для fast_ack и sequenced
Nonce 8 bytes Случайный u64 для корреляции. 0 для fast_ack и sequenced
Seq 8 bytes Sequence number. u64::MAX = нет (None на wire)

PacketType

Значение Имя Поведение
0 Standard Трекинг acks. Retry при timeout.
1 Durable WAL спилл при переполнении памяти
2 Invulnerable Зарезервирован для sync fsync persistence (не реализован)

Data Section

block-beta
  columns 3
  a["Payload Length<br/>4 bytes LE"]
  b["Type Byte<br/>1 byte"]
  c["Payload<br/>(variable)"]
Type Byte Variant Payload
0 Message Raw bytes (application data)
1 Ack u64 LE (nonce или seq)
2 Nack u64 LE (nonce или seq)
3 Ping Пусто
4 Connect ConnectionType (1 byte) + queue name (UTF-8)

ConnectionType: 0 = Publisher, 1 = Subscriber.

Framing (PacketCodec)

PacketCodec реализует tokio_util::codec::{Encoder, Decoder} для framing поверх TCP.

Decoder: считывает первые 4 байта — длину payload. Ждёт, пока в буфере не накопится 4 + payload_len + 25 байт (25 = PacketType + Timestamp + Nonce + Seq). Вызывает Packet::deserialize_from_bytes() — O(1) для Bytes (Arc bump, без копирования данных).

Encoder: вызывает packet.serialize_into(&mut BytesMut) — запись напрямую в буфер кодека, без промежуточных аллокаций.

Размер минимального пакета (Ping): 4 (len) + 1 (type) + 1 (pkt_type) + 8 (ts) + 8 (nonce) + 8 (seq) = 30 bytes.