Skip to content

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?;

Конструктор:

  1. TCP connect
  2. Spawn run_connection() с callback handle_incoming
  3. Spawn retry_task
  4. Отправляет 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).

pub struct Wal {
    path: PathBuf,
    file: Mutex<File>,  // tokio async Mutex
}

Формат записи: Length(4 bytes LE) | SerializedPacket(Length bytes).

Каждая запись завершается sync_data() (fsync). WalEntry { offset, len } хранится в PacketMetadata для последующего чтения.