Сервер¶
Конфигурация¶
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 подключается:
- Приходит
Data::Connect(Publisher, queue)→ сервер запоминаетclient_queue - Отправляет
fast_ack
Subscriber подключается:
- Приходит
Data::Connect(Subscriber, queue) - Сервер создаёт
SubscriberId(атомарный счётчик) - Создаёт
SubscriberWindowсwindow_sizeиз конфига - Spawn'ит две фоновые таски:
subscriber_enqueue_task— переносит пакеты изenqueue_inboxв sliding windowsubscriber_retry_task— проверяет timeout каждые 500ms- Регистрирует в
DashMap<queue, Vec<SubscriberInfo>> - Отправляет
fast_ack
Message приходит от publisher:
- Сервер немедленно шлёт
fast_ackпаблишеру - Оборачивает пакет в
Arc<Packet> - Push в
enqueue_inboxкаждого подписчика на очереди - Логирует каждые 100K сообщений
Ack/Nack от subscriber:
Ack(seq)→window.handle_ack(seq)→ убирает из inflight, двигает base, дренит pendingNack(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: