use super::*;
fn queued(key: &str, delivery: usize) -> QueuedAdmission<usize> {
QueuedAdmission {
key: key.to_owned(),
recorded_at: "2026-08-17T00:00:00Z".to_owned(),
admission: PendingTurnAdmission::User {
text: format!("text-{key}"),
metadata: serde_json::json!({}),
},
delivery,
}
}
#[test]
fn fifo_and_queued_duplicate_keep_original_sequence() {
let mailbox = Mailbox::new();
let sender = mailbox.sender();
assert_eq!(
sender.push(queued("a", 10)),
PushResult::Queued { sequence: 1 }
);
assert_eq!(
sender.push(queued("b", 20)),
PushResult::Queued { sequence: 2 }
);
assert_eq!(
sender.push(queued("a", 99)),
PushResult::Duplicate { sequence: 1 }
);
let drained = mailbox.drain_through(mailbox.watermark());
let observed: Vec<_> = drained
.iter()
.map(|entry| (entry.sequence, entry.item.key.as_str(), entry.item.delivery))
.collect();
assert_eq!(observed, [(1, "a", 10), (2, "b", 20)]);
}
#[test]
fn captured_watermark_excludes_later_push() {
let mailbox = Mailbox::new();
let sender = mailbox.sender();
assert_eq!(
sender.push(queued("a", 1)),
PushResult::Queued { sequence: 1 }
);
let watermark = mailbox.watermark();
assert_eq!(
sender.push(queued("b", 2)),
PushResult::Queued { sequence: 2 }
);
let first = mailbox.drain_through(watermark);
assert_eq!(
first
.iter()
.map(|entry| (entry.sequence, entry.item.delivery))
.collect::<Vec<_>>(),
[(1, 1)]
);
let later = mailbox.drain_through(mailbox.watermark());
assert_eq!(
later
.iter()
.map(|entry| (entry.sequence, entry.item.delivery))
.collect::<Vec<_>>(),
[(2, 2)]
);
}
#[test]
fn privately_drained_duplicate_and_restoration_preserve_order() {
let mailbox = Mailbox::new();
let sender = mailbox.sender();
sender.push(queued("a", 1));
sender.push(queued("b", 2));
let drained = mailbox.drain_through(mailbox.watermark());
assert_eq!(
sender.push(queued("a", 99)),
PushResult::Duplicate { sequence: 1 }
);
assert_eq!(
sender.push(queued("c", 3)),
PushResult::Queued { sequence: 3 }
);
mailbox.restore_front(drained);
let restored = mailbox.drain_through(mailbox.watermark());
let observed: Vec<_> = restored
.iter()
.map(|entry| (entry.sequence, entry.item.key.as_str(), entry.item.delivery))
.collect();
assert_eq!(observed, [(1, "a", 1), (2, "b", 2), (3, "c", 3)]);
}
#[test]
fn exact_acknowledgement_permits_key_at_later_sequence() {
let mailbox = Mailbox::new();
let sender = mailbox.sender();
assert_eq!(
sender.push(queued("key", 1)),
PushResult::Queued { sequence: 1 }
);
let first = mailbox.drain_through(mailbox.watermark()).pop().unwrap();
mailbox.acknowledge(first.sequence, &first.item.key);
assert_eq!(
sender.push(queued("key", 2)),
PushResult::Queued { sequence: 2 }
);
let second = mailbox.drain_through(mailbox.watermark()).pop().unwrap();
assert_eq!((second.sequence, second.item.delivery), (2, 2));
mailbox.acknowledge(second.sequence, &second.item.key);
}
#[tokio::test(flavor = "multi_thread")]
async fn coalesced_notify_and_concurrent_producers_preserve_sequences_and_payloads() {
let mailbox = Mailbox::new();
let sender = mailbox.sender();
sender.push(queued("first", 100));
sender.push(queued("second", 200));
mailbox.notified().await;
const N: usize = 32;
let mut tasks = Vec::new();
for delivery in 0..N {
let producer = sender.clone();
tasks.push(tokio::spawn(async move {
let result = producer.push(queued(&format!("key-{delivery}"), delivery));
let PushResult::Queued { sequence } = result else {
panic!("unique producer key was duplicated");
};
(sequence, delivery)
}));
}
let mut assigned = Vec::new();
for task in tasks {
assigned.push(task.await.unwrap());
}
assigned.sort_unstable();
assert_eq!(
assigned.iter().map(|entry| entry.0).collect::<Vec<_>>(),
(3..=N as u64 + 2).collect::<Vec<_>>()
);
let drained = mailbox.drain_through(mailbox.watermark());
let observed: Vec<_> = drained
.iter()
.map(|entry| (entry.sequence, entry.item.delivery))
.collect();
let mut expected = vec![(1, 100), (2, 200)];
expected.extend(assigned);
assert_eq!(observed, expected);
}