Skip to main content

queuey_core/
memory.rs

1//! In-memory backend for tests and local development.
2//!
3//! Contract:
4//! * `MemoryBackend::new()`; `Clone`-able handle over shared `Arc` state.
5//! * `publish` with `delay` uses `tokio::time::sleep` in a spawned task (so it is
6//!   compatible with `tokio::time::pause()` / `advance()` in tests).
7//! * `consume` honours `prefetch` (at most `prefetch` un-acked deliveries per consumer).
8//! * `retry` re-publishes `next` after `delay`, then acks the original.
9//! * `defer` holds `next` for `delay` and then inserts it *by priority*, then acks the
10//!   original. The hold is a `tokio::time::sleep` too, so it is virtual-time friendly.
11//! * Every enqueue (plain, delayed or deferred) goes through one priority insertion:
12//!   the envelope lands behind every pending envelope whose `priority` is greater than
13//!   or equal to its own and ahead of the rest, so equal priorities stay FIFO. A hold
14//!   ends inside that same critical section, so no envelope is ever counted by both
15//!   `deferred` and `pending`.
16//! * `dead_letter` moves the envelope into an inspectable `dead_letters(queue)` list with reason.
17//! * Test helpers: `pending(queue) -> usize`, `deferred(queue) -> usize`,
18//!   `dead_letters(queue) -> Vec<(Envelope, String)>`, `acked(queue) -> Vec<Envelope>`.
19//! * `close` ends all consumer streams. Afterwards `declare`, `publish`, `defer`,
20//!   `retry` and `dead_letter` all fail with [`Error::ShutDown`]; a plain `ack` still
21//!   records. A delayed publish or a deferral whose timer fires *after* the close is
22//!   dropped, exactly like a broker that went away mid-wait.
23//!
24//! Implementation notes:
25//! * Each consumer owns a [`tokio::sync::Semaphore`] with `prefetch` permits. A permit
26//!   is moved into every [`Delivery`] handed out and released on ack / retry /
27//!   dead-letter, which is what caps the number of outstanding deliveries. The
28//!   semaphore is dropped from the queue when the consumer task ends.
29//! * Dropping a stream requeues every message that was produced for it but not handed
30//!   out yet, so a message can never vanish just because a consumer went away.
31//! * A delivery that is dropped without being settled releases its permit but the
32//!   message is *not* requeued (unlike a real broker). Tests should settle deliveries.
33
34use 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
53/// State of one queue inside the backend.
54struct QueueState {
55    /// Config from the first `declare` (or a default one, for implicit queues).
56    config: QueueConfig,
57    /// Messages waiting to be handed to a consumer.
58    pending: VecDeque<Envelope>,
59    /// Everything that was acked, in ack order (retries count as an ack of the original).
60    acked: Vec<Envelope>,
61    /// Everything that was dead-lettered, with its reason.
62    dead: Vec<(Envelope, String)>,
63    /// One semaphore per live consumer; closing it stops that consumer.
64    consumers: Vec<Arc<Semaphore>>,
65    /// Deferred envelopes whose hold has not elapsed yet. Stands in for the hold
66    /// queues a real broker would own.
67    held: usize,
68    /// Bumped whenever `pending` grows or the backend closes, to wake consumers.
69    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    /// Queue `envelope` by priority: behind every pending envelope that is at least as
87    /// important, ahead of the rest. Equal priorities therefore stay FIFO, and a
88    /// deferred job (top priority) overtakes the normal backlog (priority `0`).
89    ///
90    /// The single insertion point for `publish` and `defer` alike, so priority
91    /// ordering is one code path.
92    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/// Everything the cloned handles share.
107#[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/// An in-memory [`Backend`]. Cloning gives another handle onto the same state.
122///
123/// ```
124/// # use queuey_core::MemoryBackend;
125/// let backend = MemoryBackend::new();
126/// let same_state = backend.clone();
127/// assert_eq!(same_state.pending("nothing.here"), 0);
128/// ```
129#[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    /// Create an empty backend with no queues declared.
146    pub fn new() -> Self {
147        Self::default()
148    }
149
150    fn lock(&self) -> MutexGuard<'_, Shared> {
151        // A panicking test must not poison every later assertion.
152        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    /// Number of messages waiting in `queue`.
160    ///
161    /// Messages that were handed to a consumer but not yet settled are *not* counted,
162    /// and neither are delayed publishes whose delay has not elapsed.
163    pub fn pending(&self, queue: &str) -> usize {
164        self.lock().queues.get(queue).map_or(0, |q| q.pending.len())
165    }
166
167    /// Envelopes that were acked on `queue`, in ack order.
168    ///
169    /// A successful [`Delivery::retry`] acks the original envelope, so it shows up here
170    /// too (the rescheduled copy arrives separately with `attempt + 1`).
171    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    /// Number of deferred envelopes still sitting in `queue`'s hold.
180    ///
181    /// These are the ones [`Backend::defer`] (or [`Delivery::defer`]) accepted but
182    /// whose delay has not elapsed: they are not `pending` yet and no consumer can see
183    /// them. The count drops back to zero as each hold expires or, if the backend
184    /// was closed in the meantime, as each held envelope is dropped.
185    pub fn deferred(&self, queue: &str) -> usize {
186        self.lock().queues.get(queue).map_or(0, |q| q.held)
187    }
188
189    /// Envelopes that were dead-lettered on `queue`, with the reason given.
190    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    /// Names of every declared (or implicitly created) queue.
199    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    /// Config recorded for `queue` by the first `declare` call, if any.
206    pub fn queue_config(&self, queue: &str) -> Option<QueueConfig> {
207        self.lock().queues.get(queue).map(|q| q.config.clone())
208    }
209
210    /// Whether [`Backend::close`] has been called.
211    pub fn is_closed(&self) -> bool {
212        self.lock().closed
213    }
214
215    /// Number of live consumers registered on `queue`.
216    ///
217    /// Internal state, exposed to this crate's tests so they can synchronise on a
218    /// consumer having started (or gone) instead of sleeping.
219    #[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    /// Whether [`Backend::close`] has been called, on a handle-less state pointer.
228    fn is_state_closed(state: &Arc<Mutex<Shared>>) -> bool {
229        Self::lock_shared(state).closed
230    }
231
232    /// Put an envelope on its queue right away, by priority, and wake any idle
233    /// consumer.
234    ///
235    /// `hold`, when given, is the hold this envelope is leaving. Its decrement happens
236    /// under the *same* lock as the insertion, so there is no window in which
237    /// [`MemoryBackend::deferred`] and [`MemoryBackend::pending`] both count it.
238    ///
239    /// A closed backend accepts nothing, exactly like [`Backend::publish`], but the
240    /// hold is still released, so the counter drops either way.
241    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    /// Enqueue `envelope` after `delay` (immediately when the delay is zero).
255    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        // The deadline is taken now, not when the spawned task is first polled, so the
261        // delay is measured from the publish call even under `tokio::time::pause()`.
262        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    /// Hold `envelope` for `delay`, then enqueue it by priority.
271    ///
272    /// The stand-in for a broker's hold queue: while it waits, the envelope is
273    /// invisible to consumers and counted by [`MemoryBackend::deferred`]. A zero delay
274    /// enqueues straight away, so the counter never moves. If the backend is closed
275    /// before the hold expires the envelope is dropped, like a delayed publish.
276    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        // Same reasoning as `schedule`: the deadline is taken now, not when the
282        // spawned task is first polled.
283        let deadline = tokio::time::Instant::now() + delay;
284        Self::lock_shared(state).queue_mut(&envelope.queue).held += 1;
285        // The guard decrements again however the hold ends: it is handed to the
286        // enqueue on expiry (one lock, no double count), and otherwise decrements on
287        // drop, when the task is aborted or the runtime goes away.
288        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
300/// Keeps [`MemoryBackend::deferred`] honest: one live guard per envelope in hold.
301struct HoldGuard {
302    state: Arc<Mutex<Shared>>,
303    /// The queue to decrement, taken once the decrement has happened.
304    queue: Option<String>,
305}
306
307impl HoldGuard {
308    /// Decrement under a lock the caller already holds, and disarm the `Drop`.
309    ///
310    /// This is what lets the hold end in the same critical section that enqueues the
311    /// envelope, instead of one lock later.
312    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
328/// One held envelope has left `queue`'s hold.
329fn 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            // Idempotent: re-declaring never drops pending messages or consumers.
344            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                // Dropping `tx` ends the stream immediately.
387                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                // Wakes consumers waiting for a free prefetch slot.
410                semaphore.close();
411            }
412            // Wakes consumers waiting for a message.
413            queue.wake();
414        }
415        Ok(())
416    }
417}
418
419/// Drops a consumer's semaphore from its queue when the consumer task ends, however
420/// it ends, so `QueueState::consumers` never grows with dead consumers.
421struct 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
436/// Pulls messages for one consumer, respecting its prefetch window.
437async 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        // Blocks while `prefetch` deliveries are outstanding; errors once closed.
452        // A dropped stream ends the task right away, so its semaphore does not linger.
453        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                    // `changed()` resolves immediately if a message arrived since the
477                    // last check, so no wakeup can be lost here.
478                    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            // Consumer dropped the stream: put the message back for someone else.
496            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
506/// A message handed to one consumer. Holds a prefetch permit until it is settled.
507struct 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        // A retry is a publish, so a closed backend rejects it and the original stays
548        // unacked, just like `publish`.
549        if MemoryBackend::is_state_closed(&self.state) {
550            return Err(Error::ShutDown);
551        }
552        // Schedule first, ack second: never lose the message.
553        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        // A deferral is a publish too, so a closed backend rejects it and the original
560        // stays unacked.
561        if MemoryBackend::is_state_closed(&self.state) {
562            return Err(Error::ShutDown);
563        }
564        // Hold first, ack second: never lose the message. Dropping `self` afterwards
565        // releases the prefetch permit, exactly as `retry` does.
566        MemoryBackend::hold(&self.state, next, delay);
567        self.record_ack();
568        Ok(())
569    }
570}
571
572/// Adapts the consumer channel to the [`Stream`] the [`Backend`] trait returns.
573struct 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    /// Requeues everything that was produced for this consumer but never handed to
589    /// it, the way a broker returns unacked messages when a channel closes. Without
590    /// this, dropping a stream would silently swallow the messages in flight.
591    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        // Front, in reverse, so the original order survives.
608        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        // Let the consumer task park on the empty queue.
710        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        // Third delivery is withheld: two are outstanding.
739        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        // Sleeping (rather than `advance`) lets the paused clock actually fire the
779        // scheduled timer, because the runtime parks in between.
780        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        // The original is acked straight away, the copy is still in flight.
803        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    /// Publishes `name` with an explicit priority, bypassing `Envelope::new`'s zero.
832    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    /// The names still waiting on `queue`, in the order a consumer would see them.
840    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", // 10, in publish order
872                "middle",   // 5
873                "normal-1", "normal-2", "normal-3", // 0, in publish order
874            ]
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        // Invisible while it waits: not pending, but accounted for.
905        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    /// `held + pending` for `queue`, read under one lock so the pair is a consistent
951    /// snapshot rather than two observations with a gap in between.
952    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        // Real time, real threads: the holds expire while this loop watches. Each
976        // envelope must be either in hold or pending, never both, so the sum can only
977        // ever be `HELD`.
978        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        // The original is acked straight away; the copy is in hold.
1052        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        // A refused deferral never acks the original, and nothing is held.
1110        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        // Only one consumer exists at this point, so the message lands on `busy`.
1163        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        // Even a consumer with an outstanding delivery is released.
1170        assert!(busy.next().await.is_none());
1171        assert!(backend.is_closed());
1172        // Settling after close still works and is recorded.
1173        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        // A refused retry never acks the original, and nothing is rescheduled.
1209        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        // Each consumer task notices the closed channel and deregisters itself.
1244        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        // Wait until the consumer has taken both messages off the queue but nobody
1300        // has pulled them off the stream: that is the state the requeue has to cover.
1301        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        // Dropping the stream hands them straight back, in their original order.
1310        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}