1use 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#[derive(Debug)]
67enum JobOutcome {
68 Success,
70 Decode(String),
72 Retryable(String),
74 Fatal(String),
76 Deferred {
78 delay: Duration,
80 reason: String,
82 },
83}
84
85#[async_trait]
87trait ErasedHandler: Send + Sync + 'static {
88 fn policy(&self) -> RetryPolicy;
90
91 fn defer_priority(&self) -> u8;
94
95 async fn run(&self, envelope: &Envelope, timeout: Option<Duration>) -> JobOutcome;
97}
98
99struct 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 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 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
164pub 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 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 pub fn handle(&self) -> WorkerHandle {
208 self.handle.clone()
209 }
210
211 pub fn queues(&self) -> &[QueueConfig] {
213 &self.queues
214 }
215
216 pub fn concurrency(&self) -> usize {
218 self.concurrency
219 }
220
221 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 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 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 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 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 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 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
384async 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
415enum ConsumerEvent {
417 Delivery(Box<dyn Delivery>),
418 Failed(Error),
419 Ended(String),
420}
421
422async 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 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
520fn 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
532pub 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 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 pub fn queues(mut self, queues: &[Q]) -> Self {
563 self.queues = Some(queues.to_vec());
564 self
565 }
566
567 pub fn concurrency(mut self, concurrency: usize) -> Self {
571 self.concurrency = Some(concurrency);
572 self
573 }
574
575 pub fn job_timeout(mut self, timeout: Duration) -> Self {
577 self.job_timeout = Some(timeout);
578 self
579 }
580
581 pub fn close_backend_on_shutdown(mut self, close: bool) -> Self {
589 self.close_backend = close;
590 self
591 }
592
593 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#[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 pub fn shutdown(&self) {
653 self.shutdown.send_replace(true);
654 }
655
656 pub fn is_shutdown(&self) -> bool {
658 *self.shutdown.borrow()
659 }
660
661 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 #[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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 stop(handle, task).await.unwrap();
1505
1506 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 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 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 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 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 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 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 {
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}