use super::*;
use mdns_proto::{ServiceUpdate, event::ServiceRenamed};
fn renamed(n: &str) -> ServiceUpdate {
ServiceUpdate::Renamed(ServiceRenamed::new(
mdns_proto::Name::try_from_str(n).unwrap(),
))
}
#[test]
fn mailbox_coalesces_established_and_renamed_by_kind() {
let mut mb = ServiceMailbox::new();
mb.push_update(ServiceUpdate::Established);
mb.push_update(ServiceUpdate::Established);
mb.push_update(renamed("a-1._ipp._tcp.local."));
mb.push_update(renamed("a-2._ipp._tcp.local."));
mb.push_update(renamed("a-3._ipp._tcp.local."));
assert_eq!(
mb.updates.len(),
2,
"one Established + one (latest) Renamed"
);
assert!(matches!(
mb.drain(),
Drained::Update(ServiceUpdate::Established)
));
match mb.drain() {
Drained::Update(ServiceUpdate::Renamed(r)) => {
assert!(r.new_name().as_str().contains("a-3"))
}
other => panic!("expected the latest Renamed; got {other:?}"),
}
assert!(matches!(mb.drain(), Drained::Empty));
}
#[test]
fn mailbox_preserves_post_rename_established() {
let mut mb = ServiceMailbox::new();
mb.push_update(ServiceUpdate::Established); mb.push_update(ServiceUpdate::Established); mb.push_update(renamed("svc-2._ipp._tcp.local.")); mb.push_update(ServiceUpdate::Established); assert_eq!(
mb.updates.len(),
2,
"latest Renamed + the trailing post-rename Established"
);
match mb.drain() {
Drained::Update(ServiceUpdate::Renamed(r)) => {
assert!(r.new_name().as_str().contains("svc-2"))
}
other => panic!("expected Renamed(svc-2) first; got {other:?}"),
}
assert!(
matches!(mb.drain(), Drained::Update(ServiceUpdate::Established)),
"the post-rename Established must survive, after the Renamed"
);
assert!(matches!(mb.drain(), Drained::Empty));
}
#[test]
fn mailbox_bounds_non_terminal_backlog_dropping_oldest() {
let mut mb = ServiceMailbox::new();
for i in 0..(SERVICE_UPDATE_CAPACITY + 64) {
mb.push_update(renamed(&format!("svc-{i}._ipp._tcp.local.")));
}
assert_eq!(
mb.updates.len(),
1,
"rename churn coalesces to a single pending Renamed, well within the cap"
);
}
#[test]
fn mailbox_hard_cap_drops_oldest() {
let mut mb = ServiceMailbox::new();
for i in 0..(SERVICE_UPDATE_CAPACITY as u32 + 5) {
mb.bounded_push_back(renamed(&format!("svc-{i}._ipp._tcp.local.")));
}
assert_eq!(mb.updates.len(), SERVICE_UPDATE_CAPACITY);
match mb.drain() {
Drained::Update(ServiceUpdate::Renamed(r)) => {
assert!(
r.new_name().as_str().contains("svc-5"),
"oldest survivor is svc-5"
)
}
other => panic!("expected a Renamed at the head; got {other:?}"),
}
}
#[test]
fn mailbox_terminal_reserved_under_non_terminal_pressure() {
let mut mb = ServiceMailbox::new();
for i in 0..(SERVICE_UPDATE_CAPACITY + 64) {
mb.bounded_push_back(renamed(&format!("svc-{i}._ipp._tcp.local.")));
}
mb.set_terminal(ServiceUpdate::Conflict);
let mut non_terminal = 0usize;
let mut got_terminal = false;
loop {
match mb.drain() {
Drained::Update(ServiceUpdate::Conflict) => got_terminal = true,
Drained::Update(_) => non_terminal += 1,
Drained::Ended | Drained::Empty => break,
}
}
assert_eq!(non_terminal, SERVICE_UPDATE_CAPACITY);
assert!(
got_terminal,
"terminal must survive non-terminal backpressure"
);
}
#[test]
fn mailbox_set_terminal_is_idempotent_first_wins() {
let mut mb = ServiceMailbox::new();
mb.set_terminal(ServiceUpdate::Conflict);
mb.set_terminal(ServiceUpdate::HostConflict);
assert!(matches!(
mb.drain(),
Drained::Update(ServiceUpdate::Conflict)
));
assert!(matches!(mb.drain(), Drained::Ended));
}
#[test]
fn mailbox_routes_terminal_pushed_as_update_to_reserved_slot() {
let mut mb = ServiceMailbox::new();
mb.push_update(ServiceUpdate::Established);
mb.push_update(ServiceUpdate::HostConflict);
assert_eq!(mb.updates.len(), 1, "only Established is in the ring");
assert!(matches!(
mb.drain(),
Drained::Update(ServiceUpdate::Established)
));
assert!(matches!(
mb.drain(),
Drained::Update(ServiceUpdate::HostConflict)
));
assert!(matches!(mb.drain(), Drained::Ended));
}
#[test]
fn mailbox_drains_updates_then_terminal_then_ends() {
let mut mb = ServiceMailbox::new();
mb.push_update(ServiceUpdate::Established);
mb.push_update(renamed("svc-1._ipp._tcp.local."));
mb.set_terminal(ServiceUpdate::Conflict);
assert!(matches!(
mb.drain(),
Drained::Update(ServiceUpdate::Established)
));
assert!(matches!(
mb.drain(),
Drained::Update(ServiceUpdate::Renamed(_))
));
assert!(matches!(
mb.drain(),
Drained::Update(ServiceUpdate::Conflict)
));
assert!(matches!(mb.drain(), Drained::Ended));
assert!(matches!(mb.drain(), Drained::Ended));
}
#[tokio::test]
async fn doorbell_wakes_parked_consumer_for_full_batch() {
let (mailbox, doorbell_tx, doorbell_rx) = new_service_mailbox();
let mb_consumer = Arc::clone(&mailbox);
let consumer = tokio::spawn(async move {
let mut updates = 0usize;
let mut got_terminal = false;
loop {
let drained = lock(&mb_consumer).drain();
match drained {
Drained::Update(ServiceUpdate::Conflict) => got_terminal = true,
Drained::Update(_) => updates += 1,
Drained::Ended => break,
Drained::Empty => {
if doorbell_rx.recv().await.is_err() {
break;
}
}
}
}
(updates, got_terminal)
});
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
{
let mut mb = lock(&mailbox);
mb.push_update(ServiceUpdate::Established);
mb.push_update(renamed("svc-1._ipp._tcp.local."));
mb.set_terminal(ServiceUpdate::Conflict);
}
let _ = doorbell_tx.try_send(());
let (updates, got_terminal) = consumer.await.expect("consumer task panicked");
assert_eq!(updates, 2);
assert!(got_terminal);
}