Skip to main content

queuey_core/
worker.rs

1//! Worker runtime.
2//!
3//! Contract:
4//! * `Worker::<Q, B>::builder(backend)` -> `WorkerBuilder`.
5//! * `builder.handler(h)` registers a `JobHandler` whose `Job::Queue == Q`;
6//!   duplicate `Job::NAME` -> `Error::DuplicateHandler` at `build()`.
7//! * `builder.queues(&[Q::A, Q::B])` restricts consumption to a subset (default: all).
8//! * `builder.concurrency(n)` caps in-flight jobs per worker process (default: sum of prefetch).
9//! * `builder.close_backend_on_shutdown(bool)` (default `false`) closes the shared backend
10//!   when `run()` returns.
11//! * `builder.build().await?` declares queues and returns a `Worker`.
12//! * `worker.run()` consumes until `WorkerHandle::shutdown()` is called; graceful:
13//!   stops the consumers first, then processes everything already pulled from the broker,
14//!   then waits for the in-flight jobs. The backend is left open unless
15//!   `close_backend_on_shutdown(true)` was set.
16//! * `worker.handle()` -> `WorkerHandle` (Clone) for shutdown from elsewhere, and for
17//!   `WorkerHandle::settle_failures()`.
18//!
19//! Per-delivery algorithm:
20//! 1. Look up handler by `envelope.job_type`; if none -> `dead_letter("no handler")`.
21//! 2. Decode payload; on failure -> `dead_letter("decode error")`.
22//! 3. Build `JobContext`, run handler with `tokio::time::timeout` if configured.
23//! 4. Ok -> `ack`.
24//! 5. `JobError::Fatal` -> `dead_letter`.
25//! 6. `JobError::Deferred{delay, reason}` -> logged at `INFO`, then
26//!    `defer(env.deferred(priority), delay)` where `priority` is the job queue's
27//!    `max_priority.unwrap_or(0)`. Nothing failed: `attempt` is unchanged, the retry
28//!    policy is never consulted, only `deferrals` grows.
29//! 7. `JobError::Retryable` -> policy = job override or queue default;
30//!    `policy.decide(attempt)`; `Retry{delay}` -> `retry(env.next_attempt(), delay)`,
31//!    `GiveUp` -> `dead_letter("max attempts")`. `next_attempt` puts `priority` back to
32//!    `0`, so a job that deferred earlier does not keep jumping the backlog on retries.
33//! 8. Handler panics are caught (`catch_unwind` via spawned task JoinError) and treated as Retryable.
34//!
35//! Every step emits `tracing` events with job_id / job_type / attempt fields.
36
37use std::{
38    collections::HashMap,
39    marker::PhantomData,
40    sync::{
41        Arc,
42        atomic::{AtomicU64, Ordering},
43    },
44    time::Duration,
45};
46
47use async_trait::async_trait;
48use futures::StreamExt;
49use tokio::{
50    sync::{Semaphore, mpsc, watch},
51    task::JoinSet,
52};
53use tracing::{Instrument, debug, error, info, warn};
54
55use crate::{
56    backend::{Backend, Delivery, DeliveryStream},
57    envelope::{Envelope, now_ms},
58    error::{Error, JobError, Result},
59    handler::{JobContext, JobHandler},
60    job::Job,
61    queue::{QueueConfig, QueueSet},
62    retry::{RetryDecision, RetryPolicy},
63};
64
65/// What happened while running one job.
66#[derive(Debug)]
67enum JobOutcome {
68    /// The handler returned `Ok(())`.
69    Success,
70    /// The payload could not be deserialized into the handler's job type.
71    Decode(String),
72    /// Transient failure (including panics and timeouts); consult the retry policy.
73    Retryable(String),
74    /// Permanent failure; dead-letter without consulting the retry policy.
75    Fatal(String),
76    /// Not a failure: hold the job for `delay` and run it again, attempt unchanged.
77    Deferred {
78        /// How long the job asked to wait.
79        delay: Duration,
80        /// Why it was deferred; for logs only.
81        reason: String,
82    },
83}
84
85/// Type-erased [`JobHandler`], so handlers for different job types can live in one map.
86#[async_trait]
87trait ErasedHandler: Send + Sync + 'static {
88    /// Effective policy: the job override if there is one, else the queue default.
89    fn policy(&self) -> RetryPolicy;
90
91    /// Priority a deferred envelope of this job is republished with: the highest level
92    /// its queue supports, or `0` when the queue is not a priority queue.
93    fn defer_priority(&self) -> u8;
94
95    /// Decode and run the job, never panicking and never returning an error type.
96    async fn run(&self, envelope: &Envelope, timeout: Option<Duration>) -> JobOutcome;
97}
98
99/// Adapter from a concrete [`JobHandler`] to [`ErasedHandler`].
100struct ErasedJobHandler<H: JobHandler> {
101    handler: Arc<H>,
102}
103
104#[async_trait]
105impl<H: JobHandler> ErasedHandler for ErasedJobHandler<H> {
106    fn policy(&self) -> RetryPolicy {
107        <H::Job as Job>::retry_policy().unwrap_or_else(|| <H::Job as Job>::QUEUE.config().retry)
108    }
109
110    fn defer_priority(&self) -> u8 {
111        <H::Job as Job>::QUEUE.config().max_priority.unwrap_or(0)
112    }
113
114    async fn run(&self, envelope: &Envelope, timeout: Option<Duration>) -> JobOutcome {
115        let job = match envelope.decode::<H::Job>() {
116            Ok(job) => job,
117            Err(e) => return JobOutcome::Decode(e.to_string()),
118        };
119
120        let ctx = JobContext {
121            job_id: envelope.job_id,
122            job_type: <H::Job as Job>::NAME,
123            queue: <H::Job as Job>::QUEUE.name(),
124            attempt: envelope.attempt,
125            max_attempts: self.policy().max_attempts,
126            deferrals: envelope.deferrals,
127            priority: envelope.priority,
128            age: Duration::from_millis(now_ms().saturating_sub(envelope.enqueued_at_ms)),
129        };
130
131        // Spawning turns a handler panic into a `JoinError` instead of unwinding the
132        // worker, and gives us something to abort when the job times out.
133        let handler = self.handler.clone();
134        let mut task = tokio::spawn(async move { handler.handle(job, ctx).await });
135
136        let joined = match timeout {
137            Some(limit) => match tokio::time::timeout(limit, &mut task).await {
138                Ok(joined) => joined,
139                Err(_) => {
140                    task.abort();
141                    // Wait for the cancellation to land: the concurrency permit is
142                    // released when this function returns, so the job must be gone by
143                    // then or the cap could be exceeded.
144                    let _ = task.await;
145                    return JobOutcome::Retryable(format!("job timed out after {limit:?}"));
146                }
147            },
148            None => (&mut task).await,
149        };
150
151        match joined {
152            Ok(Ok(())) => JobOutcome::Success,
153            Ok(Err(JobError::Retryable(e))) => JobOutcome::Retryable(e.to_string()),
154            Ok(Err(JobError::Fatal(e))) => JobOutcome::Fatal(e.to_string()),
155            Ok(Err(JobError::Deferred { delay, reason })) => JobOutcome::Deferred { delay, reason },
156            Err(e) if e.is_panic() => JobOutcome::Retryable("handler panicked".to_owned()),
157            Err(e) => JobOutcome::Retryable(format!("handler task failed: {e}")),
158        }
159    }
160}
161
162type HandlerMap = HashMap<&'static str, Arc<dyn ErasedHandler>>;
163
164/// Consumes one or more queues of `Q` and dispatches jobs to registered handlers.
165///
166/// Build one with [`Worker::builder`]; run it with [`Worker::run`].
167pub struct Worker<Q: QueueSet, B: Backend> {
168    backend: Arc<B>,
169    handlers: Arc<HandlerMap>,
170    queues: Vec<QueueConfig>,
171    concurrency: usize,
172    job_timeout: Option<Duration>,
173    close_backend: bool,
174    handle: WorkerHandle,
175    _q: PhantomData<fn(Q)>,
176}
177
178impl<Q: QueueSet, B: Backend> std::fmt::Debug for Worker<Q, B> {
179    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
180        f.debug_struct("Worker")
181            .field(
182                "queues",
183                &self.queues.iter().map(|c| &c.name).collect::<Vec<_>>(),
184            )
185            .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
186            .field("concurrency", &self.concurrency)
187            .field("job_timeout", &self.job_timeout)
188            .field("close_backend_on_shutdown", &self.close_backend)
189            .finish()
190    }
191}
192
193impl<Q: QueueSet, B: Backend> Worker<Q, B> {
194    /// Start configuring a worker for queue set `Q` on `backend`.
195    pub fn builder(backend: Arc<B>) -> WorkerBuilder<Q, B> {
196        WorkerBuilder {
197            backend,
198            handlers: Vec::new(),
199            queues: None,
200            concurrency: None,
201            job_timeout: None,
202            close_backend: false,
203        }
204    }
205
206    /// A clonable handle that can shut this worker down from anywhere.
207    pub fn handle(&self) -> WorkerHandle {
208        self.handle.clone()
209    }
210
211    /// The queues this worker consumes.
212    pub fn queues(&self) -> &[QueueConfig] {
213        &self.queues
214    }
215
216    /// Maximum number of jobs this worker runs at the same time.
217    pub fn concurrency(&self) -> usize {
218        self.concurrency
219    }
220
221    /// Consume and dispatch until [`WorkerHandle::shutdown`] is called.
222    ///
223    /// Shutdown is graceful, and in this order:
224    ///
225    /// 1. the consumers stop pulling new deliveries from the broker;
226    /// 2. everything they already pulled is processed normally (a delivery that has
227    ///    been handed over is never dropped unsettled);
228    /// 3. the jobs already running are awaited.
229    ///
230    /// The backend is *not* closed unless
231    /// [`WorkerBuilder::close_backend_on_shutdown`] was set: it is usually an
232    /// `Arc` shared with a [`crate::Producer`] that outlives the worker.
233    ///
234    /// Returns [`Error::ConsumerStopped`] if a consumer stream ends on its own (which
235    /// means the backend went away), or the first backend error seen on a stream.
236    pub async fn run(self) -> Result<()> {
237        let Worker {
238            backend,
239            handlers,
240            queues,
241            concurrency,
242            job_timeout,
243            close_backend,
244            handle,
245            ..
246        } = self;
247
248        info!(
249            queues = queues.len(),
250            concurrency,
251            handlers = handlers.len(),
252            "worker starting"
253        );
254
255        let mut outcome: Result<()> = Ok(());
256        let (tx, mut rx) = mpsc::channel::<ConsumerEvent>(queues.len().max(1));
257        // Tells the consumer tasks to stop pulling. Separate from the user-facing
258        // shutdown signal because every exit from the dispatch loop stops them.
259        let (stop_tx, stop_rx) = watch::channel(false);
260        let mut consumers = JoinSet::new();
261        for config in &queues {
262            match backend.consume(config).await {
263                Ok(stream) => {
264                    consumers.spawn(consume_into(
265                        stream,
266                        tx.clone(),
267                        config.name.clone(),
268                        stop_rx.clone(),
269                    ));
270                }
271                Err(e) => {
272                    error!(queue = %config.name, error = %e, "failed to start consumer");
273                    outcome = Err(e);
274                    break;
275                }
276            }
277        }
278        drop(tx);
279        drop(stop_rx);
280
281        let permits = Arc::new(Semaphore::new(concurrency));
282        let mut in_flight: JoinSet<()> = JoinSet::new();
283        let mut shutdown = handle.subscribe();
284        let settle_failures = handle.settle_failures.clone();
285
286        while outcome.is_ok() {
287            if *shutdown.borrow_and_update() {
288                break;
289            }
290
291            // Take a concurrency slot before pulling, so a delivery is never held
292            // while we wait for capacity.
293            let permit = tokio::select! {
294                biased;
295                _ = shutdown.changed() => break,
296                permit = permits.clone().acquire_owned() => match permit {
297                    Ok(permit) => permit,
298                    Err(_) => break,
299                },
300            };
301
302            let event = tokio::select! {
303                biased;
304                _ = shutdown.changed() => break,
305                event = rx.recv() => event,
306            };
307
308            match event {
309                // Every consumer task is gone.
310                None => break,
311                Some(ConsumerEvent::Ended(queue)) => {
312                    error!(%queue, "consumer stream ended unexpectedly");
313                    outcome = Err(Error::ConsumerStopped(queue));
314                }
315                Some(ConsumerEvent::Failed(e)) => {
316                    error!(error = %e, "consumer stream failed");
317                    outcome = Err(e);
318                }
319                Some(ConsumerEvent::Delivery(delivery)) => {
320                    // Reap finished jobs so the JoinSet does not grow without bound.
321                    while in_flight.try_join_next().is_some() {}
322                    in_flight.spawn(process(
323                        delivery,
324                        handlers.clone(),
325                        job_timeout,
326                        settle_failures.clone(),
327                        permit,
328                    ));
329                }
330            }
331        }
332
333        // Stop the consumers *before* the channel is drained, so nothing else can be
334        // pulled from a stream, and drain only afterwards: a delivery that reached the
335        // channel has already left the broker's hands and must still be settled.
336        debug!("worker draining");
337        stop_tx.send_replace(true);
338
339        loop {
340            let Ok(permit) = permits.clone().acquire_owned().await else {
341                break;
342            };
343            match rx.recv().await {
344                // Every consumer task has finished and the channel is empty.
345                None => break,
346                Some(ConsumerEvent::Delivery(delivery)) => {
347                    while in_flight.try_join_next().is_some() {}
348                    in_flight.spawn(process(
349                        delivery,
350                        handlers.clone(),
351                        job_timeout,
352                        settle_failures.clone(),
353                        permit,
354                    ));
355                }
356                Some(ConsumerEvent::Failed(e)) => {
357                    error!(error = %e, "consumer stream failed while draining");
358                    if outcome.is_ok() {
359                        outcome = Err(e);
360                    }
361                }
362                Some(ConsumerEvent::Ended(queue)) => {
363                    debug!(%queue, "consumer stream ended while draining");
364                }
365            }
366        }
367        while consumers.join_next().await.is_some() {}
368        while in_flight.join_next().await.is_some() {}
369
370        if close_backend && let Err(e) = backend.close().await {
371            warn!(error = %e, "backend close failed");
372            if outcome.is_ok() {
373                outcome = Err(e);
374            }
375        }
376        info!(
377            settle_failures = settle_failures.load(Ordering::Relaxed),
378            "worker stopped"
379        );
380        outcome
381    }
382}
383
384/// Forwards one queue's deliveries into the dispatch channel until told to stop.
385///
386/// Pulling and forwarding are inseparable: once an item has been taken off the
387/// stream it is always sent, so no delivery is dropped unsettled at shutdown.
388async fn consume_into(
389    mut stream: DeliveryStream,
390    tx: mpsc::Sender<ConsumerEvent>,
391    queue: String,
392    mut stop: watch::Receiver<bool>,
393) {
394    loop {
395        if *stop.borrow_and_update() {
396            return;
397        }
398        let item = tokio::select! {
399            biased;
400            _ = stop.changed() => return,
401            item = stream.next() => item,
402        };
403        let Some(item) = item else { break };
404        let event = match item {
405            Ok(delivery) => ConsumerEvent::Delivery(delivery),
406            Err(e) => ConsumerEvent::Failed(e),
407        };
408        if tx.send(event).await.is_err() {
409            return;
410        }
411    }
412    let _ = tx.send(ConsumerEvent::Ended(queue)).await;
413}
414
415/// One message from a per-queue consumer task to the dispatch loop.
416enum ConsumerEvent {
417    Delivery(Box<dyn Delivery>),
418    Failed(Error),
419    Ended(String),
420}
421
422/// Runs one delivery inside a span carrying the job identity.
423async fn process(
424    delivery: Box<dyn Delivery>,
425    handlers: Arc<HandlerMap>,
426    timeout: Option<Duration>,
427    settle_failures: Arc<AtomicU64>,
428    _permit: tokio::sync::OwnedSemaphorePermit,
429) {
430    let envelope = delivery.envelope().clone();
431    let span = tracing::info_span!(
432        "job",
433        job_id = %envelope.job_id,
434        job_type = %envelope.job_type,
435        queue = %envelope.queue,
436        attempt = envelope.attempt,
437    );
438    dispatch(delivery, envelope, handlers, timeout, &settle_failures)
439        .instrument(span)
440        .await;
441}
442
443async fn dispatch(
444    delivery: Box<dyn Delivery>,
445    envelope: Envelope,
446    handlers: Arc<HandlerMap>,
447    timeout: Option<Duration>,
448    failures: &AtomicU64,
449) {
450    let Some(handler) = handlers.get(envelope.job_type.as_str()).cloned() else {
451        warn!("no handler registered; dead-lettering");
452        settle(
453            delivery
454                .dead_letter(&format!("no handler for job type `{}`", envelope.job_type))
455                .await,
456            "dead_letter",
457            failures,
458        );
459        return;
460    };
461
462    debug!("running job");
463    match handler.run(&envelope, timeout).await {
464        JobOutcome::Success => {
465            info!("job succeeded");
466            settle(delivery.ack().await, "ack", failures);
467        }
468        JobOutcome::Decode(e) => {
469            error!(error = %e, "payload did not match the handler's job type");
470            settle(
471                delivery.dead_letter(&format!("decode error: {e}")).await,
472                "dead_letter",
473                failures,
474            );
475        }
476        JobOutcome::Fatal(e) => {
477            error!(error = %e, "job failed fatally");
478            settle(
479                delivery.dead_letter(&format!("fatal error: {e}")).await,
480                "dead_letter",
481                failures,
482            );
483        }
484        JobOutcome::Deferred { delay, reason } => {
485            // Nothing failed: the attempt counter stays put and the retry policy is
486            // never consulted. Only `deferrals` grows, and the envelope comes back at
487            // the front of its queue.
488            let priority = handler.defer_priority();
489            info!(?delay, deferrals = envelope.deferrals + 1, %reason, "job deferred");
490            settle(
491                delivery.defer(envelope.deferred(priority), delay).await,
492                "defer",
493                failures,
494            );
495        }
496        JobOutcome::Retryable(e) => {
497            let policy = handler.policy();
498            match policy.decide(envelope.attempt) {
499                RetryDecision::Retry { delay } => {
500                    warn!(error = %e, ?delay, next_attempt = envelope.attempt + 1, "job failed; retrying");
501                    settle(
502                        delivery.retry(envelope.next_attempt(), delay).await,
503                        "retry",
504                        failures,
505                    );
506                }
507                RetryDecision::GiveUp => {
508                    error!(error = %e, max_attempts = policy.max_attempts, "job failed; giving up");
509                    let reason = format!(
510                        "max attempts ({}) exhausted after attempt {}: {e}",
511                        policy.max_attempts, envelope.attempt
512                    );
513                    settle(delivery.dead_letter(&reason).await, "dead_letter", failures);
514                }
515            }
516        }
517    }
518}
519
520/// Log and count (but never propagate) a failure to settle a delivery.
521///
522/// A settle failure means the broker still owns the message: it will be redelivered
523/// once this consumer's channel or connection goes away, so the count is a signal of
524/// duplicate work ahead rather than of lost work.
525fn settle(result: Result<()>, what: &'static str, failures: &AtomicU64) {
526    if let Err(e) = result {
527        failures.fetch_add(1, Ordering::Relaxed);
528        error!(error = %e, operation = what, "failed to settle delivery");
529    }
530}
531
532/// Configures a [`Worker`]. Created by [`Worker::builder`].
533pub struct WorkerBuilder<Q: QueueSet, B: Backend> {
534    backend: Arc<B>,
535    handlers: Vec<(&'static str, Arc<dyn ErasedHandler>)>,
536    queues: Option<Vec<Q>>,
537    concurrency: Option<usize>,
538    job_timeout: Option<Duration>,
539    close_backend: bool,
540}
541
542impl<Q: QueueSet, B: Backend> WorkerBuilder<Q, B> {
543    /// Register a handler.
544    ///
545    /// The `H::Job: Job<Queue = Q>` bound is what makes the worker type-safe: a handler
546    /// for a job belonging to another queue set does not compile.
547    pub fn handler<H>(mut self, handler: H) -> Self
548    where
549        H: JobHandler,
550        H::Job: Job<Queue = Q>,
551    {
552        self.handlers.push((
553            <H::Job as Job>::NAME,
554            Arc::new(ErasedJobHandler {
555                handler: Arc::new(handler),
556            }),
557        ));
558        self
559    }
560
561    /// Consume only these queues instead of every queue in `Q`. Duplicates are ignored.
562    pub fn queues(mut self, queues: &[Q]) -> Self {
563        self.queues = Some(queues.to_vec());
564        self
565    }
566
567    /// Cap the number of jobs running at the same time.
568    ///
569    /// Defaults to the sum of the prefetch of the consumed queues. `0` is treated as `1`.
570    pub fn concurrency(mut self, concurrency: usize) -> Self {
571        self.concurrency = Some(concurrency);
572        self
573    }
574
575    /// Abort a handler that runs longer than `timeout` and treat it as a retryable failure.
576    pub fn job_timeout(mut self, timeout: Duration) -> Self {
577        self.job_timeout = Some(timeout);
578        self
579    }
580
581    /// Close the backend when [`Worker::run`] returns. Defaults to `false`.
582    ///
583    /// The backend is an `Arc` that is normally shared with a [`crate::Producer`] (and
584    /// possibly other workers), so closing it is the caller's decision: a worker that
585    /// closed it on its own would take those down with it. Turn this on for a process
586    /// whose only job is to run this worker, or call `backend.close()` yourself once
587    /// `run()` has returned.
588    pub fn close_backend_on_shutdown(mut self, close: bool) -> Self {
589        self.close_backend = close;
590        self
591    }
592
593    /// Declare the queues and produce the [`Worker`].
594    ///
595    /// Fails with [`Error::DuplicateHandler`] if two handlers claim the same `Job::NAME`.
596    pub async fn build(self) -> Result<Worker<Q, B>> {
597        let mut handlers: HandlerMap = HashMap::with_capacity(self.handlers.len());
598        for (name, handler) in self.handlers {
599            if handlers.insert(name, handler).is_some() {
600                return Err(Error::DuplicateHandler(name.to_owned()));
601            }
602        }
603
604        let selected: Vec<Q> = self.queues.unwrap_or_else(|| Q::all().to_vec());
605        let mut queues: Vec<QueueConfig> = Vec::with_capacity(selected.len());
606        for queue in selected {
607            let config = queue.config();
608            if !queues.iter().any(|c| c.name == config.name) {
609                queues.push(config);
610            }
611        }
612
613        self.backend.declare(&queues).await?;
614
615        let concurrency = self
616            .concurrency
617            .unwrap_or_else(|| queues.iter().map(|c| usize::from(c.prefetch)).sum())
618            .max(1);
619
620        Ok(Worker {
621            backend: self.backend,
622            handlers: Arc::new(handlers),
623            queues,
624            concurrency,
625            job_timeout: self.job_timeout,
626            close_backend: self.close_backend,
627            handle: WorkerHandle::new(),
628            _q: PhantomData,
629        })
630    }
631}
632
633/// Remote control for a running [`Worker`]. Cheap to clone and send across tasks.
634#[derive(Clone)]
635pub struct WorkerHandle {
636    shutdown: Arc<watch::Sender<bool>>,
637    settle_failures: Arc<AtomicU64>,
638}
639
640impl WorkerHandle {
641    fn new() -> Self {
642        Self {
643            shutdown: Arc::new(watch::channel(false).0),
644            settle_failures: Arc::new(AtomicU64::new(0)),
645        }
646    }
647
648    /// Ask the worker to stop. Returns immediately; [`Worker::run`] finishes the jobs
649    /// that are already running before it returns.
650    ///
651    /// Calling this more than once, or before [`Worker::run`] starts, is fine.
652    pub fn shutdown(&self) {
653        self.shutdown.send_replace(true);
654    }
655
656    /// Whether [`WorkerHandle::shutdown`] has been called.
657    pub fn is_shutdown(&self) -> bool {
658        *self.shutdown.borrow()
659    }
660
661    /// How many times the worker failed to ack / retry / dead-letter a delivery.
662    ///
663    /// Each failure is also logged at `ERROR`. The job itself ran; only telling the
664    /// broker about the outcome failed, so the message is still owned by the broker
665    /// and will be redelivered. A non-zero (and especially a growing) count means
666    /// duplicate processing, and usually a sick connection. Alert on it.
667    pub fn settle_failures(&self) -> u64 {
668        self.settle_failures.load(Ordering::Relaxed)
669    }
670
671    fn subscribe(&self) -> watch::Receiver<bool> {
672        self.shutdown.subscribe()
673    }
674}
675
676impl std::fmt::Debug for WorkerHandle {
677    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
678        f.debug_struct("WorkerHandle")
679            .field("shutdown", &self.is_shutdown())
680            .field("settle_failures", &self.settle_failures())
681            .finish()
682    }
683}
684
685#[cfg(test)]
686mod tests {
687    use super::*;
688    use crate::{
689        handler::FnHandler,
690        memory::MemoryBackend,
691        producer::Producer,
692        test_support::{Greet, Nudge, Orphan, Ping, Stubborn, TestQueues},
693    };
694    use std::sync::{
695        Mutex,
696        atomic::{AtomicUsize, Ordering},
697    };
698    use tokio::task::JoinHandle;
699
700    const ALPHA: &str = "test.alpha";
701    const BETA: &str = "test.beta";
702    const GAMMA: &str = "test.gamma";
703
704    /// Records what handlers saw, so assertions do not depend on logging.
705    #[derive(Default)]
706    struct Recorder {
707        attempts: AtomicUsize,
708        running: AtomicUsize,
709        max_running: AtomicUsize,
710    }
711
712    impl Recorder {
713        fn enter(&self) -> usize {
714            let running = self.running.fetch_add(1, Ordering::SeqCst) + 1;
715            self.max_running.fetch_max(running, Ordering::SeqCst);
716            self.attempts.fetch_add(1, Ordering::SeqCst) + 1
717        }
718        fn leave(&self) {
719            self.running.fetch_sub(1, Ordering::SeqCst);
720        }
721        fn attempts(&self) -> usize {
722            self.attempts.load(Ordering::SeqCst)
723        }
724        fn max_running(&self) -> usize {
725            self.max_running.load(Ordering::SeqCst)
726        }
727    }
728
729    fn backend() -> Arc<MemoryBackend> {
730        Arc::new(MemoryBackend::new())
731    }
732
733    /// A [`MemoryBackend`] whose deliveries can never be settled, so the worker's
734    /// settle-failure path can be exercised.
735    struct Unsettleable(Arc<MemoryBackend>);
736
737    #[async_trait]
738    impl Backend for Unsettleable {
739        async fn declare(&self, queues: &[QueueConfig]) -> Result<()> {
740            self.0.declare(queues).await
741        }
742
743        async fn publish(&self, envelope: &Envelope, delay: Option<Duration>) -> Result<()> {
744            self.0.publish(envelope, delay).await
745        }
746
747        async fn defer(&self, envelope: &Envelope, delay: Duration) -> Result<()> {
748            self.0.defer(envelope, delay).await
749        }
750
751        async fn consume(&self, queue: &QueueConfig) -> Result<DeliveryStream> {
752            let stream = self.0.consume(queue).await?;
753            Ok(Box::pin(stream.map(|item| {
754                item.map(|delivery| Box::new(NeverSettles(delivery)) as Box<dyn Delivery>)
755            })))
756        }
757
758        async fn close(&self) -> Result<()> {
759            self.0.close().await
760        }
761    }
762
763    /// Every settle attempt fails; the inner delivery is dropped unsettled.
764    struct NeverSettles(Box<dyn Delivery>);
765
766    #[async_trait]
767    impl Delivery for NeverSettles {
768        fn envelope(&self) -> &Envelope {
769            self.0.envelope()
770        }
771
772        async fn ack(self: Box<Self>) -> Result<()> {
773            Err(Error::backend(std::io::Error::other("channel is gone")))
774        }
775
776        async fn dead_letter(self: Box<Self>, _reason: &str) -> Result<()> {
777            Err(Error::backend(std::io::Error::other("channel is gone")))
778        }
779
780        async fn retry(self: Box<Self>, _next: Envelope, _delay: Duration) -> Result<()> {
781            Err(Error::backend(std::io::Error::other("channel is gone")))
782        }
783
784        async fn defer(self: Box<Self>, _next: Envelope, _delay: Duration) -> Result<()> {
785            Err(Error::backend(std::io::Error::other("channel is gone")))
786        }
787    }
788
789    async fn producer(backend: &Arc<MemoryBackend>) -> Producer<TestQueues, MemoryBackend> {
790        Producer::<TestQueues, _>::new(backend.clone())
791            .await
792            .unwrap()
793    }
794
795    fn start(worker: Worker<TestQueues, MemoryBackend>) -> (WorkerHandle, JoinHandle<Result<()>>) {
796        let handle = worker.handle();
797        (handle, tokio::spawn(worker.run()))
798    }
799
800    /// Polls `cond` while letting paused time (and therefore retry delays) advance.
801    async fn wait_for(mut cond: impl FnMut() -> bool) {
802        for _ in 0..1_000 {
803            if cond() {
804                return;
805            }
806            tokio::time::sleep(Duration::from_millis(100)).await;
807        }
808        panic!("condition was never met");
809    }
810
811    /// Polls `cond` without advancing the clock, for state that only needs other
812    /// tasks to be given a turn (a consumer registering itself, say).
813    async fn wait_for_tasks(mut cond: impl FnMut() -> bool) {
814        for _ in 0..10_000 {
815            if cond() {
816                return;
817            }
818            tokio::task::yield_now().await;
819        }
820        panic!("condition was never met");
821    }
822
823    async fn stop(handle: WorkerHandle, task: JoinHandle<Result<()>>) -> Result<()> {
824        handle.shutdown();
825        task.await.expect("worker task panicked")
826    }
827
828    #[tokio::test(start_paused = true)]
829    async fn successful_job_is_acked() {
830        let backend = backend();
831        let producer = producer(&backend).await;
832        let seen = Arc::new(Recorder::default());
833
834        let worker = {
835            let seen = seen.clone();
836            Worker::<TestQueues, _>::builder(backend.clone())
837                .handler(FnHandler::<Greet, _>::new(
838                    move |job: Greet, ctx: JobContext| {
839                        let seen = seen.clone();
840                        async move {
841                            seen.enter();
842                            seen.leave();
843                            assert_eq!(job.name, "ada");
844                            assert_eq!(ctx.attempt, 1);
845                            assert_eq!(ctx.max_attempts, 3);
846                            assert_eq!(ctx.job_type, Greet::NAME);
847                            assert_eq!(ctx.queue, ALPHA);
848                            Ok(())
849                        }
850                    },
851                ))
852                .build()
853                .await
854                .unwrap()
855        };
856        let (handle, task) = start(worker);
857
858        let id = producer.enqueue(&Greet::new("ada")).await.unwrap();
859        wait_for(|| backend.acked(ALPHA).len() == 1).await;
860
861        assert_eq!(backend.acked(ALPHA)[0].job_id, id);
862        assert!(backend.dead_letters(ALPHA).is_empty());
863        assert_eq!(seen.attempts(), 1);
864        assert_eq!(handle.settle_failures(), 0);
865        stop(handle, task).await.unwrap();
866        // The backend is shared with the producer, so the worker leaves it alone.
867        assert!(!backend.is_closed());
868    }
869
870    #[tokio::test(start_paused = true)]
871    async fn retryable_failure_is_retried_and_then_succeeds() {
872        let backend = backend();
873        let producer = producer(&backend).await;
874        let seen = Arc::new(Recorder::default());
875
876        let worker = {
877            let seen = seen.clone();
878            Worker::<TestQueues, _>::builder(backend.clone())
879                .handler(FnHandler::<Greet, _>::new(
880                    move |_job: Greet, ctx: JobContext| {
881                        let seen = seen.clone();
882                        async move {
883                            let n = seen.enter();
884                            seen.leave();
885                            if n == 1 {
886                                assert_eq!(ctx.attempt, 1);
887                                assert!(!ctx.is_last_attempt());
888                                Err(JobError::retryable_msg("flaky"))
889                            } else {
890                                assert_eq!(ctx.attempt, 2);
891                                Ok(())
892                            }
893                        }
894                    },
895                ))
896                .build()
897                .await
898                .unwrap()
899        };
900        let (handle, task) = start(worker);
901
902        producer.enqueue(&Greet::new("ada")).await.unwrap();
903        wait_for(|| seen.attempts() == 2).await;
904        wait_for(|| backend.acked(ALPHA).len() == 2).await;
905
906        // The first ack is the retried original, the second the successful attempt.
907        let acked = backend.acked(ALPHA);
908        assert_eq!(acked[0].attempt, 1);
909        assert_eq!(acked[1].attempt, 2);
910        assert_eq!(acked[0].job_id, acked[1].job_id);
911        assert!(backend.dead_letters(ALPHA).is_empty());
912        stop(handle, task).await.unwrap();
913    }
914
915    #[tokio::test(start_paused = true)]
916    async fn retryable_failure_dead_letters_once_attempts_are_exhausted() {
917        let backend = backend();
918        let producer = producer(&backend).await;
919        let seen = Arc::new(Recorder::default());
920
921        let worker = {
922            let seen = seen.clone();
923            Worker::<TestQueues, _>::builder(backend.clone())
924                .handler(FnHandler::<Greet, _>::new(
925                    move |_job: Greet, _ctx: JobContext| {
926                        let seen = seen.clone();
927                        async move {
928                            seen.enter();
929                            seen.leave();
930                            Err(JobError::retryable_msg("always down"))
931                        }
932                    },
933                ))
934                .build()
935                .await
936                .unwrap()
937        };
938        let (handle, task) = start(worker);
939
940        producer.enqueue(&Greet::new("ada")).await.unwrap();
941        wait_for(|| !backend.dead_letters(ALPHA).is_empty()).await;
942
943        // Queue policy is three attempts.
944        assert_eq!(seen.attempts(), 3);
945        let dead = backend.dead_letters(ALPHA);
946        assert_eq!(dead.len(), 1);
947        assert_eq!(dead[0].0.attempt, 3);
948        assert!(
949            dead[0].1.contains("max attempts"),
950            "reason was {:?}",
951            dead[0].1
952        );
953        assert!(dead[0].1.contains("always down"));
954        stop(handle, task).await.unwrap();
955    }
956
957    #[tokio::test(start_paused = true)]
958    async fn fatal_failure_dead_letters_immediately() {
959        let backend = backend();
960        let producer = producer(&backend).await;
961        let seen = Arc::new(Recorder::default());
962
963        let worker = {
964            let seen = seen.clone();
965            Worker::<TestQueues, _>::builder(backend.clone())
966                .handler(FnHandler::<Greet, _>::new(
967                    move |_job: Greet, _ctx: JobContext| {
968                        let seen = seen.clone();
969                        async move {
970                            seen.enter();
971                            seen.leave();
972                            Err(JobError::fatal_msg("bad input"))
973                        }
974                    },
975                ))
976                .build()
977                .await
978                .unwrap()
979        };
980        let (handle, task) = start(worker);
981
982        producer.enqueue(&Greet::new("ada")).await.unwrap();
983        wait_for(|| !backend.dead_letters(ALPHA).is_empty()).await;
984
985        // One attempt only, even though the queue allows three.
986        assert_eq!(seen.attempts(), 1);
987        let dead = backend.dead_letters(ALPHA);
988        assert_eq!(dead[0].0.attempt, 1);
989        assert!(dead[0].1.contains("fatal"), "reason was {:?}", dead[0].1);
990        assert!(backend.acked(ALPHA).is_empty());
991        stop(handle, task).await.unwrap();
992    }
993
994    #[tokio::test(start_paused = true)]
995    async fn job_retry_policy_overrides_the_queue_policy() {
996        let backend = backend();
997        let producer = producer(&backend).await;
998        let seen = Arc::new(Recorder::default());
999
1000        let worker = {
1001            let seen = seen.clone();
1002            Worker::<TestQueues, _>::builder(backend.clone())
1003                .handler(FnHandler::<Stubborn, _>::new(
1004                    move |_job: Stubborn, ctx: JobContext| {
1005                        let seen = seen.clone();
1006                        async move {
1007                            seen.enter();
1008                            // The queue says one attempt, the job says two.
1009                            assert_eq!(ctx.max_attempts, 2);
1010                            seen.leave();
1011                            Err(JobError::retryable_msg("nope"))
1012                        }
1013                    },
1014                ))
1015                .build()
1016                .await
1017                .unwrap()
1018        };
1019        let (handle, task) = start(worker);
1020
1021        producer.enqueue(&Stubborn { id: 1 }).await.unwrap();
1022        wait_for(|| !backend.dead_letters(BETA).is_empty()).await;
1023
1024        assert_eq!(seen.attempts(), 2);
1025        assert!(backend.dead_letters(BETA)[0].1.contains("max attempts"));
1026        stop(handle, task).await.unwrap();
1027    }
1028
1029    #[tokio::test(start_paused = true)]
1030    async fn queue_policy_without_retries_dead_letters_on_first_failure() {
1031        let backend = backend();
1032        let producer = producer(&backend).await;
1033        let seen = Arc::new(Recorder::default());
1034
1035        let worker = {
1036            let seen = seen.clone();
1037            Worker::<TestQueues, _>::builder(backend.clone())
1038                .handler(FnHandler::<Ping, _>::new(
1039                    move |_job: Ping, ctx: JobContext| {
1040                        let seen = seen.clone();
1041                        async move {
1042                            seen.enter();
1043                            assert_eq!(ctx.max_attempts, 1);
1044                            assert!(ctx.is_last_attempt());
1045                            seen.leave();
1046                            Err(JobError::retryable_msg("nope"))
1047                        }
1048                    },
1049                ))
1050                .build()
1051                .await
1052                .unwrap()
1053        };
1054        let (handle, task) = start(worker);
1055
1056        producer.enqueue(&Ping { seq: 1 }).await.unwrap();
1057        wait_for(|| !backend.dead_letters(BETA).is_empty()).await;
1058        assert_eq!(seen.attempts(), 1);
1059        stop(handle, task).await.unwrap();
1060    }
1061
1062    #[tokio::test(start_paused = true)]
1063    async fn job_without_a_handler_is_dead_lettered() {
1064        let backend = backend();
1065        let producer = producer(&backend).await;
1066
1067        let worker = Worker::<TestQueues, _>::builder(backend.clone())
1068            .handler(FnHandler::<Greet, _>::new(
1069                |_job: Greet, _ctx: JobContext| async move { Ok(()) },
1070            ))
1071            .build()
1072            .await
1073            .unwrap();
1074        let (handle, task) = start(worker);
1075
1076        producer.enqueue(&Orphan { id: 9 }).await.unwrap();
1077        wait_for(|| !backend.dead_letters(ALPHA).is_empty()).await;
1078
1079        let dead = backend.dead_letters(ALPHA);
1080        assert_eq!(dead[0].0.job_type, Orphan::NAME);
1081        assert!(
1082            dead[0].1.contains("no handler"),
1083            "reason was {:?}",
1084            dead[0].1
1085        );
1086        stop(handle, task).await.unwrap();
1087    }
1088
1089    #[tokio::test(start_paused = true)]
1090    async fn malformed_payload_is_dead_lettered() {
1091        let backend = backend();
1092        let seen = Arc::new(Recorder::default());
1093
1094        let worker = {
1095            let seen = seen.clone();
1096            Worker::<TestQueues, _>::builder(backend.clone())
1097                .handler(FnHandler::<Greet, _>::new(
1098                    move |_job: Greet, _ctx: JobContext| {
1099                        let seen = seen.clone();
1100                        async move {
1101                            seen.enter();
1102                            seen.leave();
1103                            Ok(())
1104                        }
1105                    },
1106                ))
1107                .build()
1108                .await
1109                .unwrap()
1110        };
1111        let (handle, task) = start(worker);
1112
1113        let mut envelope = Envelope::new(&Greet::new("ada")).unwrap();
1114        envelope.payload = serde_json::json!({ "not_a_name": 42 });
1115        backend.publish(&envelope, None).await.unwrap();
1116
1117        wait_for(|| !backend.dead_letters(ALPHA).is_empty()).await;
1118        let dead = backend.dead_letters(ALPHA);
1119        assert!(
1120            dead[0].1.contains("decode error"),
1121            "reason was {:?}",
1122            dead[0].1
1123        );
1124        assert_eq!(
1125            seen.attempts(),
1126            0,
1127            "handler must not run on a decode failure"
1128        );
1129        stop(handle, task).await.unwrap();
1130    }
1131
1132    #[tokio::test(start_paused = true)]
1133    async fn panicking_handler_is_retried() {
1134        let backend = backend();
1135        let producer = producer(&backend).await;
1136        let seen = Arc::new(Recorder::default());
1137
1138        let worker = {
1139            let seen = seen.clone();
1140            Worker::<TestQueues, _>::builder(backend.clone())
1141                .handler(FnHandler::<Greet, _>::new(
1142                    move |_job: Greet, _ctx: JobContext| {
1143                        let seen = seen.clone();
1144                        async move {
1145                            let n = seen.enter();
1146                            seen.leave();
1147                            if n == 1 {
1148                                panic!("handler exploded on purpose");
1149                            }
1150                            Ok(())
1151                        }
1152                    },
1153                ))
1154                .build()
1155                .await
1156                .unwrap()
1157        };
1158        let (handle, task) = start(worker);
1159
1160        producer.enqueue(&Greet::new("ada")).await.unwrap();
1161        wait_for(|| seen.attempts() == 2).await;
1162        wait_for(|| backend.acked(ALPHA).len() == 2).await;
1163
1164        assert!(backend.dead_letters(ALPHA).is_empty());
1165        // The worker itself survived the panic.
1166        assert!(!task.is_finished());
1167        stop(handle, task).await.unwrap();
1168    }
1169
1170    #[tokio::test(start_paused = true)]
1171    async fn job_timeout_is_a_retryable_failure() {
1172        let backend = backend();
1173        let producer = producer(&backend).await;
1174        let seen = Arc::new(Recorder::default());
1175
1176        let worker = {
1177            let seen = seen.clone();
1178            Worker::<TestQueues, _>::builder(backend.clone())
1179                .job_timeout(Duration::from_secs(5))
1180                .handler(FnHandler::<Greet, _>::new(
1181                    move |_job: Greet, _ctx: JobContext| {
1182                        let seen = seen.clone();
1183                        async move {
1184                            let n = seen.enter();
1185                            if n == 1 {
1186                                tokio::time::sleep(Duration::from_secs(3600)).await;
1187                            }
1188                            seen.leave();
1189                            Ok(())
1190                        }
1191                    },
1192                ))
1193                .build()
1194                .await
1195                .unwrap()
1196        };
1197        let (handle, task) = start(worker);
1198
1199        producer.enqueue(&Greet::new("ada")).await.unwrap();
1200        wait_for(|| backend.acked(ALPHA).len() == 2).await;
1201
1202        assert_eq!(seen.attempts(), 2);
1203        assert!(backend.dead_letters(ALPHA).is_empty());
1204        stop(handle, task).await.unwrap();
1205    }
1206
1207    #[tokio::test(start_paused = true)]
1208    async fn deferred_job_runs_again_with_the_same_attempt_and_top_priority() {
1209        let backend = backend();
1210        let producer = producer(&backend).await;
1211        let seen = Arc::new(Recorder::default());
1212        let second_run = Arc::new(Mutex::new(None::<JobContext>));
1213
1214        let worker = {
1215            let (seen, second_run) = (seen.clone(), second_run.clone());
1216            Worker::<TestQueues, _>::builder(backend.clone())
1217                .handler(FnHandler::<Greet, _>::new(
1218                    move |_j: Greet, ctx: JobContext| {
1219                        let (seen, second_run) = (seen.clone(), second_run.clone());
1220                        async move {
1221                            let n = seen.enter();
1222                            seen.leave();
1223                            if n == 1 {
1224                                assert_eq!(ctx.deferrals, 0);
1225                                assert_eq!(ctx.priority, 0);
1226                                Err(JobError::deferred(Duration::from_secs(1)))
1227                            } else {
1228                                *second_run.lock().unwrap() = Some(ctx);
1229                                Ok(())
1230                            }
1231                        }
1232                    },
1233                ))
1234                .build()
1235                .await
1236                .unwrap()
1237        };
1238        let (handle, task) = start(worker);
1239
1240        let id = producer.enqueue(&Greet::new("limited")).await.unwrap();
1241        wait_for(|| seen.attempts() == 2).await;
1242        wait_for(|| backend.acked(ALPHA).len() == 2).await;
1243
1244        let ctx = second_run.lock().unwrap().clone().expect("second run");
1245        assert_eq!(ctx.attempt, 1, "a deferral does not spend an attempt");
1246        assert_eq!(ctx.deferrals, 1);
1247        assert_eq!(ctx.priority, 10, "alpha's default ten levels");
1248        assert_eq!(ctx.job_id, id);
1249
1250        // First ack is the deferred original, second the successful re-delivery.
1251        let acked = backend.acked(ALPHA);
1252        assert_eq!((acked[0].attempt, acked[0].deferrals), (1, 0));
1253        assert_eq!((acked[1].attempt, acked[1].deferrals), (1, 1));
1254        assert_eq!(acked[1].priority, 10);
1255        assert!(backend.dead_letters(ALPHA).is_empty());
1256        assert_eq!(handle.settle_failures(), 0, "the defer settled cleanly");
1257        assert_eq!(backend.deferred(ALPHA), 0);
1258        stop(handle, task).await.unwrap();
1259    }
1260
1261    #[tokio::test(start_paused = true)]
1262    async fn a_retry_after_a_deferral_drops_back_to_priority_zero() {
1263        let backend = backend();
1264        let producer = producer(&backend).await;
1265        let seen = Arc::new(Recorder::default());
1266        let third_run = Arc::new(Mutex::new(None::<JobContext>));
1267
1268        // Run 1 defers, run 2 (back at the top of the queue) fails transiently, run 3
1269        // is the retry: it must have lost the deferral's priority.
1270        let worker = {
1271            let (seen, third_run) = (seen.clone(), third_run.clone());
1272            Worker::<TestQueues, _>::builder(backend.clone())
1273                .handler(FnHandler::<Greet, _>::new(
1274                    move |_j: Greet, ctx: JobContext| {
1275                        let (seen, third_run) = (seen.clone(), third_run.clone());
1276                        async move {
1277                            let n = seen.enter();
1278                            seen.leave();
1279                            match n {
1280                                1 => Err(JobError::deferred(Duration::from_secs(1))),
1281                                2 => {
1282                                    assert_eq!(ctx.priority, 10, "the deferral came back first");
1283                                    Err(JobError::retryable_msg("flaky"))
1284                                }
1285                                _ => {
1286                                    *third_run.lock().unwrap() = Some(ctx);
1287                                    Ok(())
1288                                }
1289                            }
1290                        }
1291                    },
1292                ))
1293                .build()
1294                .await
1295                .unwrap()
1296        };
1297        let (handle, task) = start(worker);
1298
1299        producer.enqueue(&Greet::new("limited")).await.unwrap();
1300        wait_for(|| seen.attempts() == 3).await;
1301        wait_for(|| backend.acked(ALPHA).len() == 3).await;
1302
1303        let ctx = third_run.lock().unwrap().clone().expect("third run");
1304        assert_eq!(ctx.attempt, 2, "the retry spent an attempt");
1305        assert_eq!(ctx.deferrals, 1, "the deferral history is carried forward");
1306        assert_eq!(ctx.priority, 0, "a retry must not keep jumping the backlog");
1307
1308        // The envelope the backend saw agrees with what the handler was told.
1309        let acked = backend.acked(ALPHA);
1310        assert_eq!(
1311            (acked[2].attempt, acked[2].deferrals, acked[2].priority),
1312            (2, 1, 0)
1313        );
1314        assert!(backend.dead_letters(ALPHA).is_empty());
1315        assert_eq!(handle.settle_failures(), 0);
1316        stop(handle, task).await.unwrap();
1317    }
1318
1319    #[tokio::test(start_paused = true)]
1320    async fn deferral_on_a_queue_without_priorities_stays_at_zero() {
1321        let backend = backend();
1322        let producer = producer(&backend).await;
1323        let seen = Arc::new(Recorder::default());
1324        let second_run = Arc::new(Mutex::new(None::<JobContext>));
1325
1326        let worker = {
1327            let (seen, second_run) = (seen.clone(), second_run.clone());
1328            Worker::<TestQueues, _>::builder(backend.clone())
1329                .handler(FnHandler::<Nudge, _>::new(
1330                    move |_j: Nudge, ctx: JobContext| {
1331                        let (seen, second_run) = (seen.clone(), second_run.clone());
1332                        async move {
1333                            let n = seen.enter();
1334                            seen.leave();
1335                            if n == 1 {
1336                                Err(JobError::deferred_msg(
1337                                    Duration::from_secs(2),
1338                                    "rate limited",
1339                                ))
1340                            } else {
1341                                *second_run.lock().unwrap() = Some(ctx);
1342                                Ok(())
1343                            }
1344                        }
1345                    },
1346                ))
1347                .build()
1348                .await
1349                .unwrap()
1350        };
1351        let (handle, task) = start(worker);
1352
1353        producer.enqueue(&Nudge { id: 1 }).await.unwrap();
1354        wait_for(|| seen.attempts() == 2).await;
1355
1356        let ctx = second_run.lock().unwrap().clone().expect("second run");
1357        assert_eq!(ctx.priority, 0, "gamma is not a priority queue");
1358        assert_eq!(ctx.deferrals, 1);
1359        assert_eq!(ctx.attempt, 1);
1360        assert!(backend.dead_letters(GAMMA).is_empty());
1361        assert_eq!(handle.settle_failures(), 0);
1362        stop(handle, task).await.unwrap();
1363    }
1364
1365    #[tokio::test(start_paused = true)]
1366    async fn deferral_never_consults_the_retry_policy() {
1367        let backend = backend();
1368        let producer = producer(&backend).await;
1369        let seen = Arc::new(Recorder::default());
1370
1371        // Beta allows a single attempt: a *failure* here dead-letters at once. A
1372        // deferral must not, however often it happens.
1373        let worker = {
1374            let seen = seen.clone();
1375            Worker::<TestQueues, _>::builder(backend.clone())
1376                .handler(FnHandler::<Ping, _>::new(
1377                    move |_j: Ping, ctx: JobContext| {
1378                        let seen = seen.clone();
1379                        async move {
1380                            let n = seen.enter();
1381                            seen.leave();
1382                            assert_eq!(ctx.max_attempts, 1);
1383                            assert_eq!(ctx.attempt, 1);
1384                            assert!(ctx.is_last_attempt());
1385                            assert_eq!(ctx.deferrals as usize, n - 1);
1386                            if n < 4 {
1387                                Err(JobError::deferred(Duration::from_secs(1)))
1388                            } else {
1389                                Ok(())
1390                            }
1391                        }
1392                    },
1393                ))
1394                .build()
1395                .await
1396                .unwrap()
1397        };
1398        let (handle, task) = start(worker);
1399
1400        producer.enqueue(&Ping { seq: 1 }).await.unwrap();
1401        wait_for(|| seen.attempts() == 4).await;
1402        wait_for(|| backend.acked(BETA).len() == 4).await;
1403
1404        assert!(
1405            backend.dead_letters(BETA).is_empty(),
1406            "three deferrals on a one-attempt queue must not dead-letter"
1407        );
1408        assert!(backend.acked(BETA).iter().all(|e| e.attempt == 1));
1409        assert_eq!(
1410            backend.acked(BETA).last().unwrap().deferrals,
1411            3,
1412            "only the deferral counter moved"
1413        );
1414        assert_eq!(handle.settle_failures(), 0);
1415        stop(handle, task).await.unwrap();
1416    }
1417
1418    #[tokio::test(start_paused = true)]
1419    async fn a_job_in_hold_survives_the_worker_shutting_down() {
1420        let backend = backend();
1421        let producer = producer(&backend).await;
1422        let seen = Arc::new(Recorder::default());
1423
1424        let worker = {
1425            let seen = seen.clone();
1426            Worker::<TestQueues, _>::builder(backend.clone())
1427                .handler(FnHandler::<Greet, _>::new(
1428                    move |_j: Greet, _c: JobContext| {
1429                        let seen = seen.clone();
1430                        async move {
1431                            seen.enter();
1432                            seen.leave();
1433                            Err(JobError::deferred(Duration::from_secs(3600)))
1434                        }
1435                    },
1436                ))
1437                .build()
1438                .await
1439                .unwrap()
1440        };
1441        let (handle, task) = start(worker);
1442
1443        let id = producer.enqueue(&Greet::new("held")).await.unwrap();
1444        wait_for(|| backend.acked(ALPHA).len() == 1).await;
1445        assert_eq!(backend.deferred(ALPHA), 1);
1446
1447        // The worker goes away while the job is still in hold.
1448        stop(handle, task).await.unwrap();
1449        assert_eq!(seen.attempts(), 1);
1450        assert_eq!(backend.deferred(ALPHA), 1, "still held, not lost");
1451        assert_eq!(backend.pending(ALPHA), 0);
1452
1453        // The backend is shared and stays open, so the hold still expires.
1454        tokio::time::sleep(Duration::from_secs(3601)).await;
1455        assert_eq!(backend.deferred(ALPHA), 0);
1456        assert_eq!(backend.pending(ALPHA), 1);
1457        let waiting = backend.acked(ALPHA);
1458        assert_eq!(waiting[0].job_id, id);
1459        assert!(backend.dead_letters(ALPHA).is_empty());
1460    }
1461
1462    #[tokio::test(start_paused = true)]
1463    async fn a_deferred_job_beats_work_enqueued_while_it_waited() {
1464        /// Runs one `Greet` at a time, recording the order, and defers the job named
1465        /// `"limited"` the first time it sees it.
1466        async fn worker(
1467            backend: Arc<MemoryBackend>,
1468            order: Arc<Mutex<Vec<String>>>,
1469            seen: Arc<Recorder>,
1470        ) -> Worker<TestQueues, MemoryBackend> {
1471            Worker::<TestQueues, _>::builder(backend)
1472                .queues(&[TestQueues::Alpha])
1473                .concurrency(1)
1474                .handler(FnHandler::<Greet, _>::new(
1475                    move |job: Greet, ctx: JobContext| {
1476                        let (order, seen) = (order.clone(), seen.clone());
1477                        async move {
1478                            seen.enter();
1479                            seen.leave();
1480                            order.lock().unwrap().push(job.name.clone());
1481                            if job.name == "limited" && ctx.deferrals == 0 {
1482                                Err(JobError::deferred(Duration::from_secs(30)))
1483                            } else {
1484                                Ok(())
1485                            }
1486                        }
1487                    },
1488                ))
1489                .build()
1490                .await
1491                .unwrap()
1492        }
1493
1494        let backend = backend();
1495        let producer = producer(&backend).await;
1496        let order = Arc::new(Mutex::new(Vec::<String>::new()));
1497        let seen = Arc::new(Recorder::default());
1498
1499        let first = worker(backend.clone(), order.clone(), seen.clone()).await;
1500        let (handle, task) = start(first);
1501        producer.enqueue(&Greet::new("limited")).await.unwrap();
1502        wait_for(|| backend.deferred(ALPHA) == 1).await;
1503        // Nothing consumes from here on, so the queue order is observable.
1504        stop(handle, task).await.unwrap();
1505
1506        // Two ordinary jobs pile up while the deferred one is still in hold, and only
1507        // then does its delay elapse.
1508        producer.enqueue(&Greet::new("backlog-1")).await.unwrap();
1509        producer.enqueue(&Greet::new("backlog-2")).await.unwrap();
1510        assert_eq!(backend.pending(ALPHA), 2);
1511        wait_for(|| backend.deferred(ALPHA) == 0).await;
1512        assert_eq!(backend.pending(ALPHA), 3);
1513
1514        let second = worker(backend.clone(), order.clone(), seen.clone()).await;
1515        let (handle, task) = start(second);
1516        wait_for(|| seen.attempts() == 4).await;
1517
1518        assert_eq!(
1519            *order.lock().unwrap(),
1520            vec!["limited", "limited", "backlog-1", "backlog-2"],
1521            "the deferred job must overtake the backlog that built up"
1522        );
1523        assert_eq!(handle.settle_failures(), 0);
1524        stop(handle, task).await.unwrap();
1525    }
1526
1527    #[tokio::test(start_paused = true)]
1528    async fn a_failed_defer_is_counted_as_a_settle_failure() {
1529        let inner = backend();
1530        let backend = Arc::new(Unsettleable(inner.clone()));
1531        let producer = Producer::<TestQueues, _>::new(backend.clone())
1532            .await
1533            .unwrap();
1534
1535        let worker = Worker::<TestQueues, _>::builder(backend.clone())
1536            .handler(FnHandler::<Greet, _>::new(
1537                |_j: Greet, _c: JobContext| async move {
1538                    Err(JobError::deferred(Duration::from_secs(1)))
1539                },
1540            ))
1541            .build()
1542            .await
1543            .unwrap();
1544        let handle = worker.handle();
1545        let task = tokio::spawn(worker.run());
1546
1547        producer.enqueue(&Greet::new("ada")).await.unwrap();
1548        wait_for(|| handle.settle_failures() == 1).await;
1549
1550        assert!(inner.acked(ALPHA).is_empty());
1551        assert_eq!(inner.deferred(ALPHA), 0);
1552        handle.shutdown();
1553        task.await.expect("worker task panicked").unwrap();
1554    }
1555
1556    #[tokio::test]
1557    async fn duplicate_handler_fails_the_build() {
1558        let backend = backend();
1559        let err = Worker::<TestQueues, _>::builder(backend)
1560            .handler(FnHandler::<Greet, _>::new(
1561                |_j: Greet, _c: JobContext| async move { Ok(()) },
1562            ))
1563            .handler(FnHandler::<Greet, _>::new(
1564                |_j: Greet, _c: JobContext| async move { Ok(()) },
1565            ))
1566            .build()
1567            .await
1568            .unwrap_err();
1569        match err {
1570            Error::DuplicateHandler(name) => assert_eq!(name, Greet::NAME),
1571            other => panic!("unexpected error: {other}"),
1572        }
1573    }
1574
1575    #[tokio::test(start_paused = true)]
1576    async fn queues_restricts_consumption_to_the_chosen_subset() {
1577        let backend = backend();
1578        let producer = producer(&backend).await;
1579        let alpha_seen = Arc::new(Recorder::default());
1580        let beta_seen = Arc::new(Recorder::default());
1581
1582        let worker = {
1583            let (a, b) = (alpha_seen.clone(), beta_seen.clone());
1584            Worker::<TestQueues, _>::builder(backend.clone())
1585                .queues(&[TestQueues::Alpha, TestQueues::Alpha])
1586                .handler(FnHandler::<Greet, _>::new(
1587                    move |_j: Greet, _c: JobContext| {
1588                        let a = a.clone();
1589                        async move {
1590                            a.enter();
1591                            a.leave();
1592                            Ok(())
1593                        }
1594                    },
1595                ))
1596                .handler(FnHandler::<Ping, _>::new(
1597                    move |_j: Ping, _c: JobContext| {
1598                        let b = b.clone();
1599                        async move {
1600                            b.enter();
1601                            b.leave();
1602                            Ok(())
1603                        }
1604                    },
1605                ))
1606                .build()
1607                .await
1608                .unwrap()
1609        };
1610        assert_eq!(worker.queues().len(), 1, "duplicates must be collapsed");
1611        let (handle, task) = start(worker);
1612
1613        producer.enqueue(&Ping { seq: 1 }).await.unwrap();
1614        producer.enqueue(&Greet::new("ada")).await.unwrap();
1615        wait_for(|| alpha_seen.attempts() == 1).await;
1616
1617        assert_eq!(beta_seen.attempts(), 0);
1618        // The beta message is still sitting on its queue, untouched.
1619        assert_eq!(backend.pending(BETA), 1);
1620        stop(handle, task).await.unwrap();
1621    }
1622
1623    #[tokio::test(start_paused = true)]
1624    async fn graceful_shutdown_waits_for_in_flight_jobs() {
1625        let backend = backend();
1626        let producer = producer(&backend).await;
1627        let started = Arc::new(Recorder::default());
1628        let finished = Arc::new(Recorder::default());
1629
1630        let worker = {
1631            let (s, f) = (started.clone(), finished.clone());
1632            Worker::<TestQueues, _>::builder(backend.clone())
1633                .handler(FnHandler::<Greet, _>::new(
1634                    move |_j: Greet, _c: JobContext| {
1635                        let (s, f) = (s.clone(), f.clone());
1636                        async move {
1637                            s.enter();
1638                            tokio::time::sleep(Duration::from_secs(30)).await;
1639                            f.enter();
1640                            f.leave();
1641                            s.leave();
1642                            Ok(())
1643                        }
1644                    },
1645                ))
1646                .build()
1647                .await
1648                .unwrap()
1649        };
1650        let (handle, task) = start(worker);
1651
1652        producer.enqueue(&Greet::new("slow")).await.unwrap();
1653        wait_for(|| started.attempts() == 1).await;
1654        assert_eq!(finished.attempts(), 0);
1655
1656        handle.shutdown();
1657        assert!(handle.is_shutdown());
1658        task.await.expect("worker task panicked").unwrap();
1659
1660        assert_eq!(
1661            finished.attempts(),
1662            1,
1663            "shutdown must not abandon a running job"
1664        );
1665        assert_eq!(backend.acked(ALPHA).len(), 1);
1666        assert!(!backend.is_closed(), "closing the backend is opt-in");
1667    }
1668
1669    #[tokio::test(start_paused = true)]
1670    async fn shutdown_processes_the_deliveries_it_already_pulled() {
1671        const PUBLISHED: usize = 6;
1672
1673        let backend = backend();
1674        let producer = producer(&backend).await;
1675        let seen = Arc::new(Recorder::default());
1676
1677        // Prefetch 4 on alpha with room for one job at a time: the consumer keeps
1678        // pulling while the dispatch loop is busy, so deliveries pile up in between.
1679        let worker = {
1680            let seen = seen.clone();
1681            Worker::<TestQueues, _>::builder(backend.clone())
1682                .queues(&[TestQueues::Alpha])
1683                .concurrency(1)
1684                .handler(FnHandler::<Greet, _>::new(
1685                    move |_j: Greet, _c: JobContext| {
1686                        let seen = seen.clone();
1687                        async move {
1688                            seen.enter();
1689                            tokio::time::sleep(Duration::from_secs(5)).await;
1690                            seen.leave();
1691                            Ok(())
1692                        }
1693                    },
1694                ))
1695                .build()
1696                .await
1697                .unwrap()
1698        };
1699        let (handle, task) = start(worker);
1700
1701        for i in 0..PUBLISHED {
1702            producer
1703                .enqueue(&Greet::new(&format!("job-{i}")))
1704                .await
1705                .unwrap();
1706        }
1707        // One job is running; the rest are buffered in the worker, in the backend's
1708        // consumer, or still on the queue.
1709        wait_for(|| seen.attempts() == 1).await;
1710
1711        handle.shutdown();
1712        task.await.expect("worker task panicked").unwrap();
1713
1714        let acked = backend.acked(ALPHA).len();
1715        let pending = backend.pending(ALPHA);
1716        assert_eq!(
1717            acked + pending,
1718            PUBLISHED,
1719            "{acked} acked + {pending} pending must account for every message"
1720        );
1721        assert!(
1722            acked >= 2,
1723            "buffered deliveries must be processed, not dropped; only {acked} were"
1724        );
1725        assert_eq!(acked, seen.attempts(), "every job that ran was settled");
1726        assert!(backend.dead_letters(ALPHA).is_empty());
1727    }
1728
1729    #[tokio::test(start_paused = true)]
1730    async fn closing_the_backend_on_shutdown_is_opt_in() {
1731        let backend = backend();
1732
1733        let worker = Worker::<TestQueues, _>::builder(backend.clone())
1734            .handler(FnHandler::<Greet, _>::new(
1735                |_j: Greet, _c: JobContext| async move { Ok(()) },
1736            ))
1737            .build()
1738            .await
1739            .unwrap();
1740        let (handle, task) = start(worker);
1741        stop(handle, task).await.unwrap();
1742        assert!(!backend.is_closed(), "default must leave the backend alone");
1743
1744        // Still usable afterwards: this is what a shared `Producer` relies on.
1745        producer(&backend)
1746            .await
1747            .enqueue(&Greet::new("after"))
1748            .await
1749            .unwrap();
1750
1751        let worker = Worker::<TestQueues, _>::builder(backend.clone())
1752            .close_backend_on_shutdown(true)
1753            .handler(FnHandler::<Greet, _>::new(
1754                |_j: Greet, _c: JobContext| async move { Ok(()) },
1755            ))
1756            .build()
1757            .await
1758            .unwrap();
1759        assert!(format!("{worker:?}").contains("close_backend_on_shutdown: true"));
1760        let (handle, task) = start(worker);
1761        stop(handle, task).await.unwrap();
1762        assert!(backend.is_closed());
1763    }
1764
1765    #[tokio::test(start_paused = true)]
1766    async fn failures_to_settle_are_counted() {
1767        let inner = backend();
1768        let backend = Arc::new(Unsettleable(inner.clone()));
1769        let producer = Producer::<TestQueues, _>::new(backend.clone())
1770            .await
1771            .unwrap();
1772
1773        let worker = Worker::<TestQueues, _>::builder(backend.clone())
1774            .handler(FnHandler::<Greet, _>::new(
1775                |_j: Greet, _c: JobContext| async move { Ok(()) },
1776            ))
1777            .build()
1778            .await
1779            .unwrap();
1780        let handle = worker.handle();
1781        assert_eq!(handle.settle_failures(), 0);
1782        let task = tokio::spawn(worker.run());
1783
1784        producer.enqueue(&Greet::new("ada")).await.unwrap();
1785        wait_for(|| handle.settle_failures() == 1).await;
1786
1787        // The handler ran, but the broker never heard about it.
1788        assert!(inner.acked(ALPHA).is_empty());
1789        assert!(format!("{handle:?}").contains("settle_failures: 1"));
1790        handle.shutdown();
1791        task.await.expect("worker task panicked").unwrap();
1792    }
1793
1794    #[tokio::test(start_paused = true)]
1795    async fn concurrency_is_capped() {
1796        let backend = backend();
1797        let producer = producer(&backend).await;
1798        let seen = Arc::new(Recorder::default());
1799
1800        let worker = {
1801            let seen = seen.clone();
1802            Worker::<TestQueues, _>::builder(backend.clone())
1803                .concurrency(2)
1804                .handler(FnHandler::<Greet, _>::new(
1805                    move |_j: Greet, _c: JobContext| {
1806                        let seen = seen.clone();
1807                        async move {
1808                            seen.enter();
1809                            tokio::time::sleep(Duration::from_secs(1)).await;
1810                            seen.leave();
1811                            Ok(())
1812                        }
1813                    },
1814                ))
1815                .build()
1816                .await
1817                .unwrap()
1818        };
1819        assert_eq!(worker.concurrency(), 2);
1820        let (handle, task) = start(worker);
1821
1822        for i in 0..8 {
1823            producer
1824                .enqueue(&Greet::new(&format!("job-{i}")))
1825                .await
1826                .unwrap();
1827        }
1828        wait_for(|| backend.acked(ALPHA).len() == 8).await;
1829
1830        assert_eq!(seen.attempts(), 8);
1831        assert!(
1832            seen.max_running() <= 2,
1833            "ran {} jobs at once",
1834            seen.max_running()
1835        );
1836        stop(handle, task).await.unwrap();
1837    }
1838
1839    #[tokio::test(start_paused = true)]
1840    async fn default_concurrency_is_the_sum_of_prefetch() {
1841        let backend = backend();
1842        let worker = Worker::<TestQueues, _>::builder(backend)
1843            .handler(FnHandler::<Greet, _>::new(
1844                |_j: Greet, _c: JobContext| async move { Ok(()) },
1845            ))
1846            .build()
1847            .await
1848            .unwrap();
1849        // Alpha prefetch 4 + Beta prefetch 2 + Gamma prefetch 1.
1850        assert_eq!(worker.concurrency(), 7);
1851        assert_eq!(worker.queues().len(), 3);
1852    }
1853
1854    #[tokio::test(start_paused = true)]
1855    async fn lost_backend_ends_the_run_with_an_error() {
1856        let backend = backend();
1857        let worker = Worker::<TestQueues, _>::builder(backend.clone())
1858            .handler(FnHandler::<Greet, _>::new(
1859                |_j: Greet, _c: JobContext| async move { Ok(()) },
1860            ))
1861            .build()
1862            .await
1863            .unwrap();
1864        let (_handle, task) = start(worker);
1865
1866        // Simulate the connection going away underneath the worker, once both
1867        // consumers have actually registered with the backend.
1868        {
1869            let backend = backend.clone();
1870            wait_for_tasks(move || {
1871                backend.consumer_count(ALPHA) == 1 && backend.consumer_count(BETA) == 1
1872            })
1873            .await;
1874        }
1875        backend.close().await.unwrap();
1876
1877        match task.await.expect("worker task panicked") {
1878            Err(Error::ConsumerStopped(queue)) => assert!(queue.starts_with("test.")),
1879            other => panic!("unexpected result: {other:?}"),
1880        }
1881    }
1882
1883    #[tokio::test(start_paused = true)]
1884    async fn shutdown_before_run_stops_immediately() {
1885        let backend = backend();
1886        let worker = Worker::<TestQueues, _>::builder(backend.clone())
1887            .handler(FnHandler::<Greet, _>::new(
1888                |_j: Greet, _c: JobContext| async move { Ok(()) },
1889            ))
1890            .build()
1891            .await
1892            .unwrap();
1893        let handle = worker.handle();
1894        handle.shutdown();
1895        worker.run().await.unwrap();
1896        assert!(!backend.is_closed());
1897        assert!(format!("{handle:?}").contains("shutdown: true"));
1898    }
1899
1900    #[tokio::test(start_paused = true)]
1901    async fn handlers_of_different_job_types_are_routed_independently() {
1902        let backend = backend();
1903        let producer = producer(&backend).await;
1904        let greets = Arc::new(Recorder::default());
1905        let pings = Arc::new(Recorder::default());
1906
1907        let worker = {
1908            let (g, p) = (greets.clone(), pings.clone());
1909            Worker::<TestQueues, _>::builder(backend.clone())
1910                .handler(FnHandler::<Greet, _>::new(
1911                    move |_j: Greet, _c: JobContext| {
1912                        let g = g.clone();
1913                        async move {
1914                            g.enter();
1915                            g.leave();
1916                            Ok(())
1917                        }
1918                    },
1919                ))
1920                .handler(FnHandler::<Ping, _>::new(
1921                    move |_j: Ping, _c: JobContext| {
1922                        let p = p.clone();
1923                        async move {
1924                            p.enter();
1925                            p.leave();
1926                            Ok(())
1927                        }
1928                    },
1929                ))
1930                .build()
1931                .await
1932                .unwrap()
1933        };
1934        let (handle, task) = start(worker);
1935
1936        producer.enqueue(&Greet::new("ada")).await.unwrap();
1937        producer.enqueue(&Ping { seq: 1 }).await.unwrap();
1938        producer.enqueue(&Ping { seq: 2 }).await.unwrap();
1939
1940        wait_for(|| greets.attempts() == 1 && pings.attempts() == 2).await;
1941        assert_eq!(backend.acked(ALPHA).len(), 1);
1942        assert_eq!(backend.acked(BETA).len(), 2);
1943        stop(handle, task).await.unwrap();
1944    }
1945}