use std::{
collections::VecDeque,
sync::{Arc, Mutex, MutexGuard},
};
use mdns_proto::{ServiceHandle, ServiceUpdate};
use crate::{command::Command, error::CancelError};
pub(crate) const SERVICE_UPDATE_CAPACITY: usize = 16;
pub(crate) struct ServiceMailbox {
updates: VecDeque<ServiceUpdate>,
terminal: Option<ServiceUpdate>,
terminal_delivered: bool,
}
#[cfg_attr(test, derive(Debug))]
enum Drained {
Update(ServiceUpdate),
Ended,
Empty,
}
impl ServiceMailbox {
fn new() -> Self {
Self {
updates: VecDeque::new(),
terminal: None,
terminal_delivered: false,
}
}
pub(crate) fn push_update(&mut self, upd: ServiceUpdate) {
if upd.is_conflict() || upd.is_host_conflict() {
self.set_terminal(upd);
return;
}
if upd.is_renamed() {
self.updates.retain(|u| !u.is_renamed());
self.bounded_push_back(upd);
return;
}
if upd.is_established() {
self.updates.retain(|u| !u.is_established());
self.bounded_push_back(upd);
return;
}
self.bounded_push_back(upd);
}
fn bounded_push_back(&mut self, upd: ServiceUpdate) {
if self.updates.len() >= SERVICE_UPDATE_CAPACITY {
self.updates.pop_front();
}
self.updates.push_back(upd);
}
pub(crate) fn set_terminal(&mut self, terminal: ServiceUpdate) {
if self.terminal.is_none() && !self.terminal_delivered {
self.terminal = Some(terminal);
}
}
#[cfg(test)]
pub(crate) fn non_terminal_len(&self) -> usize {
self.updates.len()
}
#[cfg(test)]
pub(crate) fn drain_for_test(&mut self) -> Option<ServiceUpdate> {
match self.drain() {
Drained::Update(upd) => Some(upd),
Drained::Ended | Drained::Empty => None,
}
}
#[cfg(test)]
pub(crate) fn fill_non_terminal_to_cap_for_test(&mut self) {
use mdns_proto::event::ServiceRenamed;
self.updates.clear();
for i in 0..SERVICE_UPDATE_CAPACITY {
self
.updates
.push_back(ServiceUpdate::Renamed(ServiceRenamed::new(
mdns_proto::Name::try_from_str(&format!("fill-{i}._ipp._tcp.local.")).unwrap(),
)));
}
}
fn drain(&mut self) -> Drained {
if let Some(upd) = self.updates.pop_front() {
Drained::Update(upd)
} else if let Some(terminal) = self.terminal.take() {
self.terminal_delivered = true;
Drained::Update(terminal)
} else if self.terminal_delivered {
Drained::Ended
} else {
Drained::Empty
}
}
}
fn lock(mailbox: &Mutex<ServiceMailbox>) -> MutexGuard<'_, ServiceMailbox> {
mailbox
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
pub struct Service {
handle: ServiceHandle,
mailbox: Arc<Mutex<ServiceMailbox>>,
doorbell: async_channel::Receiver<()>,
cmd: async_channel::Sender<Command>,
}
impl Service {
pub(crate) fn new(
handle: ServiceHandle,
mailbox: Arc<Mutex<ServiceMailbox>>,
doorbell: async_channel::Receiver<()>,
cmd: async_channel::Sender<Command>,
) -> Self {
Self {
handle,
mailbox,
doorbell,
cmd,
}
}
#[inline]
pub const fn handle(&self) -> ServiceHandle {
self.handle
}
pub async fn next(&self) -> Option<ServiceUpdate> {
loop {
match lock(&self.mailbox).drain() {
Drained::Update(upd) => return Some(upd),
Drained::Ended => return None,
Drained::Empty => {}
}
if self.doorbell.recv().await.is_err() {
return match lock(&self.mailbox).drain() {
Drained::Update(upd) => Some(upd),
_ => None,
};
}
}
}
pub async fn unregister(self) -> Result<(), CancelError> {
self
.cmd
.send(Command::UnregisterService {
handle: self.handle,
})
.await
.map_err(|_| CancelError::DriverGone)?;
Ok(())
}
}
impl Drop for Service {
fn drop(&mut self) {
let _ = self.cmd.try_send(Command::UnregisterService {
handle: self.handle,
});
}
}
pub(crate) fn new_service_mailbox() -> (
Arc<Mutex<ServiceMailbox>>,
async_channel::Sender<()>,
async_channel::Receiver<()>,
) {
let mailbox = Arc::new(Mutex::new(ServiceMailbox::new()));
let (doorbell_tx, doorbell_rx) = async_channel::bounded(1);
(mailbox, doorbell_tx, doorbell_rx)
}
#[cfg(test)]
mod tests;