Publisher¶
API¶
#[async_trait]
pub trait Publish {
async fn publish(&self, packet: Packet);
async fn forget_publish(&self, packet: Packet);
async fn wait_for_ack(&self, id: u64, timeout: Duration) -> Result<(), Box<dyn Error>>;
}
Конструкторы¶
// Стандартный (WAL: "publisher-{queue}.wal", memory limit: 100MB)
let pub = Publisher::new(address, Connection::publisher(queue)).await?;
// С кастомным конфигом
let pub = Publisher::new_with_config(
address,
Connection::publisher(queue),
"/tmp/my.wal",
200 * 1024 * 1024, // 200MB in-memory limit
).await?;
Конструктор:
- TCP connect
- Spawn
run_connection()с callbackhandle_incoming - Spawn
retry_task - Отправляет
Packet::connect()и ждёт ack
Возвращает Arc<Publisher>.
Методы публикации¶
publish(packet):
PacketType::Standard→ делегирует вforget_publish()(без трекинга)-
PacketType::Durable/Invulnerable: -
Push в outbox (отправка)
- Сохранение в
internal_queue(HashMap по nonce)- Если
in_memory_bytes + packet_size <= in_memory_limit→ в память - Иначе → WAL spill (append на диск с fsync)
- Если
- Insert в
pending_acksс текущим timestamp
forget_publish(packet):
- Просто
outbox.push(packet). Без трекинга, без retry.
wait_for_ack(nonce, timeout):
- Spin-wait: каждую 1ms проверяет, удалён ли nonce из
pending_acks - Timeout →
Err("Ack wait timed out")
Обработка входящих (callback в poll_fn)¶
Callback handle_incoming вызывается синхронно в фазе READ:
| Входящий пакет | Действие |
|---|---|
Ack(nonce) |
Удалить из pending_acks и internal_queue. Декрементировать in_memory_bytes |
Nack(nonce) |
Если пакет в памяти: re-push в outbox + reset timestamp. WAL-backed: оставить для retry_task |
Retry task¶
Отдельная таска, запускается при создании Publisher:
- Интервал проверки: 100ms
- Таймаут: 10s (
DEFAULT_TIMEOUT) - Для каждого nonce в
pending_acksсelapsed >= 10s: - Читает пакет из
internal_queue(или WAL) - Re-push в outbox
- Reset timestamp
- Если metadata отсутствует — удаляет из
pending_acks
Ограничения (TODO):
- Нет max retry count — бесконечные retries
- Нет exponential backoff — фиксированный интервал
WAL (Write-Ahead Log)¶
Publisher использует WAL для spill Durable пакетов, когда in_memory_bytes превышает лимит (default 100MB).
Формат записи: Length(4 bytes LE) | SerializedPacket(Length bytes).
Каждая запись завершается sync_data() (fsync). WalEntry { offset, len } хранится в PacketMetadata для последующего чтения.