Skip to content

Сервер

Конфигурация

pub struct ServerConfig {
    pub max_subscriber_retries: u32,        // default: 3
    pub subscriber_retry_timeout: Duration, // default: 10s
    pub subscriber_window_size: u64,        // default: 128
}

Все параметры задаются при старте. Hot reload не поддерживается.

Запуск

// Программный запуск
let handle = start_server("0.0.0.0:7427", ServerConfig::default()).await?;
println!("Listening on {}", handle.addr());

// Для тестов: случайный порт
let handle = start_server_random_port(ServerConfig::default()).await?;

// Shutdown
handle.shutdown().await;

ServerHandle содержит oneshot::Sender для сигнала shutdown и JoinHandle для ожидания завершения. Drop отправляет сигнал shutdown автоматически.

Внутренние структуры

classDiagram
  ServerData "1" --> "*" SubscriberInfo : queue name → subscribers
  SubscriberInfo "1" --> "1" SubscriberWindow : sliding window

  class ServerData {
    +DashMap~String, Vec~SubscriberInfo~~ subscribers
    +ServerConfig config
  }
  class SubscriberInfo {
    +SubscriberId id
    +Arc~NotifyChannel~Arc~Packet~~~ enqueue_inbox
    +Arc~Mutex~SubscriberWindow~~ window
  }
  class SubscriberWindow {
    +u64 next_seq
    +u64 window_base
    +u64 window_size
    +HashMap~u64, InflightEntry~ inflight
    +VecDeque~Arc~Packet~~ pending_send
    +Arc~NotifyChannel~Packet~~ outbox
  }
Поле Назначение
ServerData.subscribers queue name → подписчики
SubscriberInfo.id u64, атомарный счётчик
SubscriberInfo.enqueue_inbox Входящие сообщения
SubscriberWindow.next_seq Следующий seq для назначения
SubscriberWindow.window_base Первый неподтверждённый seq
SubscriberWindow.window_size Макс. inflight (default 128)
SubscriberWindow.inflight seq → пакет + timestamp + retries
SubscriberWindow.pending_send Очередь при полном окне
SubscriberWindow.outbox Для отправки подписчику

Connection flow

При подключении клиента сервер spawn'ит handle_connection() — свой poll_fn loop.

Publisher подключается:

  1. Приходит Data::Connect(Publisher, queue) → сервер запоминает client_queue
  2. Отправляет fast_ack

Subscriber подключается:

  1. Приходит Data::Connect(Subscriber, queue)
  2. Сервер создаёт SubscriberId (атомарный счётчик)
  3. Создаёт SubscriberWindow с window_size из конфига
  4. Spawn'ит две фоновые таски:
  5. subscriber_enqueue_task — переносит пакеты из enqueue_inbox в sliding window
  6. subscriber_retry_task — проверяет timeout каждые 500ms
  7. Регистрирует в DashMap<queue, Vec<SubscriberInfo>>
  8. Отправляет fast_ack

Message приходит от publisher:

  1. Сервер немедленно шлёт fast_ack паблишеру
  2. Оборачивает пакет в Arc<Packet>
  3. Push в enqueue_inbox каждого подписчика на очереди
  4. Логирует каждые 100K сообщений

Ack/Nack от subscriber:

  • Ack(seq)window.handle_ack(seq) → убирает из inflight, двигает base, дренит pending
  • Nack(seq)window.handle_nack(seq, max_retries) → resend или drop

Disconnect subscriber:

  • Удаляет из DashMap по SubscriberId

Sliding Window Protocol

Per-subscriber ordered delivery с backpressure.

block-beta
  columns 8
  s0["…"] s3["3"] s4["4"] s5["5"] s6["6"] s7["7"] s8["8"] s9["…"]
  space   b1<["window_base"]>(up) space:3 b2<["window_base + window_size"]>(up) space:2

  classDef inflight fill:#1f6feb,stroke:#8ab4f8,color:#fff
  classDef pending  fill:#5a6570,stroke:#9aa4b0,color:#fff
  class s3,s4,s5,s6 inflight
  class s7,s8 pending

Синие ячейки — inflight, серые — pending_send.

Метод Описание
enqueue(packet) Если окно не полное: назначает next_seq, добавляет в inflight, возвращает sequenced packet. Если полное: кладёт в pending_send, возвращает None
handle_ack(seq) Убирает из inflight. Вызывает advance_base() + drain_pending()
handle_nack(seq, max_retries) Если retries < max: resend (reset sent_at). Если retries >= max: drop из inflight, advance + drain
handle_timeout(max_retries, timeout) Находит все inflight с sent_at.elapsed() >= timeout, вызывает handle_nack для каждого
advance_base() Двигает window_base мимо подтверждённых seq (заполняет "дыры")
drain_pending() Переносит из pending_send в inflight, пока есть место в окне

Все методы синхронные — вызываются из poll_fn контекста, где .await недоступен.

Per-subscriber background tasks

subscriber_enqueue_task:

loop {
    batch = enqueue_inbox.recv().await     // ждёт пакеты от publisher'ов
    lock window
    for pkt in batch:
        if let Some(sequenced) = window.enqueue(pkt):
            outbox.push(sequenced)         // отправить подписчику
}

subscriber_retry_task:

loop {
    sleep(500ms)
    lock window
    retries = window.handle_timeout(max_retries, retry_timeout)
    for pkt in retries:
        outbox.push(pkt)                   // переотправить подписчику
}

DashMap: подводные камни

DashMap шардирует записи по хэшу ключа. Каждый шард — отдельный RwLock.

Антипаттерн: удержание shard lock через .await:

// ПЛОХО: lock держится через await
for entry in subscribers.iter() {
    entry.send(pkt).await;  // ← deadlock при нагрузке
}

Паттерн Sharp: collect-then-process:

// ХОРОШО: быстрый sync iter, потом работа без lock
let infos: Vec<_> = subscribers.get(&queue)
    .map(|v| v.clone())
    .unwrap_or_default();
// lock уже отпущен
for info in infos {
    info.enqueue_inbox.push(pkt.clone());
}