Writing a queue driver
A backend-agnostic scaffold for a queue driver — the traits you implement, the inventory you drain, with storage-specific code marked as your own.
The nest-rs-queue
crate is a seam: a JobProducer trait on the push side, a link-time
ProcessMethod inventory + a wire envelope on the consume side, and a
Transport slot the framework drives at boot. The seam names no storage
and no worker runtime — Redis, NATS JetStream, SQS, RabbitMQ, an
in-process channel, a Postgres queue table all plug in the same way.
This page gives you the shape of a driver, not its implementation:
the traits you implement, the registry you drain, with // your code here markers where the storage-specific work lands. The reference impl
is nest-rs-redis
(oxana on Redis) — read it once when a placeholder is unclear.
Install
Section titled “Install”cargo add nest-rs --features queueThe contract without a backend: redis would pull the first-party producer and worker, which is what you are replacing.
What you’ll implement
Section titled “What you’ll implement”- A
JobProducerimpl: wrap the wire envelope, push to your storage. - A
Transportimpl: a receive loop that calls the port’sconsume::attemptand translates its outcome into your backend’s ack / retry / dead-letter. - Three modules, three shapes — the substrate (
<Vendor>Module::for_root), the producer binding (bare), the worker binding (for_root).
A driver lives in a nest-rs-<technology> crate that depends on
nest-rs-core, nest-rs-queue, plus whatever your storage and runtime
need. An app picks a driver by importing its module instead of
nest-rs-redis.
The wire envelope
Section titled “The wire envelope”The one binding rule across backends — the producer wraps, the macro-emitted consumer unwraps, the storage in between only moves JSON.
{ "v": 1, "payload": <user job as JSON>, "traceparent": "00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01", "tracestate": "rojo=00f067aa0ba902b7", "actor_id": "01a0…"}You never write that shape by hand: envelope::seal(payload) builds it,
and envelope::open(value) takes the trace context back off, leaving the
{ v, payload } envelope the macro decodes. A queue is the one hop that
crosses a process, so a hand-rolled wrapper silently costs every job
its link back to the request that enqueued it — see
Correlation.
The macro on the #[process] side validates the version and fails closed
with a clear error on mismatch, so a producer that wraps a different
shape surfaces immediately.
Step 1 — The producer
Section titled “Step 1 — The producer”JobProducer::push_json is the only method the framework asks a
producer to implement. Per-job options, batching, scheduled jobs,
priority hints — those are your driver’s own API on the same struct.
use nest_rs::queue::async_trait;use nest_rs::queue::QueueError;
/// Your driver's connection handle. The producer side of a queue.#[derive(Clone)]pub struct YourConnection { // your storage handle — a Redis pool, an SQS client, a NATS // connection, a channel sender, a Postgres pool...}
#[async_trait]impl nest_rs::queue::JobProducer for YourConnection { async fn push_json(&self, queue: &str, payload: serde_json::Value) -> Result<(), QueueError> { // (1) Wrap in the wire envelope — required for the // macro-emitted consumer to decode it. `seal` stamps the // version and seals the ambient trace context into the value. let envelope = nest_rs::queue::envelope::seal(payload);
// (2) Push the envelope onto your storage for `queue`. // Your code: `self.client.send(queue, envelope).await?`, // `RPUSH`, `Publish`, `INSERT INTO jobs (...) VALUES (...)`, // `SendMessage`, ... Ok(()) }}The typed producer.push::<Job>(queue, job).await call sites see is a
blanket impl on top of push_json — no extra code on your side.
Step 2 — The worker (Transport)
Section titled “Step 2 — The worker (Transport)”This is the meat, and it is smaller than it used to be. The worker is a
nest_rs::core::Transport: the framework calls configure with the assembled
container at boot, then serve with a cancellation token. Two calls into the port do everything that is not your storage’s.
consume::discover drains the #[process] inventory — module-gated,
duplicate-checked, announced. consume::attempt runs one attempt: it opens
the envelope, continues or mints the trace, opens the queue.job span and
the ambient scope, catches a panic, classifies the outcome, and files the
events and the nest_rs::operation line.
You never restate any of that: an adapter that opens a queue.job span of its
own has taken semantics it does not own, and the conformance suite says so.
What your code owns is the transport: how envelopes arrive (block-pop,
subscribe, long-poll, ReceiveMessage), how retries and visibility timeouts
wire up, and the translation of an Attempt into your backend’s vocabulary.
Those are storage decisions.
One contract is not a storage decision: a #[process] method runs
one job at a time. Serialize per method in
your receive loop, the way nest-rs-redis does, so a handler behaves the
same whichever driver is underneath it.
use nest_rs::core::{Container, Transport};use nest_rs::queue::async_trait;use nest_rs::queue::consume::{self, Attempt};use nest_rs::queue::ProcessMethod;use tokio_util::sync::CancellationToken;
pub struct YourWorker { // `serve` takes only a cancellation token, so everything it needs is // captured in `configure`: the reachable methods and the container. methods: Vec<&'static ProcessMethod>, container: Option<Container>,}
#[async_trait]impl Transport for YourWorker { async fn configure(&mut self, container: &Container) -> anyhow::Result<()> { // (1) Which methods this app serves is the port's answer. self.methods = consume::discover(container)?; // Fail fast here if your connection is not in the container — name the // module that opens it, the way `RedisWorker` names `RedisModule`. self.container = Some(container.clone()); Ok(()) }
async fn serve(self: Box<Self>, cancel: CancellationToken) -> anyhow::Result<()> { let container = self.container.expect("configure runs before serve");
// (2) One receive loop per method, serialized. Your code: replace // `your_pop_from_queue(...)` with your storage's receive primitive — // BRPOP, subscribe, ReceiveMessage, NOTIFY/LISTEN, a channel `recv`… // Respect `cancel` for graceful shutdown. for method in self.methods { let container = container.clone(); let cancel = cancel.clone(); tokio::spawn(async move { loop { tokio::select! { _ = cancel.cancelled() => break, (job_id, attempt_no, payload) = /* your_pop_from_queue(method.queue) */ => { // (3) The attempt is the port's — envelope, trace, span, // panic, classification, events and line included. match consume::attempt(method, payload, job_id, attempt_no, container.clone()).await { // (4) What you decide: how each outcome maps onto your // backend — ack / retry-later / dead-letter. Attempt::Ok => { /* ack */ } Attempt::Retry(error) => { /* nak / re-queue within `method.retries` */ let _ = error; } Attempt::DeadLetter(error) => { /* term / dead-letter now */ let _ = error; } } } } } }); }
Ok(()) }}Attempt::Retry is the user method’s Err — re-run it within
method.retries. Attempt::DeadLetter is deterministic (an undeserializable
payload, a pipe rejection, a missing provider, a panic): re-running would burn
the budget on a payload that cannot succeed, so dead-letter it at once. Inside
attempt the macro-emitted handler already did the serde_json::from_value to
the user’s job type, the Container::get of the #[injectable] host and the
run_in_job_context wrap. Your code is the transport. The port is the
dispatcher.
Step 3 — The modules
Section titled “Step 3 — The modules”Three modules, three shapes — the same three every adapter in the framework
has, so an app that has read nest-rs-redis’s composition root has read yours:
| Shape | Yours | Role | Variables |
|---|---|---|---|
<Vendor>Module::for_root(cfg) | NatsModule::for_root(None) | opens the one connection every binding in your crate shares | NESTRS_NATS__* |
<Vendor><Port>Module (bare) | NatsQueueModule | binds Arc<dyn JobProducer> over that connection | — |
<Vendor><Port>Module::for_root(cfg) | NatsWorkerModule::for_root(None) | contributes the Transport; its own settings, if any | NESTRS_NATS__WORKER__* |
The layout follows: src/config.rs, src/connection.rs, src/module.rs at the
crate root; queue/producer.rs + queue/module.rs; worker/consumer.rs +
worker/config.rs + worker/module.rs. A namespace is the stem of its path —
nats at the root, nats__worker in the binding folder — which is what keeps
your variables findable from your types and your types from your variables.
use std::sync::Arc;
use nest_rs::config::{ConfigModule, ConfigSetup};use nest_rs::core::{ContainerBuilder, DynamicModule, Module, TransportContribution};use nest_rs::queue::JobProducer;
/// The substrate: opens `YourConnection` from `YourConfig` (env + pinned).pub struct YourModule;
impl YourModule { pub fn for_root(config: impl Into<Option<YourConfig>>) -> YourSetup { YourSetup { pinned: config.into() } }}
pub struct YourSetup { pinned: Option<YourConfig>,}
impl DynamicModule for YourSetup { fn collect(&self, builder: ContainerBuilder) -> ContainerBuilder { let builder = ConfigModule::provide_feature(self.pinned.clone(), builder); builder.provide_factory::<YourConnection, _, _>(|container| async move { let config = container.get::<YourConfig>().expect("resolved above"); YourConnection::connect(&config).await }) }}
/// The producer binding — bare. Declares it runs *after* the connection's/// factory, so `imports` order stays a readability choice.pub struct YourQueueModule;
impl Module for YourQueueModule { fn register(builder: ContainerBuilder) -> ContainerBuilder { builder } fn collect(builder: ContainerBuilder) -> ContainerBuilder { builder.provide_factory_dyn_after::<YourQueueProducer, dyn JobProducer, YourConnection, _, _>( |container| async move { let conn = container.get::<YourConnection>() .ok_or_else(|| anyhow::anyhow!("import YourModule::for_root(None) beside YourQueueModule"))?; Ok(YourQueueProducer::new((*conn).clone())) }, |producer| Arc::new(producer) as Arc<dyn JobProducer>, ) }}
/// The consumer binding. Contributes the `Transport`; `for_root` because it/// owns settings of its own (`YourWorkerConfig`, namespace `<vendor>__worker`).pub struct YourWorkerModule;
impl YourWorkerModule { pub fn for_root(config: impl Into<Option<YourWorkerConfig>>) -> YourWorkerSetup { ConfigModule::setup(config) }}
pub type YourWorkerSetup = ConfigSetup<YourWorkerModule, YourWorkerConfig>;
impl Module for YourWorkerModule { fn register(builder: ContainerBuilder) -> ContainerBuilder { builder.provide_meta(TransportContribution { name: "YourWorker", build: |_| Ok(Box::new(YourWorker { methods: Vec::new(), container: None })), }) }}Every field of YourConfig is settable both via NESTRS_<NAMESPACE>__<KEY>
and pinned in code — the framework-wide dual-path rule.
Step 4 — Using the integration
Section titled “Step 4 — Using the integration”Application code stays put. The swap is three lines in the app’s
module.rs, and they read like nest-rs-redis’s:
// Producer-only API binary:#[module(imports = [YourModule::for_root(None), YourQueueModule, /* features */])]pub struct ApiModule;
// Worker binary (consumes too):#[module(imports = [ YourModule::for_root(None), YourQueueModule, YourWorkerModule::for_root(None), /* features that ship #[processor]s */])]pub struct WorkerModule;The #[processor] impl ... blocks in user code are unchanged — same
impl, same #[process(queue = <Name>Queue)], same method signatures — and the
features that enqueue inject Arc<dyn JobProducer>, so they never knew which
backend they were on.
Ephemeral vs durable backends
Section titled “Ephemeral vs durable backends”The same three traits cover both shapes:
- Ephemeral (in-memory channels, in-process pub/sub): jobs vanish on restart, no cross-process delivery. Useful for tests and local CLIs where the queue’s lifetime equals the process’s.
- Durable (Redis, Postgres, SQS, RabbitMQ, NATS JetStream): retries, visibility timeouts, dead-letter queues. Your driver decides how to expose those — the framework takes no position.
No tradeoff is hidden. Both shapes plug into the same JobProducer and
Transport.
What you don’t get for free
Section titled “What you don’t get for free”The framework gives you the contract — Job, Processor,
JobHandler, WIRE_FORMAT_VERSION, the Transport slot, the
JobContext ambient executor for handlers that touch a DB.
Retry semantics, batching, priorities, scheduled jobs, ordering,
visibility timeouts, dead-letter routing — those are your driver’s
contract with its storage. Decide them deliberately, document them in
the driver’s README.
Reference
Section titled “Reference”crates/nest-rs-redis/— the reference impl; read itsconnection.rs+module.rs(the shared connection),queue/producer.rs+queue/module.rs(the producer) andworker/consumer.rs+worker/module.rs(the transport) whenever a placeholder above needs a concrete shape.crates/nest-rs-queue/README.md— the extension contract, theProcessMethodregistry, the envelope spec.- Database driver — the analogous scaffold for the data layer, same shape, same seams.
Going further
Section titled “Going further”- Wiring — the producer/worker module split a driver plugs into.
- Retries and failure — the retry budget your driver owns.
- Database driver — the analogous data-layer scaffold.