Files
lanspread/crates/lanspread-mdns/src/lib.rs
T
ddidderr 60fd7ba0c2 feat(peer)!: cut over to authenticated catalog sharing
Replace address-only trust and pushed peer state with installation identities,
SPKI-pinned QUIC, candidate-only discovery, and bounded responder-owned
protocol-8 pulls. The runtime now owns each network generation and all admitted
work through shutdown.

Add exact bundled content identities, reproducible manifest publishing,
capability-confined downloads, streaming BLAKE3 verification, quarantine and
retry, and crash-recoverable download and install transactions. Ship generated
fixture catalogs and fail closed when production manifests are absent.

The Tauri backend exposes durable sharing policy, redacted identity state, and
attempt-keyed transfer snapshots. Frontend consumption follows in the next
commit. Repository-wide test certificates and protocol-7 paths are removed.

BREAKING CHANGE: peers must use protocol 8 and exact catalog content artifacts;
protocol-7 frames and shared-certificate identities are no longer accepted.

Test Plan:
- `just test` -- passed on the completed stack (708 workspace tests)
- `just clippy` -- passed on the completed stack
- `just build` -- passed with fixture catalogs on the completed stack
- `just catalog-check-production` -- failed closed because the external
  production manifest corpus is absent
- `git diff --cached --check` -- passed
2026-08-10 13:59:18 +02:00

537 lines
16 KiB
Rust

#![allow(clippy::missing_errors_doc, clippy::missing_panics_doc)]
use std::{
collections::HashMap,
net::SocketAddr,
thread,
time::{Duration, Instant},
};
use eyre::{WrapErr as _, bail};
pub use mdns_sd::DaemonEvent;
use mdns_sd::{
DaemonStatus,
Error as MdnsError,
Receiver,
ResolvedService,
ServiceDaemon,
ServiceEvent,
ServiceInfo,
UnregisterStatus,
};
pub const LANSPREAD_SERVICE_TYPE: &str = "_lanspread._udp.local.";
pub type MdnsMonitor = Receiver<DaemonEvent>;
const DAEMON_COMMAND_RETRY_DELAY: Duration = Duration::from_millis(1);
#[derive(Debug, PartialEq, Eq)]
enum UnregisterOutcome {
Removed,
AlreadyAbsent,
}
struct DaemonOwner {
daemon: Option<ServiceDaemon>,
}
impl DaemonOwner {
fn new(daemon: ServiceDaemon) -> Self {
Self {
daemon: Some(daemon),
}
}
fn daemon(&self) -> &ServiceDaemon {
self.daemon
.as_ref()
.expect("mDNS daemon is available until its owner is closed")
}
fn is_closed(&self) -> bool {
self.daemon.is_none()
}
fn close(&mut self) -> eyre::Result<()> {
let Some(daemon) = self.daemon.as_ref() else {
return Ok(());
};
shutdown_and_wait(daemon)?;
self.daemon.take();
Ok(())
}
}
impl Drop for DaemonOwner {
fn drop(&mut self) {
if let Err(err) = self.close() {
log::error!("Failed to stop mDNS daemon during cleanup: {err:#}");
}
}
}
pub struct MdnsAdvertiser {
daemon: DaemonOwner,
service_info: ServiceInfo,
pub monitor: Receiver<DaemonEvent>,
}
impl MdnsAdvertiser {
pub fn new(
service_type: &str,
instance_name: &str,
address: SocketAddr,
properties: Option<HashMap<String, String>>,
) -> eyre::Result<Self> {
let daemon = ServiceDaemon::new()?;
Self::new_with_daemon(daemon, service_type, instance_name, address, properties)
}
fn new_with_daemon(
daemon: ServiceDaemon,
service_type: &str,
instance_name: &str,
address: SocketAddr,
properties: Option<HashMap<String, String>>,
) -> eyre::Result<Self> {
let daemon = DaemonOwner::new(daemon);
let host_name = format!("{}.local.", address.ip());
let service_info = ServiceInfo::new(
service_type,
instance_name,
&host_name,
address.ip(),
address.port(),
properties,
)?;
let monitor = daemon.daemon().monitor()?;
// Register the service
daemon.daemon().register(service_info.clone())?;
Ok(Self {
daemon,
service_info,
monitor,
})
}
/// Unregisters the service and waits for the daemon's shutdown acknowledgement.
pub fn close(mut self) -> eyre::Result<()> {
self.close_inner()
}
fn close_inner(&mut self) -> eyre::Result<()> {
if self.daemon.is_closed() {
return Ok(());
}
let unregister_result =
unregister_and_wait(self.daemon.daemon(), self.service_info.get_fullname()).map(|_| ());
let shutdown_result = self.daemon.close();
combine_cleanup_results(unregister_result, shutdown_result)
}
}
impl Drop for MdnsAdvertiser {
fn drop(&mut self) {
if let Err(err) = self.close_inner() {
log::error!("Failed to close mDNS advertiser during cleanup: {err:#}");
}
}
}
pub struct MdnsBrowser {
daemon: DaemonOwner,
receiver: Receiver<ServiceEvent>,
service_type: String,
}
#[derive(Debug, Clone)]
pub struct MdnsService {
pub addr: SocketAddr,
pub fullname: String,
pub hostname: String,
pub properties: HashMap<String, String>,
}
#[derive(Debug, Clone)]
pub enum MdnsServicePoll {
Service(MdnsService),
Timeout,
Closed,
}
impl MdnsBrowser {
pub fn new(service_type: &str) -> eyre::Result<Self> {
let daemon = ServiceDaemon::new()?;
Self::new_with_daemon(daemon, service_type)
}
fn new_with_daemon(daemon: ServiceDaemon, service_type: &str) -> eyre::Result<Self> {
let daemon = DaemonOwner::new(daemon);
let receiver = daemon.daemon().browse(service_type)?;
Ok(Self {
daemon,
receiver,
service_type: service_type.to_string(),
})
}
/// Stops browsing and waits for the daemon's shutdown acknowledgement.
pub fn close(mut self) -> eyre::Result<()> {
self.close_inner()
}
pub fn next_service(
&self,
ignore_addr: Option<SocketAddr>,
) -> eyre::Result<Option<MdnsService>> {
loop {
match self.receiver.recv() {
Ok(event) => {
if let Some(service) = self.service_from_event(event, ignore_addr) {
return Ok(Some(service));
}
}
Err(err) => {
log::error!("mDNS browse channel closed: {err}");
return Ok(None);
}
}
}
}
pub fn next_service_timeout(
&self,
ignore_addr: Option<SocketAddr>,
timeout: Duration,
) -> eyre::Result<MdnsServicePoll> {
let deadline = Instant::now() + timeout;
loop {
let remaining = deadline.saturating_duration_since(Instant::now());
if remaining.is_zero() {
return Ok(MdnsServicePoll::Timeout);
}
match self.receiver.recv_timeout(remaining) {
Ok(event) => {
if let Some(service) = self.service_from_event(event, ignore_addr) {
return Ok(MdnsServicePoll::Service(service));
}
}
Err(err) if self.receiver.is_disconnected() => {
log::error!("mDNS browse channel closed: {err}");
return Ok(MdnsServicePoll::Closed);
}
Err(err) => {
log::trace!("mDNS browse timeout: {err}");
return Ok(MdnsServicePoll::Timeout);
}
}
}
}
pub fn next_address(
&self,
ignore_addr: Option<SocketAddr>,
) -> eyre::Result<Option<SocketAddr>> {
Ok(self.next_service(ignore_addr)?.map(|service| service.addr))
}
fn service_from_event(
&self,
event: ServiceEvent,
ignore_addr: Option<SocketAddr>,
) -> Option<MdnsService> {
match event {
ServiceEvent::ServiceResolved(info) => self.service_from_resolved(&info, ignore_addr),
other_event => {
log::trace!("mdns unrelated event: {other_event:?}");
None
}
}
}
fn service_from_resolved(
&self,
info: &ResolvedService,
ignore_addr: Option<SocketAddr>,
) -> Option<MdnsService> {
log::trace!("mdns ServiceResolved event: {info:?}");
if info.ty_domain != self.service_type {
log::trace!(
"Got mDNS with uninteresting service type: {} (expected: {})",
info.ty_domain,
self.service_type,
);
return None;
}
let mut ignored_match = false;
for address in info.get_addresses() {
let addr = SocketAddr::new(address.to_ip_addr(), info.get_port());
if ignore_addr.is_some_and(|ignore| ignore == addr) {
ignored_match = true;
log::trace!("Ignoring mDNS advertisement for local server at {addr}");
continue;
}
log::info!("Found server at {addr}");
let properties = info.get_properties().clone().into_property_map_str();
return Some(MdnsService {
addr,
fullname: info.get_fullname().to_string(),
hostname: info.get_hostname().to_string(),
properties,
});
}
if ignored_match {
log::trace!(
"Only saw ignored mDNS advertisements (probably ourselves) for {:?}",
info.get_fullname()
);
return None;
}
log::error!("No address found in mDNS response: {info:?}");
None
}
fn close_inner(&mut self) -> eyre::Result<()> {
self.daemon.close()
}
}
impl Drop for MdnsBrowser {
fn drop(&mut self) {
if let Err(err) = self.close_inner() {
log::error!("Failed to close mDNS browser during cleanup: {err:#}");
}
}
}
fn unregister_and_wait(daemon: &ServiceDaemon, fullname: &str) -> eyre::Result<UnregisterOutcome> {
let receiver = loop {
match daemon.unregister(fullname) {
Ok(receiver) => break receiver,
Err(MdnsError::Again) => thread::sleep(DAEMON_COMMAND_RETRY_DELAY),
Err(err) => return Err(err).wrap_err("failed to request mDNS service unregister"),
}
};
match receiver
.recv()
.wrap_err("mDNS unregister response channel closed")?
{
UnregisterStatus::OK => Ok(UnregisterOutcome::Removed),
// Cleanup is idempotent: the desired postcondition already holds if
// the daemon no longer has this service in its registration table.
UnregisterStatus::NotFound => Ok(UnregisterOutcome::AlreadyAbsent),
}
}
fn shutdown_and_wait(daemon: &ServiceDaemon) -> eyre::Result<()> {
let receiver = loop {
match daemon.shutdown() {
Ok(receiver) => break receiver,
Err(MdnsError::Again) => thread::sleep(DAEMON_COMMAND_RETRY_DELAY),
Err(MdnsError::DaemonShutdown) => return Ok(()),
// `mdns-sd` enqueues the command before waking its daemon socket. A
// wake failure can therefore return an error even though shutdown
// is already queued. Retry until the daemon confirms it stopped.
Err(err) => {
log::warn!("Failed to signal mDNS daemon shutdown; retrying: {err}");
thread::sleep(DAEMON_COMMAND_RETRY_DELAY);
}
}
};
match receiver.recv() {
Ok(DaemonStatus::Shutdown) => Ok(()),
Ok(status) => bail!("mDNS daemon returned unexpected shutdown status: {status:?}"),
Err(_) => wait_for_reported_shutdown(daemon),
}
}
fn wait_for_reported_shutdown(daemon: &ServiceDaemon) -> eyre::Result<()> {
loop {
let receiver = match daemon.status() {
Ok(receiver) => receiver,
Err(MdnsError::Again) => {
thread::sleep(DAEMON_COMMAND_RETRY_DELAY);
continue;
}
Err(MdnsError::DaemonShutdown) => return Ok(()),
Err(err) => {
log::warn!("Failed to query mDNS daemon shutdown status; retrying: {err}");
thread::sleep(DAEMON_COMMAND_RETRY_DELAY);
continue;
}
};
match receiver.recv() {
Ok(DaemonStatus::Shutdown) => return Ok(()),
Ok(DaemonStatus::Running) | Err(_) => {
thread::sleep(DAEMON_COMMAND_RETRY_DELAY);
}
Ok(status) => bail!("mDNS daemon returned unexpected status: {status:?}"),
}
}
}
fn combine_cleanup_results(
first: eyre::Result<()>,
shutdown: eyre::Result<()>,
) -> eyre::Result<()> {
match (first, shutdown) {
(Ok(()), Ok(())) => Ok(()),
(Err(err), Ok(())) | (Ok(()), Err(err)) => Err(err),
(Err(first_err), Err(shutdown_err)) => bail!(
"mDNS cleanup failed: {first_err:#}; daemon shutdown also failed: {shutdown_err:#}"
),
}
}
pub fn discover_service(
service_type: &str,
ignore_addr: Option<SocketAddr>,
) -> eyre::Result<SocketAddr> {
// Currently unused; kept for potential one-off discovery callers that just need a single address.
let browser = MdnsBrowser::new(service_type)?;
let result = browser
.next_address(ignore_addr)
.and_then(|address| address.ok_or_else(|| eyre::eyre!("No server found.")));
let shutdown_result = browser.close();
match (result, shutdown_result) {
(Ok(address), Ok(())) => Ok(address),
(Err(err), Ok(())) | (Ok(_), Err(err)) => Err(err),
(Err(discovery_err), Err(shutdown_err)) => bail!(
"mDNS discovery failed: {discovery_err:#}; daemon shutdown also failed: {shutdown_err:#}"
),
}
}
#[cfg(test)]
mod tests {
use std::{net::SocketAddr, time::Duration};
use mdns_sd::{DaemonStatus, ServiceDaemon};
use super::{
MdnsAdvertiser,
MdnsBrowser,
UnregisterOutcome,
shutdown_and_wait,
unregister_and_wait,
};
// mdns-sd rejects service names longer than 15 bytes asynchronously. Keep
// the lifecycle fixture valid so it exercises a real registration.
const TEST_SERVICE_TYPE: &str = "_ls-lifecycle._udp.local.";
#[test]
fn advertiser_close_waits_for_unregister_and_daemon_shutdown() {
let daemon = ServiceDaemon::new_with_port(0).expect("test daemon should start");
let probe = daemon.clone();
let advertiser = MdnsAdvertiser::new_with_daemon(
daemon,
TEST_SERVICE_TYPE,
"advertiser-close",
SocketAddr::from(([127, 0, 0, 1], 41_234)),
None,
)
.expect("test advertiser should register");
advertiser.close().expect("advertiser should close cleanly");
assert_daemon_stopped(&probe);
}
#[test]
fn advertiser_drop_waits_for_unregister_and_daemon_shutdown() {
let daemon = ServiceDaemon::new_with_port(0).expect("test daemon should start");
let probe = daemon.clone();
let advertiser = MdnsAdvertiser::new_with_daemon(
daemon,
TEST_SERVICE_TYPE,
"advertiser-drop",
SocketAddr::from(([127, 0, 0, 1], 41_235)),
None,
)
.expect("test advertiser should register");
drop(advertiser);
assert_daemon_stopped(&probe);
}
#[test]
fn lifecycle_fixture_reaches_the_daemon_registration_table() {
let daemon = ServiceDaemon::new_with_port(0).expect("test daemon should start");
let advertiser = MdnsAdvertiser::new_with_daemon(
daemon,
TEST_SERVICE_TYPE,
"registered-fixture",
SocketAddr::from(([127, 0, 0, 1], 41_236)),
None,
)
.expect("test advertiser should register");
let outcome = unregister_and_wait(
advertiser.daemon.daemon(),
advertiser.service_info.get_fullname(),
)
.expect("test registration should be removable");
assert_eq!(outcome, UnregisterOutcome::Removed);
advertiser.close().expect("advertiser should close cleanly");
}
#[test]
fn unregister_is_idempotent_when_service_is_already_absent() {
let daemon = ServiceDaemon::new_with_port(0).expect("test daemon should start");
let outcome = unregister_and_wait(&daemon, "absent._ls-lifecycle._udp.local.")
.expect("an already-absent service satisfies the cleanup postcondition");
assert_eq!(outcome, UnregisterOutcome::AlreadyAbsent);
shutdown_and_wait(&daemon).expect("test daemon should stop");
}
#[test]
fn browser_close_waits_for_daemon_shutdown() {
let daemon = ServiceDaemon::new_with_port(0).expect("test daemon should start");
let probe = daemon.clone();
let browser = MdnsBrowser::new_with_daemon(daemon, TEST_SERVICE_TYPE)
.expect("test browser should start");
browser.close().expect("browser should close cleanly");
assert_daemon_stopped(&probe);
}
#[test]
fn browser_drop_waits_for_daemon_shutdown() {
let daemon = ServiceDaemon::new_with_port(0).expect("test daemon should start");
let probe = daemon.clone();
let browser = MdnsBrowser::new_with_daemon(daemon, TEST_SERVICE_TYPE)
.expect("test browser should start");
drop(browser);
assert_daemon_stopped(&probe);
}
fn assert_daemon_stopped(daemon: &ServiceDaemon) {
let status = daemon
.status()
.expect("daemon status should remain observable")
.recv_timeout(Duration::from_secs(1))
.expect("daemon status should arrive");
assert_eq!(status, DaemonStatus::Shutdown);
}
}