Expand description
Worker runtime.
Contract:
Worker::<Q, B>::builder(backend)->WorkerBuilder.builder.handler(h)registers aJobHandlerwhoseJob::Queue == Q; duplicateJob::NAME->Error::DuplicateHandleratbuild().builder.queues(&[Q::A, Q::B])restricts consumption to a subset (default: all).builder.concurrency(n)caps in-flight jobs per worker process (default: sum of prefetch).builder.close_backend_on_shutdown(bool)(defaultfalse) closes the shared backend whenrun()returns.builder.build().await?declares queues and returns aWorker.worker.run()consumes untilWorkerHandle::shutdown()is called; graceful: stops the consumers first, then processes everything already pulled from the broker, then waits for the in-flight jobs. The backend is left open unlessclose_backend_on_shutdown(true)was set.worker.handle()->WorkerHandle(Clone) for shutdown from elsewhere, and forWorkerHandle::settle_failures().
Per-delivery algorithm:
- Look up handler by
envelope.job_type; if none ->dead_letter("no handler"). - Decode payload; on failure ->
dead_letter("decode error"). - Build
JobContext, run handler withtokio::time::timeoutif configured. - Ok ->
ack. JobError::Fatal->dead_letter.JobError::Deferred{delay, reason}-> logged atINFO, thendefer(env.deferred(priority), delay)wherepriorityis the job queue’smax_priority.unwrap_or(0). Nothing failed:attemptis unchanged, the retry policy is never consulted, onlydeferralsgrows.JobError::Retryable-> policy = job override or queue default;policy.decide(attempt);Retry{delay}->retry(env.next_attempt(), delay),GiveUp->dead_letter("max attempts").next_attemptputspriorityback to0, so a job that deferred earlier does not keep jumping the backlog on retries.- Handler panics are caught (
catch_unwindvia spawned task JoinError) and treated as Retryable.
Every step emits tracing events with job_id / job_type / attempt fields.
Structs§
- Worker
- Consumes one or more queues of
Qand dispatches jobs to registered handlers. - Worker
Builder - Configures a
Worker. Created byWorker::builder. - Worker
Handle - Remote control for a running
Worker. Cheap to clone and send across tasks.