use std::time::{Duration, Instant};
use futures::StreamExt;
use lapin::{
Channel, Connection, ConnectionProperties,
options::{
BasicGetOptions, BasicPublishOptions, ConfirmSelectOptions, QueueDeclareOptions,
QueueDeleteOptions,
},
types::{AMQPValue, FieldTable},
};
use queuey_core::{Backend, Delivery, Envelope, Error, QueueConfig};
use queuey_rabbitmq::{RabbitMqBackend, RabbitMqOptions};
use uuid::Uuid;
fn unique_queue(label: &str) -> String {
format!("aq-test.{label}.{}", Uuid::new_v4().simple())
}
fn family(queue: &str) -> [String; 2] {
[queue.to_owned(), format!("{queue}.dead")]
}
fn hold_queue(queue: &str, delay: Duration) -> String {
format!("{queue}.deferred.{}", delay.as_millis())
}
fn envelope(queue: &str, attempt: u32) -> Envelope {
Envelope {
job_id: Uuid::new_v4(),
job_type: "queuey_rabbitmq::tests::Probe".to_owned(),
queue: queue.to_owned(),
attempt,
enqueued_at_ms: 1_700_000_000_000,
deferrals: 0,
priority: 0,
payload: serde_json::json!({ "probe": true, "n": attempt }),
}
}
async fn control(url: &str) -> Connection {
Connection::connect(url, ConnectionProperties::default())
.await
.expect("control connection")
}
async fn ready_count(control: &Connection, queue: &str) -> u32 {
let channel = control.create_channel().await.expect("channel");
let declared = channel
.queue_declare(
queue.into(),
QueueDeclareOptions {
passive: true,
..Default::default()
},
FieldTable::default(),
)
.await
.unwrap_or_else(|error| panic!("passive declare of `{queue}` failed: {error}"));
let count = declared.message_count();
let _ = channel.close(200, "OK".into()).await;
count
}
async fn queue_exists(control: &Connection, queue: &str) -> bool {
let channel = control.create_channel().await.expect("channel");
let exists = channel
.queue_declare(
queue.into(),
QueueDeclareOptions {
passive: true,
..Default::default()
},
FieldTable::default(),
)
.await
.is_ok();
let _ = channel.close(200, "OK".into()).await;
exists
}
async fn await_count(control: &Connection, queue: &str, expected: u32) {
let deadline = std::time::Instant::now() + Duration::from_secs(10);
let mut last = u32::MAX;
while std::time::Instant::now() < deadline {
last = ready_count(control, queue).await;
if last == expected {
return;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
panic!("`{queue}` held {last} ready messages, expected {expected}");
}
async fn publish_raw(channel: &Channel, queue: &str, body: &[u8]) {
channel
.confirm_select(ConfirmSelectOptions::default())
.await
.expect("confirm_select");
let confirmation = channel
.basic_publish(
"".into(),
queue.into(),
BasicPublishOptions {
mandatory: true,
immediate: false,
},
body,
lapin::BasicProperties::default(),
)
.await
.expect("publish")
.await
.expect("confirm");
assert!(
confirmation.is_ack() && confirmation.take_message().is_none(),
"raw publish to `{queue}` was not routed anywhere"
);
}
async fn cleanup(control: &Connection, queue: &str) {
delete_queues(control, &family(queue)).await;
}
async fn delete_queues(control: &Connection, names: &[String]) {
let channel = control.create_channel().await.expect("channel");
for name in names {
let _ = channel
.queue_delete(name.as_str().into(), QueueDeleteOptions::default())
.await;
}
let _ = channel.close(200, "OK".into()).await;
}
async fn redeclare_raw(
control: &Connection,
queue: &str,
args: FieldTable,
) -> Result<(), lapin::Error> {
let channel = control.create_channel().await.expect("channel");
let result = channel
.queue_declare(
queue.into(),
QueueDeclareOptions {
durable: true,
..Default::default()
},
args,
)
.await
.map(|_| ());
let _ = channel.close(200, "OK".into()).await;
result
}
fn header_string(headers: &FieldTable, key: &str) -> Option<String> {
match headers.inner().get(key)? {
AMQPValue::LongString(value) => Some(value.to_string()),
AMQPValue::ShortString(value) => Some(value.to_string()),
_ => None,
}
}
fn header_u32(headers: &FieldTable, key: &str) -> Option<u32> {
match headers.inner().get(key)? {
AMQPValue::ShortShortInt(value) => u32::try_from(*value).ok(),
AMQPValue::ShortShortUInt(value) => Some(u32::from(*value)),
AMQPValue::ShortInt(value) => u32::try_from(*value).ok(),
AMQPValue::ShortUInt(value) => Some(u32::from(*value)),
AMQPValue::LongInt(value) => u32::try_from(*value).ok(),
AMQPValue::LongUInt(value) => Some(*value),
AMQPValue::LongLongInt(value) => u32::try_from(*value).ok(),
_ => None,
}
}
async fn next_delivery(
stream: &mut queuey_core::DeliveryStream,
within: Duration,
) -> Box<dyn Delivery> {
match tokio::time::timeout(within, stream.next()).await {
Ok(Some(Ok(delivery))) => delivery,
Ok(Some(Err(error))) => panic!("consumer stream yielded an error: {error}"),
Ok(None) => panic!("consumer stream ended unexpectedly"),
Err(_) => panic!("no delivery within {within:?}"),
}
}
async fn expect_idle(stream: &mut queuey_core::DeliveryStream, within: Duration) {
match tokio::time::timeout(within, stream.next()).await {
Err(_) => {}
Ok(Some(Ok(delivery))) => panic!(
"unexpected delivery of job {} (attempt {})",
delivery.envelope().job_id,
delivery.envelope().attempt
),
Ok(Some(Err(error))) => panic!("consumer stream yielded an error: {error}"),
Ok(None) => panic!("consumer stream ended unexpectedly"),
}
}
#[tokio::test]
async fn declare_is_idempotent() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let queue = unique_queue("declare");
let config = QueueConfig::new(queue.clone());
let backend = RabbitMqBackend::connect(&url).await.expect("connect");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&config))
.await
.expect("first declare");
backend
.declare(std::slice::from_ref(&config))
.await
.expect("second declare with identical arguments");
for name in family(&queue) {
assert_eq!(ready_count(&control, &name).await, 0, "queue {name}");
}
let other = RabbitMqBackend::connect(&url).await.expect("connect");
other
.declare(std::slice::from_ref(&config))
.await
.expect("third declare");
other.close().await.expect("close");
cleanup(&control, &queue).await;
backend.close().await.expect("close");
}
#[tokio::test]
async fn publish_then_consume_and_ack() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let queue = unique_queue("roundtrip");
let config = QueueConfig::new(queue.clone());
let backend = RabbitMqBackend::connect(&url).await.expect("connect");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare");
let sent = envelope(&queue, 1);
backend.publish(&sent, None).await.expect("publish");
await_count(&control, &queue, 1).await;
let mut stream = backend.consume(&config).await.expect("consume");
let delivery = next_delivery(&mut stream, Duration::from_secs(5)).await;
assert_eq!(delivery.envelope(), &sent, "envelope round-tripped intact");
delivery.ack().await.expect("ack");
drop(stream);
backend.close().await.expect("close");
await_count(&control, &queue, 0).await;
cleanup(&control, &queue).await;
}
#[tokio::test]
async fn prefetch_limits_unacked_deliveries() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let queue = unique_queue("prefetch");
let config = QueueConfig::new(queue.clone()).prefetch(2);
let backend = RabbitMqBackend::connect(&url).await.expect("connect");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare");
for attempt in 1..=5 {
backend
.publish(&envelope(&queue, attempt), None)
.await
.expect("publish");
}
await_count(&control, &queue, 5).await;
let mut stream = backend.consume(&config).await.expect("consume");
let first = next_delivery(&mut stream, Duration::from_secs(5)).await;
let second = next_delivery(&mut stream, Duration::from_secs(5)).await;
expect_idle(&mut stream, Duration::from_millis(750)).await;
first.ack().await.expect("ack");
let third = next_delivery(&mut stream, Duration::from_secs(5)).await;
expect_idle(&mut stream, Duration::from_millis(750)).await;
second.ack().await.expect("ack");
third.ack().await.expect("ack");
drop(stream);
backend.close().await.expect("close");
cleanup(&control, &queue).await;
}
#[tokio::test]
async fn delayed_publish_waits_for_the_delay() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let queue = unique_queue("delayed");
let config = QueueConfig::new(queue.clone());
let backend = RabbitMqBackend::connect(&url).await.expect("connect");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare");
let mut stream = backend.consume(&config).await.expect("consume");
let sent = envelope(&queue, 1);
let delay = Duration::from_secs(2);
backend
.publish(&sent, Some(delay))
.await
.expect("delayed publish");
let hold = hold_queue(&queue, delay);
await_count(&control, &hold, 1).await;
expect_idle(&mut stream, Duration::from_millis(1_200)).await;
let delivery = next_delivery(&mut stream, Duration::from_secs(8)).await;
assert_eq!(delivery.envelope(), &sent);
assert_eq!(
*delivery.envelope(),
sent,
"a delayed publish comes back byte for byte"
);
delivery.ack().await.expect("ack");
drop(stream);
backend.close().await.expect("close");
cleanup(&control, &queue).await;
delete_queues(&control, &[hold]).await;
}
#[tokio::test]
async fn retry_redelivers_with_the_next_attempt() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let queue = unique_queue("retry");
let config = QueueConfig::new(queue.clone());
let backend = RabbitMqBackend::connect(&url).await.expect("connect");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare");
let sent = envelope(&queue, 1);
backend.publish(&sent, None).await.expect("publish");
let mut stream = backend.consume(&config).await.expect("consume");
let first = next_delivery(&mut stream, Duration::from_secs(5)).await;
assert_eq!(first.envelope().attempt, 1);
let next = first.envelope().next_attempt();
let started = Instant::now();
first
.retry(next.clone(), Duration::from_millis(500))
.await
.expect("retry");
let hold = hold_queue(&queue, Duration::from_secs(1));
assert_eq!(ready_count(&control, &hold).await, 1, "held in `{hold}`");
expect_idle(&mut stream, Duration::from_millis(250)).await;
let second = next_delivery(&mut stream, Duration::from_secs(8)).await;
let elapsed = started.elapsed();
assert!(
elapsed >= Duration::from_millis(900),
"retry came back after {elapsed:?}; rounding up to the granularity means never early"
);
assert_eq!(second.envelope().job_id, sent.job_id, "same job id");
assert_eq!(second.envelope().attempt, 2, "attempt was incremented");
assert_eq!(second.envelope(), &next);
second.ack().await.expect("ack");
drop(stream);
backend.close().await.expect("close");
await_count(&control, &queue, 0).await;
await_count(&control, &hold, 0).await;
cleanup(&control, &queue).await;
delete_queues(&control, &[hold]).await;
}
#[tokio::test]
async fn a_short_retry_is_not_stuck_behind_a_long_one() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let queue = unique_queue("retry-hol");
let config = QueueConfig::new(queue.clone());
let backend = RabbitMqBackend::connect(&url).await.expect("connect");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare");
let slow = envelope(&queue, 1);
let fast = envelope(&queue, 1);
backend.publish(&slow, None).await.expect("publish");
backend.publish(&fast, None).await.expect("publish");
let mut stream = backend.consume(&config).await.expect("consume");
let first = next_delivery(&mut stream, Duration::from_secs(5)).await;
let second = next_delivery(&mut stream, Duration::from_secs(5)).await;
assert_eq!(first.envelope(), &slow, "FIFO on the work queue");
assert_eq!(second.envelope(), &fast);
let slow_next = slow.next_attempt();
let fast_next = fast.next_attempt();
let long_delay = Duration::from_secs(6);
let short_delay = Duration::from_secs(1);
let started = Instant::now();
first
.retry(slow_next.clone(), long_delay)
.await
.expect("retry (long)");
second
.retry(fast_next.clone(), short_delay)
.await
.expect("retry (short)");
let long_hold = hold_queue(&queue, long_delay);
let short_hold = hold_queue(&queue, short_delay);
assert_eq!(ready_count(&control, &long_hold).await, 1);
assert_eq!(ready_count(&control, &short_hold).await, 1);
let delivery = next_delivery(&mut stream, Duration::from_secs(4)).await;
let elapsed = started.elapsed();
assert_eq!(
delivery.envelope(),
&fast_next,
"the short retry returns first"
);
assert!(
elapsed < Duration::from_secs(4),
"the short retry took {elapsed:?}; it was stuck behind the long one"
);
assert!(
elapsed >= Duration::from_millis(900),
"the short retry came back after only {elapsed:?}"
);
delivery.ack().await.expect("ack");
assert_eq!(ready_count(&control, &long_hold).await, 1);
let delivery = next_delivery(&mut stream, Duration::from_secs(10)).await;
assert_eq!(delivery.envelope(), &slow_next);
assert!(
started.elapsed() >= Duration::from_millis(5_900),
"the long retry came back early"
);
delivery.ack().await.expect("ack");
drop(stream);
backend.close().await.expect("close");
await_count(&control, &queue, 0).await;
cleanup(&control, &queue).await;
delete_queues(&control, &[long_hold, short_hold]).await;
}
#[tokio::test]
async fn retries_and_deferrals_share_a_hold_queue_but_return_at_their_own_priority() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let queue = unique_queue("retry-shared-hold");
let config = QueueConfig::new(queue.clone());
let backend = RabbitMqBackend::connect(&url).await.expect("connect");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare");
let retried = envelope(&queue, 1).next_attempt();
let deferred = envelope(&queue, 1).deferred(10);
assert_eq!(retried.priority, 0);
let delay = Duration::from_secs(1);
backend
.publish(&retried, Some(delay))
.await
.expect("delayed publish");
backend.defer(&deferred, delay).await.expect("defer");
let hold = hold_queue(&queue, delay);
assert_eq!(
ready_count(&control, &hold).await,
2,
"both wait in `{hold}`"
);
await_count(&control, &queue, 2).await;
let mut stream = backend.consume(&config).await.expect("consume");
let first = next_delivery(&mut stream, Duration::from_secs(5)).await;
assert_eq!(
first.envelope(),
&deferred,
"the deferral overtakes the retry"
);
first.ack().await.expect("ack");
let second = next_delivery(&mut stream, Duration::from_secs(5)).await;
assert_eq!(second.envelope(), &retried);
second.ack().await.expect("ack");
drop(stream);
backend.close().await.expect("close");
cleanup(&control, &queue).await;
delete_queues(&control, &[hold]).await;
}
#[tokio::test]
async fn retry_granularity_is_separate_from_deferral_granularity() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let queue = unique_queue("retry-granularity");
let config = QueueConfig::new(queue.clone());
let backend = RabbitMqBackend::with_options(
&url,
RabbitMqOptions::default()
.retry_granularity(Duration::from_secs(5))
.deferred_granularity(Duration::from_secs(1)),
)
.await
.expect("connect");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare");
let delay = Duration::from_millis(1_200);
backend
.publish(&envelope(&queue, 1), Some(delay))
.await
.expect("delayed publish");
backend
.defer(&envelope(&queue, 1).deferred(10), delay)
.await
.expect("defer");
let retry_hold = backend.deferred_queue_name(&queue, 5_000);
let defer_hold = backend.deferred_queue_name(&queue, 2_000);
assert_eq!(
ready_count(&control, &retry_hold).await,
1,
"`{retry_hold}`"
);
assert_eq!(
ready_count(&control, &defer_hold).await,
1,
"`{defer_hold}`"
);
backend.close().await.expect("close");
cleanup(&control, &queue).await;
delete_queues(&control, &[retry_hold, defer_hold]).await;
}
#[tokio::test]
async fn retrying_onto_an_undeclared_queue_is_refused() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let queue = unique_queue("retry-undeclared");
let backend = RabbitMqBackend::connect(&url).await.expect("connect");
let control = control(&url).await;
let error = backend
.publish(&envelope(&queue, 1), Some(Duration::from_secs(1)))
.await
.expect_err("a delayed publish onto an undeclared queue must fail");
assert!(
matches!(&error, Error::UnknownQueue(name) if name == &queue),
"expected UnknownQueue, got {error:?}"
);
let hold = hold_queue(&queue, Duration::from_secs(1));
assert!(
!queue_exists(&control, &hold).await,
"`{hold}` must not have been created"
);
backend.close().await.expect("close");
}
#[tokio::test]
async fn dead_letter_moves_the_message_with_headers() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let queue = unique_queue("dead");
let dead = format!("{queue}.dead");
let config = QueueConfig::new(queue.clone());
let backend = RabbitMqBackend::connect(&url).await.expect("connect");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare");
let sent = envelope(&queue, 3);
backend.publish(&sent, None).await.expect("publish");
let mut stream = backend.consume(&config).await.expect("consume");
let delivery = next_delivery(&mut stream, Duration::from_secs(5)).await;
delivery
.dead_letter("handler returned Fatal")
.await
.expect("dead_letter");
drop(stream);
backend.close().await.expect("close");
await_count(&control, &queue, 0).await;
await_count(&control, &dead, 1).await;
let channel = control.create_channel().await.expect("channel");
let message = channel
.basic_get(dead.as_str().into(), BasicGetOptions { no_ack: true })
.await
.expect("basic_get")
.expect("a dead-lettered message");
assert_eq!(
Envelope::from_bytes(&message.delivery.data).expect("decode"),
sent
);
let headers = message
.delivery
.properties
.headers()
.clone()
.expect("headers");
assert_eq!(
header_string(&headers, "x-death-reason").as_deref(),
Some("handler returned Fatal")
);
assert_eq!(
header_string(&headers, "x-original-queue").as_deref(),
Some(queue.as_str())
);
assert_eq!(header_u32(&headers, "x-attempts"), Some(3));
let _ = channel.close(200, "OK".into()).await;
cleanup(&control, &queue).await;
}
#[tokio::test]
async fn publishing_to_a_missing_queue_is_an_error() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let queue = unique_queue("undeclared");
let backend = RabbitMqBackend::connect(&url).await.expect("connect");
let control = control(&url).await;
let error = backend
.publish(&envelope(&queue, 1), None)
.await
.expect_err("publishing to a queue that does not exist must fail");
let text = error.to_string();
assert!(
text.contains(&queue),
"error did not name the queue: {text}"
);
assert!(
text.contains("unroutable"),
"error did not explain the return: {text}"
);
let config = QueueConfig::new(queue.clone());
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare");
backend
.publish(&envelope(&queue, 1), None)
.await
.expect("publish after declaring");
await_count(&control, &queue, 1).await;
backend.close().await.expect("close");
cleanup(&control, &queue).await;
}
#[tokio::test]
async fn dead_letter_rejects_when_dead_letter_queues_are_disabled() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let queue = unique_queue("dead-nodlq");
let dead = format!("{queue}.dead");
let config = QueueConfig::new(queue.clone());
let backend = RabbitMqBackend::with_options(
&url,
RabbitMqOptions::default().declare_dead_letter_queues(false),
)
.await
.expect("connect");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare");
backend
.publish(&envelope(&queue, 3), None)
.await
.expect("publish");
let mut stream = backend.consume(&config).await.expect("consume");
let delivery = next_delivery(&mut stream, Duration::from_secs(5)).await;
delivery
.dead_letter("handler returned Fatal")
.await
.expect("dead_letter must succeed by rejecting");
drop(stream);
backend.close().await.expect("close");
await_count(&control, &queue, 0).await;
assert!(
!queue_exists(&control, &dead).await,
"`{dead}` must not exist when dead-letter queues are disabled"
);
cleanup(&control, &queue).await;
}
#[tokio::test]
async fn malformed_body_is_skipped_without_stalling() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let queue = unique_queue("malformed");
let dead = format!("{queue}.dead");
let config = QueueConfig::new(queue.clone());
let backend = RabbitMqBackend::connect(&url).await.expect("connect");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare");
let channel = control.create_channel().await.expect("channel");
publish_raw(&channel, &queue, b"{ this is not an envelope ]").await;
await_count(&control, &queue, 1).await;
let sent = envelope(&queue, 1);
backend.publish(&sent, None).await.expect("publish");
await_count(&control, &queue, 2).await;
let mut stream = backend.consume(&config).await.expect("consume");
let delivery = next_delivery(&mut stream, Duration::from_secs(5)).await;
assert_eq!(delivery.envelope(), &sent);
delivery.ack().await.expect("ack");
expect_idle(&mut stream, Duration::from_millis(500)).await;
drop(stream);
backend.close().await.expect("close");
await_count(&control, &queue, 0).await;
await_count(&control, &dead, 1).await;
let message = channel
.basic_get(dead.as_str().into(), BasicGetOptions { no_ack: true })
.await
.expect("basic_get")
.expect("the malformed body");
assert_eq!(message.delivery.data, b"{ this is not an envelope ]");
let headers = message
.delivery
.properties
.headers()
.clone()
.expect("headers");
assert_eq!(
header_string(&headers, "x-death-reason").as_deref(),
Some("malformed envelope")
);
assert_eq!(
header_string(&headers, "x-original-queue").as_deref(),
Some(queue.as_str())
);
let _ = channel.close(200, "OK".into()).await;
cleanup(&control, &queue).await;
}
#[tokio::test]
async fn malformed_body_is_rejected_when_dead_letter_queues_are_disabled() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let queue = unique_queue("malformed-nodlq");
let config = QueueConfig::new(queue.clone());
let backend = RabbitMqBackend::with_options(
&url,
RabbitMqOptions::default().declare_dead_letter_queues(false),
)
.await
.expect("connect");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare");
let channel = control.create_channel().await.expect("channel");
publish_raw(&channel, &queue, b"nonsense").await;
await_count(&control, &queue, 1).await;
let mut stream = backend.consume(&config).await.expect("consume");
expect_idle(&mut stream, Duration::from_secs(2)).await;
drop(stream);
backend.close().await.expect("close");
await_count(&control, &queue, 0).await;
let _ = channel.close(200, "OK".into()).await;
cleanup(&control, &queue).await;
}
#[tokio::test]
async fn close_ends_the_consumer_stream() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let queue = unique_queue("close");
let config = QueueConfig::new(queue.clone());
let backend = RabbitMqBackend::connect(&url).await.expect("connect");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare");
let mut stream = backend.consume(&config).await.expect("consume");
expect_idle(&mut stream, Duration::from_millis(250)).await;
backend.close().await.expect("close");
let ended = tokio::time::timeout(Duration::from_secs(10), async {
while let Some(item) = stream.next().await {
if let Ok(delivery) = item {
panic!(
"unexpected delivery after close: {}",
delivery.envelope().job_id
);
}
}
})
.await;
assert!(ended.is_ok(), "consumer stream did not end after close()");
cleanup(&control, &queue).await;
}
#[tokio::test]
async fn deferred_job_reappears_after_ttl() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let queue = unique_queue("defer-ttl");
let config = QueueConfig::new(queue.clone());
let backend = RabbitMqBackend::connect(&url).await.expect("connect");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare");
let hold = backend.deferred_queue_name(&queue, 1_000);
assert_eq!(hold, format!("{queue}.deferred.1000"));
let held = envelope(&queue, 2).deferred(10);
let started = Instant::now();
backend
.defer(&held, Duration::from_secs(1))
.await
.expect("defer");
assert!(
queue_exists(&control, &hold).await,
"`{hold}` must exist while the job is held"
);
assert_eq!(ready_count(&control, &hold).await, 1);
let mut stream = backend.consume(&config).await.expect("consume");
expect_idle(&mut stream, Duration::from_millis(600)).await;
let delivery = next_delivery(&mut stream, Duration::from_secs(8)).await;
let elapsed = started.elapsed();
assert!(
elapsed >= Duration::from_millis(900),
"job came back after {elapsed:?}, expected to wait about a second"
);
assert_eq!(delivery.envelope(), &held);
assert_eq!(delivery.envelope().deferrals, 1);
assert_eq!(delivery.envelope().priority, 10);
assert_eq!(
delivery.envelope().attempt,
2,
"a deferral is not a failed attempt"
);
delivery.ack().await.expect("ack");
drop(stream);
backend.close().await.expect("close");
tokio::time::sleep_until((started + Duration::from_millis(3_500)).into()).await;
assert!(
!queue_exists(&control, &hold).await,
"`{hold}` must have expired once it went unused for ttl + grace"
);
cleanup(&control, &queue).await;
delete_queues(&control, &[hold]).await;
}
#[tokio::test]
async fn deferred_job_is_consumed_before_the_backlog() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let queue = unique_queue("defer-priority");
let config = QueueConfig::new(queue.clone());
let backend = RabbitMqBackend::connect(&url).await.expect("connect");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare");
let first = envelope(&queue, 1);
let second = envelope(&queue, 2);
backend.publish(&first, None).await.expect("publish");
backend.publish(&second, None).await.expect("publish");
await_count(&control, &queue, 2).await;
let held = envelope(&queue, 7).deferred(10);
backend
.defer(&held, Duration::from_secs(1))
.await
.expect("defer");
await_count(&control, &queue, 3).await;
let mut stream = backend.consume(&config).await.expect("consume");
let deferred = next_delivery(&mut stream, Duration::from_secs(5)).await;
assert_eq!(
deferred.envelope(),
&held,
"the deferred job must overtake the backlog"
);
deferred.ack().await.expect("ack");
let one = next_delivery(&mut stream, Duration::from_secs(5)).await;
assert_eq!(one.envelope(), &first);
one.ack().await.expect("ack");
let two = next_delivery(&mut stream, Duration::from_secs(5)).await;
assert_eq!(two.envelope(), &second);
two.ack().await.expect("ack");
drop(stream);
backend.close().await.expect("close");
await_count(&control, &queue, 0).await;
cleanup(&control, &queue).await;
delete_queues(&control, &[format!("{queue}.deferred.1000")]).await;
}
#[tokio::test]
async fn queue_without_priorities_still_defers() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let queue = unique_queue("defer-nopriority");
let config = QueueConfig::new(queue.clone()).max_priority(0);
assert_eq!(config.max_priority, None);
let backend = RabbitMqBackend::connect(&url).await.expect("connect");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare must succeed without priorities");
redeclare_raw(&control, &queue, FieldTable::default())
.await
.expect("`q` must have been declared without any arguments");
let held = envelope(&queue, 1).deferred(0);
backend
.defer(&held, Duration::from_secs(1))
.await
.expect("defer");
let mut stream = backend.consume(&config).await.expect("consume");
let delivery = next_delivery(&mut stream, Duration::from_secs(8)).await;
assert_eq!(delivery.envelope(), &held);
assert_eq!(delivery.envelope().priority, 0);
assert_eq!(delivery.envelope().deferrals, 1);
delivery.ack().await.expect("ack");
drop(stream);
backend.close().await.expect("close");
cleanup(&control, &queue).await;
delete_queues(&control, &[format!("{queue}.deferred.1000")]).await;
}
#[tokio::test]
async fn redeclaring_with_different_max_priority_is_an_error() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let queue = unique_queue("defer-redeclare");
let with_priorities = QueueConfig::new(queue.clone()).max_priority(10);
let without = QueueConfig::new(queue.clone()).max_priority(0);
let backend = RabbitMqBackend::connect(&url).await.expect("connect");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&with_priorities))
.await
.expect("declare");
let error = tokio::time::timeout(
Duration::from_secs(10),
backend.declare(std::slice::from_ref(&without)),
)
.await
.expect("redeclaring must fail rather than hang")
.expect_err("redeclaring with different arguments must be an error");
let text = error.to_string();
assert!(
text.contains("PRECONDITION_FAILED"),
"error did not explain the refusal: {text}"
);
let other = RabbitMqBackend::connect(&url).await.expect("connect");
other
.declare(std::slice::from_ref(&without))
.await
.expect_err("a fresh backend must be refused too");
other
.publish(&envelope(&queue, 1), None)
.await
.expect("publishing must still work after a refused declaration");
other.close().await.expect("close");
backend
.publish(&envelope(&queue, 2), None)
.await
.expect("publishing must still work after a refused declaration");
await_count(&control, &queue, 2).await;
backend
.declare(std::slice::from_ref(&with_priorities))
.await
.expect("the queue still has its original arguments");
backend.close().await.expect("close");
cleanup(&control, &queue).await;
}
#[tokio::test]
async fn delivery_defer_acks_the_original_before_holding_the_next() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let queue = unique_queue("defer-delivery");
let config = QueueConfig::new(queue.clone());
let backend = RabbitMqBackend::connect(&url).await.expect("connect");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare");
let sent = envelope(&queue, 1);
backend.publish(&sent, None).await.expect("publish");
let mut stream = backend.consume(&config).await.expect("consume");
let delivery = next_delivery(&mut stream, Duration::from_secs(5)).await;
let held = delivery.envelope().deferred(10);
delivery
.defer(held.clone(), Duration::from_secs(2))
.await
.expect("defer");
let hold = backend.deferred_queue_name(&queue, 2_000);
assert_eq!(ready_count(&control, &hold).await, 1, "held in `{hold}`");
drop(stream);
backend.close().await.expect("close");
assert_eq!(
ready_count(&control, &queue).await,
0,
"the original must have been acked, not requeued"
);
tokio::time::sleep(Duration::from_millis(500)).await;
assert_eq!(ready_count(&control, &queue).await, 0, "still nothing due");
await_count(&control, &queue, 1).await;
let channel = control.create_channel().await.expect("channel");
let message = channel
.basic_get(queue.as_str().into(), BasicGetOptions { no_ack: true })
.await
.expect("basic_get")
.expect("the deferred job");
assert_eq!(
Envelope::from_bytes(&message.delivery.data).expect("decode"),
held
);
assert_eq!(*message.delivery.properties.priority(), Some(10));
let headers = message
.delivery
.properties
.headers()
.clone()
.expect("headers");
assert_eq!(header_u32(&headers, "x-deferrals"), Some(1));
assert_eq!(header_u32(&headers, "x-attempt"), Some(1));
let _ = channel.close(200, "OK".into()).await;
cleanup(&control, &queue).await;
delete_queues(&control, &[hold]).await;
}
#[tokio::test]
async fn hold_queue_with_foreign_arguments_makes_defer_fail_and_leaves_the_original_unacked() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let queue = unique_queue("defer-foreign");
let config = QueueConfig::new(queue.clone());
let backend = RabbitMqBackend::connect(&url).await.expect("connect");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare");
let hold = backend.deferred_queue_name(&queue, 1_000);
let mut foreign = FieldTable::default();
foreign.insert("x-message-ttl".into(), AMQPValue::LongLongInt(1_000));
redeclare_raw(&control, &hold, foreign)
.await
.expect("pre-declaring the hold queue with foreign arguments");
let sent = envelope(&queue, 1);
backend.publish(&sent, None).await.expect("publish");
let mut stream = backend.consume(&config).await.expect("consume");
let delivery = next_delivery(&mut stream, Duration::from_secs(5)).await;
let held = delivery.envelope().deferred(10);
let error = delivery
.defer(held, Duration::from_secs(1))
.await
.expect_err("deferring into a hold queue with foreign arguments must fail");
let text = error.to_string();
assert!(
text.contains("PRECONDITION_FAILED"),
"error did not explain the refusal: {text}"
);
assert_eq!(ready_count(&control, &hold).await, 0);
let extra = envelope(&queue, 9);
backend
.publish(&extra, None)
.await
.expect("publishing must still work after a refused hold queue declaration");
drop(stream);
backend.close().await.expect("close");
await_count(&control, &queue, 2).await;
cleanup(&control, &queue).await;
delete_queues(&control, &[hold]).await;
}
#[tokio::test]
async fn deferring_onto_an_undeclared_queue_is_refused() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let queue = unique_queue("defer-undeclared");
let backend = RabbitMqBackend::connect(&url).await.expect("connect");
let control = control(&url).await;
let held = envelope(&queue, 1).deferred(10);
let error = backend
.defer(&held, Duration::from_secs(1))
.await
.expect_err("deferring onto a queue this backend never declared must fail");
assert!(
matches!(&error, Error::UnknownQueue(name) if name == &queue),
"expected UnknownQueue, got {error:?}"
);
let hold = backend.deferred_queue_name(&queue, 1_000);
assert!(
!queue_exists(&control, &hold).await,
"`{hold}` must not have been created"
);
let config = QueueConfig::new(queue.clone());
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare");
backend
.defer(&held, Duration::from_secs(1))
.await
.expect("deferring must work once the queue has been declared");
assert_eq!(ready_count(&control, &hold).await, 1);
backend.close().await.expect("close");
cleanup(&control, &queue).await;
delete_queues(&control, &[hold]).await;
}
#[tokio::test]
async fn deferral_past_the_cap_is_refused() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let queue = unique_queue("defer-too-long");
let config = QueueConfig::new(queue.clone());
let backend = RabbitMqBackend::connect(&url).await.expect("connect");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare");
let held = envelope(&queue, 1).deferred(10);
let thirty_days = Duration::from_secs(30 * 86_400);
let error = backend
.defer(&held, thirty_days)
.await
.expect_err("a deferral past the cap must fail");
let text = error.to_string();
assert!(
text.contains("longer than a hold queue can wait"),
"error did not explain the cap: {text}"
);
assert_eq!(ready_count(&control, &queue).await, 0);
let absurd = backend.deferred_queue_name(&queue, 2_592_000_000);
assert!(!queue_exists(&control, &absurd).await, "`{absurd}` exists");
let hold = backend.deferred_queue_name(&queue, 1_000);
backend
.defer(&held, Duration::from_secs(1))
.await
.expect("a delay within the cap must still work");
assert_eq!(ready_count(&control, &hold).await, 1);
backend.close().await.expect("close");
cleanup(&control, &queue).await;
delete_queues(&control, &[hold]).await;
}
struct BrokerProxy {
addr: std::net::SocketAddr,
cut: tokio::sync::broadcast::Sender<()>,
}
impl BrokerProxy {
async fn start(upstream: String) -> Self {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind the proxy");
let addr = listener.local_addr().expect("proxy address");
let (cut, _) = tokio::sync::broadcast::channel(16);
let accepts = cut.clone();
tokio::spawn(async move {
loop {
let Ok((mut inbound, _)) = listener.accept().await else {
return;
};
let upstream = upstream.clone();
let mut stop = accepts.subscribe();
tokio::spawn(async move {
let Ok(mut outbound) = tokio::net::TcpStream::connect(&upstream).await else {
return;
};
tokio::select! {
_ = tokio::io::copy_bidirectional(&mut inbound, &mut outbound) => {}
_ = stop.recv() => {}
}
});
}
});
Self { addr, cut }
}
fn cut(&self) {
let _ = self.cut.send(());
}
fn url(&self, direct: &str) -> String {
let (head, _, tail) = split_url(direct);
format!("{head}{}{tail}", self.addr)
}
}
fn split_url(url: &str) -> (&str, &str, &str) {
let scheme_end = url.find("://").expect("an amqp:// url") + "://".len();
let rest = &url[scheme_end..];
let authority_end = rest.find('/').unwrap_or(rest.len());
let authority = &rest[..authority_end];
let tail = &rest[authority_end..];
match authority.rfind('@') {
Some(at) => (&url[..scheme_end + at + 1], &authority[at + 1..], tail),
None => (&url[..scheme_end], authority, tail),
}
}
fn upstream_addr(url: &str) -> String {
let (_, hostport, _) = split_url(url);
if hostport.contains(':') {
hostport.to_owned()
} else {
format!("{hostport}:5672")
}
}
async fn await_disconnected(backend: &RabbitMqBackend) {
let deadline = Instant::now() + Duration::from_secs(10);
while Instant::now() < deadline {
if !backend.is_connected() {
return;
}
tokio::time::sleep(Duration::from_millis(25)).await;
}
panic!("the backend never noticed its connection was cut");
}
#[tokio::test]
async fn publishing_reconnects_after_the_connection_drops() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let proxy = BrokerProxy::start(upstream_addr(&url)).await;
let queue = unique_queue("reconnect-publish");
let config = QueueConfig::new(queue.clone());
let backend = RabbitMqBackend::connect(&proxy.url(&url))
.await
.expect("connect through the proxy");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare");
backend
.publish(&envelope(&queue, 1), None)
.await
.expect("publish before the cut");
await_count(&control, &queue, 1).await;
proxy.cut();
await_disconnected(&backend).await;
backend
.publish(&envelope(&queue, 2), None)
.await
.expect("publish after the cut must reconnect, not fail");
await_count(&control, &queue, 2).await;
assert!(
backend.is_connected(),
"the backend is back on a connection"
);
cleanup(&control, &queue).await;
backend.close().await.expect("close");
}
#[tokio::test]
async fn a_consumer_resubscribes_after_the_connection_drops() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let proxy = BrokerProxy::start(upstream_addr(&url)).await;
let queue = unique_queue("reconnect-consume");
let config = QueueConfig::new(queue.clone());
let backend = RabbitMqBackend::connect(&proxy.url(&url))
.await
.expect("connect through the proxy");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare");
let mut stream = backend.consume(&config).await.expect("consume");
backend
.publish(&envelope(&queue, 1), None)
.await
.expect("publish");
let first = next_delivery(&mut stream, Duration::from_secs(5)).await;
assert_eq!(first.envelope().attempt, 1);
first.ack().await.expect("ack before the cut");
proxy.cut();
await_disconnected(&backend).await;
backend
.publish(&envelope(&queue, 2), None)
.await
.expect("publish after the cut");
let mut resubscribed = false;
for _ in 0..3 {
let delivery = next_delivery(&mut stream, Duration::from_secs(15)).await;
let attempt = delivery.envelope().attempt;
delivery.ack().await.expect("ack after the cut");
if attempt == 2 {
resubscribed = true;
break;
}
assert_eq!(
attempt, 1,
"only the job that was in flight across the cut can be redelivered"
);
}
assert!(
resubscribed,
"the resubscribed consumer never delivered the job published after the cut"
);
drop(stream);
cleanup(&control, &queue).await;
backend.close().await.expect("close");
}
#[tokio::test]
async fn a_reconnect_redeclares_the_queues_it_had_declared() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let proxy = BrokerProxy::start(upstream_addr(&url)).await;
let queue = unique_queue("reconnect-redeclare");
let config = QueueConfig::new(queue.clone());
let backend = RabbitMqBackend::connect(&proxy.url(&url))
.await
.expect("connect through the proxy");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare");
proxy.cut();
await_disconnected(&backend).await;
delete_queues(&control, &family(&queue)).await;
assert!(!queue_exists(&control, &queue).await, "the queue is gone");
backend
.publish(&envelope(&queue, 1), None)
.await
.expect("publish after the cut must find a re-declared queue");
await_count(&control, &queue, 1).await;
assert!(
queue_exists(&control, &backend.dead_queue_name(&queue)).await,
"the dead-letter queue is re-declared too"
);
cleanup(&control, &queue).await;
backend.close().await.expect("close");
}
#[tokio::test]
async fn reconnection_can_be_turned_off() {
let Some(url) = std::env::var("AMQP_URL").ok() else {
eprintln!("skipping: AMQP_URL not set");
return;
};
let proxy = BrokerProxy::start(upstream_addr(&url)).await;
let queue = unique_queue("reconnect-disabled");
let config = QueueConfig::new(queue.clone());
let backend =
RabbitMqBackend::with_options(&proxy.url(&url), RabbitMqOptions::default().reconnect(None))
.await
.expect("connect through the proxy");
let control = control(&url).await;
backend
.declare(std::slice::from_ref(&config))
.await
.expect("declare");
let mut stream = backend.consume(&config).await.expect("consume");
proxy.cut();
await_disconnected(&backend).await;
match tokio::time::timeout(Duration::from_secs(10), stream.next()).await {
Ok(None) => {}
Ok(Some(Ok(_))) => panic!("a delivery arrived on a connection that was cut"),
Ok(Some(Err(_))) => {}
Err(_) => panic!("the consumer stream neither ended nor errored"),
}
let error = backend
.publish(&envelope(&queue, 1), None)
.await
.expect_err("publishing must fail rather than reconnect");
assert!(
error.to_string().contains("reconnection is disabled"),
"unexpected error: {error}"
);
drop(stream);
cleanup(&control, &queue).await;
let _ = backend.close().await;
}