Producing jobs
How an HTTP handler enqueues — the typed queue handle, the JobProducer publisher, and the service-as-choke-point rule.
Pushing a job is a one-line call — queue.push_to::<AudioQueue>(job).await?
— but where that line lives decides whether the feature stays auditable.
The rule is the same one that holds Repo calls behind a service:
controllers, resolvers, and gateways ask the service to enqueue. The
service names the queue; the transport adapter doesn’t.
The payload
Section titled “The payload”A job is any Serialize + DeserializeOwned + Clone + Send + Sync + 'static
type — the marker nest_rs::queue::Job is auto-implemented for every
type that satisfies those bounds. No derive, no trait to write. The
convention is to keep payloads small: a reference to the work, not the
work itself.
use nest_rs::queue::{QueueName, queue};use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Serialize, Deserialize)]pub struct TranscodeCommand { pub file: String,}
#[queue(name = "audio", job = TranscodeCommand)]pub struct AudioQueue;
pub const AUDIO_QUEUE: &str = <AudioQueue as QueueName>::NAME;Plain serde derives, not #[input] — the one place the wire-DTO
shorthand is the wrong tool. #[input] carries deny_unknown_fields, which is right at an edge a client
controls and wrong on a producer↔worker contract. A producer one deploy
ahead that adds a field would have every job it enqueues refused by the
older worker. And a decode failure aborts rather than retries, so the job
dead-letters on its first attempt.
Adding a field to a payload has to be survivable, and that means accepting one
you do not know yet.
A queue payload is named for the boundary it crosses: an imperative
“do X” handled by one processor is a Command (command.rs at the
feature port); a published past-tense fact with potentially many
consumers is an Event (event.rs). The HTTP body that triggers
the enqueue is a separate TranscodeDto — each boundary speaks its
own vocabulary.
Producer and consumer agree on both facts by importing AudioQueue:
the wire name and the payload type travel together, so a renamed queue or
a changed payload breaks the build on whichever side did not follow. The
#[queue] marker sits at the port beside the payload, not in the
consumer-side queue/ adapter — the producer is one directory up, and a
marker it cannot see is a marker it cannot use.
The service is the choke point
Section titled “The service is the choke point”A producer is any code that asks the service to enqueue. The HTTP controller, the GraphQL resolver, the scheduled producer — all of them delegate. They do not hold the producer and do not know the queue name.
use std::sync::Arc;use std::time::Duration;
use anyhow::Result;use nest_rs::core::injectable;use nest_rs::queue::{JobProducer, JobProducerExt};
use super::command::{AudioQueue, TranscodeCommand};
#[injectable]pub struct AudioService { #[inject] queue: Arc<dyn JobProducer>,}
impl AudioService { pub async fn enqueue_transcode(&self, file: String) -> Result<()> { self.queue .push_to::<AudioQueue>(TranscodeCommand { file: file.clone() }) .await?; tracing::info!(target: "features::audio", %file, "enqueued transcode job"); Ok(()) }
pub async fn transcode(&self, file: &str) -> Result<()> { tokio::time::sleep(Duration::from_millis(300)).await; tracing::info!(target: "features::audio", file, "transcoded"); Ok(()) }}push_to::<Q> takes both the wire name and the accepted payload from
Q, so enqueueing onto the wrong queue, or with the wrong payload, does
not compile. push_to comes from JobProducerExt, a blanket extension
over any producer — import the trait, and it is available on
Arc<dyn JobProducer> and on a backend’s concrete producer alike. It is async
because Redis is async; that is the whole producer surface from a
feature’s point of view. Inside the call, nest-rs-redis wraps your
payload in the wire envelope ({ "v": 1, "payload": ... }) so any other
backend can drain it, and seals your trace context into it so the job
runs under the trace that enqueued it. You never see either — see
Correlation.
The HTTP adapter
Section titled “The HTTP adapter”The controller is a thin transport. It accepts a TranscodeDto — the
REST body, a distinct type from the queue’s TranscodeCommand — calls
the service, returns the wire response. No queue knowledge of its own.
Declare the DTO with #[input]: a Json<T> handler argument requires
schemars::JsonSchema (see extractors), and #[input]
derives it along with Deserialize and Validate. A hand-derived DTO needs
#[derive(schemars::JsonSchema)] and the schemars dependency to go with it.
use nest_rs::http::input;
#[input]pub struct TranscodeDto { pub file: String,}use std::sync::Arc;
use nest_rs::http::{controller, routes};use nest_rs::http::poem::http::StatusCode;use nest_rs::http::poem::web::Json;use nest_rs::http::poem::{Error, Result};
use crate::audio::{AudioService, TranscodeDto};
#[controller(path = "/audio")]#[use_guards(AuthnGuard, AuthzGuard)]pub struct AudioController { #[inject] svc: Arc<AudioService>,}
#[routes]impl AudioController { #[post("/transcode")] async fn transcode(&self, body: Json<TranscodeDto>) -> Result<Json<TranscodeDto>> { let job = body.0; self.svc .enqueue_transcode(job.file.clone()) .await .map_err(|e| Error::from_string(e.to_string(), StatusCode::INTERNAL_SERVER_ERROR))?; Ok(Json(job)) }}Swap this for a #[resolver] mutation, a #[gateway] message, or a
#[scheduled] method on a timer — the controller is the only file that
changes. The service stays the single audited entry point, every
transport hits it through enqueue_transcode.
Why this rule
Section titled “Why this rule”Three reasons the push_to call hides behind a service:
- One audit point per queue.
AudioQueueis named in exactly one place per side. A grep forenqueue_transcodeshows every producer. - Pre-flight logic survives a swap of transport. If enqueueing needs an ability check, an idempotency key, a tracing field — the service holds it. Adding a second producer doesn’t duplicate that.
- The consumer can change without touching producers. Renaming the
queue, splitting
audiointoaudio.fastandaudio.slow, moving from Redis to SQS — one file changes, every adapter follows. And because the name lives on the#[queue]type, a rename that misses a call site is a compile error rather than a queue nobody drains.
The same shape the rest of the framework uses for Repo, just on the
push side. The service is the contract; transports adapt.
Reference
Section titled “Reference”crates/features/src/audio/command.rs—TranscodeCommand,AudioQueue.crates/features/src/audio/service.rs—enqueue_transcode.crates/features/src/audio/http/controller.rs— the HTTP producer.crates/features/src/audio/schedule/tasks.rs— the scheduled producer (a timer that calls the same service).crates/nest-rs-queue/src/queue_name.rs— theQueueNametrait the#[queue]macro implements.crates/nest-rs-queue/src/producer.rs—JobProducer,JobProducerExt::push_to.crates/nest-rs-redis/src/connection.rs— the RedisJobProducer.
Going further
Section titled “Going further”- Wiring — producer-only vs worker apps sharing a features crate.
- Retries and failure — what happens once a worker drains the job.
- Queue — the section landing and worker overview.