Skip to main content

Crate queuey_rabbitmq

Crate queuey_rabbitmq 

Source
Expand description

RabbitMQ backend for queuey, built on lapin.

RabbitMqBackend implements queuey_core::Backend: it owns one AMQP connection, publishes with publisher confirms, and hands the worker runtime a stream of RabbitMqDelivery values.

use std::{sync::Arc, time::Duration};

use queuey_core::{Backend, QueueConfig};
use queuey_rabbitmq::{RabbitMqBackend, RabbitMqOptions};

let backend = RabbitMqBackend::with_options(
    "amqp://guest:guest@localhost:5672/%2f",
    RabbitMqOptions::default().retry_granularity(Duration::from_secs(5)),
)
.await?;

let emails = QueueConfig::new("myapp.emails").prefetch(10);
backend.declare(std::slice::from_ref(&emails)).await?;

let backend = Arc::new(backend);
// ... hand `backend` to a `Producer` / `Worker` ...
backend.close().await?;

§Topology

Each logical queue q is backed by two long-lived broker queues, q and q.dead, plus a short-lived hold queue q.deferred.{ttl_ms} per distinct delay. Every wait, whether a retry backoff, a delayed enqueue or a deferral, happens in a hold queue. See topology for the exact arguments.

§Why hold queues, and not one wait queue with per-message expirations

RabbitMQ only expires the message at the head of a classic queue. In a shared wait queue, a message with a five-minute expiration at the head holds back every one-second expiration queued behind it, and exponential backoff produces exactly that mix of delays. So instead the delay is part of the queue name, the wait is the queue-wide x-message-ttl, and every message in q.deferred.30000 expires in publish order. A short wait is never stuck behind a long one, because the two live in different queues.

Delays are rounded up to a granularity to bound how many hold queues exist at once: RabbitMqOptions::retry_granularity for retries and Producer::enqueue_after, RabbitMqOptions::deferred_granularity for deferrals, both 1s by default. The hold queue is declared on demand right before each publish: an idle hold queue deletes itself one TTL after the last publish to it (x-expires = 2 * TTL), and every declare resets that timer.

§Retry versus deferral

Both wait in the same hold queues. They differ in what happens when the job is back on q:

  • A retry (Delivery::retry, and a delayed Backend::publish) carries priority 0 and joins the back of the queue like any other message. Its attempt counter has been incremented.
  • A deferral (Backend::defer, Delivery::defer) is what a 429 Too Many Requests with Retry-After: 30 needs: the job did not fail, must not burn an attempt, and must run ahead of the backlog when it returns. q is declared with x-max-priority from QueueConfig::max_priority (default Some(10)), every publish carries the envelope’s priority, and a deferred envelope carries the queue’s top level, so it is served before everything that piled up meanwhile. “Ahead of the backlog” means ahead of what is still on the queue: a consumer with prefetch N already holds up to N backlog messages, and the returning deferral is first among what is left.

§What a hold requires

  • The queue must have been declared through this backend, in this process. Otherwise the hold queue’s durability and the queue it dead-letters back to would be guesses, and a TTL expiry into a queue that does not exist is discarded silently by the broker. Unlike a mandatory publish, nothing comes back and nothing is logged. Holding onto an unknown queue is Error::UnknownQueue instead. Producer::new and WorkerBuilder::build declare the queue set; Producer::new_undeclared deliberately does not, so a producer built that way can enqueue, but not enqueue with a delay or defer.
  • The delay must fit. It is capped at MAX_DEFERRAL_MS, about 24.8 days. That is half of what a 32-bit millisecond TTL can express, because the hold queue’s x-expires is twice its TTL. A longer delay is refused rather than clamped: releasing a job early is the one thing a hold promises not to do. Rounding up to the granularity happens first, so a delay just under the cap can be refused too.

Both failures happen before anything is acked, so from Delivery::retry and Delivery::defer they leave the original message unacknowledged and the broker redelivers it.

§Upgrading from the q.retry wait queue

Earlier versions declared a q.retry queue per work queue and published retries into it with a per-message expiration. This version neither declares nor uses it. Nothing needs migrating: messages still waiting in an existing q.retry expire back onto q on their own, because the dead-letter routing is an argument of that queue, and workers running the old version keep declaring it themselves. Delete q.retry once it is empty and no old worker is left. RabbitMqOptions::retry_suffix is gone with it; RabbitMqOptions::retry_granularity is the retry tunable now.

§Breaking topology change

x-max-priority is a declaration argument, and RabbitMQ refuses to change the arguments of a queue that already exists: the declaration is answered with PRECONDITION_FAILED, which closes the channel and surfaces here as an error from declare.

A q created before this feature has no x-max-priority, so declaring it again with the default config will fail. Either:

  • drain and delete q, then let this backend redeclare it. Deferred jobs then come back ahead of the backlog; or
  • set max_priority = 0 on the queue’s QueueConfig (or #[queue(max_priority = 0)]), which declares q exactly as before. Deferral still works, it just returns jobs FIFO instead of ahead of the queue.

q.dead and hold queues are unchanged, so only q is affected.

§Guarantees

  • Every publish (enqueue, retry, defer, dead-letter) is mandatory and confirmed by the broker before it is reported as successful. A routing key that matches no queue is returned by the broker and reported as an error rather than passing for a confirmed publish.
  • Delivery::retry, Delivery::defer and Delivery::dead_letter publish first and ack second, and skip the ack entirely when the publish fails, so a job is never lost. At worst it is redelivered. Different delays never block each other: each waits in its own hold queue. With RabbitMqOptions::declare_dead_letter_queues off, dead_letter rejects the delivery instead of publishing to a q.dead this backend does not own.
  • Messages whose body is not a valid queuey_core::Envelope are moved aside and logged, never surfaced as a stream error, so one poison message cannot stall a consumer.

§Not in scope for v1

Reconnection. When the connection drops, consumer streams end and further calls fail; supervising and rebuilding the backend is the caller’s job.

Re-exports§

pub use lapin;

Modules§

codec
Pure mapping between Envelope and AMQP BasicProperties / headers.
topology
Pure helpers describing the RabbitMQ topology used by this backend.

Structs§

RabbitMqBackend
A Backend backed by a single RabbitMQ connection.
RabbitMqDelivery
One RabbitMQ message, decoded into an Envelope.
RabbitMqOptions
Configuration for RabbitMqBackend::with_options.