Skip to content

Архитектура

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)