use crate::config::Config;
use socket2::{Domain, Protocol, Socket, Type};
use std::net::{Ipv4Addr, SocketAddrV4};
use std::sync::Arc;
use std::time::Duration;
use tokio::net::UdpSocket;
const MULTICAST_ADDR: Ipv4Addr = Ipv4Addr::new(224, 0, 0, 167);
const MULTICAST_TTL: u32 = 1;
pub async fn run(config: Config) -> anyhow::Result<()> {
let discovery_port = config.server.port + 1;
let multicast_target = format!("{}:{}", MULTICAST_ADDR, discovery_port);
let recv_socket = {
let socket = Socket::new(Domain::IPV4, Type::DGRAM, Some(Protocol::UDP))?;
socket.set_reuse_address(true)?;
#[cfg(unix)]
socket.set_reuse_port(true)?;
socket.bind(&socket2::SockAddr::from(SocketAddrV4::new(
Ipv4Addr::UNSPECIFIED,
discovery_port,
)))?;
socket.join_multicast_v4(&MULTICAST_ADDR, &Ipv4Addr::UNSPECIFIED)?;
socket.set_multicast_ttl_v4(MULTICAST_TTL)?;
socket.set_multicast_loop_v4(true)?;
let std_sock: std::net::UdpSocket = socket.into();
std_sock.set_nonblocking(true)?;
UdpSocket::from_std(std_sock)?
};
tracing::info!(
"Discovery service listening on UDP {} (multicast {})",
discovery_port,
MULTICAST_ADDR
);
let announce_socket = {
let socket = Socket::new(Domain::IPV4, Type::DGRAM, Some(Protocol::UDP))?;
socket.bind(&socket2::SockAddr::from(SocketAddrV4::new(
Ipv4Addr::UNSPECIFIED,
0, )))?;
socket.set_multicast_ttl_v4(MULTICAST_TTL)?;
socket.set_multicast_loop_v4(true)?;
let std_sock: std::net::UdpSocket = socket.into();
std_sock.set_nonblocking(true)?;
UdpSocket::from_std(std_sock)?
};
let announce = serde_json::json!({
"type": "announce",
"name": config.service.name,
"tcp_port": config.server.port,
});
let announce_bytes = Arc::new(serde_json::to_vec(&announce)?);
let broadcast_bytes = announce_bytes.clone();
let broadcast_target = multicast_target.clone();
tokio::spawn(async move {
send_burst(&announce_socket, &broadcast_bytes, &broadcast_target).await;
loop {
tokio::time::sleep(Duration::from_secs(5)).await;
if let Err(e) = announce_socket.send_to(&broadcast_bytes, &broadcast_target).await {
tracing::warn!("Failed to multicast announcement: {}", e);
} else {
tracing::debug!("Sent announcement to {}", broadcast_target);
}
}
});
let mut buf = [0u8; 512];
loop {
match recv_socket.recv_from(&mut buf).await {
Ok((_len, src)) => {
tracing::debug!("Discovery request from {}", src);
if let Err(e) = recv_socket.send_to(&announce_bytes, src).await {
tracing::warn!("Failed to send discovery response to {}: {}", src, e);
}
}
Err(e) => {
tracing::error!("Discovery socket error: {}", e);
}
}
}
}
async fn send_burst(socket: &UdpSocket, data: &[u8], target: &str) {
let delays_ms = [0u64, 100, 600]; for &delay_ms in &delays_ms {
if delay_ms > 0 {
tokio::time::sleep(Duration::from_millis(delay_ms)).await;
}
if let Err(e) = socket.send_to(data, target).await {
tracing::warn!("Burst announcement failed: {}", e);
return; }
}
tracing::debug!("Announcement burst complete");
}