use event_listener::Event;
use std::{
collections::HashMap,
rc::Rc,
time::{Duration, Instant as StdInstant, SystemTime},
};
use core::{
cell::RefCell,
net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr, SocketAddrV4, SocketAddrV6},
};
use hick_trace::*;
use mdns_proto::{
CacheEntry, CollectedAnswer, Endpoint as ProtoEp, EndpointConfig, EndpointEventEntry,
QueryHandle, QueryUpdate, ServiceHandle, ServiceRoute, ServiceUpdate, WithdrawalSend,
WithdrawalToken, query::Query as ProtoQuery, service::Service as ProtoSvc, transmit::Transmit,
};
use rand::{SeedableRng, rngs::StdRng};
use slab::Slab;
use crate::{
service::ServiceMailbox,
socket::{RecvMeta, Socket},
};
#[cfg(test)]
mod tests;
pub(crate) const MAX_TRANSMIT_CREDITS_PER_PASS: usize = 64;
pub(crate) const MDNS_V4_DST: SocketAddr =
SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::new(224, 0, 0, 251), 5353));
pub(crate) const MDNS_V6_DST: SocketAddr = SocketAddr::V6(SocketAddrV6::new(
Ipv6Addr::new(0xff02, 0, 0, 0, 0, 0, 0, 0x00fb),
5353,
0,
0,
));
pub(crate) fn is_mdns_multicast_dst(dst: SocketAddr) -> bool {
use hick_udp::constants::MDNS_PORT;
matches!(dst, SocketAddr::V4(a) if a.ip().is_multicast() && a.port() == MDNS_PORT)
|| matches!(dst, SocketAddr::V6(a) if a.ip().is_multicast() && a.port() == MDNS_PORT)
}
pub(crate) type ProtoEndpoint = ProtoEp<
StdInstant,
StdRng,
Slab<CacheEntry<StdInstant>>,
Slab<ServiceRoute>,
Slab<ProtoQuery<StdInstant, Slab<CollectedAnswer>, Slab<QueryUpdate>>>,
Slab<EndpointEventEntry>,
Slab<CollectedAnswer>,
Slab<QueryUpdate>,
>;
pub(crate) type ProtoService = ProtoSvc<StdInstant, Slab<Transmit>, Slab<ServiceUpdate>>;
#[derive(Clone)]
pub(crate) struct LocalNotify {
inner: Rc<Event>,
}
impl LocalNotify {
pub(crate) fn new() -> Self {
Self {
inner: Rc::new(Event::new()),
}
}
pub(crate) fn notify(&self) {
self.inner.notify(usize::MAX);
}
pub(crate) async fn listen(&self) {
self.inner.listen().await;
}
}
#[derive(Clone, Copy)]
pub(crate) enum TransmitOrigin {
Service(ServiceHandle),
Query(QueryHandle),
}
pub(crate) const MAX_CONSECUTIVE_ENCODE_ERRORS: u8 = 3;
pub(crate) struct ServiceCtx {
pub(crate) proto: ProtoService,
pub(crate) mailbox: Rc<RefCell<ServiceMailbox>>,
pub(crate) cancelled: bool,
pub(crate) encode_failures: u8,
pub(crate) errored: bool,
}
pub(crate) struct QueryCtx {
pub(crate) last_seq: u64,
pub(crate) cancelled: bool,
pub(crate) errored: bool,
pub(crate) terminal_wake_pending: bool,
}
pub(crate) struct State {
pub(crate) endpoint: ProtoEndpoint,
pub(crate) services: HashMap<ServiceHandle, ServiceCtx>,
pub(crate) queries: HashMap<QueryHandle, QueryCtx>,
pub(crate) recent_sends: crate::selfsend::SelfSends,
pub(crate) completed_withdrawals: Vec<ServiceHandle>,
pub(crate) svc_handle_scratch: Vec<ServiceHandle>,
pub(crate) query_handle_scratch: Vec<QueryHandle>,
pub(crate) bound_interface: u32,
pub(crate) local_subnets: Vec<(IpAddr, u8)>,
pub(crate) max_payload: usize,
pub(crate) max_recv: usize,
#[cfg(feature = "stats")]
pub(crate) stats: std::sync::Arc<stats::Stats>,
}
impl State {
pub(crate) fn new(cfg: EndpointConfig, max_payload: usize, max_recv: usize) -> Self {
let rng = StdRng::from_rng(&mut rand::rng());
let endpoint = ProtoEndpoint::try_new(cfg, rng);
#[cfg(feature = "stats")]
let stats = endpoint.stats_handle();
Self {
endpoint,
services: HashMap::new(),
queries: HashMap::new(),
recent_sends: Vec::new(),
completed_withdrawals: Vec::new(),
svc_handle_scratch: Vec::new(),
query_handle_scratch: Vec::new(),
bound_interface: 0,
local_subnets: Vec::new(),
max_payload,
max_recv,
#[cfg(feature = "stats")]
stats,
}
}
pub(crate) fn register_service(
&mut self,
spec: mdns_proto::ServiceSpec,
now: StdInstant,
mailbox: Rc<RefCell<ServiceMailbox>>,
) -> Result<ServiceHandle, mdns_proto::error::RegisterServiceError> {
let (handle, svc) = self
.endpoint
.try_register_service::<Slab<_>, Slab<_>>(spec, now)?;
self.services.insert(
handle,
ServiceCtx {
proto: svc,
mailbox,
cancelled: false,
encode_failures: 0,
errored: false,
},
);
Ok(handle)
}
#[cfg(test)]
pub(crate) fn test_register_service(
&mut self,
spec: mdns_proto::ServiceSpec,
now: StdInstant,
) -> Result<ServiceHandle, mdns_proto::error::RegisterServiceError> {
use crate::service::new_service_mailbox;
let mailbox = new_service_mailbox();
self.register_service(spec, now, mailbox)
}
pub(crate) fn start_query(
&mut self,
spec: mdns_proto::QuerySpec,
now: StdInstant,
) -> Result<QueryHandle, mdns_proto::error::StartQueryError> {
let h = self.endpoint.try_start_query(spec, now)?;
self.queries.insert(
h,
QueryCtx {
last_seq: 0,
cancelled: false,
errored: false,
terminal_wake_pending: false,
},
);
Ok(h)
}
pub(crate) fn flag_query_cancelled(&mut self, h: QueryHandle) {
if let Some(q) = self.queries.get_mut(&h) {
q.cancelled = true;
}
}
pub(crate) fn flag_service_unregistered(&mut self, h: ServiceHandle) {
if let Some(s) = self.services.get_mut(&h) {
s.cancelled = true;
}
}
pub(crate) fn begin_service_withdrawal(&mut self, handle: ServiceHandle, now: StdInstant) {
let (snap, handoff) = match self.services.get_mut(&handle) {
Some(ctx) => {
ctx.errored = true;
let handoff = ctx.proto.take_rename_goodbye_handoff();
(ctx.proto.withdrawal_snapshot(), handoff)
}
None => return,
};
if let Some(handoff) = handoff {
self.endpoint.enqueue_rename_withdrawal(handoff, now, true);
}
self.endpoint.begin_withdrawal(handle, snap, now);
}
pub(crate) fn poll_one_withdrawal(
&mut self,
now: StdInstant,
scratch: &mut [u8],
) -> Option<(SocketAddr, usize, WithdrawalToken)> {
self.endpoint.poll_withdrawal_transmit(now, scratch)
}
pub(crate) fn note_withdrawal_result(
&mut self,
token: WithdrawalToken,
now: StdInstant,
v4: WithdrawalSend,
v6: WithdrawalSend,
) {
self.endpoint.note_withdrawal_result(token, now, v4, v6);
}
pub(crate) fn drain_completed_withdrawals(&mut self, now: StdInstant) -> bool {
self.completed_withdrawals.clear();
self
.endpoint
.drain_completed_withdrawals(now, &mut self.completed_withdrawals);
let mut gcd_any = false;
while let Some(handle) = self.completed_withdrawals.pop() {
if self.services.remove(&handle).is_some() {
gcd_any = true;
}
}
gcd_any
}
pub(crate) fn fire_timeouts(&mut self, now: StdInstant) {
let _ = self.endpoint.handle_timeout(now);
let Self {
endpoint, queries, ..
} = &mut *self;
for (&h, ctx) in queries.iter() {
if ctx.errored {
continue;
}
let _ = endpoint.handle_query_timeout(h, now);
}
for ctx in self.services.values_mut() {
if !ctx.cancelled && !ctx.errored {
let _ = ctx.proto.handle_timeout(now);
}
}
}
pub(crate) fn sweep_cancelled_services(&mut self, now: StdInstant) -> bool {
let cancelled: Vec<ServiceHandle> = self
.services
.iter()
.filter(|(_, ctx)| ctx.cancelled && !ctx.errored)
.map(|(h, _)| *h)
.collect();
let swept = !cancelled.is_empty();
for h in cancelled {
self.begin_service_withdrawal(h, now);
}
swept
}
pub(crate) fn push_service_updates(&mut self, now: StdInstant) -> bool {
let mut pushed_any = false;
self.svc_handle_scratch.clear();
self
.svc_handle_scratch
.extend(self.services.keys().copied());
let mut i = 0;
while i < self.svc_handle_scratch.len() {
let h = self.svc_handle_scratch[i];
i += 1;
if self.services.get(&h).is_some_and(|c| c.errored) {
continue;
}
while let Some(upd) = self.services.get_mut(&h).and_then(|c| c.proto.poll()) {
if let ServiceUpdate::Renamed(ref renamed) = upd {
let new_name = renamed.new_name().clone();
let rename_result = self.endpoint.handle_service_renamed(h, new_name);
let handoff = self
.services
.get_mut(&h)
.and_then(|c| c.proto.take_rename_goodbye_handoff());
if let Some(handoff) = handoff {
self
.endpoint
.enqueue_rename_withdrawal(handoff, now, rename_result.is_err());
}
if let Err(_e) = rename_result {
warn!(
handle = ?h,
error = ?_e,
"auto-rename collided with another local service; emitting Conflict and beginning withdrawal"
);
if let Some(ctx) = self.services.get(&h) {
ctx
.mailbox
.borrow_mut()
.set_terminal(ServiceUpdate::Conflict);
}
self.begin_service_withdrawal(h, now);
pushed_any = true;
break;
}
}
let is_terminal = upd.is_conflict() || upd.is_host_conflict();
if let Some(ctx) = self.services.get(&h) {
ctx.mailbox.borrow_mut().push_update(upd);
}
pushed_any = true;
if is_terminal {
self.begin_service_withdrawal(h, now);
break;
}
}
}
pushed_any
}
pub(crate) fn poll_one_transmit(
&mut self,
now: StdInstant,
scratch: &mut [u8],
) -> Option<(SocketAddr, usize, TransmitOrigin)> {
self.svc_handle_scratch.clear();
self
.svc_handle_scratch
.extend(self.services.keys().copied());
let mut i = 0;
while i < self.svc_handle_scratch.len() {
let h = self.svc_handle_scratch[i];
i += 1;
{
let Some(ctx) = self.services.get_mut(&h) else {
continue;
};
if ctx.cancelled || ctx.errored {
continue;
}
}
let escalated = {
let ctx = self
.services
.get_mut(&h)
.expect("handle present (just checked)");
match ctx.proto.poll_transmit(now, scratch) {
Ok(Some(t)) => {
ctx.encode_failures = 0;
return Some((t.dst(), t.size(), TransmitOrigin::Service(h)));
}
Ok(None) => {
ctx.encode_failures = 0;
false
}
Err(_e) => {
ctx.encode_failures = ctx.encode_failures.saturating_add(1);
if ctx.encode_failures >= MAX_CONSECUTIVE_ENCODE_ERRORS {
warn!(
handle = ?h,
error = ?_e,
scratch_size = scratch.len(),
consecutive_failures = ctx.encode_failures,
"Service::poll_transmit failed; escalating to Conflict and beginning withdrawal"
);
ctx
.mailbox
.borrow_mut()
.set_terminal(ServiceUpdate::Conflict);
true
} else {
false
}
}
}
};
if escalated {
self.begin_service_withdrawal(h, now);
}
}
self.query_handle_scratch.clear();
self
.query_handle_scratch
.extend(self.queries.keys().copied());
let mut i = 0;
while i < self.query_handle_scratch.len() {
let h = self.query_handle_scratch[i];
i += 1;
let Some(ctx) = self.queries.get_mut(&h) else {
continue;
};
if ctx.cancelled || ctx.errored {
continue;
}
match self.endpoint.poll_query_transmit(h, now, scratch) {
Ok(Some(t)) => return Some((t.dst(), t.size(), TransmitOrigin::Query(h))),
Ok(None) => {}
Err(_e) => {
if let Some(ctx) = self.queries.get_mut(&h) {
if !ctx.errored {
ctx.errored = true;
ctx.terminal_wake_pending = true;
}
}
self.endpoint.retire_query(h);
warn!(
handle = ?h,
error = ?_e,
scratch_size = scratch.len(),
"Query::poll_query_transmit failed to encode; retiring proto query and marking errored (Query::next will surface terminal)"
);
}
}
}
None
}
pub(crate) fn note_service_transmit_result(
&mut self,
h: ServiceHandle,
now: StdInstant,
delivered: bool,
) {
if let Some(ctx) = self.services.get_mut(&h) {
ctx.proto.note_transmit_result(now, delivered);
if delivered {
self.endpoint.note_service_advertised(
h,
ctx.proto.advertised_a_addrs(),
ctx.proto.advertised_aaaa_addrs(),
ctx.proto.advertises_instance(),
);
}
}
}
pub(crate) fn note_query_transmit_result(
&mut self,
h: QueryHandle,
now: StdInstant,
delivered: bool,
) {
self.endpoint.note_query_transmit_result(h, now, delivered);
}
#[cfg(test)]
pub(crate) fn has_pending_withdrawal(&self) -> bool {
self
.services
.values()
.any(|ctx| ctx.cancelled && !ctx.errored)
}
pub(crate) fn take_query_terminal_wakes(&mut self) -> bool {
let mut woke = false;
for ctx in self.queries.values_mut() {
if ctx.terminal_wake_pending {
ctx.terminal_wake_pending = false;
woke = true;
}
}
woke
}
pub(crate) fn next_withdrawal_deadline(&self) -> Option<StdInstant> {
self.endpoint.next_withdrawal_deadline()
}
pub(crate) fn poll_deadline(&self) -> Option<StdInstant> {
let mut best = self.endpoint.poll_timeout();
for ctx in self.services.values() {
if ctx.errored {
continue;
}
if let Some(t) = ctx.proto.poll_timeout() {
best = Some(best.map_or(t, |b| b.min(t)));
}
}
for (h, ctx) in &self.queries {
if ctx.errored {
continue;
}
if let Some(t) = self.endpoint.poll_query_timeout(*h) {
best = Some(best.map_or(t, |b| b.min(t)));
}
}
best
}
pub(crate) fn handle_datagram(&mut self, meta: &crate::socket::RecvMeta, data: &[u8]) {
let on_link = if meta.hop_limit().is_some() {
crate::onlink::is_on_link(meta.hop_limit())
} else {
crate::onlink::src_on_local_link(
&self.local_subnets,
self.bound_interface,
meta.interface_index(),
meta.peer().ip(),
)
};
if !on_link {
debug!(
src = %meta.peer(),
hop_limit = ?meta.hop_limit(),
"dropping off-link packet (RFC 6762 §11 trust boundary)"
);
#[cfg(feature = "stats")]
{
self.stats.packets_rx(1);
self.stats.bytes_rx(data.len() as u64);
self.stats.packets_dropped(1);
}
return;
}
if data.get(2).is_some_and(|b| b & 0x80 != 0)
&& meta.peer().port() != hick_udp::constants::MDNS_PORT
{
debug!(
src = %meta.peer(),
"dropping untrusted response (source port != 5353) before self-send match"
);
#[cfg(feature = "stats")]
{
self.stats.packets_rx(1);
self.stats.bytes_rx(data.len() as u64);
self.stats.packets_dropped(1);
}
return;
}
let caller_is_self = match meta.kernel_rx_time() {
Some(rx) => crate::selfsend::take_self_send(
&mut self.recent_sends,
data,
rx,
crate::selfsend::MatchMode::Ordered,
),
None => crate::selfsend::take_self_send(
&mut self.recent_sends,
data,
SystemTime::now(),
crate::selfsend::MatchMode::Degraded,
),
};
let now = StdInstant::now();
let Self {
endpoint, services, ..
} = self;
let route_events = match endpoint.handle(
now,
meta.peer(),
meta.local_ip(),
meta.interface_index(),
data,
caller_is_self,
) {
Ok(it) => it,
Err(_) => return,
};
for ev in route_events {
match ev {
Ok(mdns_proto::event::RouteEvent::ToService(ts)) => {
if let Some(ctx) = services.get_mut(&ts.handle())
&& !ctx.errored
{
ctx.proto.handle_event(ts.into_event(), now);
}
}
Ok(_) => {}
Err(_) => break,
}
}
}
}
pub(crate) struct EndpointInner {
pub(crate) state: RefCell<State>,
pub(crate) notify: LocalNotify,
pub(crate) dirty: core::cell::Cell<bool>,
}
impl EndpointInner {
pub(crate) fn new(cfg: EndpointConfig, max_payload: usize, max_recv: usize) -> Rc<Self> {
Rc::new(Self {
state: RefCell::new(State::new(cfg, max_payload, max_recv)),
notify: LocalNotify::new(),
dirty: core::cell::Cell::new(false),
})
}
#[inline]
pub(crate) fn mark_dirty(&self) {
self.dirty.set(true);
self.notify.notify();
}
}
#[cfg_attr(feature = "tracing", tracing::instrument(level = "trace", skip_all))]
pub(crate) async fn run(
inner: Rc<EndpointInner>,
sock_v4: Option<Rc<Socket>>,
sock_v6: Option<Rc<Socket>>,
) {
use futures::{FutureExt, future::Either};
let (mut scratch, max_recv) = {
let s = inner.state.borrow();
(vec![0u8; s.max_payload], s.max_recv)
};
loop {
let mut credits = MAX_TRANSMIT_CREDITS_PER_PASS;
loop {
if credits == 0 {
break;
}
credits -= 1;
let pumped = {
let mut s = inner.state.borrow_mut();
let now = StdInstant::now();
s.poll_one_transmit(now, &mut scratch)
};
let Some((dst, n, origin)) = pumped else {
break;
};
let delivered = if is_mdns_multicast_dst(dst) {
let mut sent_any = false;
if let Some(s4) = sock_v4.as_ref() {
let when = SystemTime::now();
let res = s4.send_to(&scratch[..n], MDNS_V4_DST, None).await;
if res.is_ok() {
trace!(dst = %MDNS_V4_DST, len = n, "send_to v4");
let mut state = inner.state.borrow_mut();
crate::selfsend::record_self_send(&mut state.recent_sends, &scratch[..n], when);
#[cfg(feature = "stats")]
{
state.stats.packets_tx(1);
state.stats.bytes_tx(n as u64);
}
sent_any = true;
} else {
debug!(dst = %MDNS_V4_DST, "send_to v4 failed");
#[cfg(feature = "stats")]
inner.state.borrow().stats.send_errors(1);
}
}
if let Some(s6) = sock_v6.as_ref() {
let when = SystemTime::now();
let res = s6.send_to(&scratch[..n], MDNS_V6_DST, None).await;
if res.is_ok() {
trace!(dst = %MDNS_V6_DST, len = n, "send_to v6");
let mut state = inner.state.borrow_mut();
crate::selfsend::record_self_send(&mut state.recent_sends, &scratch[..n], when);
#[cfg(feature = "stats")]
{
state.stats.packets_tx(1);
state.stats.bytes_tx(n as u64);
}
sent_any = true;
} else {
debug!(dst = %MDNS_V6_DST, "send_to v6 failed");
#[cfg(feature = "stats")]
inner.state.borrow().stats.send_errors(1);
}
}
sent_any
} else {
let sock = match dst {
SocketAddr::V4(_) => sock_v4.as_ref(),
SocketAddr::V6(_) => sock_v6.as_ref(),
};
if let Some(s) = sock {
let when = SystemTime::now();
let res = s.send_to(&scratch[..n], dst, None).await;
if res.is_ok() {
trace!(dst = %dst, len = n, "send_to");
let mut state = inner.state.borrow_mut();
crate::selfsend::record_self_send(&mut state.recent_sends, &scratch[..n], when);
#[cfg(feature = "stats")]
{
state.stats.packets_tx(1);
state.stats.bytes_tx(n as u64);
}
} else {
debug!(dst = %dst, "send_to failed");
#[cfg(feature = "stats")]
inner.state.borrow().stats.send_errors(1);
}
res.is_ok()
} else {
false
}
};
match origin {
TransmitOrigin::Service(h) => {
let mut state = inner.state.borrow_mut();
state.note_service_transmit_result(h, StdInstant::now(), delivered);
}
TransmitOrigin::Query(h) => {
let mut state = inner.state.borrow_mut();
state.note_query_transmit_result(h, StdInstant::now(), delivered);
}
}
}
let pump_budget_exhausted = credits == 0;
{
let now = StdInstant::now();
inner.state.borrow_mut().sweep_cancelled_services(now);
}
{
let now = StdInstant::now();
let pushed = inner.state.borrow_mut().push_service_updates(now);
if pushed {
inner.notify.notify();
}
}
drain_withdrawals(&inner, &sock_v4, &sock_v6, &mut scratch).await;
{
let woke = inner.state.borrow_mut().take_query_terminal_wakes();
if woke {
inner.notify.notify();
}
}
if Rc::strong_count(&inner) == 1 {
let shutdown_deadline = StdInstant::now() + Duration::from_secs(10);
loop {
drain_withdrawals(&inner, &sock_v4, &sock_v6, &mut scratch).await;
inner
.state
.borrow_mut()
.sweep_cancelled_services(StdInstant::now());
let now = StdInstant::now();
let Some(next) = ({ inner.state.borrow().next_withdrawal_deadline() }) else {
break;
};
if now >= shutdown_deadline {
debug!("shutdown withdrawal flush hit its wall-clock backstop; exiting");
break;
}
let dur = next
.saturating_duration_since(now)
.min(shutdown_deadline.saturating_duration_since(now));
if dur > Duration::ZERO {
compio::time::sleep(dur).await;
}
}
break;
}
let deadline = { inner.state.borrow().poll_deadline() };
let force_now = inner.dirty.replace(false) || pump_budget_exhausted;
let timer_fut = match (force_now, deadline) {
(true, _) => Either::Left(compio::time::sleep(Duration::ZERO)),
(false, Some(at)) => {
let dur = at.saturating_duration_since(StdInstant::now());
Either::Left(compio::time::sleep(dur))
}
(false, None) => Either::Right(core::future::pending::<()>()),
}
.fuse();
futures::pin_mut!(timer_fut);
let notify_fut = inner.notify.listen().fuse();
futures::pin_mut!(notify_fut);
let mut woke_state = false;
match (sock_v4.as_ref(), sock_v6.as_ref()) {
(Some(s4), Some(s6)) => {
let r4 = s4.recv(max_recv).fuse();
let r6 = s6.recv(max_recv).fuse();
futures::pin_mut!(r4, r6);
futures::select! {
r = r4 => { handle_recv(&inner, r); woke_state = true; }
r = r6 => { handle_recv(&inner, r); woke_state = true; }
_ = timer_fut => { inner.state.borrow_mut().fire_timeouts(StdInstant::now()); woke_state = true; }
_ = notify_fut => {}
}
}
(Some(s4), None) => {
let r4 = s4.recv(max_recv).fuse();
futures::pin_mut!(r4);
futures::select! {
r = r4 => { handle_recv(&inner, r); woke_state = true; }
_ = timer_fut => { inner.state.borrow_mut().fire_timeouts(StdInstant::now()); woke_state = true; }
_ = notify_fut => {}
}
}
(None, Some(s6)) => {
let r6 = s6.recv(max_recv).fuse();
futures::pin_mut!(r6);
futures::select! {
r = r6 => { handle_recv(&inner, r); woke_state = true; }
_ = timer_fut => { inner.state.borrow_mut().fire_timeouts(StdInstant::now()); woke_state = true; }
_ = notify_fut => {}
}
}
(None, None) => {
futures::select! {
_ = timer_fut => { inner.state.borrow_mut().fire_timeouts(StdInstant::now()); woke_state = true; }
_ = notify_fut => {}
}
}
}
if woke_state {
inner.notify.notify();
}
}
}
fn present_socket_send_outcome<T>(res: &std::io::Result<T>) -> WithdrawalSend {
match res {
Ok(_) => WithdrawalSend::Sent,
Err(_) => WithdrawalSend::Retry,
}
}
async fn drain_withdrawals(
inner: &Rc<EndpointInner>,
sock_v4: &Option<Rc<Socket>>,
sock_v6: &Option<Rc<Socket>>,
scratch: &mut [u8],
) {
loop {
let due = {
let mut s = inner.state.borrow_mut();
let now = StdInstant::now();
s.poll_one_withdrawal(now, scratch)
};
let Some((_dst, len, token)) = due else {
break;
};
let mut v4_out = WithdrawalSend::WriteOff;
let mut v6_out = WithdrawalSend::WriteOff;
if let Some(s4) = sock_v4.as_ref() {
let when = SystemTime::now();
let res = s4.send_to(&scratch[..len], MDNS_V4_DST, None).await;
v4_out = present_socket_send_outcome(&res);
match res {
Ok(_) => {
trace!(dst = %MDNS_V4_DST, len, "withdrawal send_to v4");
let mut state = inner.state.borrow_mut();
crate::selfsend::record_self_send(&mut state.recent_sends, &scratch[..len], when);
#[cfg(feature = "stats")]
{
state.stats.packets_tx(1);
state.stats.bytes_tx(len as u64);
}
}
Err(_e) => {
debug!(error = %_e, dst = %MDNS_V4_DST, "withdrawal send_to v4 failed");
#[cfg(feature = "stats")]
inner.state.borrow().stats.send_errors(1);
}
}
}
if let Some(s6) = sock_v6.as_ref() {
let when = SystemTime::now();
let res = s6.send_to(&scratch[..len], MDNS_V6_DST, None).await;
v6_out = present_socket_send_outcome(&res);
match res {
Ok(_) => {
trace!(dst = %MDNS_V6_DST, len, "withdrawal send_to v6");
let mut state = inner.state.borrow_mut();
crate::selfsend::record_self_send(&mut state.recent_sends, &scratch[..len], when);
#[cfg(feature = "stats")]
{
state.stats.packets_tx(1);
state.stats.bytes_tx(len as u64);
}
}
Err(_e) => {
debug!(error = %_e, dst = %MDNS_V6_DST, "withdrawal send_to v6 failed");
#[cfg(feature = "stats")]
inner.state.borrow().stats.send_errors(1);
}
}
}
#[cfg(feature = "stats")]
if matches!(v4_out, WithdrawalSend::Sent) || matches!(v6_out, WithdrawalSend::Sent) {
inner.state.borrow().stats.goodbyes_tx(1);
}
{
let now = StdInstant::now();
inner
.state
.borrow_mut()
.note_withdrawal_result(token, now, v4_out, v6_out);
}
}
let gcd_any = {
let now = StdInstant::now();
inner.state.borrow_mut().drain_completed_withdrawals(now)
};
if gcd_any {
inner.notify.notify();
}
}
#[inline]
fn handle_recv(inner: &Rc<EndpointInner>, r: std::io::Result<(Vec<u8>, RecvMeta)>) {
match r {
Ok((data, meta)) => {
trace!(src = %meta.peer(), len = data.len(), truncated = meta.truncated(), "recv datagram");
if meta.truncated() {
debug!(
src = %meta.peer(),
len = data.len(),
"dropping truncated (oversized) datagram before proto routing"
);
#[cfg(feature = "stats")]
{
let s = inner.state.borrow();
s.stats.packets_rx(1);
s.stats.bytes_rx(data.len() as u64);
s.stats.packets_dropped(1);
}
return;
}
let mut s = inner.state.borrow_mut();
s.handle_datagram(&meta, &data);
}
Err(_e) => {
debug!(error = %_e, "socket recv failed");
}
}
}