Skip to content

Subscriber

API

#[async_trait]
pub trait Subscribe {
    async fn subscribe(&self) -> UnboundedReceiver<Packet>;
    async fn ack(&self, packet: &Packet);
    async fn nack(&self, packet: &Packet);
}

Конструктор

let sub = Subscriber::new(address, Connection::subscriber(queue)).await?;
  1. TCP connect
  2. Создаёт unbounded mpsc channel
  3. Spawn run_connection() с callback: фильтрует Data::Message, отправляет в channel
  4. Push Packet::connect() в outbox (без ожидания ack)

Возвращает Arc<Subscriber>.

Получение сообщений

let mut rx = subscriber.subscribe().await;
// subscribe() можно вызвать только ОДИН раз — забирает Option<Receiver>

while let Some(packet) = rx.recv().await {
    // обработка
    subscriber.ack(&packet).await;
}

Ack/Nack

Два режима, выбирается автоматически:

Условие Ack отправляет Nack отправляет
packet.seq.is_some() Packet::seq_ack(seq) Packet::seq_nack(seq)
packet.seq.is_none() packet.ack() (nonce-based) packet.nack() (nonce-based)

Sequenced пакеты — от сервера с sliding window. Nonce-based — legacy или direct connections.

Отсутствие retry

Subscriber не имеет retry логики. Если сервер не получил ack (таймаут), он сам переотправит сообщение. Это deliberate design — ответственность за доставку лежит на сервере.