use alloc::{collections::VecDeque, vec::Vec};
use core::net::{IpAddr, Ipv4Addr, SocketAddr};
use mdns_proto::{Name, ServiceRecords, ServiceSpec};
use rand::{SeedableRng, rngs::StdRng};
use smoltcp::time::Instant as RawInstant;
use super::*;
use crate::{
SmoltcpInstant,
constants::{MDNS_SOCKET_V4, MDNS_SOCKET_V6},
udpio::{RecvMeta, SendError, UdpIo},
};
#[derive(Default)]
struct MockUdp {
inbound: VecDeque<(Vec<u8>, RecvMeta)>,
sent: Vec<(SocketAddr, Vec<u8>)>,
v4_fail: Option<SendError>,
v6_fail: Option<SendError>,
capacity: Option<usize>,
}
impl UdpIo for MockUdp {
fn try_recv(&mut self, buf: &mut [u8]) -> Option<RecvMeta> {
let (data, mut meta) = self.inbound.pop_front()?;
let n = data.len().min(buf.len());
buf[..n].copy_from_slice(&data[..n]);
meta.len = n;
Some(meta)
}
fn try_send(&mut self, buf: &[u8], dst: SocketAddr) -> Result<(), SendError> {
if let Some(err) = if dst.is_ipv4() {
self.v4_fail
} else {
self.v6_fail
} {
return Err(err);
}
if let Some(slots) = self.capacity.as_mut() {
if *slots == 0 {
return Err(SendError::Busy);
}
*slots -= 1;
}
self.sent.push((dst, buf.to_vec()));
Ok(())
}
}
fn at(micros: i64) -> SmoltcpInstant {
SmoltcpInstant(RawInstant::from_micros(micros))
}
fn sample_spec() -> ServiceSpec {
let service_type = Name::try_from_str("_ipp._tcp.local.").unwrap();
let instance = Name::try_from_str("Test._ipp._tcp.local.").unwrap();
let host = Name::try_from_str("test.local.").unwrap();
let mut records = ServiceRecords::new(service_type, instance, host, 631, 120);
records.add_a(Ipv4Addr::new(192, 168, 1, 10));
ServiceSpec::new(records)
}
fn spec_for(service_type: &str, instance: &str, host: &str, addr: Ipv4Addr) -> ServiceSpec {
let mut records = ServiceRecords::new(
Name::try_from_str(service_type).unwrap(),
Name::try_from_str(instance).unwrap(),
Name::try_from_str(host).unwrap(),
631,
120,
);
records.add_a(addr);
ServiceSpec::new(records)
}
#[test]
fn registering_a_service_emits_a_probe_to_the_mdns_group() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(1));
engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp::default();
let mut scratch = [0u8; 1500];
for micros in [0, 250_000, 500_000, 1_000_000, 2_000_000] {
engine.pump(at(micros), &mut io, &mut scratch);
}
assert!(
io.sent
.iter()
.any(|(dst, _)| *dst == MDNS_SOCKET_V4 || *dst == MDNS_SOCKET_V6),
"expected at least one probe to an mDNS group; sent dsts = {:?}",
io.sent.iter().map(|(d, _)| *d).collect::<Vec<_>>()
);
}
#[test]
fn a_goodbye_with_no_socket_on_any_family_writes_off_without_error() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(101));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp::default();
let mut scratch = [0u8; 1500];
for micros in [
0, 250_000, 500_000, 750_000, 1_000_000, 1_500_000, 2_000_000, 3_000_000, 4_000_000,
] {
engine.pump(at(micros), &mut io, &mut scratch);
}
io.v4_fail = Some(SendError::Unsupported);
io.v6_fail = Some(SendError::Unsupported);
io.sent.clear();
engine.unregister_service(handle, at(5_000_000));
for micros in [5_000_000, 5_250_001, 5_500_001, 5_750_001, 6_000_001] {
engine.pump(at(micros), &mut io, &mut scratch);
}
assert!(
io.sent.is_empty(),
"nothing can leave when every family is Unsupported; sent = {:?}",
io.sent.iter().map(|(d, _)| *d).collect::<Vec<_>>()
);
}
#[cfg(feature = "stats")]
#[test]
fn stats_handle_exposes_the_shared_counter() {
let engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(102));
let s = engine.stats_handle();
assert!(Arc::strong_count(&s) >= 2);
}
fn build_ptr_query(qname: &Name) -> Vec<u8> {
use mdns_proto::wire::{Header, MessageBuilder, ResourceClass, ResourceType};
let mut buf = [0u8; 512];
let mut b: MessageBuilder<'_, 0> = MessageBuilder::try_new(&mut buf, Header::new()).unwrap();
b.push_question(qname, ResourceType::Ptr, ResourceClass::In, false)
.unwrap();
let n = b.finish().unwrap();
buf[..n].to_vec()
}
fn unicast_reply_scenario(seed: u64, v4_fail: Option<SendError>) -> Vec<(SocketAddr, Vec<u8>)> {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(seed));
engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp::default();
let mut scratch = [0u8; 1500];
for micros in [
0, 250_000, 500_000, 750_000, 1_000_000, 1_500_000, 2_000_000, 3_000_000, 4_000_000,
] {
engine.pump(at(micros), &mut io, &mut scratch);
}
io.sent.clear();
io.v4_fail = v4_fail;
let querier = SocketAddr::from((Ipv4Addr::new(192, 168, 1, 50), 6000));
io.inbound.push_back((
build_ptr_query(&Name::try_from_str("_ipp._tcp.local.").unwrap()),
RecvMeta {
src: querier,
local: Some(MDNS_SOCKET_V4.ip()),
hop_limit: None,
len: 0,
},
));
for micros in [5_000_000, 5_250_000, 5_500_000] {
engine.pump(at(micros), &mut io, &mut scratch);
}
io.sent
}
#[test]
fn a_legacy_unicast_query_gets_a_unicast_reply() {
let querier = SocketAddr::from((Ipv4Addr::new(192, 168, 1, 50), 6000));
let sent = unicast_reply_scenario(201, None);
assert!(
sent.iter().any(|(dst, _)| *dst == querier),
"expected a unicast reply to the legacy querier; sent = {:?}",
sent.iter().map(|(d, _)| *d).collect::<Vec<_>>()
);
}
#[test]
fn a_unicast_reply_too_large_is_handled_without_panicking() {
let querier = SocketAddr::from((Ipv4Addr::new(192, 168, 1, 50), 6000));
let sent = unicast_reply_scenario(202, Some(SendError::TooLarge));
assert!(sent.iter().all(|(dst, _)| *dst != querier));
}
#[test]
fn a_unicast_reply_busy_is_best_effort_not_fatal() {
let querier = SocketAddr::from((Ipv4Addr::new(192, 168, 1, 50), 6000));
let sent = unicast_reply_scenario(203, Some(SendError::Busy));
assert!(sent.iter().all(|(dst, _)| *dst != querier));
}
#[test]
fn unregistering_an_announced_service_emits_a_goodbye() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(2));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp::default();
let mut scratch = [0u8; 1500];
for micros in [
0, 250_000, 500_000, 750_000, 1_000_000, 1_500_000, 2_000_000, 3_000_000, 4_000_000,
] {
engine.pump(at(micros), &mut io, &mut scratch);
}
io.sent.clear();
engine.unregister_service(handle, at(5_000_000));
for micros in [5_000_000, 5_000_001, 5_250_001, 5_500_001, 5_750_001] {
engine.pump(at(micros), &mut io, &mut scratch);
}
assert!(
!io.sent.is_empty(),
"unregistering an announced service should emit a §10.1 goodbye burst"
);
}
fn pump_schedule() -> [i64; 10] {
[
0, 250_000, 500_000, 750_000, 1_000_000, 1_500_000, 2_000_000, 3_000_000, 4_000_000, 5_000_000,
]
}
#[test]
fn v6_only_node_advertises_via_multicast_fan_out() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(4));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp {
v4_fail: Some(SendError::Unsupported),
..Default::default()
};
let mut scratch = [0u8; 1500];
let mut established = false;
for micros in pump_schedule() {
engine.pump(at(micros), &mut io, &mut scratch);
while let Some(update) = engine.poll_service_update(handle) {
established |= matches!(update, ServiceUpdate::Established);
}
}
assert!(
established,
"a v6-only node must still reach Established via the v6 group"
);
assert!(!io.sent.is_empty(), "expected real sends to the v6 group");
assert!(
io.sent.iter().all(|(dst, _)| *dst == MDNS_SOCKET_V6),
"v6-only: every queued send must target the v6 group; got {:?}",
io.sent.iter().map(|(d, _)| *d).collect::<Vec<_>>()
);
}
#[test]
fn no_reachable_group_does_not_falsely_advance() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(5));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp {
v4_fail: Some(SendError::Unsupported),
v6_fail: Some(SendError::Unsupported),
..Default::default()
};
let mut scratch = [0u8; 1500];
let mut established = false;
for micros in pump_schedule() {
engine.pump(at(micros), &mut io, &mut scratch);
while let Some(update) = engine.poll_service_update(handle) {
established |= matches!(update, ServiceUpdate::Established);
}
}
assert!(
!established,
"a service must NOT reach Established when no datagram is ever queued"
);
assert!(
io.sent.is_empty(),
"no send should be recorded when both families are blocked"
);
}
#[test]
fn goodbye_budget_is_not_consumed_while_transport_is_busy() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(6));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp::default();
let mut scratch = [0u8; 1500];
for micros in [
0, 250_000, 500_000, 750_000, 1_000_000, 1_500_000, 2_000_000, 3_000_000, 4_000_000,
] {
engine.pump(at(micros), &mut io, &mut scratch);
}
engine.unregister_service(handle, at(5_000_000));
io.v4_fail = Some(SendError::Busy);
io.v6_fail = Some(SendError::Busy);
io.sent.clear();
for micros in [5_000_000, 5_250_001, 5_500_001, 5_750_001, 6_000_001] {
engine.pump(at(micros), &mut io, &mut scratch);
}
assert!(
io.sent.is_empty(),
"no goodbye should be recorded while busy"
);
assert!(
engine.services.contains_key(&handle),
"an all-busy withdrawal must not complete (its budget is re-armed, not spent), \
so the driver slot is still held"
);
io.v4_fail = None;
io.v6_fail = None;
engine.pump(at(6_250_001), &mut io, &mut scratch);
assert!(
io.sent.iter().any(|(_, d)| datagram_kind(d) == Some(true)),
"the TTL=0 goodbye must go out once the transport frees"
);
}
fn datagram_kind(bytes: &[u8]) -> Option<bool> {
use mdns_proto::wire::MessageReader;
let reader = MessageReader::try_parse(bytes).ok()?;
let mut saw_answer = false;
let mut saw_zero_ttl = false;
for rec in reader.answers().flatten() {
saw_answer = true;
if rec.ttl() == 0 {
saw_zero_ttl = true;
}
}
if !saw_answer {
return None;
}
Some(saw_zero_ttl)
}
#[test]
fn same_name_replacement_is_rejected_until_withdrawal_completes() {
let cfg = EndpointConfig::new().with_probe_unique_names(false);
let mut engine: Engine<SmoltcpInstant, StdRng> = Engine::new(cfg, StdRng::seed_from_u64(101));
let mut io = MockUdp::default();
let mut scratch = [0u8; 1500];
let a = engine.register_service(sample_spec(), at(0)).unwrap();
let mut established = false;
let mut t = 0i64;
for _ in 0..16 {
engine.pump(at(t), &mut io, &mut scratch);
while let Some(u) = engine.poll_service_update(a) {
established |= matches!(u, ServiceUpdate::Established);
}
t += 250_000;
}
assert!(
established,
"service A must reach Established before withdrawal"
);
engine.unregister_service(a, at(t));
let rejected = engine.register_service(sample_spec(), at(t + 1));
assert!(
matches!(
rejected,
Err(RegisterServiceError::NameAlreadyRegistered(_))
),
"a same-name registration must be rejected while the withdrawal holds the \
name; got {rejected:?}"
);
io.sent.clear();
let mut completed = false;
for _ in 0..32 {
t += 250_000;
engine.pump(at(t), &mut io, &mut scratch);
if !engine.services.contains_key(&a) {
completed = true;
break;
}
}
assert!(
completed,
"the withdrawal must complete (route freed + driver slot GC'd) on a working \
transport"
);
assert!(
io.sent.iter().any(|(_, d)| datagram_kind(d) == Some(true)),
"the withdrawal must emit a TTL=0 §10.1 goodbye; sent kinds = {:?}",
io.sent
.iter()
.map(|(_, d)| datagram_kind(d))
.collect::<Vec<_>>()
);
engine
.register_service(sample_spec(), at(t))
.expect("the same name must be re-registerable once the withdrawal completes");
}
#[test]
fn unregister_then_discard_with_unread_update_gc_s_the_slot() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(202));
let mut io = MockUdp::default();
let mut scratch = [0u8; 1500];
let a = engine.register_service(sample_spec(), at(0)).unwrap();
engine
.services
.get_mut(&a)
.unwrap()
.push_update(ServiceUpdate::Established);
engine.unregister_service(a, at(1));
let mut t = 1i64;
let mut gcd = false;
for _ in 0..4 {
t += 250_000;
engine.pump(at(t), &mut io, &mut scratch);
if !engine.services.contains_key(&a) {
gcd = true;
break;
}
}
assert!(
gcd,
"an unregistered service with an unread update must be GC'd (caller_gone), \
not deferred forever and leaked"
);
}
#[test]
fn flooded_conflict_updates_are_coalesced_and_bounded() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(7));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
let slot = engine.services.get_mut(&handle).unwrap();
for _ in 0..1000 {
slot.push_update(ServiceUpdate::HostConflict);
}
assert_eq!(
slot.updates.len(),
1,
"repeated HostConflict must coalesce to one queued update"
);
for _ in 0..1000 {
slot.push_update(ServiceUpdate::HostConflict);
slot.push_update(ServiceUpdate::Conflict);
}
assert!(
slot.updates.len() <= MAX_SERVICE_UPDATES,
"the update backlog must stay capped; got {}",
slot.updates.len()
);
}
#[test]
fn a_partial_fan_out_confirms_and_latches_goodbye_ownership() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(8));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp {
v6_fail: Some(SendError::Busy),
..Default::default()
};
let mut scratch = [0u8; 1500];
let mut established = false;
for micros in [
0, 250_000, 500_000, 750_000, 1_000_000, 1_500_000, 2_000_000, 2_500_000, 3_000_000, 4_000_000,
] {
engine.pump(at(micros), &mut io, &mut scratch);
while let Some(update) = engine.poll_service_update(handle) {
established |= matches!(update, ServiceUpdate::Established);
}
}
assert!(
established,
"a partial fan-out must confirm on the reachable family, not stall on the \
transiently-busy one"
);
assert!(
io.sent.iter().all(|(dst, _)| *dst == MDNS_SOCKET_V4),
"only v4 should carry sends while v6 is busy; got {:?}",
io.sent.iter().map(|(d, _)| *d).collect::<Vec<_>>()
);
engine.unregister_service(handle, at(4_500_000));
io.sent.clear();
engine.pump(at(4_500_001), &mut io, &mut scratch);
assert!(
io.sent
.iter()
.any(|(dst, d)| *dst == MDNS_SOCKET_V4 && datagram_kind(d) == Some(true)),
"a v4-only advertisement must still latch goodbye ownership, so the \
withdrawal emits a TTL=0 goodbye to v4; sent = {:?}",
io.sent
.iter()
.map(|(dst, d)| (*dst, datagram_kind(d)))
.collect::<Vec<_>>()
);
io.v6_fail = None;
io.sent.clear();
for micros in [4_750_001, 5_000_001] {
engine.pump(at(micros), &mut io, &mut scratch);
}
assert!(
io.sent.iter().any(|(dst, _)| *dst == MDNS_SOCKET_V6),
"the goodbye must reach v6 once it recovers"
);
}
fn build_conflict_srv_authority(instance_str: &str) -> Vec<u8> {
use mdns_proto::wire::{Header, MessageBuilder};
let mut buf = [0u8; 512];
let mut b = MessageBuilder::<'_, 32>::try_new(&mut buf, Header::new()).unwrap();
let name = Name::try_from_str(instance_str).unwrap();
let target = Name::try_from_str("rival-host.local.").unwrap();
b.push_srv_authority(&name, 120, 0, 0, 9999, &target)
.unwrap();
let n = b.finish().unwrap();
buf[..n].to_vec()
}
fn build_conflict_a_authority(host_str: &str, addr: [u8; 4]) -> Vec<u8> {
use mdns_proto::wire::{Header, MessageBuilder};
let mut buf = [0u8; 512];
let mut b = MessageBuilder::<'_, 32>::try_new(&mut buf, Header::new()).unwrap();
let name = Name::try_from_str(host_str).unwrap();
b.push_a_authority(&name, 120, Ipv4Addr::from(addr))
.unwrap();
let n = b.finish().unwrap();
buf[..n].to_vec()
}
#[test]
fn a_constrained_transport_does_not_starve_either_family() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(22));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp::default();
let mut scratch = [0u8; 1500];
let mut established = false;
let mut t = 0i64;
for _ in 0..40 {
t += 250_000;
io.capacity = Some(1);
engine.pump(at(t), &mut io, &mut scratch);
while let Some(update) = engine.poll_service_update(handle) {
established |= matches!(update, ServiceUpdate::Established);
}
}
assert!(
established,
"the service must still reach Established on a one-slot transport"
);
let hit_v4 = io.sent.iter().any(|(dst, _)| *dst == MDNS_SOCKET_V4);
let hit_v6 = io.sent.iter().any(|(dst, _)| *dst == MDNS_SOCKET_V6);
assert!(
hit_v4 && hit_v6,
"both families must receive sends on a constrained transport, not just the \
one that wins a fixed order; v4={hit_v4} v6={hit_v6}"
);
}
#[test]
fn a_constrained_transport_drains_a_withdrawal_after_each_family_gets_a_round() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(23));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp::default();
let mut scratch = [0u8; 1500];
for micros in [
0, 250_000, 500_000, 750_000, 1_000_000, 1_500_000, 2_000_000, 3_000_000, 4_000_000,
] {
engine.pump(at(micros), &mut io, &mut scratch);
}
engine.unregister_service(handle, at(5_000_000));
io.sent.clear();
let mut t = 5_000_000i64;
let mut completed = false;
for _ in 0..16 {
t += 250_000;
io.capacity = Some(1);
engine.pump(at(t), &mut io, &mut scratch);
while engine.poll_service_update(handle).is_some() {}
if !engine.services.contains_key(&handle) {
completed = true;
break;
}
}
assert!(
completed,
"the withdrawal must drain via the endpoint resend schedule on a one-slot \
transport, not linger"
);
let v4 = io.sent.iter().filter(|(d, _)| *d == MDNS_SOCKET_V4).count();
let v6 = io.sent.iter().filter(|(d, _)| *d == MDNS_SOCKET_V6).count();
assert!(
v4 >= 1 && v6 >= 1,
"each reachable family must receive at least one goodbye on a constrained \
transport; v4={v4} v6={v6}"
);
}
#[test]
fn default_setup_processes_rx_without_hop_limit_or_subnets() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(47));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp::default();
let mut scratch = [0u8; 1500];
for micros in pump_schedule() {
engine.pump(at(micros), &mut io, &mut scratch);
while engine.poll_service_update(handle).is_some() {}
}
let conflict = build_conflict_srv_authority("Test._ipp._tcp.local.");
let mut t = 6_000_000i64;
let mut reacted = false;
for _ in 0..16 {
io.inbound.push_back((
conflict.clone(),
RecvMeta {
src: SocketAddr::from((Ipv4Addr::new(192, 168, 1, 200), 5353)),
local: Some(MDNS_SOCKET_V4.ip()),
hop_limit: None,
len: 0,
},
));
engine.pump(at(t), &mut io, &mut scratch);
t += 250_000;
while let Some(u) = engine.poll_service_update(handle) {
reacted |= matches!(u, ServiceUpdate::Renamed(_) | ServiceUpdate::Conflict);
}
if reacted {
break;
}
}
assert!(
reacted,
"a default node (hop_limit None, no subnets) must PROCESS inbound mDNS — the §11 \
gate dropping everything would leave it deaf to queries, answers, and conflicts"
);
}
#[test]
fn default_setup_rejects_off_link_unicast() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(59));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp::default();
let mut scratch = [0u8; 1500];
for micros in pump_schedule() {
engine.pump(at(micros), &mut io, &mut scratch);
while engine.poll_service_update(handle).is_some() {}
}
let conflict = build_conflict_srv_authority("Test._ipp._tcp.local.");
let mut t = 6_000_000i64;
let mut reacted = false;
for _ in 0..16 {
io.inbound.push_back((
conflict.clone(),
RecvMeta {
src: SocketAddr::from((Ipv4Addr::new(192, 168, 1, 200), 5353)),
local: Some(IpAddr::V4(Ipv4Addr::new(192, 168, 1, 10))),
hop_limit: None,
len: 0,
},
));
engine.pump(at(t), &mut io, &mut scratch);
t += 250_000;
while let Some(u) = engine.poll_service_update(handle) {
reacted |= matches!(u, ServiceUpdate::Renamed(_) | ServiceUpdate::Conflict);
}
}
assert!(
!reacted,
"off-link unicast must NOT drive a conflict rename when no hop-limit or subnet \
vouches for it — only link-scoped multicast is trusted by default"
);
}
#[test]
fn proto_emitted_host_conflict_retires_and_gcs_the_smoltcp_service() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(83));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp::default();
let mut scratch = [0u8; 1500];
let mut established = false;
for micros in pump_schedule() {
engine.pump(at(micros), &mut io, &mut scratch);
while let Some(u) = engine.poll_service_update(handle) {
established |= matches!(u, ServiceUpdate::Established);
}
}
assert!(
established,
"service must reach Established before the host conflict"
);
let conflict = build_conflict_a_authority("test.local.", [10, 0, 0, 99]);
let mut t = 6_000_000i64;
let mut retired = false;
for _ in 0..16 {
io.inbound.push_back((
conflict.clone(),
RecvMeta {
src: SocketAddr::from((Ipv4Addr::new(192, 168, 1, 200), 5353)),
local: Some(MDNS_SOCKET_V4.ip()),
hop_limit: None,
len: 0,
},
));
engine.pump(at(t), &mut io, &mut scratch);
t += 250_000;
if engine
.services
.get(&handle)
.map(|s| s.errored)
.unwrap_or(false)
{
retired = true;
break;
}
}
assert!(
retired,
"a proto-emitted HostConflict must begin the endpoint-owned withdrawal (errored)"
);
let mut saw_host_conflict = false;
while let Some(u) = engine.poll_service_update(handle) {
saw_host_conflict |= u.is_host_conflict();
}
assert!(
saw_host_conflict,
"the HostConflict terminal must reach the caller via poll_service_update"
);
let mut gced = false;
for _ in 0..64 {
t += 250_000;
engine.pump(at(t), &mut io, &mut scratch);
if !engine.services.contains_key(&handle) {
gced = true;
break;
}
}
assert!(
gced,
"the withdrawn slot must be GC'd after the §10.1 goodbye completes"
);
}
#[test]
fn rx_drain_is_capped_per_pump_with_immediate_repump() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(53));
let mut io = MockUdp::default();
let mut scratch = [0u8; 1500];
let pkt = build_conflict_srv_authority("Whatever._ipp._tcp.local.");
let flood = MAX_RX_PER_PUMP + 10;
for _ in 0..flood {
io.inbound.push_back((
pkt.clone(),
RecvMeta {
src: SocketAddr::from((Ipv4Addr::new(192, 168, 1, 200), 5353)),
local: Some(MDNS_SOCKET_V4.ip()),
hop_limit: None,
len: 0,
},
));
}
let now = at(1_000_000);
let deadline = engine.pump(now, &mut io, &mut scratch);
assert_eq!(
io.inbound.len(),
flood - MAX_RX_PER_PUMP,
"one pump must drain at most MAX_RX_PER_PUMP datagrams, leaving the rest buffered"
);
assert_eq!(
deadline,
Some(now),
"a capped RX drain must request an immediate re-pump (deadline = now)"
);
engine.pump(at(1_000_001), &mut io, &mut scratch);
assert!(
io.inbound.is_empty(),
"the follow-up pump drains the remaining buffered datagrams"
);
}
#[test]
fn an_oversized_service_is_not_advertised_so_it_is_never_unwithdrawable() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(30));
let mut records = ServiceRecords::new(
Name::try_from_str("_ipp._tcp.local.").unwrap(),
Name::try_from_str("Huge._ipp._tcp.local.").unwrap(),
Name::try_from_str("huge.local.").unwrap(),
631,
120,
);
for i in 0..400u16 {
records.add_aaaa(core::net::Ipv6Addr::new(0x2001, 0xdb8, 0, 0, 0, 0, 0, i));
}
let handle = engine
.register_service(ServiceSpec::new(records), at(0))
.unwrap();
let mut io = MockUdp::default();
let mut scratch = [0u8; 12_000];
let mut established = false;
for micros in pump_schedule() {
engine.pump(at(micros), &mut io, &mut scratch);
while let Some(u) = engine.poll_service_update(handle) {
established |= matches!(u, ServiceUpdate::Established);
}
}
assert!(
!established,
"an oversized service must not reach Established (it cannot be encoded \
within the §17 ceiling, even with a larger caller scratch)"
);
io.sent.clear();
engine.unregister_service(handle, at(6_000_000));
for micros in [6_000_001, 6_250_001, 6_500_001] {
engine.pump(at(micros), &mut io, &mut scratch);
}
assert!(
io.sent.iter().all(|(_, d)| datagram_kind(d) != Some(true)),
"an oversized service that never advertised must not emit any TTL=0 goodbye; \
sent kinds = {:?}",
io.sent
.iter()
.map(|(_, d)| datagram_kind(d))
.collect::<Vec<_>>()
);
}
#[test]
fn permanently_failing_family_does_not_stall_the_healthy_one() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(15));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp {
v6_fail: Some(SendError::Busy),
..Default::default()
};
let mut scratch = [0u8; 1500];
let mut established = false;
let mut t = 0;
for _ in 0..80 {
t += 250_000;
engine.pump(at(t), &mut io, &mut scratch);
while let Some(update) = engine.poll_service_update(handle) {
established |= matches!(update, ServiceUpdate::Established);
}
}
assert!(
established,
"a healthy v4 family must reach Established despite a permanently-failing v6"
);
assert!(
io.sent.iter().all(|(dst, _)| *dst == MDNS_SOCKET_V4),
"only v4 should carry real sends; got {:?}",
io.sent.iter().map(|(d, _)| *d).collect::<Vec<_>>()
);
}
#[test]
fn own_multicast_loopback_is_not_treated_as_conflict() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(9));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp::default();
let mut scratch = [0u8; 1500];
for micros in pump_schedule() {
engine.pump(at(micros), &mut io, &mut scratch);
}
let (_, datagram) = io.sent.last().cloned().expect("a datagram was sent");
io.inbound.push_back((
datagram,
RecvMeta {
src: SocketAddr::from((Ipv4Addr::new(10, 0, 0, 99), 5353)),
local: None,
hop_limit: Some(255),
len: 0,
},
));
engine.pump(at(5_000_001), &mut io, &mut scratch);
let mut conflict = false;
while let Some(update) = engine.poll_service_update(handle) {
conflict |= matches!(
update,
ServiceUpdate::Conflict | ServiceUpdate::HostConflict
);
}
assert!(
!conflict,
"our own looped-back multicast must not be seen as a conflicting peer"
);
}
#[test]
fn actionable_updates_survive_conflict_flood() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(10));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
let slot = engine.services.get_mut(&handle).unwrap();
slot.push_update(ServiceUpdate::Established);
for _ in 0..1000 {
slot.push_update(ServiceUpdate::HostConflict);
slot.push_update(ServiceUpdate::Conflict);
}
assert!(
slot
.updates
.iter()
.any(|u| matches!(u, ServiceUpdate::Established)),
"the Established transition must not be evicted by conflict noise"
);
assert!(
slot.updates.len() <= MAX_SERVICE_UPDATES,
"the backlog must stay bounded; got {}",
slot.updates.len()
);
}
#[test]
fn busy_goodbye_is_held_then_force_completed_at_the_ceiling() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(11));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp::default();
let mut scratch = [0u8; 1500];
for micros in [
0, 250_000, 500_000, 750_000, 1_000_000, 1_500_000, 2_000_000, 3_000_000, 4_000_000,
] {
engine.pump(at(micros), &mut io, &mut scratch);
}
while engine.poll_service_update(handle).is_some() {}
engine.unregister_service(handle, at(5_000_000));
io.v4_fail = Some(SendError::Busy);
io.v6_fail = Some(SendError::Busy);
for micros in [5_250_001, 5_500_001, 6_000_001, 6_500_001] {
engine.pump(at(micros), &mut io, &mut scratch);
while engine.poll_service_update(handle).is_some() {}
}
assert!(
engine.services.contains_key(&handle),
"a never-delivered withdrawal must be HELD (route reserved + slot present) \
within the 2 s anti-pin ceiling"
);
engine.pump(at(7_500_001), &mut io, &mut scratch);
assert!(
!engine.services.contains_key(&handle),
"an undeliverable withdrawal must be force-completed at its anti-pin ceiling"
);
}
#[test]
fn loopback_detected_across_a_large_send_burst() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(14));
let mut handles = Vec::new();
for i in 0..8u8 {
let instance = alloc::format!("Dev{i}._ipp._tcp.local.");
let host = alloc::format!("dev{i}.local.");
handles.push(
engine
.register_service(
spec_for(
"_ipp._tcp.local.",
&instance,
&host,
Ipv4Addr::new(192, 168, 1, 10 + i),
),
at(0),
)
.unwrap(),
);
}
let mut io = MockUdp::default();
let mut scratch = [0u8; 1500];
for micros in [0, 250_000, 500_000] {
engine.pump(at(micros), &mut io, &mut scratch);
}
assert!(
io.sent.len() > 4,
"expected a burst larger than any small fixed ring; got {}",
io.sent.len()
);
let first = io.sent.first().cloned().expect("a probe was sent");
io.inbound.push_back((
first.1,
RecvMeta {
src: SocketAddr::from((Ipv4Addr::new(10, 0, 0, 99), 5353)),
local: None,
hop_limit: Some(255),
len: 0,
},
));
engine.pump(at(750_000), &mut io, &mut scratch);
let mut conflict = false;
for h in &handles {
while let Some(u) = engine.poll_service_update(*h) {
conflict |= matches!(u, ServiceUpdate::Conflict | ServiceUpdate::HostConflict);
}
}
assert!(
!conflict,
"the oldest self-send in a large burst must still be loopback-detected"
);
}
#[test]
fn send_multicast_confirms_when_any_family_queues() {
let mut tx = Multicaster::<SmoltcpInstant>::new();
let mut partial = MockUdp {
v6_fail: Some(SendError::Busy),
..Default::default()
};
let (outcome, fanout) = tx.send_multicast(&mut partial, b"a-multicast-datagram", at(0));
assert!(
matches!(outcome, MulticastOutcome::Delivered),
"v4 queued + v6 transiently busy must confirm (>= 1 socket succeeded)"
);
assert_eq!(
fanout.sent_count(),
1,
"v4 queued, v6 busy: exactly 1 datagram on the wire"
);
assert!(
matches!(fanout.v4, FamilySend::Sent(_)),
"v4 must have sent"
);
assert!(matches!(fanout.v6, FamilySend::Busy), "v6 must be Busy");
let mut all_busy = MockUdp {
v4_fail: Some(SendError::Busy),
v6_fail: Some(SendError::Busy),
..Default::default()
};
let (outcome_busy, fanout_busy) =
tx.send_multicast(&mut all_busy, b"a-multicast-datagram", at(0));
assert!(
matches!(outcome_busy, MulticastOutcome::Retry),
"both families busy: nothing on the link, so retry rather than confirm or retire"
);
assert_eq!(
fanout_busy.sent_count(),
0,
"both families busy: no datagrams on the wire"
);
}
#[test]
fn a_permanently_too_large_send_retires_the_service() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(31));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp {
v4_fail: Some(SendError::TooLarge),
v6_fail: Some(SendError::TooLarge),
..Default::default()
};
let mut scratch = [0u8; 1500];
let mut conflict = false;
let mut established = false;
for micros in pump_schedule() {
engine.pump(at(micros), &mut io, &mut scratch);
while let Some(u) = engine.poll_service_update(handle) {
conflict |= matches!(u, ServiceUpdate::Conflict | ServiceUpdate::HostConflict);
established |= matches!(u, ServiceUpdate::Established);
}
}
assert!(
conflict,
"a permanently-too-large send must retire the service with an actionable update"
);
assert!(
!established,
"a service whose datagrams can never be sent must not reach Established"
);
assert!(
io.sent.is_empty(),
"nothing is ever queued when every send is permanently too large"
);
}
#[test]
fn a_too_large_family_does_not_retire_while_the_other_may_recover() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(33));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp {
v4_fail: Some(SendError::TooLarge), v6_fail: Some(SendError::Busy), ..Default::default()
};
let mut scratch = [0u8; 1500];
let mut conflict = false;
let mut established = false;
let mut t = 0i64;
for _ in 0..40 {
t += 250_000;
engine.pump(at(t), &mut io, &mut scratch);
while let Some(u) = engine.poll_service_update(handle) {
conflict |= matches!(u, ServiceUpdate::Conflict | ServiceUpdate::HostConflict);
established |= matches!(u, ServiceUpdate::Established);
}
}
assert!(
!conflict,
"a TooLarge family must not retire the service while the other (Busy) may \
still recover"
);
assert!(
!established,
"cannot advertise while v6 is busy and v4 is permanently too large"
);
io.v6_fail = None;
for ms in 41..=64i64 {
engine.pump(at(ms * 250_000), &mut io, &mut scratch);
while let Some(u) = engine.poll_service_update(handle) {
established |= matches!(u, ServiceUpdate::Established);
}
}
assert!(
established,
"once v6 recovers the service advertises on it — it was never retired"
);
}
#[test]
fn established_is_observable_on_the_pump_that_confirms_it() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(32));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp::default();
let mut scratch = [0u8; 1500];
let mut established = false;
let mut settled = false;
let mut t = 0i64;
for _ in 0..40 {
t += 250_000;
let deadline = engine.pump(at(t), &mut io, &mut scratch);
while let Some(u) = engine.poll_service_update(handle) {
established |= matches!(u, ServiceUpdate::Established);
}
if deadline.is_some_and(|d| d >= at(t + 30_000_000)) {
settled = true;
break;
}
}
assert!(
settled,
"the service should have reached its re-announce deadline"
);
assert!(
established,
"Established must be surfaced on the pump that confirms the final \
announcement, not stranded until the distant re-announce"
);
}
#[test]
fn a_query_exposes_collected_answers_via_the_public_api() {
let mut responder: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(40));
responder.register_service(sample_spec(), at(0)).unwrap();
let mut rio = MockUdp::default();
let mut scratch = [0u8; 1500];
for micros in pump_schedule() {
responder.pump(at(micros), &mut rio, &mut scratch);
}
let (_, announcement) = rio
.sent
.iter()
.rev()
.find(|(dst, _)| *dst == MDNS_SOCKET_V4)
.cloned()
.expect("the responder must have multicast an announcement");
let mut querier: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(41));
let q = querier
.start_query(
QuerySpec::new(
Name::try_from_str("_ipp._tcp.local.").unwrap(),
mdns_proto::wire::ResourceType::Ptr,
),
at(0),
)
.unwrap();
let mut qio = MockUdp::default();
qio.inbound.push_back((
announcement,
RecvMeta {
src: SocketAddr::from((Ipv4Addr::new(10, 0, 0, 5), 5353)),
local: None,
hop_limit: Some(255),
len: 0,
},
));
for micros in pump_schedule() {
querier.pump(at(micros), &mut qio, &mut scratch);
}
let answers = querier.collected_answers(q).count();
assert!(
answers >= 1,
"a query's collected answers must be readable via the public API; got {answers}"
);
assert!(
querier.query_accepted_count(q).unwrap_or(0) >= 1,
"query_accepted_count must reflect the accepted answer"
);
}
#[test]
fn a_query_that_can_never_send_surfaces_a_terminal_update() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(42));
let q = engine
.start_query(
QuerySpec::new(
Name::try_from_str("_ipp._tcp.local.").unwrap(),
mdns_proto::wire::ResourceType::Ptr,
),
at(0),
)
.unwrap();
let mut io = MockUdp {
v4_fail: Some(SendError::TooLarge),
v6_fail: Some(SendError::TooLarge),
..Default::default()
};
let mut scratch = [0u8; 1500];
let mut terminal = false;
for micros in pump_schedule() {
engine.pump(at(micros), &mut io, &mut scratch);
while let Some(u) = engine.poll_query_update(q) {
terminal |= matches!(u, QueryUpdate::Timeout | QueryUpdate::Done);
}
}
assert!(
terminal,
"a query that can never send must surface a terminal update, not hang silently"
);
}
#[test]
fn a_retired_query_freezes_answers_and_emits_no_second_terminal() {
let mut responder: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(43));
responder.register_service(sample_spec(), at(0)).unwrap();
let mut rio = MockUdp::default();
let mut scratch = [0u8; 1500];
for micros in pump_schedule() {
responder.pump(at(micros), &mut rio, &mut scratch);
}
let (_, announcement) = rio
.sent
.iter()
.rev()
.find(|(d, _)| *d == MDNS_SOCKET_V4)
.cloned()
.expect("the responder must have announced");
let mut querier: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(44));
let q = querier
.start_query(
QuerySpec::new(
Name::try_from_str("_ipp._tcp.local.").unwrap(),
mdns_proto::wire::ResourceType::Ptr,
),
at(0),
)
.unwrap();
let mut qio = MockUdp {
v4_fail: Some(SendError::TooLarge),
v6_fail: Some(SendError::TooLarge),
..Default::default()
};
let mut terminals = 0;
for micros in pump_schedule() {
querier.pump(at(micros), &mut qio, &mut scratch);
while let Some(u) = querier.poll_query_update(q) {
if matches!(u, QueryUpdate::Timeout | QueryUpdate::Done) {
terminals += 1;
}
}
}
assert_eq!(
terminals, 1,
"a retired query surfaces exactly one terminal"
);
assert_eq!(
querier.collected_answers(q).count(),
0,
"a retired query collected nothing (it never sent)"
);
qio.inbound.push_back((
announcement,
RecvMeta {
src: SocketAddr::from((Ipv4Addr::new(10, 0, 0, 7), 5353)),
local: None,
hop_limit: Some(255),
len: 0,
},
));
let mut t = 100_000_000i64;
for _ in 0..10 {
t += 250_000;
querier.pump(at(t), &mut qio, &mut scratch);
while let Some(u) = querier.poll_query_update(q) {
if matches!(u, QueryUpdate::Timeout | QueryUpdate::Done) {
terminals += 1;
}
}
}
assert_eq!(
terminals, 1,
"no SECOND terminal after a late response to a retired query"
);
assert_eq!(
querier.collected_answers(q).count(),
0,
"a late response must be frozen — collected_answers unchanged after the terminal"
);
}
#[cfg(feature = "stats")]
#[test]
fn stats_withdrawal_dual_stack_counts_rounds_and_per_family_datagrams() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(1005));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp::default();
let mut scratch = [0u8; 1500];
for micros in pump_schedule() {
engine.pump(at(micros), &mut io, &mut scratch);
}
engine.unregister_service(handle, at(5_000_000));
let snap_before = engine.stats();
io.sent.clear();
let mut t = 5_000_000i64;
let mut completed = false;
for _ in 0..16 {
t += 250_000;
engine.pump(at(t), &mut io, &mut scratch);
if engine.stats().services_active == 0 {
completed = true;
break;
}
}
assert!(completed, "the withdrawal must drain on dual-stack");
let snap_after = engine.stats();
let v4 = io.sent.iter().filter(|(d, _)| *d == MDNS_SOCKET_V4).count();
let v6 = io.sent.iter().filter(|(d, _)| *d == MDNS_SOCKET_V6).count();
assert!(
v4 >= 1 && v6 >= 1,
"both families must carry goodbyes; v4={v4} v6={v6}"
);
assert_eq!(
v4, v6,
"dual-stack: each round fans to both families equally"
);
let rounds = v4 as u64;
assert_eq!(
snap_after.goodbyes_tx - snap_before.goodbyes_tx,
rounds,
"goodbyes_tx must count one per delivered round (== {rounds})"
);
assert_eq!(
snap_after.packets_tx - snap_before.packets_tx,
(v4 + v6) as u64,
"packets_tx delta must equal per-family goodbye datagrams"
);
assert_eq!(
snap_after.send_errors - snap_before.send_errors,
0,
"dual-stack healthy: send_errors must be 0"
);
}
#[cfg(feature = "stats")]
#[test]
fn stats_withdrawal_v6_busy_until_recovery_not_freed_before_v6_sends() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(2006));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp::default();
let mut scratch = [0u8; 1500];
for micros in pump_schedule() {
engine.pump(at(micros), &mut io, &mut scratch);
}
while engine.poll_service_update(handle).is_some() {}
engine.unregister_service(handle, at(5_000_000)); io.sent.clear();
io.v6_fail = Some(SendError::Busy);
for micros in [5_250_001, 5_500_001, 5_750_001, 6_000_001] {
engine.pump(at(micros), &mut io, &mut scratch);
while engine.poll_service_update(handle).is_some() {}
}
assert!(
engine.services.contains_key(&handle),
"a withdrawal whose v6 family is still busy must NOT be freed before the \
2 s ceiling — v6 peers still hold the records"
);
let v6_before = io.sent.iter().filter(|(d, _)| *d == MDNS_SOCKET_V6).count();
assert_eq!(
v6_before, 0,
"no v6 goodbye can have reached the wire while v6 was busy; got {v6_before}"
);
assert!(
io.sent.iter().any(|(d, _)| *d == MDNS_SOCKET_V4),
"v4 must have emitted its TTL=0 goodbyes while v6 was busy"
);
io.v6_fail = None;
let mut completed = false;
for micros in [6_250_001, 6_500_001, 6_750_001, 6_900_001] {
engine.pump(at(micros), &mut io, &mut scratch);
while engine.poll_service_update(handle).is_some() {}
if !engine.services.contains_key(&handle) {
completed = true;
break;
}
}
assert!(
completed,
"once v6 recovers and sends its goodbyes the withdrawal completes (before \
the 2 s ceiling)"
);
let v6_after = io.sent.iter().filter(|(d, _)| *d == MDNS_SOCKET_V6).count();
assert!(
v6_after >= 1,
"v6 must have emitted at least one TTL=0 goodbye after recovery; got {v6_after}"
);
}
#[cfg(feature = "stats")]
#[test]
fn stats_goodbye_both_families_failed_no_goodbyes_tx() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(1004));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp::default();
let mut scratch = [0u8; 1500];
for micros in pump_schedule() {
engine.pump(at(micros), &mut io, &mut scratch);
}
engine.unregister_service(handle, at(5_500_000));
io.v4_fail = Some(SendError::TooLarge);
io.v6_fail = Some(SendError::TooLarge);
let snap_before = engine.stats();
io.sent.clear();
engine.pump(at(6_500_000), &mut io, &mut scratch);
let snap_after = engine.stats();
assert_eq!(
io.sent.len(),
0,
"no datagrams should be sent when both families fail"
);
assert_eq!(
snap_after.goodbyes_tx - snap_before.goodbyes_tx,
0,
"goodbyes_tx must be 0 when nothing ever goes on the wire; delta={}",
snap_after.goodbyes_tx - snap_before.goodbyes_tx
);
let errors_delta = snap_after.send_errors - snap_before.send_errors;
assert!(
errors_delta >= 2,
"both families TooLarge must bump send_errors at least once each; delta={errors_delta}"
);
}
#[cfg(feature = "stats")]
#[test]
fn stats_multicast_tx_partial_failure_counted_per_family() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(1006));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp {
v6_fail: Some(SendError::TooLarge),
..Default::default()
};
let mut scratch = [0u8; 1500];
let snap_before = engine.stats();
for micros in [0, 250_000, 500_000, 750_000, 1_000_000] {
engine.pump(at(micros), &mut io, &mut scratch);
}
let _ = handle;
let snap_after = engine.stats();
let v4_sent = io.sent.iter().filter(|(d, _)| *d == MDNS_SOCKET_V4).count();
assert!(
snap_after.packets_tx > snap_before.packets_tx,
"v4 probes must increment packets_tx"
);
assert_eq!(
snap_after.packets_tx - snap_before.packets_tx,
v4_sent as u64,
"packets_tx delta must equal v4 sends only; delta={}, v4_sent={v4_sent}",
snap_after.packets_tx - snap_before.packets_tx
);
assert_eq!(
snap_after.send_errors - snap_before.send_errors,
v4_sent as u64,
"send_errors delta must equal v4_sent (one v6-TooLarge per fan-out); \
errors_delta={}, v4_sent={v4_sent}",
snap_after.send_errors - snap_before.send_errors
);
}
#[cfg(feature = "stats")]
#[test]
fn stats_multicast_sent_plus_failed_send_errors_exact() {
let mut tx: Multicaster<SmoltcpInstant> = Multicaster::new();
let mut io = MockUdp {
v6_fail: Some(SendError::TooLarge),
..Default::default()
};
let data = b"probe-datagram";
let (outcome, fanout) = tx.send_multicast(&mut io, data, at(0));
assert!(
matches!(outcome, MulticastOutcome::Delivered),
"v4 Sent + v6 TooLarge must yield Delivered"
);
assert_eq!(
fanout.failed_count(),
1,
"exactly one family (v6) must be Failed; failed_count={}",
fanout.failed_count()
);
assert_eq!(
fanout.sent_count(),
1,
"exactly one family (v4) must be Sent; sent_count={}",
fanout.sent_count()
);
assert_eq!(
fanout.failed_count(),
1,
"send_errors delta must be 1 (v6 failure must not be dropped by Delivered arm)"
);
}
#[cfg(feature = "stats")]
#[test]
fn stats_multicast_failed_plus_busy_send_errors_exact() {
let mut tx: Multicaster<SmoltcpInstant> = Multicaster::new();
let mut io = MockUdp {
v4_fail: Some(SendError::TooLarge),
v6_fail: Some(SendError::Busy),
..Default::default()
};
let data = b"probe-datagram";
let (outcome, fanout) = tx.send_multicast(&mut io, data, at(0));
assert!(
matches!(outcome, MulticastOutcome::Retry),
"v4 Failed + v6 Busy must yield Retry (v6 Busy keeps things alive)"
);
assert_eq!(
fanout.failed_count(),
1,
"only v4 is Failed; failed_count must be 1, got {}",
fanout.failed_count()
);
assert!(
!matches!(fanout.v6, FamilySend::Failed),
"v6 Busy must not be mapped to Failed"
);
assert_eq!(
fanout.failed_count(),
1,
"send_errors delta must be 1 (Failed only), never 2 (Busy must not count)"
);
}
#[cfg(feature = "stats")]
#[test]
fn stats_unicast_busy_does_not_increment_send_errors() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(2001));
let _handle = engine.register_service(sample_spec(), at(0)).unwrap();
let mut io = MockUdp {
capacity: Some(0),
..Default::default()
};
let mut scratch = [0u8; 1500];
let snap_before = engine.stats();
engine.pump(at(0), &mut io, &mut scratch);
let snap_after = engine.stats();
assert_eq!(
snap_after.send_errors - snap_before.send_errors,
0,
"Busy (capacity=0) must not increment send_errors; delta={}",
snap_after.send_errors - snap_before.send_errors
);
}
#[cfg(feature = "stats")]
#[test]
fn stats_unicast_too_large_increments_send_errors() {
let mut io = MockUdp {
v4_fail: Some(SendError::TooLarge),
..Default::default()
};
let unicast_dst: SocketAddr = "192.168.1.100:5353".parse().unwrap();
let result = io.try_send(b"unicast-reply", unicast_dst);
assert!(
matches!(result, Err(SendError::TooLarge)),
"MockUdp with v4_fail=TooLarge must return TooLarge for IPv4 unicast"
);
let errors: u64 = match result {
Ok(()) => 0,
Err(SendError::TooLarge) => 1,
Err(SendError::Busy) | Err(SendError::Unsupported) => 0,
};
assert_eq!(
errors, 1,
"TooLarge unicast must count as send_errors=1; got {errors}"
);
}
#[cfg(feature = "stats")]
#[test]
fn stats_unicast_unsupported_does_not_increment_send_errors() {
let mut io = MockUdp {
v4_fail: Some(SendError::Unsupported),
..Default::default()
};
let unicast_dst: SocketAddr = "192.168.1.100:5353".parse().unwrap();
let result = io.try_send(b"unicast-reply", unicast_dst);
assert!(
matches!(result, Err(SendError::Unsupported)),
"MockUdp with v4_fail=Unsupported must return Unsupported for IPv4 unicast"
);
let errors: u64 = match result {
Ok(()) => 0,
Err(SendError::TooLarge) => 1,
Err(SendError::Busy) | Err(SendError::Unsupported) => 0,
};
assert_eq!(
errors, 0,
"Unsupported unicast must not count as send_errors; got {errors}"
);
}
#[cfg(feature = "stats")]
#[test]
fn stats_off_link_datagram_counts_rx_bytes_and_dropped() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(9001));
let mut io = MockUdp::default();
let mut scratch = [0u8; 1500];
let pkt = build_conflict_srv_authority("Test._ipp._tcp.local.");
let pkt_len = pkt.len();
io.inbound.push_back((
pkt,
RecvMeta {
src: SocketAddr::from((Ipv4Addr::new(192, 168, 2, 1), 5353)),
local: Some(IpAddr::V4(Ipv4Addr::new(224, 0, 0, 251))),
hop_limit: Some(1),
len: pkt_len,
},
));
let snap_before = engine.stats();
engine.pump(at(0), &mut io, &mut scratch);
let snap_after = engine.stats();
assert_eq!(
snap_after.packets_rx - snap_before.packets_rx,
1,
"an off-link datagram WAS received → packets_rx must rise by 1"
);
assert_eq!(
snap_after.bytes_rx - snap_before.bytes_rx,
pkt_len as u64,
"off-link datagram bytes_rx must rise by the datagram length"
);
assert_eq!(
snap_after.packets_dropped - snap_before.packets_dropped,
1,
"an off-link datagram must increment packets_dropped by 1"
);
}
#[cfg(feature = "stats")]
#[test]
fn stats_oversized_zero_len_marker_counts_rx_and_dropped() {
use std::net::{IpAddr, Ipv4Addr, SocketAddr};
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(42));
let mut io = MockUdp::default();
let mut scratch = [0u8; 1500];
io.inbound.push_back((
vec![],
RecvMeta {
src: SocketAddr::from((Ipv4Addr::new(192, 168, 1, 5), 5353)),
local: Some(IpAddr::V4(Ipv4Addr::new(224, 0, 0, 251))),
hop_limit: Some(255),
len: 0,
},
));
let snap_before = engine.stats();
engine.pump(at(0), &mut io, &mut scratch);
let snap_after = engine.stats();
assert_eq!(
snap_after.packets_rx - snap_before.packets_rx,
1,
"a zero-length (oversized) marker WAS consumed → packets_rx must rise by 1"
);
assert_eq!(
snap_after.packets_dropped - snap_before.packets_dropped,
1,
"a zero-length marker is an unusable datagram → packets_dropped must rise by 1"
);
assert_eq!(
snap_after.bytes_rx, snap_before.bytes_rx,
"bytes_rx must not change (oversized payload is lost before the zero-len marker)"
);
}
#[cfg(feature = "stats")]
#[test]
fn encode_failure_retirement_frees_proto_route_and_decrements_services_active() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(99));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
assert_eq!(
engine.stats().services_active,
1,
"services_active must be 1 after registration"
);
let mut io = MockUdp::default();
let mut scratch_tiny = [0u8; 1];
let mut got_conflict = false;
for micros in [0i64, 100_000, 200_000, 300_000, 400_000] {
engine.pump(at(micros), &mut io, &mut scratch_tiny);
while let Some(u) = engine.poll_service_update(handle) {
got_conflict |= matches!(u, ServiceUpdate::Conflict | ServiceUpdate::HostConflict);
}
if got_conflict {
break;
}
}
assert!(
got_conflict,
"encode failure must surface Conflict to the caller (poll_service_update)"
);
assert_eq!(
engine.stats().services_active,
0,
"services_active must be 0 after encode-failure retirement (proto route freed)"
);
engine
.register_service(sample_spec(), at(500_000))
.expect("same service name must be re-registerable after encode-failure retirement");
assert_eq!(
engine.stats().services_active,
1,
"services_active must be 1 again after re-registration"
);
}
#[cfg(feature = "stats")]
#[test]
fn multi_service_encode_failure_frees_route_even_with_sibling_transmit() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(200));
let handle_a = engine.register_service(sample_spec(), at(0)).unwrap();
let handle_b = engine
.register_service(
spec_for(
"_ipp._tcp.local.",
"Sibling._ipp._tcp.local.",
"sibling.local.",
Ipv4Addr::new(192, 168, 1, 11),
),
at(0),
)
.unwrap();
assert_eq!(
engine.stats().services_active,
2,
"both services registered: services_active must be 2"
);
let mut io = MockUdp::default();
let mut tiny = [0u8; 1];
let mut got_conflict_a = false;
let mut got_conflict_b = false;
for i in 0..30i64 {
let t = at(i * 100_000);
engine.pump(t, &mut io, &mut tiny);
while let Some(u) = engine.poll_service_update(handle_a) {
if matches!(u, ServiceUpdate::Conflict | ServiceUpdate::HostConflict) {
got_conflict_a = true;
}
}
while let Some(u) = engine.poll_service_update(handle_b) {
if matches!(u, ServiceUpdate::Conflict | ServiceUpdate::HostConflict) {
got_conflict_b = true;
}
}
if got_conflict_a && got_conflict_b {
break;
}
}
assert!(
got_conflict_a,
"A's Conflict must be surfaced via poll_service_update"
);
assert!(
got_conflict_b,
"B's Conflict must be surfaced via poll_service_update"
);
assert_eq!(
engine.stats().services_active,
0,
"services_active must be 0 after both services are retired by encode failure \
(each begins + completes an empty withdrawal; no route leak)"
);
engine
.register_service(sample_spec(), at(3_000_000))
.expect("A's name must be re-registerable after in-iteration unregister (fix)");
engine
.register_service(
spec_for(
"_ipp._tcp.local.",
"Sibling._ipp._tcp.local.",
"sibling.local.",
Ipv4Addr::new(192, 168, 1, 11),
),
at(3_000_000),
)
.expect("B's name must be re-registerable after in-iteration unregister (fix)");
assert_eq!(
engine.stats().services_active,
2,
"services_active must be 2 after re-registering both A and B"
);
}
#[cfg(feature = "stats")]
#[test]
fn send_too_large_retirement_frees_proto_route_and_decrements_services_active() {
let mut engine: Engine<SmoltcpInstant, StdRng> =
Engine::new(EndpointConfig::new(), StdRng::seed_from_u64(100));
let handle = engine.register_service(sample_spec(), at(0)).unwrap();
assert_eq!(
engine.stats().services_active,
1,
"services_active must be 1 after registration"
);
let mut io = MockUdp {
v4_fail: Some(SendError::TooLarge),
v6_fail: Some(SendError::TooLarge),
..Default::default()
};
let mut scratch = [0u8; 1500];
let mut got_conflict = false;
for micros in pump_schedule() {
engine.pump(at(micros), &mut io, &mut scratch);
while let Some(u) = engine.poll_service_update(handle) {
got_conflict |= matches!(u, ServiceUpdate::Conflict | ServiceUpdate::HostConflict);
}
if got_conflict {
break;
}
}
assert!(
got_conflict,
"permanently-too-large sends must surface Conflict (retire_origin path)"
);
assert_eq!(
engine.stats().services_active,
0,
"services_active must be 0 after retire_origin (proto route freed)"
);
engine
.register_service(sample_spec(), at(10_000_000))
.expect("same service name must be re-registerable after retire_origin");
assert_eq!(
engine.stats().services_active,
1,
"services_active must be 1 again after re-registration"
);
}