1use std::{
35 collections::{HashMap, VecDeque},
36 pin::Pin,
37 sync::{Arc, Mutex, MutexGuard},
38 task::{Context, Poll},
39 time::Duration,
40};
41
42use async_trait::async_trait;
43use futures::Stream;
44use tokio::sync::{OwnedSemaphorePermit, Semaphore, mpsc, watch};
45
46use crate::{
47 backend::{Backend, Delivery, DeliveryStream},
48 envelope::Envelope,
49 error::{Error, Result},
50 queue::QueueConfig,
51};
52
53struct QueueState {
55 config: QueueConfig,
57 pending: VecDeque<Envelope>,
59 acked: Vec<Envelope>,
61 dead: Vec<(Envelope, String)>,
63 consumers: Vec<Arc<Semaphore>>,
65 held: usize,
68 signal: watch::Sender<u64>,
70}
71
72impl QueueState {
73 fn new(config: QueueConfig) -> Self {
74 let (signal, _) = watch::channel(0);
75 Self {
76 config,
77 pending: VecDeque::new(),
78 acked: Vec::new(),
79 dead: Vec::new(),
80 consumers: Vec::new(),
81 held: 0,
82 signal,
83 }
84 }
85
86 fn insert_by_priority(&mut self, envelope: Envelope) {
93 let position = self
94 .pending
95 .iter()
96 .rposition(|pending| pending.priority >= envelope.priority)
97 .map_or(0, |last| last + 1);
98 self.pending.insert(position, envelope);
99 }
100
101 fn wake(&self) {
102 self.signal.send_modify(|v| *v = v.wrapping_add(1));
103 }
104}
105
106#[derive(Default)]
108struct Shared {
109 queues: HashMap<String, QueueState>,
110 closed: bool,
111}
112
113impl Shared {
114 fn queue_mut(&mut self, name: &str) -> &mut QueueState {
115 self.queues
116 .entry(name.to_owned())
117 .or_insert_with(|| QueueState::new(QueueConfig::new(name)))
118 }
119}
120
121#[derive(Clone, Default)]
130pub struct MemoryBackend {
131 state: Arc<Mutex<Shared>>,
132}
133
134impl std::fmt::Debug for MemoryBackend {
135 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
136 let state = self.lock();
137 f.debug_struct("MemoryBackend")
138 .field("queues", &state.queues.keys().collect::<Vec<_>>())
139 .field("closed", &state.closed)
140 .finish()
141 }
142}
143
144impl MemoryBackend {
145 pub fn new() -> Self {
147 Self::default()
148 }
149
150 fn lock(&self) -> MutexGuard<'_, Shared> {
151 self.state.lock().unwrap_or_else(|e| e.into_inner())
153 }
154
155 fn lock_shared(state: &Arc<Mutex<Shared>>) -> MutexGuard<'_, Shared> {
156 state.lock().unwrap_or_else(|e| e.into_inner())
157 }
158
159 pub fn pending(&self, queue: &str) -> usize {
164 self.lock().queues.get(queue).map_or(0, |q| q.pending.len())
165 }
166
167 pub fn acked(&self, queue: &str) -> Vec<Envelope> {
172 self.lock()
173 .queues
174 .get(queue)
175 .map(|q| q.acked.clone())
176 .unwrap_or_default()
177 }
178
179 pub fn deferred(&self, queue: &str) -> usize {
186 self.lock().queues.get(queue).map_or(0, |q| q.held)
187 }
188
189 pub fn dead_letters(&self, queue: &str) -> Vec<(Envelope, String)> {
191 self.lock()
192 .queues
193 .get(queue)
194 .map(|q| q.dead.clone())
195 .unwrap_or_default()
196 }
197
198 pub fn queue_names(&self) -> Vec<String> {
200 let mut names: Vec<_> = self.lock().queues.keys().cloned().collect();
201 names.sort();
202 names
203 }
204
205 pub fn queue_config(&self, queue: &str) -> Option<QueueConfig> {
207 self.lock().queues.get(queue).map(|q| q.config.clone())
208 }
209
210 pub fn is_closed(&self) -> bool {
212 self.lock().closed
213 }
214
215 #[cfg(test)]
220 pub(crate) fn consumer_count(&self, queue: &str) -> usize {
221 self.lock()
222 .queues
223 .get(queue)
224 .map_or(0, |q| q.consumers.len())
225 }
226
227 fn is_state_closed(state: &Arc<Mutex<Shared>>) -> bool {
229 Self::lock_shared(state).closed
230 }
231
232 fn enqueue_now(state: &Arc<Mutex<Shared>>, envelope: Envelope, hold: Option<HoldGuard>) {
242 let mut shared = Self::lock_shared(state);
243 if let Some(hold) = hold {
244 hold.release(&mut shared);
245 }
246 if shared.closed {
247 return;
248 }
249 let queue = shared.queue_mut(&envelope.queue);
250 queue.insert_by_priority(envelope);
251 queue.wake();
252 }
253
254 fn schedule(state: &Arc<Mutex<Shared>>, envelope: Envelope, delay: Duration) {
256 if delay.is_zero() {
257 Self::enqueue_now(state, envelope, None);
258 return;
259 }
260 let deadline = tokio::time::Instant::now() + delay;
263 let state = state.clone();
264 tokio::spawn(async move {
265 tokio::time::sleep_until(deadline).await;
266 Self::enqueue_now(&state, envelope, None);
267 });
268 }
269
270 fn hold(state: &Arc<Mutex<Shared>>, envelope: Envelope, delay: Duration) {
277 if delay.is_zero() {
278 Self::enqueue_now(state, envelope, None);
279 return;
280 }
281 let deadline = tokio::time::Instant::now() + delay;
284 Self::lock_shared(state).queue_mut(&envelope.queue).held += 1;
285 let guard = HoldGuard {
289 state: state.clone(),
290 queue: Some(envelope.queue.clone()),
291 };
292 let state = state.clone();
293 tokio::spawn(async move {
294 tokio::time::sleep_until(deadline).await;
295 Self::enqueue_now(&state, envelope, Some(guard));
296 });
297 }
298}
299
300struct HoldGuard {
302 state: Arc<Mutex<Shared>>,
303 queue: Option<String>,
305}
306
307impl HoldGuard {
308 fn release(mut self, shared: &mut Shared) {
313 if let Some(queue) = self.queue.take() {
314 release_hold(shared, &queue);
315 }
316 }
317}
318
319impl Drop for HoldGuard {
320 fn drop(&mut self) {
321 if let Some(queue) = self.queue.take() {
322 let mut shared = MemoryBackend::lock_shared(&self.state);
323 release_hold(&mut shared, &queue);
324 }
325 }
326}
327
328fn release_hold(shared: &mut Shared, queue: &str) {
330 if let Some(queue) = shared.queues.get_mut(queue) {
331 queue.held = queue.held.saturating_sub(1);
332 }
333}
334
335#[async_trait]
336impl Backend for MemoryBackend {
337 async fn declare(&self, queues: &[QueueConfig]) -> Result<()> {
338 let mut shared = self.lock();
339 if shared.closed {
340 return Err(Error::ShutDown);
341 }
342 for config in queues {
343 shared
345 .queues
346 .entry(config.name.clone())
347 .or_insert_with(|| QueueState::new(config.clone()));
348 }
349 Ok(())
350 }
351
352 async fn publish(&self, envelope: &Envelope, delay: Option<Duration>) -> Result<()> {
353 if self.is_closed() {
354 return Err(Error::ShutDown);
355 }
356 Self::schedule(&self.state, envelope.clone(), delay.unwrap_or_default());
357 Ok(())
358 }
359
360 async fn defer(&self, envelope: &Envelope, delay: Duration) -> Result<()> {
361 if self.is_closed() {
362 return Err(Error::ShutDown);
363 }
364 Self::hold(&self.state, envelope.clone(), delay);
365 Ok(())
366 }
367
368 async fn consume(&self, queue: &QueueConfig) -> Result<DeliveryStream> {
369 let (tx, rx) = mpsc::unbounded_channel();
370 let stream: DeliveryStream = Box::pin(DeliveryReceiver {
371 rx,
372 state: self.state.clone(),
373 queue: queue.name.clone(),
374 });
375
376 let permits = if queue.prefetch == 0 {
377 Semaphore::MAX_PERMITS
378 } else {
379 usize::from(queue.prefetch)
380 };
381 let semaphore = Arc::new(Semaphore::new(permits));
382
383 let signal = {
384 let mut shared = self.lock();
385 if shared.closed {
386 return Ok(stream);
388 }
389 let state = shared.queue_mut(&queue.name);
390 state.consumers.push(semaphore.clone());
391 state.signal.subscribe()
392 };
393
394 tokio::spawn(consumer_loop(
395 self.state.clone(),
396 queue.name.clone(),
397 semaphore,
398 tx,
399 signal,
400 ));
401 Ok(stream)
402 }
403
404 async fn close(&self) -> Result<()> {
405 let mut shared = self.lock();
406 shared.closed = true;
407 for queue in shared.queues.values_mut() {
408 for semaphore in queue.consumers.drain(..) {
409 semaphore.close();
411 }
412 queue.wake();
414 }
415 Ok(())
416 }
417}
418
419struct ConsumerGuard {
422 state: Arc<Mutex<Shared>>,
423 queue: String,
424 semaphore: Arc<Semaphore>,
425}
426
427impl Drop for ConsumerGuard {
428 fn drop(&mut self) {
429 let mut shared = MemoryBackend::lock_shared(&self.state);
430 if let Some(queue) = shared.queues.get_mut(&self.queue) {
431 queue.consumers.retain(|s| !Arc::ptr_eq(s, &self.semaphore));
432 }
433 }
434}
435
436async fn consumer_loop(
438 state: Arc<Mutex<Shared>>,
439 queue: String,
440 semaphore: Arc<Semaphore>,
441 tx: mpsc::UnboundedSender<Result<Box<dyn Delivery>>>,
442 mut signal: watch::Receiver<u64>,
443) {
444 let _guard = ConsumerGuard {
445 state: state.clone(),
446 queue: queue.clone(),
447 semaphore: semaphore.clone(),
448 };
449
450 loop {
451 let permit = tokio::select! {
454 biased;
455 () = tx.closed() => return,
456 permit = semaphore.clone().acquire_owned() => match permit {
457 Ok(permit) => permit,
458 Err(_) => return,
459 },
460 };
461
462 let envelope = loop {
463 let next = {
464 let mut shared = MemoryBackend::lock_shared(&state);
465 if shared.closed {
466 return;
467 }
468 shared
469 .queues
470 .get_mut(&queue)
471 .and_then(|q| q.pending.pop_front())
472 };
473 match next {
474 Some(envelope) => break envelope,
475 None => {
476 tokio::select! {
479 biased;
480 () = tx.closed() => return,
481 changed = signal.changed() => if changed.is_err() {
482 return;
483 },
484 }
485 }
486 }
487 };
488
489 let delivery = MemoryDelivery {
490 state: state.clone(),
491 envelope: envelope.clone(),
492 _permit: permit,
493 };
494 if tx.send(Ok(Box::new(delivery))).is_err() {
495 let mut shared = MemoryBackend::lock_shared(&state);
497 if let Some(q) = shared.queues.get_mut(&queue) {
498 q.pending.push_front(envelope);
499 q.wake();
500 }
501 return;
502 }
503 }
504}
505
506struct MemoryDelivery {
508 state: Arc<Mutex<Shared>>,
509 envelope: Envelope,
510 _permit: OwnedSemaphorePermit,
511}
512
513impl MemoryDelivery {
514 fn record_ack(&self) {
515 let mut shared = MemoryBackend::lock_shared(&self.state);
516 shared
517 .queue_mut(&self.envelope.queue)
518 .acked
519 .push(self.envelope.clone());
520 }
521}
522
523#[async_trait]
524impl Delivery for MemoryDelivery {
525 fn envelope(&self) -> &Envelope {
526 &self.envelope
527 }
528
529 async fn ack(self: Box<Self>) -> Result<()> {
530 self.record_ack();
531 Ok(())
532 }
533
534 async fn dead_letter(self: Box<Self>, reason: &str) -> Result<()> {
535 let mut shared = MemoryBackend::lock_shared(&self.state);
536 if shared.closed {
537 return Err(Error::ShutDown);
538 }
539 shared
540 .queue_mut(&self.envelope.queue)
541 .dead
542 .push((self.envelope.clone(), reason.to_owned()));
543 Ok(())
544 }
545
546 async fn retry(self: Box<Self>, next: Envelope, delay: Duration) -> Result<()> {
547 if MemoryBackend::is_state_closed(&self.state) {
550 return Err(Error::ShutDown);
551 }
552 MemoryBackend::schedule(&self.state, next, delay);
554 self.record_ack();
555 Ok(())
556 }
557
558 async fn defer(self: Box<Self>, next: Envelope, delay: Duration) -> Result<()> {
559 if MemoryBackend::is_state_closed(&self.state) {
562 return Err(Error::ShutDown);
563 }
564 MemoryBackend::hold(&self.state, next, delay);
567 self.record_ack();
568 Ok(())
569 }
570}
571
572struct DeliveryReceiver {
574 rx: mpsc::UnboundedReceiver<Result<Box<dyn Delivery>>>,
575 state: Arc<Mutex<Shared>>,
576 queue: String,
577}
578
579impl Stream for DeliveryReceiver {
580 type Item = Result<Box<dyn Delivery>>;
581
582 fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
583 self.get_mut().rx.poll_recv(cx)
584 }
585}
586
587impl Drop for DeliveryReceiver {
588 fn drop(&mut self) {
592 self.rx.close();
593 let mut unread = Vec::new();
594 while let Ok(item) = self.rx.try_recv() {
595 if let Ok(delivery) = item {
596 unread.push(delivery.envelope().clone());
597 }
598 }
599 if unread.is_empty() {
600 return;
601 }
602 let mut shared = MemoryBackend::lock_shared(&self.state);
603 if shared.closed {
604 return;
605 }
606 let queue = shared.queue_mut(&self.queue);
607 for envelope in unread.into_iter().rev() {
609 queue.pending.push_front(envelope);
610 }
611 queue.wake();
612 }
613}
614
615#[cfg(test)]
616mod tests {
617 use super::*;
618 use crate::{
619 job::Job,
620 queue::QueueSet,
621 test_support::{Greet, Ping, TestQueues},
622 };
623 use futures::{StreamExt, future::poll_immediate};
624
625 fn alpha() -> QueueConfig {
626 TestQueues::Alpha.config()
627 }
628
629 async fn declared() -> MemoryBackend {
630 let backend = MemoryBackend::new();
631 backend.declare(&[alpha()]).await.unwrap();
632 backend
633 }
634
635 async fn publish(backend: &MemoryBackend, name: &str) -> Envelope {
636 let envelope = Envelope::new(&Greet::new(name)).unwrap();
637 backend.publish(&envelope, None).await.unwrap();
638 envelope
639 }
640
641 async fn next_delivery(stream: &mut DeliveryStream) -> Box<dyn Delivery> {
642 stream
643 .next()
644 .await
645 .expect("stream ended")
646 .expect("delivery error")
647 }
648
649 #[tokio::test]
650 async fn declare_is_idempotent() {
651 let backend = declared().await;
652 publish(&backend, "a").await;
653
654 backend.declare(&[alpha(), alpha()]).await.unwrap();
655 backend.declare(&[alpha()]).await.unwrap();
656
657 assert_eq!(backend.pending("test.alpha"), 1);
658 assert_eq!(backend.queue_names(), vec!["test.alpha".to_owned()]);
659 assert_eq!(backend.queue_config("test.alpha").unwrap().prefetch, 4);
660 }
661
662 #[tokio::test]
663 async fn declares_every_queue_of_a_queue_set() {
664 let backend = MemoryBackend::new();
665 let configs: Vec<_> = TestQueues::all().iter().map(|q| q.config()).collect();
666 backend.declare(&configs).await.unwrap();
667 assert_eq!(
668 backend.queue_names(),
669 vec![
670 "test.alpha".to_owned(),
671 "test.beta".to_owned(),
672 "test.gamma".to_owned()
673 ]
674 );
675 }
676
677 #[tokio::test]
678 async fn publish_then_consume_is_fifo() {
679 let backend = declared().await;
680 for name in ["a", "b", "c"] {
681 publish(&backend, name).await;
682 }
683 assert_eq!(backend.pending("test.alpha"), 3);
684
685 let mut stream = backend.consume(&alpha()).await.unwrap();
686 for name in ["a", "b", "c"] {
687 let delivery = next_delivery(&mut stream).await;
688 assert_eq!(
689 delivery.envelope().decode::<Greet>().unwrap(),
690 Greet::new(name)
691 );
692 assert_eq!(delivery.envelope().attempt, 1);
693 delivery.ack().await.unwrap();
694 }
695
696 assert_eq!(backend.pending("test.alpha"), 0);
697 let acked: Vec<_> = backend
698 .acked("test.alpha")
699 .iter()
700 .map(|e| e.decode::<Greet>().unwrap().name)
701 .collect();
702 assert_eq!(acked, vec!["a", "b", "c"]);
703 }
704
705 #[tokio::test]
706 async fn consumer_started_before_publish_receives_messages() {
707 let backend = declared().await;
708 let mut stream = backend.consume(&alpha()).await.unwrap();
709 tokio::task::yield_now().await;
711
712 publish(&backend, "late").await;
713 let delivery = next_delivery(&mut stream).await;
714 assert_eq!(
715 delivery.envelope().decode::<Greet>().unwrap(),
716 Greet::new("late")
717 );
718 delivery.ack().await.unwrap();
719 }
720
721 #[tokio::test]
722 async fn prefetch_caps_outstanding_deliveries() {
723 let backend = declared().await;
724 let config = QueueConfig::new("test.alpha").prefetch(2);
725 for name in ["a", "b", "c"] {
726 publish(&backend, name).await;
727 }
728
729 let mut stream = backend.consume(&config).await.unwrap();
730 let first = next_delivery(&mut stream).await;
731 let second = next_delivery(&mut stream).await;
732 assert_eq!(first.envelope().decode::<Greet>().unwrap(), Greet::new("a"));
733 assert_eq!(
734 second.envelope().decode::<Greet>().unwrap(),
735 Greet::new("b")
736 );
737
738 assert!(poll_immediate(stream.next()).await.is_none());
740 assert_eq!(backend.pending("test.alpha"), 1);
741
742 first.ack().await.unwrap();
743 let third = next_delivery(&mut stream).await;
744 assert_eq!(third.envelope().decode::<Greet>().unwrap(), Greet::new("c"));
745 assert_eq!(backend.pending("test.alpha"), 0);
746
747 second.ack().await.unwrap();
748 third.ack().await.unwrap();
749 assert_eq!(backend.acked("test.alpha").len(), 3);
750 }
751
752 #[tokio::test]
753 async fn zero_prefetch_means_unlimited() {
754 let backend = declared().await;
755 for name in ["a", "b", "c"] {
756 publish(&backend, name).await;
757 }
758 let mut stream = backend
759 .consume(&QueueConfig::new("test.alpha").prefetch(0))
760 .await
761 .unwrap();
762 let mut held = Vec::new();
763 for _ in 0..3 {
764 held.push(next_delivery(&mut stream).await);
765 }
766 assert_eq!(held.len(), 3);
767 }
768
769 #[tokio::test(start_paused = true)]
770 async fn delayed_publish_is_invisible_until_the_delay_elapses() {
771 let backend = declared().await;
772 let envelope = Envelope::new(&Greet::new("later")).unwrap();
773 backend
774 .publish(&envelope, Some(Duration::from_secs(30)))
775 .await
776 .unwrap();
777
778 tokio::time::sleep(Duration::from_secs(29)).await;
781 assert_eq!(backend.pending("test.alpha"), 0);
782
783 tokio::time::sleep(Duration::from_secs(2)).await;
784 assert_eq!(backend.pending("test.alpha"), 1);
785
786 let mut stream = backend.consume(&alpha()).await.unwrap();
787 let delivery = next_delivery(&mut stream).await;
788 assert_eq!(delivery.envelope().job_id, envelope.job_id);
789 delivery.ack().await.unwrap();
790 }
791
792 #[tokio::test(start_paused = true)]
793 async fn retry_redelivers_with_a_higher_attempt_after_the_delay() {
794 let backend = declared().await;
795 let original = publish(&backend, "flaky").await;
796 let mut stream = backend.consume(&alpha()).await.unwrap();
797
798 let delivery = next_delivery(&mut stream).await;
799 let next = delivery.envelope().next_attempt();
800 delivery.retry(next, Duration::from_secs(10)).await.unwrap();
801
802 assert_eq!(backend.acked("test.alpha").len(), 1);
804 assert_eq!(backend.pending("test.alpha"), 0);
805 tokio::time::advance(Duration::from_secs(5)).await;
806 assert!(poll_immediate(stream.next()).await.is_none());
807
808 tokio::time::advance(Duration::from_secs(6)).await;
809 let redelivered = next_delivery(&mut stream).await;
810 assert_eq!(redelivered.envelope().attempt, 2);
811 assert_eq!(redelivered.envelope().job_id, original.job_id);
812 redelivered.ack().await.unwrap();
813 assert_eq!(backend.acked("test.alpha").len(), 2);
814 }
815
816 #[tokio::test]
817 async fn retry_without_delay_is_immediate() {
818 let backend = declared().await;
819 publish(&backend, "now").await;
820 let mut stream = backend.consume(&alpha()).await.unwrap();
821
822 let delivery = next_delivery(&mut stream).await;
823 let next = delivery.envelope().next_attempt();
824 delivery.retry(next, Duration::ZERO).await.unwrap();
825
826 let redelivered = next_delivery(&mut stream).await;
827 assert_eq!(redelivered.envelope().attempt, 2);
828 redelivered.ack().await.unwrap();
829 }
830
831 async fn publish_with_priority(backend: &MemoryBackend, name: &str, priority: u8) -> Envelope {
833 let mut envelope = Envelope::new(&Greet::new(name)).unwrap();
834 envelope.priority = priority;
835 backend.publish(&envelope, None).await.unwrap();
836 envelope
837 }
838
839 fn pending_names(backend: &MemoryBackend, queue: &str) -> Vec<String> {
841 backend
842 .lock()
843 .queues
844 .get(queue)
845 .map(|q| {
846 q.pending
847 .iter()
848 .map(|e| e.decode::<Greet>().unwrap().name)
849 .collect()
850 })
851 .unwrap_or_default()
852 }
853
854 #[tokio::test]
855 async fn publish_orders_by_priority_and_keeps_each_level_fifo() {
856 let backend = declared().await;
857 for (name, priority) in [
858 ("normal-1", 0),
859 ("urgent-1", 10),
860 ("normal-2", 0),
861 ("urgent-2", 10),
862 ("middle", 5),
863 ("normal-3", 0),
864 ] {
865 publish_with_priority(&backend, name, priority).await;
866 }
867
868 assert_eq!(
869 pending_names(&backend, "test.alpha"),
870 vec![
871 "urgent-1", "urgent-2", "middle", "normal-1", "normal-2", "normal-3", ]
875 );
876 }
877
878 #[tokio::test]
879 async fn a_higher_priority_publish_overtakes_the_whole_backlog() {
880 let backend = declared().await;
881 for name in ["a", "b", "c"] {
882 publish(&backend, name).await;
883 }
884 publish_with_priority(&backend, "jumper", 1).await;
885
886 let mut stream = backend.consume(&alpha()).await.unwrap();
887 let first = next_delivery(&mut stream).await;
888 assert_eq!(
889 first.envelope().decode::<Greet>().unwrap(),
890 Greet::new("jumper")
891 );
892 first.ack().await.unwrap();
893 }
894
895 #[tokio::test(start_paused = true)]
896 async fn defer_holds_the_message_and_then_lands_it_ahead_of_the_backlog() {
897 let backend = declared().await;
898 for name in ["backlog-1", "backlog-2"] {
899 publish(&backend, name).await;
900 }
901 let held = Envelope::new(&Greet::new("held")).unwrap().deferred(10);
902 backend.defer(&held, Duration::from_secs(30)).await.unwrap();
903
904 assert_eq!(backend.pending("test.alpha"), 2);
906 assert_eq!(backend.deferred("test.alpha"), 1);
907 tokio::time::sleep(Duration::from_secs(29)).await;
908 assert_eq!(backend.pending("test.alpha"), 2);
909 assert_eq!(backend.deferred("test.alpha"), 1);
910
911 tokio::time::sleep(Duration::from_secs(2)).await;
912 assert_eq!(backend.deferred("test.alpha"), 0, "the hold has expired");
913 assert_eq!(
914 pending_names(&backend, "test.alpha"),
915 vec!["held", "backlog-1", "backlog-2"],
916 "a deferred job comes back in front"
917 );
918
919 let mut stream = backend.consume(&alpha()).await.unwrap();
920 let first = next_delivery(&mut stream).await;
921 assert_eq!(first.envelope().job_id, held.job_id);
922 assert_eq!(first.envelope().deferrals, 1);
923 assert_eq!(first.envelope().priority, 10);
924 assert_eq!(first.envelope().attempt, 1, "a deferral is not an attempt");
925 first.ack().await.unwrap();
926 }
927
928 #[tokio::test(start_paused = true)]
929 async fn deferred_counts_every_message_in_hold() {
930 let backend = declared().await;
931 assert_eq!(backend.deferred("test.alpha"), 0);
932 assert_eq!(backend.deferred("never.declared"), 0);
933
934 for _ in 0..3 {
935 backend
936 .defer(
937 &Envelope::new(&Greet::new("held")).unwrap(),
938 Duration::from_secs(5),
939 )
940 .await
941 .unwrap();
942 }
943 assert_eq!(backend.deferred("test.alpha"), 3);
944
945 tokio::time::sleep(Duration::from_secs(6)).await;
946 assert_eq!(backend.deferred("test.alpha"), 0);
947 assert_eq!(backend.pending("test.alpha"), 3);
948 }
949
950 fn held_plus_pending(backend: &MemoryBackend, queue: &str) -> usize {
953 backend
954 .lock()
955 .queues
956 .get(queue)
957 .map_or(0, |q| q.held + q.pending.len())
958 }
959
960 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
961 async fn a_hold_is_never_counted_twice_while_it_expires() {
962 const HELD: usize = 64;
963
964 let backend = declared().await;
965 for i in 0..HELD {
966 backend
967 .defer(
968 &Envelope::new(&Greet::new(&i.to_string())).unwrap(),
969 Duration::from_millis(5),
970 )
971 .await
972 .unwrap();
973 }
974
975 for _ in 0..1_000_000 {
979 let total = held_plus_pending(&backend, "test.alpha");
980 assert!(
981 total <= HELD,
982 "an envelope was counted as held and pending at once ({total} > {HELD})"
983 );
984 if backend.pending("test.alpha") == HELD {
985 break;
986 }
987 tokio::task::yield_now().await;
988 }
989
990 assert_eq!(backend.pending("test.alpha"), HELD);
991 assert_eq!(backend.deferred("test.alpha"), 0);
992 }
993
994 #[tokio::test]
995 async fn a_zero_delay_deferral_never_enters_the_hold() {
996 let backend = declared().await;
997 backend
998 .defer(&Envelope::new(&Greet::new("now")).unwrap(), Duration::ZERO)
999 .await
1000 .unwrap();
1001 assert_eq!(backend.deferred("test.alpha"), 0);
1002 assert_eq!(backend.pending("test.alpha"), 1);
1003 }
1004
1005 #[tokio::test]
1006 async fn defer_after_close_reports_the_shutdown() {
1007 let backend = declared().await;
1008 backend.close().await.unwrap();
1009
1010 assert!(matches!(
1011 backend
1012 .defer(
1013 &Envelope::new(&Greet::new("x")).unwrap(),
1014 Duration::from_secs(1)
1015 )
1016 .await,
1017 Err(Error::ShutDown)
1018 ));
1019 assert_eq!(backend.deferred("test.alpha"), 0);
1020 assert_eq!(backend.pending("test.alpha"), 0);
1021 }
1022
1023 #[tokio::test(start_paused = true)]
1024 async fn a_deferral_that_lands_after_close_is_dropped() {
1025 let backend = declared().await;
1026 backend
1027 .defer(
1028 &Envelope::new(&Greet::new("held")).unwrap(),
1029 Duration::from_secs(30),
1030 )
1031 .await
1032 .unwrap();
1033 assert_eq!(backend.deferred("test.alpha"), 1);
1034 backend.close().await.unwrap();
1035
1036 tokio::time::sleep(Duration::from_secs(31)).await;
1037 assert_eq!(backend.pending("test.alpha"), 0);
1038 assert_eq!(backend.deferred("test.alpha"), 0);
1039 }
1040
1041 #[tokio::test(start_paused = true)]
1042 async fn delivery_defer_acks_the_original_and_reschedules_it() {
1043 let backend = declared().await;
1044 let original = publish(&backend, "rate-limited").await;
1045 let mut stream = backend.consume(&alpha()).await.unwrap();
1046
1047 let delivery = next_delivery(&mut stream).await;
1048 let next = delivery.envelope().deferred(10);
1049 delivery.defer(next, Duration::from_secs(30)).await.unwrap();
1050
1051 assert_eq!(backend.acked("test.alpha").len(), 1);
1053 assert_eq!(backend.acked("test.alpha")[0].job_id, original.job_id);
1054 assert_eq!(backend.acked("test.alpha")[0].deferrals, 0);
1055 assert_eq!(backend.pending("test.alpha"), 0);
1056 assert_eq!(backend.deferred("test.alpha"), 1);
1057 tokio::time::advance(Duration::from_secs(10)).await;
1058 assert!(poll_immediate(stream.next()).await.is_none());
1059
1060 tokio::time::advance(Duration::from_secs(21)).await;
1061 let redelivered = next_delivery(&mut stream).await;
1062 assert_eq!(redelivered.envelope().job_id, original.job_id);
1063 assert_eq!(redelivered.envelope().attempt, 1);
1064 assert_eq!(redelivered.envelope().deferrals, 1);
1065 assert_eq!(redelivered.envelope().priority, 10);
1066 assert_eq!(backend.deferred("test.alpha"), 0);
1067 redelivered.ack().await.unwrap();
1068 assert_eq!(backend.acked("test.alpha").len(), 2);
1069 }
1070
1071 #[tokio::test]
1072 async fn delivery_defer_frees_a_prefetch_slot() {
1073 let backend = declared().await;
1074 for name in ["a", "b"] {
1075 publish(&backend, name).await;
1076 }
1077 let mut stream = backend
1078 .consume(&QueueConfig::new("test.alpha").prefetch(1))
1079 .await
1080 .unwrap();
1081
1082 let first = next_delivery(&mut stream).await;
1083 assert!(poll_immediate(stream.next()).await.is_none());
1084 let next = first.envelope().deferred(10);
1085 first.defer(next, Duration::from_secs(600)).await.unwrap();
1086
1087 let second = next_delivery(&mut stream).await;
1088 assert_eq!(
1089 second.envelope().decode::<Greet>().unwrap(),
1090 Greet::new("b")
1091 );
1092 second.ack().await.unwrap();
1093 }
1094
1095 #[tokio::test]
1096 async fn delivery_defer_after_close_reports_the_shutdown() {
1097 let backend = declared().await;
1098 publish(&backend, "a").await;
1099 let mut stream = backend.consume(&alpha()).await.unwrap();
1100 let delivery = next_delivery(&mut stream).await;
1101
1102 backend.close().await.unwrap();
1103
1104 let next = delivery.envelope().deferred(10);
1105 assert!(matches!(
1106 delivery.defer(next, Duration::from_secs(1)).await,
1107 Err(Error::ShutDown)
1108 ));
1109 assert!(backend.acked("test.alpha").is_empty());
1111 assert_eq!(backend.deferred("test.alpha"), 0);
1112 }
1113
1114 #[tokio::test]
1115 async fn dead_letter_records_the_reason() {
1116 let backend = declared().await;
1117 let original = publish(&backend, "doomed").await;
1118 let mut stream = backend.consume(&alpha()).await.unwrap();
1119
1120 let delivery = next_delivery(&mut stream).await;
1121 delivery
1122 .dead_letter("max attempts (3) exhausted")
1123 .await
1124 .unwrap();
1125
1126 let dead = backend.dead_letters("test.alpha");
1127 assert_eq!(dead.len(), 1);
1128 assert_eq!(dead[0].0.job_id, original.job_id);
1129 assert_eq!(dead[0].1, "max attempts (3) exhausted");
1130 assert!(backend.acked("test.alpha").is_empty());
1131 assert_eq!(backend.pending("test.alpha"), 0);
1132 }
1133
1134 #[tokio::test]
1135 async fn dead_letter_frees_a_prefetch_slot() {
1136 let backend = declared().await;
1137 for name in ["a", "b"] {
1138 publish(&backend, name).await;
1139 }
1140 let mut stream = backend
1141 .consume(&QueueConfig::new("test.alpha").prefetch(1))
1142 .await
1143 .unwrap();
1144
1145 let first = next_delivery(&mut stream).await;
1146 assert!(poll_immediate(stream.next()).await.is_none());
1147 first.dead_letter("nope").await.unwrap();
1148
1149 let second = next_delivery(&mut stream).await;
1150 assert_eq!(
1151 second.envelope().decode::<Greet>().unwrap(),
1152 Greet::new("b")
1153 );
1154 second.ack().await.unwrap();
1155 }
1156
1157 #[tokio::test]
1158 async fn close_ends_all_streams() {
1159 let backend = declared().await;
1160 let mut busy = backend.consume(&alpha()).await.unwrap();
1161 publish(&backend, "held").await;
1162 let held = next_delivery(&mut busy).await;
1164 let mut idle = backend.consume(&alpha()).await.unwrap();
1165
1166 backend.close().await.unwrap();
1167
1168 assert!(idle.next().await.is_none());
1169 assert!(busy.next().await.is_none());
1171 assert!(backend.is_closed());
1172 held.ack().await.unwrap();
1174 assert_eq!(backend.acked("test.alpha").len(), 1);
1175
1176 assert!(matches!(
1177 backend
1178 .publish(&Envelope::new(&Greet::new("x")).unwrap(), None)
1179 .await,
1180 Err(Error::ShutDown)
1181 ));
1182 assert!(matches!(
1183 backend.declare(&[alpha()]).await,
1184 Err(Error::ShutDown)
1185 ));
1186 }
1187
1188 #[tokio::test]
1189 async fn retry_and_dead_letter_after_close_report_the_shutdown() {
1190 let backend = declared().await;
1191 publish(&backend, "a").await;
1192 publish(&backend, "b").await;
1193 let mut stream = backend.consume(&alpha()).await.unwrap();
1194 let first = next_delivery(&mut stream).await;
1195 let second = next_delivery(&mut stream).await;
1196
1197 backend.close().await.unwrap();
1198
1199 let next = first.envelope().next_attempt();
1200 assert!(matches!(
1201 first.retry(next, Duration::ZERO).await,
1202 Err(Error::ShutDown)
1203 ));
1204 assert!(matches!(
1205 second.dead_letter("too late").await,
1206 Err(Error::ShutDown)
1207 ));
1208 assert!(backend.acked("test.alpha").is_empty());
1210 assert!(backend.dead_letters("test.alpha").is_empty());
1211 assert_eq!(backend.pending("test.alpha"), 0);
1212 }
1213
1214 #[tokio::test]
1215 async fn a_delayed_publish_that_lands_after_close_is_dropped() {
1216 let backend = declared().await;
1217 backend
1218 .publish(
1219 &Envelope::new(&Greet::new("later")).unwrap(),
1220 Some(Duration::from_millis(1)),
1221 )
1222 .await
1223 .unwrap();
1224 backend.close().await.unwrap();
1225
1226 tokio::time::sleep(Duration::from_millis(20)).await;
1227 assert_eq!(backend.pending("test.alpha"), 0);
1228 }
1229
1230 #[tokio::test]
1231 async fn consumers_are_forgotten_when_their_streams_are_dropped() {
1232 let backend = declared().await;
1233 let streams: Vec<_> = {
1234 let mut v = Vec::new();
1235 for _ in 0..3 {
1236 v.push(backend.consume(&alpha()).await.unwrap());
1237 }
1238 v
1239 };
1240 assert_eq!(backend.consumer_count("test.alpha"), 3);
1241
1242 drop(streams);
1243 for _ in 0..1_000 {
1245 if backend.consumer_count("test.alpha") == 0 {
1246 break;
1247 }
1248 tokio::task::yield_now().await;
1249 }
1250 assert_eq!(backend.consumer_count("test.alpha"), 0);
1251 }
1252
1253 #[tokio::test]
1254 async fn consuming_after_close_yields_an_ended_stream() {
1255 let backend = declared().await;
1256 backend.close().await.unwrap();
1257 let mut stream = backend.consume(&alpha()).await.unwrap();
1258 assert!(stream.next().await.is_none());
1259 }
1260
1261 #[tokio::test]
1262 async fn multiple_consumers_share_one_queue() {
1263 let backend = declared().await;
1264 let config = QueueConfig::new("test.alpha").prefetch(1);
1265 let mut left = backend.consume(&config).await.unwrap();
1266 let mut right = backend.consume(&config).await.unwrap();
1267 tokio::task::yield_now().await;
1268
1269 for name in ["a", "b"] {
1270 publish(&backend, name).await;
1271 }
1272
1273 let one = next_delivery(&mut left).await;
1274 let two = next_delivery(&mut right).await;
1275 let names = [
1276 one.envelope().decode::<Greet>().unwrap().name,
1277 two.envelope().decode::<Greet>().unwrap().name,
1278 ];
1279 assert!(
1280 names.contains(&"a".to_owned()) && names.contains(&"b".to_owned()),
1281 "{names:?}"
1282 );
1283
1284 one.ack().await.unwrap();
1285 two.ack().await.unwrap();
1286 assert_eq!(backend.acked("test.alpha").len(), 2);
1287 }
1288
1289 #[tokio::test]
1290 async fn dropping_a_stream_returns_undelivered_messages_to_the_queue() {
1291 let backend = declared().await;
1292 let stream = backend
1293 .consume(&QueueConfig::new("test.alpha").prefetch(2))
1294 .await
1295 .unwrap();
1296 publish(&backend, "orphaned").await;
1297 publish(&backend, "also-orphaned").await;
1298
1299 for _ in 0..1_000 {
1302 if backend.pending("test.alpha") == 0 {
1303 break;
1304 }
1305 tokio::task::yield_now().await;
1306 }
1307 assert_eq!(backend.pending("test.alpha"), 0);
1308
1309 drop(stream);
1311 assert_eq!(backend.pending("test.alpha"), 2);
1312 let mut stream = backend.consume(&alpha()).await.unwrap();
1313 for name in ["orphaned", "also-orphaned"] {
1314 let delivery = next_delivery(&mut stream).await;
1315 assert_eq!(
1316 delivery.envelope().decode::<Greet>().unwrap(),
1317 Greet::new(name)
1318 );
1319 delivery.ack().await.unwrap();
1320 }
1321 }
1322
1323 #[tokio::test]
1324 async fn queues_are_independent() {
1325 let backend = MemoryBackend::new();
1326 let configs: Vec<_> = TestQueues::all().iter().map(|q| q.config()).collect();
1327 backend.declare(&configs).await.unwrap();
1328
1329 backend
1330 .publish(&Envelope::new(&Greet::new("a")).unwrap(), None)
1331 .await
1332 .unwrap();
1333 backend
1334 .publish(&Envelope::new(&Ping { seq: 1 }).unwrap(), None)
1335 .await
1336 .unwrap();
1337
1338 assert_eq!(backend.pending("test.alpha"), 1);
1339 assert_eq!(backend.pending("test.beta"), 1);
1340
1341 let mut stream = backend.consume(&TestQueues::Beta.config()).await.unwrap();
1342 let delivery = next_delivery(&mut stream).await;
1343 assert_eq!(delivery.envelope().job_type, Ping::NAME);
1344 delivery.ack().await.unwrap();
1345 assert_eq!(backend.pending("test.alpha"), 1);
1346 }
1347
1348 #[tokio::test]
1349 async fn clones_share_state() {
1350 let backend = declared().await;
1351 let clone = backend.clone();
1352 publish(&clone, "shared").await;
1353 assert_eq!(backend.pending("test.alpha"), 1);
1354 assert!(format!("{backend:?}").contains("test.alpha"));
1355 }
1356
1357 #[tokio::test]
1358 async fn delivery_is_usable_behind_a_trait_object() {
1359 let backend = declared().await;
1360 publish(&backend, "boxed").await;
1361 let mut stream = backend.consume(&alpha()).await.unwrap();
1362 let delivery: Box<dyn Delivery> = next_delivery(&mut stream).await;
1363 let envelope = delivery.envelope().clone();
1364 delivery.ack().await.unwrap();
1365 assert_eq!(backend.acked("test.alpha"), vec![envelope]);
1366 }
1367}