use std::time::Duration;
use lapin::{
options::QueueDeclareOptions,
types::{AMQPValue, FieldTable, LongString},
};
use queuey_core::QueueConfig;
pub const DEFAULT_RETRY_SUFFIX: &str = ".retry";
pub const DEFAULT_DEAD_SUFFIX: &str = ".dead";
pub const DEFAULT_DEFERRED_SUFFIX: &str = ".deferred";
pub const HEADER_DEATH_REASON: &str = "x-death-reason";
pub const HEADER_ORIGINAL_QUEUE: &str = "x-original-queue";
pub const HEADER_ATTEMPTS: &str = "x-attempts";
pub const HEADER_ATTEMPT: &str = "x-attempt";
pub const HEADER_DEFERRALS: &str = "x-deferrals";
pub const ARG_DEAD_LETTER_EXCHANGE: &str = "x-dead-letter-exchange";
pub const ARG_DEAD_LETTER_ROUTING_KEY: &str = "x-dead-letter-routing-key";
pub const ARG_MESSAGE_TTL: &str = "x-message-ttl";
pub const ARG_MAX_PRIORITY: &str = "x-max-priority";
pub const ARG_EXPIRES: &str = "x-expires";
pub const MAX_TTL_MS: u32 = u32::MAX;
pub const MAX_DEFERRAL_MS: u32 = MAX_TTL_MS / 2;
#[must_use]
pub fn retry_queue_name(queue: &str, suffix: &str) -> String {
format!("{queue}{suffix}")
}
#[must_use]
pub fn dead_queue_name(queue: &str, suffix: &str) -> String {
format!("{queue}{suffix}")
}
#[must_use]
pub fn deferred_queue_name(queue: &str, suffix: &str, ttl_ms: u32) -> String {
format!("{queue}{suffix}.{ttl_ms}")
}
#[must_use]
pub fn deferred_ttl_ms(delay: Duration, granularity: Duration) -> Option<u32> {
let max = u128::from(MAX_DEFERRAL_MS);
let step = millis_ceil(granularity).clamp(1, max);
let ttl = millis_ceil(delay)
.div_ceil(step)
.saturating_mul(step)
.max(step);
(ttl <= max).then(|| u32::try_from(ttl).unwrap_or(MAX_DEFERRAL_MS))
}
#[must_use]
pub fn queue_args(config: &QueueConfig) -> FieldTable {
let mut args = FieldTable::default();
if let Some(ttl) = config.message_ttl {
args.insert(
ARG_MESSAGE_TTL.into(),
AMQPValue::LongLongInt(ttl_millis(ttl)),
);
}
if let Some(levels) = config.max_priority {
args.insert(ARG_MAX_PRIORITY.into(), AMQPValue::ShortShortUInt(levels));
}
args
}
#[must_use]
pub fn deferred_queue_args(config: &QueueConfig, ttl_ms: u32) -> FieldTable {
let mut args = FieldTable::default();
args.insert(
ARG_MESSAGE_TTL.into(),
AMQPValue::LongLongInt(ttl_ms.into()),
);
args.insert(
ARG_DEAD_LETTER_EXCHANGE.into(),
AMQPValue::LongString(LongString::from("")),
);
args.insert(
ARG_DEAD_LETTER_ROUTING_KEY.into(),
AMQPValue::LongString(LongString::from(config.name.as_str())),
);
let expires = u64::from(ttl_ms)
.saturating_mul(2)
.min(u64::from(MAX_TTL_MS));
args.insert(
ARG_EXPIRES.into(),
AMQPValue::LongLongInt(i64::try_from(expires).unwrap_or(i64::from(MAX_TTL_MS))),
);
args
}
pub(crate) fn declare_options(durable: bool) -> QueueDeclareOptions {
QueueDeclareOptions {
passive: false,
durable,
exclusive: false,
auto_delete: false,
nowait: false,
}
}
#[must_use]
pub fn retry_queue_args(config: &QueueConfig) -> FieldTable {
let mut args = FieldTable::default();
args.insert(
ARG_DEAD_LETTER_EXCHANGE.into(),
AMQPValue::LongString(LongString::from("")),
);
args.insert(
ARG_DEAD_LETTER_ROUTING_KEY.into(),
AMQPValue::LongString(LongString::from(config.name.as_str())),
);
args
}
#[must_use]
pub fn dead_queue_args(_config: &QueueConfig) -> FieldTable {
FieldTable::default()
}
fn ttl_millis(ttl: Duration) -> i64 {
let ms = ttl
.as_nanos()
.div_ceil(1_000_000)
.clamp(1, u128::from(MAX_TTL_MS));
i64::try_from(ms).unwrap_or(i64::from(MAX_TTL_MS))
}
fn millis_ceil(duration: Duration) -> u128 {
duration.as_nanos().div_ceil(1_000_000)
}
#[cfg(test)]
mod tests {
use super::*;
fn config(name: &str) -> QueueConfig {
QueueConfig::new(name)
}
#[test]
fn retry_name_without_prefix() {
assert_eq!(
retry_queue_name("emails", DEFAULT_RETRY_SUFFIX),
"emails.retry"
);
}
#[test]
fn retry_name_with_prefix() {
assert_eq!(
retry_queue_name("myapp.emails", DEFAULT_RETRY_SUFFIX),
"myapp.emails.retry"
);
}
#[test]
fn dead_name_without_prefix() {
assert_eq!(
dead_queue_name("emails", DEFAULT_DEAD_SUFFIX),
"emails.dead"
);
}
#[test]
fn dead_name_with_prefix() {
assert_eq!(
dead_queue_name("myapp.emails", DEFAULT_DEAD_SUFFIX),
"myapp.emails.dead"
);
}
#[test]
fn custom_suffixes_are_honoured() {
assert_eq!(retry_queue_name("q", "-wait"), "q-wait");
assert_eq!(dead_queue_name("q", "-dlq"), "q-dlq");
}
fn plain(name: &str) -> QueueConfig {
QueueConfig::new(name).max_priority(0)
}
#[test]
fn queue_args_are_empty_without_ttl_or_priorities() {
let args = queue_args(&plain("emails"));
assert!(args.inner().is_empty(), "unexpected args: {args:?}");
assert!(!args.contains_key(ARG_MESSAGE_TTL));
assert!(!args.contains_key(ARG_MAX_PRIORITY));
}
#[test]
fn queue_args_carry_message_ttl() {
let cfg = plain("emails").message_ttl(Duration::from_secs(30));
let args = queue_args(&cfg);
assert_eq!(
args.inner().get(ARG_MESSAGE_TTL),
Some(&AMQPValue::LongLongInt(30_000))
);
assert_eq!(args.inner().len(), 1);
}
#[test]
fn queue_args_carry_max_priority_when_the_queue_has_priorities() {
let args = queue_args(&config("emails"));
assert_eq!(
args.inner().get(ARG_MAX_PRIORITY),
Some(&AMQPValue::ShortShortUInt(10))
);
assert_eq!(args.inner().len(), 1);
let args = queue_args(&config("emails").max_priority(255));
assert_eq!(
args.inner().get(ARG_MAX_PRIORITY),
Some(&AMQPValue::ShortShortUInt(255))
);
}
#[test]
fn queue_args_omit_max_priority_when_priorities_are_off() {
let cfg = config("emails").max_priority(0);
assert_eq!(cfg.max_priority, None);
assert!(!queue_args(&cfg).contains_key(ARG_MAX_PRIORITY));
}
#[test]
fn queue_args_carry_ttl_and_priority_together() {
let cfg = config("emails")
.message_ttl(Duration::from_secs(5))
.max_priority(3);
let args = queue_args(&cfg);
assert_eq!(
args.inner().get(ARG_MESSAGE_TTL),
Some(&AMQPValue::LongLongInt(5_000))
);
assert_eq!(
args.inner().get(ARG_MAX_PRIORITY),
Some(&AMQPValue::ShortShortUInt(3))
);
assert_eq!(args.inner().len(), 2);
}
#[test]
fn only_the_main_queue_gets_priorities() {
let cfg = config("emails");
assert!(!retry_queue_args(&cfg).contains_key(ARG_MAX_PRIORITY));
assert!(!dead_queue_args(&cfg).contains_key(ARG_MAX_PRIORITY));
assert!(!deferred_queue_args(&cfg, 1_000).contains_key(ARG_MAX_PRIORITY));
}
#[test]
fn sub_millisecond_ttl_rounds_up_to_one() {
let cfg = config("emails").message_ttl(Duration::from_nanos(1));
let args = queue_args(&cfg);
assert_eq!(
args.inner().get(ARG_MESSAGE_TTL),
Some(&AMQPValue::LongLongInt(1))
);
}
#[test]
fn zero_ttl_is_clamped_to_one_millisecond() {
let cfg = config("emails").message_ttl(Duration::ZERO);
let args = queue_args(&cfg);
assert_eq!(
args.inner().get(ARG_MESSAGE_TTL),
Some(&AMQPValue::LongLongInt(1))
);
}
#[test]
fn huge_ttl_is_clamped_to_what_rabbitmq_accepts() {
let cfg = config("emails").message_ttl(Duration::MAX);
let args = queue_args(&cfg);
assert_eq!(
args.inner().get(ARG_MESSAGE_TTL),
Some(&AMQPValue::LongLongInt(i64::from(MAX_TTL_MS)))
);
}
#[test]
fn ttl_exactly_at_the_limit_is_kept_verbatim() {
let cfg = config("emails").message_ttl(Duration::from_millis(u64::from(MAX_TTL_MS)));
let args = queue_args(&cfg);
assert_eq!(
args.inner().get(ARG_MESSAGE_TTL),
Some(&AMQPValue::LongLongInt(i64::from(MAX_TTL_MS)))
);
}
#[test]
fn ttl_one_millisecond_past_the_limit_is_clamped() {
let cfg = config("emails").message_ttl(Duration::from_millis(u64::from(MAX_TTL_MS) + 1));
let args = queue_args(&cfg);
assert_eq!(
args.inner().get(ARG_MESSAGE_TTL),
Some(&AMQPValue::LongLongInt(i64::from(MAX_TTL_MS)))
);
}
#[test]
fn retry_args_dead_letter_back_to_the_main_queue() {
let args = retry_queue_args(&config("myapp.emails"));
assert_eq!(
args.inner().get(ARG_DEAD_LETTER_EXCHANGE),
Some(&AMQPValue::LongString(LongString::from("")))
);
assert_eq!(
args.inner().get(ARG_DEAD_LETTER_ROUTING_KEY),
Some(&AMQPValue::LongString(LongString::from("myapp.emails")))
);
assert_eq!(args.inner().len(), 2);
}
#[test]
fn retry_args_never_carry_a_queue_ttl() {
let cfg = config("emails").message_ttl(Duration::from_secs(5));
let args = retry_queue_args(&cfg);
assert!(!args.contains_key(ARG_MESSAGE_TTL));
}
#[test]
fn dead_args_are_empty() {
assert!(dead_queue_args(&config("emails")).inner().is_empty());
}
const SECOND: Duration = Duration::from_secs(1);
#[test]
fn deferred_name_puts_the_ttl_in_the_queue_name() {
assert_eq!(
deferred_queue_name("emails", DEFAULT_DEFERRED_SUFFIX, 30_000),
"emails.deferred.30000"
);
assert_eq!(
deferred_queue_name("myapp.emails", DEFAULT_DEFERRED_SUFFIX, 1_000),
"myapp.emails.deferred.1000"
);
}
#[test]
fn deferred_name_honours_a_custom_suffix() {
assert_eq!(deferred_queue_name("q", "-hold", 250), "q-hold.250");
}
#[test]
fn deferred_names_differ_per_ttl_which_is_the_whole_point() {
let short = deferred_queue_name("q", DEFAULT_DEFERRED_SUFFIX, 1_000);
let long = deferred_queue_name("q", DEFAULT_DEFERRED_SUFFIX, 60_000);
assert_ne!(short, long);
}
#[test]
fn an_exact_multiple_of_the_granularity_is_kept_verbatim() {
assert_eq!(deferred_ttl_ms(SECOND, SECOND), Some(1_000));
assert_eq!(
deferred_ttl_ms(Duration::from_secs(30), SECOND),
Some(30_000)
);
assert_eq!(
deferred_ttl_ms(Duration::from_millis(750), Duration::from_millis(250)),
Some(750)
);
}
#[test]
fn a_partial_step_rounds_up_so_a_job_never_returns_early() {
assert_eq!(
deferred_ttl_ms(Duration::from_millis(29_200), SECOND),
Some(30_000)
);
assert_eq!(
deferred_ttl_ms(Duration::from_millis(1_001), SECOND),
Some(2_000)
);
assert_eq!(
deferred_ttl_ms(Duration::from_millis(1_999), SECOND),
Some(2_000)
);
assert_eq!(
deferred_ttl_ms(Duration::from_micros(1_000_001), SECOND),
Some(2_000)
);
}
#[test]
fn a_delay_below_the_granularity_becomes_one_step() {
assert_eq!(
deferred_ttl_ms(Duration::from_millis(5), SECOND),
Some(1_000)
);
assert_eq!(
deferred_ttl_ms(Duration::from_nanos(1), SECOND),
Some(1_000)
);
assert_eq!(deferred_ttl_ms(Duration::ZERO, SECOND), Some(1_000));
}
#[test]
fn a_delay_past_the_cap_is_refused_rather_than_released_early() {
assert_eq!(deferred_ttl_ms(Duration::MAX, SECOND), None);
assert_eq!(
deferred_ttl_ms(Duration::from_secs(30 * 86_400), SECOND),
None
);
assert_eq!(
deferred_ttl_ms(
Duration::from_millis(u64::from(MAX_DEFERRAL_MS) + 1),
SECOND
),
None
);
}
#[test]
fn the_cap_itself_is_accepted_at_a_matching_granularity() {
assert_eq!(
deferred_ttl_ms(
Duration::from_millis(u64::from(MAX_DEFERRAL_MS)),
Duration::from_millis(1)
),
Some(MAX_DEFERRAL_MS)
);
}
#[test]
fn rounding_up_can_itself_cross_the_cap_and_is_then_refused() {
let just_under = Duration::from_millis(u64::from(MAX_DEFERRAL_MS) - 1);
assert_eq!(deferred_ttl_ms(just_under, SECOND), None);
}
#[test]
fn a_zero_granularity_is_clamped_to_one_millisecond_rather_than_panicking() {
assert_eq!(
deferred_ttl_ms(Duration::from_millis(7), Duration::ZERO),
Some(7)
);
assert_eq!(deferred_ttl_ms(Duration::ZERO, Duration::ZERO), Some(1));
assert_eq!(
deferred_ttl_ms(Duration::from_millis(7), Duration::from_nanos(1)),
Some(7)
);
}
#[test]
fn an_absurd_granularity_is_clamped_to_the_cap_not_past_it() {
assert_eq!(
deferred_ttl_ms(SECOND, Duration::MAX),
Some(MAX_DEFERRAL_MS)
);
}
#[test]
fn deferred_args_hold_for_the_ttl_then_dead_letter_onto_the_main_queue() {
let args = deferred_queue_args(&config("myapp.emails"), 30_000);
assert_eq!(
args.inner().get(ARG_MESSAGE_TTL),
Some(&AMQPValue::LongLongInt(30_000))
);
assert_eq!(
args.inner().get(ARG_DEAD_LETTER_EXCHANGE),
Some(&AMQPValue::LongString(LongString::from("")))
);
assert_eq!(
args.inner().get(ARG_DEAD_LETTER_ROUTING_KEY),
Some(&AMQPValue::LongString(LongString::from("myapp.emails")))
);
assert_eq!(
args.inner().get(ARG_EXPIRES),
Some(&AMQPValue::LongLongInt(60_000))
);
assert_eq!(args.inner().len(), 4);
}
#[test]
fn deferred_args_depend_only_on_the_hold_queue_name() {
let one = config("q").message_ttl(Duration::from_secs(99)).prefetch(1);
let two = config("q").max_priority(0).prefetch(200);
assert_eq!(
deferred_queue_args(&one, 1_000),
deferred_queue_args(&two, 1_000)
);
}
#[test]
fn expires_is_always_strictly_greater_than_the_ttl() {
for ttl in [1_u32, 1_000, 30_000, MAX_DEFERRAL_MS] {
let args = deferred_queue_args(&config("q"), ttl);
let AMQPValue::LongLongInt(expires) = args.inner().get(ARG_EXPIRES).expect("x-expires")
else {
panic!("x-expires is not a long long int");
};
assert!(
*expires > i64::from(ttl),
"x-expires {expires} must outlive the {ttl}ms TTL"
);
assert!(
*expires <= i64::from(MAX_TTL_MS),
"x-expires {expires} is past what RabbitMQ accepts"
);
}
}
#[test]
fn expires_is_clamped_when_doubling_the_ttl_would_overflow() {
let args = deferred_queue_args(&config("q"), MAX_TTL_MS);
assert_eq!(
args.inner().get(ARG_EXPIRES),
Some(&AMQPValue::LongLongInt(i64::from(MAX_TTL_MS)))
);
}
#[test]
fn deferred_args_never_carry_a_queue_ttl_of_their_own_from_the_config() {
let cfg = config("q").message_ttl(Duration::from_secs(99));
let args = deferred_queue_args(&cfg, 2_000);
assert_eq!(
args.inner().get(ARG_MESSAGE_TTL),
Some(&AMQPValue::LongLongInt(2_000))
);
}
#[test]
fn declare_options_are_shared_and_never_auto_delete() {
let durable = declare_options(true);
assert!(durable.durable);
assert!(!durable.auto_delete);
assert!(!durable.exclusive);
assert!(!durable.passive);
assert!(!durable.nowait);
assert!(!declare_options(false).durable);
}
}