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);
}
Конструктор¶
- TCP connect
- Создаёт unbounded mpsc channel
- Spawn
run_connection()с callback: фильтруетData::Message, отправляет в channel - 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 — ответственность за доставку лежит на сервере.