use std::{
collections::{HashMap, VecDeque},
pin::Pin,
sync::{Arc, Mutex, MutexGuard},
task::{Context, Poll},
time::Duration,
};
use async_trait::async_trait;
use futures::Stream;
use tokio::sync::{OwnedSemaphorePermit, Semaphore, mpsc, watch};
use crate::{
backend::{Backend, Delivery, DeliveryStream},
envelope::Envelope,
error::{Error, Result},
queue::QueueConfig,
};
struct QueueState {
config: QueueConfig,
pending: VecDeque<Envelope>,
acked: Vec<Envelope>,
dead: Vec<(Envelope, String)>,
consumers: Vec<Arc<Semaphore>>,
held: usize,
signal: watch::Sender<u64>,
}
impl QueueState {
fn new(config: QueueConfig) -> Self {
let (signal, _) = watch::channel(0);
Self {
config,
pending: VecDeque::new(),
acked: Vec::new(),
dead: Vec::new(),
consumers: Vec::new(),
held: 0,
signal,
}
}
fn insert_by_priority(&mut self, envelope: Envelope) {
let position = self
.pending
.iter()
.rposition(|pending| pending.priority >= envelope.priority)
.map_or(0, |last| last + 1);
self.pending.insert(position, envelope);
}
fn wake(&self) {
self.signal.send_modify(|v| *v = v.wrapping_add(1));
}
}
#[derive(Default)]
struct Shared {
queues: HashMap<String, QueueState>,
closed: bool,
}
impl Shared {
fn queue_mut(&mut self, name: &str) -> &mut QueueState {
self.queues
.entry(name.to_owned())
.or_insert_with(|| QueueState::new(QueueConfig::new(name)))
}
}
#[derive(Clone, Default)]
pub struct MemoryBackend {
state: Arc<Mutex<Shared>>,
}
impl std::fmt::Debug for MemoryBackend {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let state = self.lock();
f.debug_struct("MemoryBackend")
.field("queues", &state.queues.keys().collect::<Vec<_>>())
.field("closed", &state.closed)
.finish()
}
}
impl MemoryBackend {
pub fn new() -> Self {
Self::default()
}
fn lock(&self) -> MutexGuard<'_, Shared> {
self.state.lock().unwrap_or_else(|e| e.into_inner())
}
fn lock_shared(state: &Arc<Mutex<Shared>>) -> MutexGuard<'_, Shared> {
state.lock().unwrap_or_else(|e| e.into_inner())
}
pub fn pending(&self, queue: &str) -> usize {
self.lock().queues.get(queue).map_or(0, |q| q.pending.len())
}
pub fn acked(&self, queue: &str) -> Vec<Envelope> {
self.lock()
.queues
.get(queue)
.map(|q| q.acked.clone())
.unwrap_or_default()
}
pub fn deferred(&self, queue: &str) -> usize {
self.lock().queues.get(queue).map_or(0, |q| q.held)
}
pub fn dead_letters(&self, queue: &str) -> Vec<(Envelope, String)> {
self.lock()
.queues
.get(queue)
.map(|q| q.dead.clone())
.unwrap_or_default()
}
pub fn queue_names(&self) -> Vec<String> {
let mut names: Vec<_> = self.lock().queues.keys().cloned().collect();
names.sort();
names
}
pub fn queue_config(&self, queue: &str) -> Option<QueueConfig> {
self.lock().queues.get(queue).map(|q| q.config.clone())
}
pub fn is_closed(&self) -> bool {
self.lock().closed
}
#[cfg(test)]
pub(crate) fn consumer_count(&self, queue: &str) -> usize {
self.lock()
.queues
.get(queue)
.map_or(0, |q| q.consumers.len())
}
fn is_state_closed(state: &Arc<Mutex<Shared>>) -> bool {
Self::lock_shared(state).closed
}
fn enqueue_now(state: &Arc<Mutex<Shared>>, envelope: Envelope, hold: Option<HoldGuard>) {
let mut shared = Self::lock_shared(state);
if let Some(hold) = hold {
hold.release(&mut shared);
}
if shared.closed {
return;
}
let queue = shared.queue_mut(&envelope.queue);
queue.insert_by_priority(envelope);
queue.wake();
}
fn schedule(state: &Arc<Mutex<Shared>>, envelope: Envelope, delay: Duration) {
if delay.is_zero() {
Self::enqueue_now(state, envelope, None);
return;
}
let deadline = tokio::time::Instant::now() + delay;
let state = state.clone();
tokio::spawn(async move {
tokio::time::sleep_until(deadline).await;
Self::enqueue_now(&state, envelope, None);
});
}
fn hold(state: &Arc<Mutex<Shared>>, envelope: Envelope, delay: Duration) {
if delay.is_zero() {
Self::enqueue_now(state, envelope, None);
return;
}
let deadline = tokio::time::Instant::now() + delay;
Self::lock_shared(state).queue_mut(&envelope.queue).held += 1;
let guard = HoldGuard {
state: state.clone(),
queue: Some(envelope.queue.clone()),
};
let state = state.clone();
tokio::spawn(async move {
tokio::time::sleep_until(deadline).await;
Self::enqueue_now(&state, envelope, Some(guard));
});
}
}
struct HoldGuard {
state: Arc<Mutex<Shared>>,
queue: Option<String>,
}
impl HoldGuard {
fn release(mut self, shared: &mut Shared) {
if let Some(queue) = self.queue.take() {
release_hold(shared, &queue);
}
}
}
impl Drop for HoldGuard {
fn drop(&mut self) {
if let Some(queue) = self.queue.take() {
let mut shared = MemoryBackend::lock_shared(&self.state);
release_hold(&mut shared, &queue);
}
}
}
fn release_hold(shared: &mut Shared, queue: &str) {
if let Some(queue) = shared.queues.get_mut(queue) {
queue.held = queue.held.saturating_sub(1);
}
}
#[async_trait]
impl Backend for MemoryBackend {
async fn declare(&self, queues: &[QueueConfig]) -> Result<()> {
let mut shared = self.lock();
if shared.closed {
return Err(Error::ShutDown);
}
for config in queues {
shared
.queues
.entry(config.name.clone())
.or_insert_with(|| QueueState::new(config.clone()));
}
Ok(())
}
async fn publish(&self, envelope: &Envelope, delay: Option<Duration>) -> Result<()> {
if self.is_closed() {
return Err(Error::ShutDown);
}
Self::schedule(&self.state, envelope.clone(), delay.unwrap_or_default());
Ok(())
}
async fn defer(&self, envelope: &Envelope, delay: Duration) -> Result<()> {
if self.is_closed() {
return Err(Error::ShutDown);
}
Self::hold(&self.state, envelope.clone(), delay);
Ok(())
}
async fn consume(&self, queue: &QueueConfig) -> Result<DeliveryStream> {
let (tx, rx) = mpsc::unbounded_channel();
let stream: DeliveryStream = Box::pin(DeliveryReceiver {
rx,
state: self.state.clone(),
queue: queue.name.clone(),
});
let permits = if queue.prefetch == 0 {
Semaphore::MAX_PERMITS
} else {
usize::from(queue.prefetch)
};
let semaphore = Arc::new(Semaphore::new(permits));
let signal = {
let mut shared = self.lock();
if shared.closed {
return Ok(stream);
}
let state = shared.queue_mut(&queue.name);
state.consumers.push(semaphore.clone());
state.signal.subscribe()
};
tokio::spawn(consumer_loop(
self.state.clone(),
queue.name.clone(),
semaphore,
tx,
signal,
));
Ok(stream)
}
async fn close(&self) -> Result<()> {
let mut shared = self.lock();
shared.closed = true;
for queue in shared.queues.values_mut() {
for semaphore in queue.consumers.drain(..) {
semaphore.close();
}
queue.wake();
}
Ok(())
}
}
struct ConsumerGuard {
state: Arc<Mutex<Shared>>,
queue: String,
semaphore: Arc<Semaphore>,
}
impl Drop for ConsumerGuard {
fn drop(&mut self) {
let mut shared = MemoryBackend::lock_shared(&self.state);
if let Some(queue) = shared.queues.get_mut(&self.queue) {
queue.consumers.retain(|s| !Arc::ptr_eq(s, &self.semaphore));
}
}
}
async fn consumer_loop(
state: Arc<Mutex<Shared>>,
queue: String,
semaphore: Arc<Semaphore>,
tx: mpsc::UnboundedSender<Result<Box<dyn Delivery>>>,
mut signal: watch::Receiver<u64>,
) {
let _guard = ConsumerGuard {
state: state.clone(),
queue: queue.clone(),
semaphore: semaphore.clone(),
};
loop {
let permit = tokio::select! {
biased;
() = tx.closed() => return,
permit = semaphore.clone().acquire_owned() => match permit {
Ok(permit) => permit,
Err(_) => return,
},
};
let envelope = loop {
let next = {
let mut shared = MemoryBackend::lock_shared(&state);
if shared.closed {
return;
}
shared
.queues
.get_mut(&queue)
.and_then(|q| q.pending.pop_front())
};
match next {
Some(envelope) => break envelope,
None => {
tokio::select! {
biased;
() = tx.closed() => return,
changed = signal.changed() => if changed.is_err() {
return;
},
}
}
}
};
let delivery = MemoryDelivery {
state: state.clone(),
envelope: envelope.clone(),
_permit: permit,
};
if tx.send(Ok(Box::new(delivery))).is_err() {
let mut shared = MemoryBackend::lock_shared(&state);
if let Some(q) = shared.queues.get_mut(&queue) {
q.pending.push_front(envelope);
q.wake();
}
return;
}
}
}
struct MemoryDelivery {
state: Arc<Mutex<Shared>>,
envelope: Envelope,
_permit: OwnedSemaphorePermit,
}
impl MemoryDelivery {
fn record_ack(&self) {
let mut shared = MemoryBackend::lock_shared(&self.state);
shared
.queue_mut(&self.envelope.queue)
.acked
.push(self.envelope.clone());
}
}
#[async_trait]
impl Delivery for MemoryDelivery {
fn envelope(&self) -> &Envelope {
&self.envelope
}
async fn ack(self: Box<Self>) -> Result<()> {
self.record_ack();
Ok(())
}
async fn dead_letter(self: Box<Self>, reason: &str) -> Result<()> {
let mut shared = MemoryBackend::lock_shared(&self.state);
if shared.closed {
return Err(Error::ShutDown);
}
shared
.queue_mut(&self.envelope.queue)
.dead
.push((self.envelope.clone(), reason.to_owned()));
Ok(())
}
async fn retry(self: Box<Self>, next: Envelope, delay: Duration) -> Result<()> {
if MemoryBackend::is_state_closed(&self.state) {
return Err(Error::ShutDown);
}
MemoryBackend::schedule(&self.state, next, delay);
self.record_ack();
Ok(())
}
async fn defer(self: Box<Self>, next: Envelope, delay: Duration) -> Result<()> {
if MemoryBackend::is_state_closed(&self.state) {
return Err(Error::ShutDown);
}
MemoryBackend::hold(&self.state, next, delay);
self.record_ack();
Ok(())
}
}
struct DeliveryReceiver {
rx: mpsc::UnboundedReceiver<Result<Box<dyn Delivery>>>,
state: Arc<Mutex<Shared>>,
queue: String,
}
impl Stream for DeliveryReceiver {
type Item = Result<Box<dyn Delivery>>;
fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
self.get_mut().rx.poll_recv(cx)
}
}
impl Drop for DeliveryReceiver {
fn drop(&mut self) {
self.rx.close();
let mut unread = Vec::new();
while let Ok(item) = self.rx.try_recv() {
if let Ok(delivery) = item {
unread.push(delivery.envelope().clone());
}
}
if unread.is_empty() {
return;
}
let mut shared = MemoryBackend::lock_shared(&self.state);
if shared.closed {
return;
}
let queue = shared.queue_mut(&self.queue);
for envelope in unread.into_iter().rev() {
queue.pending.push_front(envelope);
}
queue.wake();
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{
job::Job,
queue::QueueSet,
test_support::{Greet, Ping, TestQueues},
};
use futures::{StreamExt, future::poll_immediate};
fn alpha() -> QueueConfig {
TestQueues::Alpha.config()
}
async fn declared() -> MemoryBackend {
let backend = MemoryBackend::new();
backend.declare(&[alpha()]).await.unwrap();
backend
}
async fn publish(backend: &MemoryBackend, name: &str) -> Envelope {
let envelope = Envelope::new(&Greet::new(name)).unwrap();
backend.publish(&envelope, None).await.unwrap();
envelope
}
async fn next_delivery(stream: &mut DeliveryStream) -> Box<dyn Delivery> {
stream
.next()
.await
.expect("stream ended")
.expect("delivery error")
}
#[tokio::test]
async fn declare_is_idempotent() {
let backend = declared().await;
publish(&backend, "a").await;
backend.declare(&[alpha(), alpha()]).await.unwrap();
backend.declare(&[alpha()]).await.unwrap();
assert_eq!(backend.pending("test.alpha"), 1);
assert_eq!(backend.queue_names(), vec!["test.alpha".to_owned()]);
assert_eq!(backend.queue_config("test.alpha").unwrap().prefetch, 4);
}
#[tokio::test]
async fn declares_every_queue_of_a_queue_set() {
let backend = MemoryBackend::new();
let configs: Vec<_> = TestQueues::all().iter().map(|q| q.config()).collect();
backend.declare(&configs).await.unwrap();
assert_eq!(
backend.queue_names(),
vec![
"test.alpha".to_owned(),
"test.beta".to_owned(),
"test.gamma".to_owned()
]
);
}
#[tokio::test]
async fn publish_then_consume_is_fifo() {
let backend = declared().await;
for name in ["a", "b", "c"] {
publish(&backend, name).await;
}
assert_eq!(backend.pending("test.alpha"), 3);
let mut stream = backend.consume(&alpha()).await.unwrap();
for name in ["a", "b", "c"] {
let delivery = next_delivery(&mut stream).await;
assert_eq!(
delivery.envelope().decode::<Greet>().unwrap(),
Greet::new(name)
);
assert_eq!(delivery.envelope().attempt, 1);
delivery.ack().await.unwrap();
}
assert_eq!(backend.pending("test.alpha"), 0);
let acked: Vec<_> = backend
.acked("test.alpha")
.iter()
.map(|e| e.decode::<Greet>().unwrap().name)
.collect();
assert_eq!(acked, vec!["a", "b", "c"]);
}
#[tokio::test]
async fn consumer_started_before_publish_receives_messages() {
let backend = declared().await;
let mut stream = backend.consume(&alpha()).await.unwrap();
tokio::task::yield_now().await;
publish(&backend, "late").await;
let delivery = next_delivery(&mut stream).await;
assert_eq!(
delivery.envelope().decode::<Greet>().unwrap(),
Greet::new("late")
);
delivery.ack().await.unwrap();
}
#[tokio::test]
async fn prefetch_caps_outstanding_deliveries() {
let backend = declared().await;
let config = QueueConfig::new("test.alpha").prefetch(2);
for name in ["a", "b", "c"] {
publish(&backend, name).await;
}
let mut stream = backend.consume(&config).await.unwrap();
let first = next_delivery(&mut stream).await;
let second = next_delivery(&mut stream).await;
assert_eq!(first.envelope().decode::<Greet>().unwrap(), Greet::new("a"));
assert_eq!(
second.envelope().decode::<Greet>().unwrap(),
Greet::new("b")
);
assert!(poll_immediate(stream.next()).await.is_none());
assert_eq!(backend.pending("test.alpha"), 1);
first.ack().await.unwrap();
let third = next_delivery(&mut stream).await;
assert_eq!(third.envelope().decode::<Greet>().unwrap(), Greet::new("c"));
assert_eq!(backend.pending("test.alpha"), 0);
second.ack().await.unwrap();
third.ack().await.unwrap();
assert_eq!(backend.acked("test.alpha").len(), 3);
}
#[tokio::test]
async fn zero_prefetch_means_unlimited() {
let backend = declared().await;
for name in ["a", "b", "c"] {
publish(&backend, name).await;
}
let mut stream = backend
.consume(&QueueConfig::new("test.alpha").prefetch(0))
.await
.unwrap();
let mut held = Vec::new();
for _ in 0..3 {
held.push(next_delivery(&mut stream).await);
}
assert_eq!(held.len(), 3);
}
#[tokio::test(start_paused = true)]
async fn delayed_publish_is_invisible_until_the_delay_elapses() {
let backend = declared().await;
let envelope = Envelope::new(&Greet::new("later")).unwrap();
backend
.publish(&envelope, Some(Duration::from_secs(30)))
.await
.unwrap();
tokio::time::sleep(Duration::from_secs(29)).await;
assert_eq!(backend.pending("test.alpha"), 0);
tokio::time::sleep(Duration::from_secs(2)).await;
assert_eq!(backend.pending("test.alpha"), 1);
let mut stream = backend.consume(&alpha()).await.unwrap();
let delivery = next_delivery(&mut stream).await;
assert_eq!(delivery.envelope().job_id, envelope.job_id);
delivery.ack().await.unwrap();
}
#[tokio::test(start_paused = true)]
async fn retry_redelivers_with_a_higher_attempt_after_the_delay() {
let backend = declared().await;
let original = publish(&backend, "flaky").await;
let mut stream = backend.consume(&alpha()).await.unwrap();
let delivery = next_delivery(&mut stream).await;
let next = delivery.envelope().next_attempt();
delivery.retry(next, Duration::from_secs(10)).await.unwrap();
assert_eq!(backend.acked("test.alpha").len(), 1);
assert_eq!(backend.pending("test.alpha"), 0);
tokio::time::advance(Duration::from_secs(5)).await;
assert!(poll_immediate(stream.next()).await.is_none());
tokio::time::advance(Duration::from_secs(6)).await;
let redelivered = next_delivery(&mut stream).await;
assert_eq!(redelivered.envelope().attempt, 2);
assert_eq!(redelivered.envelope().job_id, original.job_id);
redelivered.ack().await.unwrap();
assert_eq!(backend.acked("test.alpha").len(), 2);
}
#[tokio::test]
async fn retry_without_delay_is_immediate() {
let backend = declared().await;
publish(&backend, "now").await;
let mut stream = backend.consume(&alpha()).await.unwrap();
let delivery = next_delivery(&mut stream).await;
let next = delivery.envelope().next_attempt();
delivery.retry(next, Duration::ZERO).await.unwrap();
let redelivered = next_delivery(&mut stream).await;
assert_eq!(redelivered.envelope().attempt, 2);
redelivered.ack().await.unwrap();
}
async fn publish_with_priority(backend: &MemoryBackend, name: &str, priority: u8) -> Envelope {
let mut envelope = Envelope::new(&Greet::new(name)).unwrap();
envelope.priority = priority;
backend.publish(&envelope, None).await.unwrap();
envelope
}
fn pending_names(backend: &MemoryBackend, queue: &str) -> Vec<String> {
backend
.lock()
.queues
.get(queue)
.map(|q| {
q.pending
.iter()
.map(|e| e.decode::<Greet>().unwrap().name)
.collect()
})
.unwrap_or_default()
}
#[tokio::test]
async fn publish_orders_by_priority_and_keeps_each_level_fifo() {
let backend = declared().await;
for (name, priority) in [
("normal-1", 0),
("urgent-1", 10),
("normal-2", 0),
("urgent-2", 10),
("middle", 5),
("normal-3", 0),
] {
publish_with_priority(&backend, name, priority).await;
}
assert_eq!(
pending_names(&backend, "test.alpha"),
vec![
"urgent-1", "urgent-2", "middle", "normal-1", "normal-2", "normal-3", ]
);
}
#[tokio::test]
async fn a_higher_priority_publish_overtakes_the_whole_backlog() {
let backend = declared().await;
for name in ["a", "b", "c"] {
publish(&backend, name).await;
}
publish_with_priority(&backend, "jumper", 1).await;
let mut stream = backend.consume(&alpha()).await.unwrap();
let first = next_delivery(&mut stream).await;
assert_eq!(
first.envelope().decode::<Greet>().unwrap(),
Greet::new("jumper")
);
first.ack().await.unwrap();
}
#[tokio::test(start_paused = true)]
async fn defer_holds_the_message_and_then_lands_it_ahead_of_the_backlog() {
let backend = declared().await;
for name in ["backlog-1", "backlog-2"] {
publish(&backend, name).await;
}
let held = Envelope::new(&Greet::new("held")).unwrap().deferred(10);
backend.defer(&held, Duration::from_secs(30)).await.unwrap();
assert_eq!(backend.pending("test.alpha"), 2);
assert_eq!(backend.deferred("test.alpha"), 1);
tokio::time::sleep(Duration::from_secs(29)).await;
assert_eq!(backend.pending("test.alpha"), 2);
assert_eq!(backend.deferred("test.alpha"), 1);
tokio::time::sleep(Duration::from_secs(2)).await;
assert_eq!(backend.deferred("test.alpha"), 0, "the hold has expired");
assert_eq!(
pending_names(&backend, "test.alpha"),
vec!["held", "backlog-1", "backlog-2"],
"a deferred job comes back in front"
);
let mut stream = backend.consume(&alpha()).await.unwrap();
let first = next_delivery(&mut stream).await;
assert_eq!(first.envelope().job_id, held.job_id);
assert_eq!(first.envelope().deferrals, 1);
assert_eq!(first.envelope().priority, 10);
assert_eq!(first.envelope().attempt, 1, "a deferral is not an attempt");
first.ack().await.unwrap();
}
#[tokio::test(start_paused = true)]
async fn deferred_counts_every_message_in_hold() {
let backend = declared().await;
assert_eq!(backend.deferred("test.alpha"), 0);
assert_eq!(backend.deferred("never.declared"), 0);
for _ in 0..3 {
backend
.defer(
&Envelope::new(&Greet::new("held")).unwrap(),
Duration::from_secs(5),
)
.await
.unwrap();
}
assert_eq!(backend.deferred("test.alpha"), 3);
tokio::time::sleep(Duration::from_secs(6)).await;
assert_eq!(backend.deferred("test.alpha"), 0);
assert_eq!(backend.pending("test.alpha"), 3);
}
fn held_plus_pending(backend: &MemoryBackend, queue: &str) -> usize {
backend
.lock()
.queues
.get(queue)
.map_or(0, |q| q.held + q.pending.len())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_hold_is_never_counted_twice_while_it_expires() {
const HELD: usize = 64;
let backend = declared().await;
for i in 0..HELD {
backend
.defer(
&Envelope::new(&Greet::new(&i.to_string())).unwrap(),
Duration::from_millis(5),
)
.await
.unwrap();
}
for _ in 0..1_000_000 {
let total = held_plus_pending(&backend, "test.alpha");
assert!(
total <= HELD,
"an envelope was counted as held and pending at once ({total} > {HELD})"
);
if backend.pending("test.alpha") == HELD {
break;
}
tokio::task::yield_now().await;
}
assert_eq!(backend.pending("test.alpha"), HELD);
assert_eq!(backend.deferred("test.alpha"), 0);
}
#[tokio::test]
async fn a_zero_delay_deferral_never_enters_the_hold() {
let backend = declared().await;
backend
.defer(&Envelope::new(&Greet::new("now")).unwrap(), Duration::ZERO)
.await
.unwrap();
assert_eq!(backend.deferred("test.alpha"), 0);
assert_eq!(backend.pending("test.alpha"), 1);
}
#[tokio::test]
async fn defer_after_close_reports_the_shutdown() {
let backend = declared().await;
backend.close().await.unwrap();
assert!(matches!(
backend
.defer(
&Envelope::new(&Greet::new("x")).unwrap(),
Duration::from_secs(1)
)
.await,
Err(Error::ShutDown)
));
assert_eq!(backend.deferred("test.alpha"), 0);
assert_eq!(backend.pending("test.alpha"), 0);
}
#[tokio::test(start_paused = true)]
async fn a_deferral_that_lands_after_close_is_dropped() {
let backend = declared().await;
backend
.defer(
&Envelope::new(&Greet::new("held")).unwrap(),
Duration::from_secs(30),
)
.await
.unwrap();
assert_eq!(backend.deferred("test.alpha"), 1);
backend.close().await.unwrap();
tokio::time::sleep(Duration::from_secs(31)).await;
assert_eq!(backend.pending("test.alpha"), 0);
assert_eq!(backend.deferred("test.alpha"), 0);
}
#[tokio::test(start_paused = true)]
async fn delivery_defer_acks_the_original_and_reschedules_it() {
let backend = declared().await;
let original = publish(&backend, "rate-limited").await;
let mut stream = backend.consume(&alpha()).await.unwrap();
let delivery = next_delivery(&mut stream).await;
let next = delivery.envelope().deferred(10);
delivery.defer(next, Duration::from_secs(30)).await.unwrap();
assert_eq!(backend.acked("test.alpha").len(), 1);
assert_eq!(backend.acked("test.alpha")[0].job_id, original.job_id);
assert_eq!(backend.acked("test.alpha")[0].deferrals, 0);
assert_eq!(backend.pending("test.alpha"), 0);
assert_eq!(backend.deferred("test.alpha"), 1);
tokio::time::advance(Duration::from_secs(10)).await;
assert!(poll_immediate(stream.next()).await.is_none());
tokio::time::advance(Duration::from_secs(21)).await;
let redelivered = next_delivery(&mut stream).await;
assert_eq!(redelivered.envelope().job_id, original.job_id);
assert_eq!(redelivered.envelope().attempt, 1);
assert_eq!(redelivered.envelope().deferrals, 1);
assert_eq!(redelivered.envelope().priority, 10);
assert_eq!(backend.deferred("test.alpha"), 0);
redelivered.ack().await.unwrap();
assert_eq!(backend.acked("test.alpha").len(), 2);
}
#[tokio::test]
async fn delivery_defer_frees_a_prefetch_slot() {
let backend = declared().await;
for name in ["a", "b"] {
publish(&backend, name).await;
}
let mut stream = backend
.consume(&QueueConfig::new("test.alpha").prefetch(1))
.await
.unwrap();
let first = next_delivery(&mut stream).await;
assert!(poll_immediate(stream.next()).await.is_none());
let next = first.envelope().deferred(10);
first.defer(next, Duration::from_secs(600)).await.unwrap();
let second = next_delivery(&mut stream).await;
assert_eq!(
second.envelope().decode::<Greet>().unwrap(),
Greet::new("b")
);
second.ack().await.unwrap();
}
#[tokio::test]
async fn delivery_defer_after_close_reports_the_shutdown() {
let backend = declared().await;
publish(&backend, "a").await;
let mut stream = backend.consume(&alpha()).await.unwrap();
let delivery = next_delivery(&mut stream).await;
backend.close().await.unwrap();
let next = delivery.envelope().deferred(10);
assert!(matches!(
delivery.defer(next, Duration::from_secs(1)).await,
Err(Error::ShutDown)
));
assert!(backend.acked("test.alpha").is_empty());
assert_eq!(backend.deferred("test.alpha"), 0);
}
#[tokio::test]
async fn dead_letter_records_the_reason() {
let backend = declared().await;
let original = publish(&backend, "doomed").await;
let mut stream = backend.consume(&alpha()).await.unwrap();
let delivery = next_delivery(&mut stream).await;
delivery
.dead_letter("max attempts (3) exhausted")
.await
.unwrap();
let dead = backend.dead_letters("test.alpha");
assert_eq!(dead.len(), 1);
assert_eq!(dead[0].0.job_id, original.job_id);
assert_eq!(dead[0].1, "max attempts (3) exhausted");
assert!(backend.acked("test.alpha").is_empty());
assert_eq!(backend.pending("test.alpha"), 0);
}
#[tokio::test]
async fn dead_letter_frees_a_prefetch_slot() {
let backend = declared().await;
for name in ["a", "b"] {
publish(&backend, name).await;
}
let mut stream = backend
.consume(&QueueConfig::new("test.alpha").prefetch(1))
.await
.unwrap();
let first = next_delivery(&mut stream).await;
assert!(poll_immediate(stream.next()).await.is_none());
first.dead_letter("nope").await.unwrap();
let second = next_delivery(&mut stream).await;
assert_eq!(
second.envelope().decode::<Greet>().unwrap(),
Greet::new("b")
);
second.ack().await.unwrap();
}
#[tokio::test]
async fn close_ends_all_streams() {
let backend = declared().await;
let mut busy = backend.consume(&alpha()).await.unwrap();
publish(&backend, "held").await;
let held = next_delivery(&mut busy).await;
let mut idle = backend.consume(&alpha()).await.unwrap();
backend.close().await.unwrap();
assert!(idle.next().await.is_none());
assert!(busy.next().await.is_none());
assert!(backend.is_closed());
held.ack().await.unwrap();
assert_eq!(backend.acked("test.alpha").len(), 1);
assert!(matches!(
backend
.publish(&Envelope::new(&Greet::new("x")).unwrap(), None)
.await,
Err(Error::ShutDown)
));
assert!(matches!(
backend.declare(&[alpha()]).await,
Err(Error::ShutDown)
));
}
#[tokio::test]
async fn retry_and_dead_letter_after_close_report_the_shutdown() {
let backend = declared().await;
publish(&backend, "a").await;
publish(&backend, "b").await;
let mut stream = backend.consume(&alpha()).await.unwrap();
let first = next_delivery(&mut stream).await;
let second = next_delivery(&mut stream).await;
backend.close().await.unwrap();
let next = first.envelope().next_attempt();
assert!(matches!(
first.retry(next, Duration::ZERO).await,
Err(Error::ShutDown)
));
assert!(matches!(
second.dead_letter("too late").await,
Err(Error::ShutDown)
));
assert!(backend.acked("test.alpha").is_empty());
assert!(backend.dead_letters("test.alpha").is_empty());
assert_eq!(backend.pending("test.alpha"), 0);
}
#[tokio::test]
async fn a_delayed_publish_that_lands_after_close_is_dropped() {
let backend = declared().await;
backend
.publish(
&Envelope::new(&Greet::new("later")).unwrap(),
Some(Duration::from_millis(1)),
)
.await
.unwrap();
backend.close().await.unwrap();
tokio::time::sleep(Duration::from_millis(20)).await;
assert_eq!(backend.pending("test.alpha"), 0);
}
#[tokio::test]
async fn consumers_are_forgotten_when_their_streams_are_dropped() {
let backend = declared().await;
let streams: Vec<_> = {
let mut v = Vec::new();
for _ in 0..3 {
v.push(backend.consume(&alpha()).await.unwrap());
}
v
};
assert_eq!(backend.consumer_count("test.alpha"), 3);
drop(streams);
for _ in 0..1_000 {
if backend.consumer_count("test.alpha") == 0 {
break;
}
tokio::task::yield_now().await;
}
assert_eq!(backend.consumer_count("test.alpha"), 0);
}
#[tokio::test]
async fn consuming_after_close_yields_an_ended_stream() {
let backend = declared().await;
backend.close().await.unwrap();
let mut stream = backend.consume(&alpha()).await.unwrap();
assert!(stream.next().await.is_none());
}
#[tokio::test]
async fn multiple_consumers_share_one_queue() {
let backend = declared().await;
let config = QueueConfig::new("test.alpha").prefetch(1);
let mut left = backend.consume(&config).await.unwrap();
let mut right = backend.consume(&config).await.unwrap();
tokio::task::yield_now().await;
for name in ["a", "b"] {
publish(&backend, name).await;
}
let one = next_delivery(&mut left).await;
let two = next_delivery(&mut right).await;
let names = [
one.envelope().decode::<Greet>().unwrap().name,
two.envelope().decode::<Greet>().unwrap().name,
];
assert!(
names.contains(&"a".to_owned()) && names.contains(&"b".to_owned()),
"{names:?}"
);
one.ack().await.unwrap();
two.ack().await.unwrap();
assert_eq!(backend.acked("test.alpha").len(), 2);
}
#[tokio::test]
async fn dropping_a_stream_returns_undelivered_messages_to_the_queue() {
let backend = declared().await;
let stream = backend
.consume(&QueueConfig::new("test.alpha").prefetch(2))
.await
.unwrap();
publish(&backend, "orphaned").await;
publish(&backend, "also-orphaned").await;
for _ in 0..1_000 {
if backend.pending("test.alpha") == 0 {
break;
}
tokio::task::yield_now().await;
}
assert_eq!(backend.pending("test.alpha"), 0);
drop(stream);
assert_eq!(backend.pending("test.alpha"), 2);
let mut stream = backend.consume(&alpha()).await.unwrap();
for name in ["orphaned", "also-orphaned"] {
let delivery = next_delivery(&mut stream).await;
assert_eq!(
delivery.envelope().decode::<Greet>().unwrap(),
Greet::new(name)
);
delivery.ack().await.unwrap();
}
}
#[tokio::test]
async fn queues_are_independent() {
let backend = MemoryBackend::new();
let configs: Vec<_> = TestQueues::all().iter().map(|q| q.config()).collect();
backend.declare(&configs).await.unwrap();
backend
.publish(&Envelope::new(&Greet::new("a")).unwrap(), None)
.await
.unwrap();
backend
.publish(&Envelope::new(&Ping { seq: 1 }).unwrap(), None)
.await
.unwrap();
assert_eq!(backend.pending("test.alpha"), 1);
assert_eq!(backend.pending("test.beta"), 1);
let mut stream = backend.consume(&TestQueues::Beta.config()).await.unwrap();
let delivery = next_delivery(&mut stream).await;
assert_eq!(delivery.envelope().job_type, Ping::NAME);
delivery.ack().await.unwrap();
assert_eq!(backend.pending("test.alpha"), 1);
}
#[tokio::test]
async fn clones_share_state() {
let backend = declared().await;
let clone = backend.clone();
publish(&clone, "shared").await;
assert_eq!(backend.pending("test.alpha"), 1);
assert!(format!("{backend:?}").contains("test.alpha"));
}
#[tokio::test]
async fn delivery_is_usable_behind_a_trait_object() {
let backend = declared().await;
publish(&backend, "boxed").await;
let mut stream = backend.consume(&alpha()).await.unwrap();
let delivery: Box<dyn Delivery> = next_delivery(&mut stream).await;
let envelope = delivery.envelope().clone();
delivery.ack().await.unwrap();
assert_eq!(backend.acked("test.alpha"), vec![envelope]);
}
}