Skip to main content

Module worker

Module worker 

Source
Expand description

Worker runtime.

Contract:

  • Worker::<Q, B>::builder(backend) -> WorkerBuilder.
  • builder.handler(h) registers a JobHandler whose Job::Queue == Q; duplicate Job::NAME -> Error::DuplicateHandler at build().
  • 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) (default false) closes the shared backend when run() returns.
  • builder.build().await? declares queues and returns a Worker.
  • worker.run() consumes until WorkerHandle::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 unless close_backend_on_shutdown(true) was set.
  • worker.handle() -> WorkerHandle (Clone) for shutdown from elsewhere, and for WorkerHandle::settle_failures().

Per-delivery algorithm:

  1. Look up handler by envelope.job_type; if none -> dead_letter("no handler").
  2. Decode payload; on failure -> dead_letter("decode error").
  3. Build JobContext, run handler with tokio::time::timeout if configured.
  4. Ok -> ack.
  5. JobError::Fatal -> dead_letter.
  6. JobError::Deferred{delay, reason} -> logged at INFO, then defer(env.deferred(priority), delay) where priority is the job queue’s max_priority.unwrap_or(0). Nothing failed: attempt is unchanged, the retry policy is never consulted, only deferrals grows.
  7. JobError::Retryable -> policy = job override or queue default; policy.decide(attempt); Retry{delay} -> retry(env.next_attempt(), delay), GiveUp -> dead_letter("max attempts"). next_attempt puts priority back to 0, so a job that deferred earlier does not keep jumping the backlog on retries.
  8. Handler panics are caught (catch_unwind via 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 Q and dispatches jobs to registered handlers.
WorkerBuilder
Configures a Worker. Created by Worker::builder.
WorkerHandle
Remote control for a running Worker. Cheap to clone and send across tasks.