Архитектура¶
Workspace¶
Cargo workspace (resolver = "3") из четырёх крейтов:
graph TD
W["sharp/ — Cargo workspace"]
W --> C["sharp/<br/>Core protocol library<br/><small>packet, codec, duplex, notify_channel</small>"]
W --> L["client/<br/>Publisher + Subscriber<br/><small>клиентская библиотека</small>"]
W --> S["server/<br/>MQ сервер<br/><small>lib + binary</small>"]
W --> B["bench/<br/>Инструмент бенчмаркинга"]
Двухуровневая модель надёжности¶
graph LR
P[Publisher] -->|"Level 1<br/>client resp."| S[Server]
S -->|"Level 2<br/>server resp."| Sub[Subscriber]
Level 1 — Publisher → Server. Клиент отвечает за доставку до сервера:
- Retry при timeout (проверка каждые 100ms, таймаут 10s)
- Немедленная переотправка при NACK
- Tracking pending acks в
IndexMap<u64, RetryQueue>
Level 2 — Server → Subscriber. Сервер отвечает за доставку подписчикам:
- Sliding window per subscriber (default size 128)
- Retry при timeout (проверка каждые 500ms, таймаут 10s)
- NACK handling с лимитом retries (default 3)
- Cleanup при disconnect подписчика
Ключевое решение: сервер ACK'ает паблишера сразу при получении. Паблишер не знает о подписчиках — это зона ответственности сервера.
Поток данных¶
sequenceDiagram
autonumber
participant P as Publisher
participant S as Server
participant Sub as Subscriber
P->>S: Connect(Pub, queue)
S-->>P: Ack
Sub->>S: Connect(Sub, queue)
S-->>Sub: Ack
P->>S: Message
S-->>P: fast_ack
S->>Sub: SequencedMessage(seq=0)
Sub-->>S: SeqAck(seq=0)
P->>S: Message
S-->>P: fast_ack
S->>Sub: SequencedMessage(seq=1)
Sub-->>S: SeqNack(seq=1)
Note over S,Sub: resend
S->>Sub: SequencedMessage(seq=1)
Sub-->>S: SeqAck(seq=1)