Skip to content

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.

Terminal window
cargo add nest-rs --features queue

The contract without a backend: redis would pull the first-party producer and worker, which is what you are replacing.

  • A JobProducer impl: wrap the wire envelope, push to your storage.
  • A Transport impl: a receive loop that calls the port’s consume::attempt and 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 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.

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.

crates/yourqueue/src/queue/producer.rs
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.

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.

crates/yourqueue/src/worker/transport.rs
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.

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:

ShapeYoursRoleVariables
<Vendor>Module::for_root(cfg)NatsModule::for_root(None)opens the one connection every binding in your crate sharesNESTRS_NATS__*
<Vendor><Port>Module (bare)NatsQueueModulebinds Arc<dyn JobProducer> over that connection—
<Vendor><Port>Module::for_root(cfg)NatsWorkerModule::for_root(None)contributes the Transport; its own settings, if anyNESTRS_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.

the three modules — substrate, producer, consumer
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.

Application code stays put. The swap is three lines in the app’s module.rs, and they read like nest-rs-redis’s:

two binaries, one integration
// 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.

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.

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.

  • crates/nest-rs-redis/ — the reference impl; read its connection.rs + module.rs (the shared connection), queue/producer.rs + queue/module.rs (the producer) and worker/consumer.rs + worker/module.rs (the transport) whenever a placeholder above needs a concrete shape.
  • crates/nest-rs-queue/README.md — the extension contract, the ProcessMethod registry, the envelope spec.
  • Database driver — the analogous scaffold for the data layer, same shape, same seams.