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.