use super::{TransportMetrics};
use crate::{
types::SecureMessage,
error::Result,
circuit_breaker::{CircuitBreaker, CircuitBreakerConfig},
};
use std::{
time::{Duration, Instant},
collections::HashMap,
net::{SocketAddr, IpAddr, Ipv4Addr},
sync::{Arc, RwLock},
};
use tracing::{info, debug, warn, error};
use tokio::{net::UdpSocket as TokioUdpSocket, sync::Mutex, time::interval};
use serde::{Serialize, Deserialize};
pub struct EnhancedMdnsTransport {
instance_name: String,
service_type: String,
local_port: u16,
entity_id: String,
multicast_socket: Arc<Mutex<TokioUdpSocket>>,
discovered_peers: Arc<RwLock<HashMap<String, EnhancedMdnsPeer>>>,
our_announcements: Vec<ServiceAnnouncement>,
config: MdnsConfig,
metrics: Arc<RwLock<TransportMetrics>>,
circuit_breaker: Arc<CircuitBreaker>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct EnhancedMdnsPeer {
pub entity_id: String,
pub instance_name: String,
pub service_type: String,
pub host_name: String,
pub addresses: Vec<IpAddr>,
pub port: u16,
pub txt_records: HashMap<String, String>,
pub priority: u16,
pub weight: u16,
pub ttl: u32,
#[serde(skip, default = "Instant::now")]
pub discovered_at: Instant,
#[serde(skip, default = "Instant::now")]
pub last_seen: Instant,
pub capabilities: Vec<String>,
pub protocol_version: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ServiceAnnouncement {
pub instance_name: String,
pub service_type: String,
pub domain: String,
pub host_name: String,
pub port: u16,
pub txt_records: HashMap<String, String>,
pub ttl: u32,
#[serde(skip, default = "Instant::now")]
pub announced_at: Instant,
}
#[derive(Debug, Clone)]
pub struct MdnsConfig {
pub multicast_addr: Ipv4Addr,
pub multicast_port: u16,
pub announce_interval: Duration,
pub query_interval: Duration,
pub default_ttl: u32,
pub discovery_timeout: Duration,
pub max_peers: usize,
pub peer_timeout: Duration,
}
#[derive(Debug, Clone)]
pub enum MdnsPacket {
Query {
questions: Vec<MdnsQuestion>,
transaction_id: u16,
},
Response {
answers: Vec<MdnsRecord>,
authorities: Vec<MdnsRecord>,
additionals: Vec<MdnsRecord>,
transaction_id: u16,
authoritative: bool,
},
}
#[derive(Debug, Clone)]
pub struct MdnsQuestion {
pub name: String,
pub qtype: u16, pub qclass: u16, }
#[derive(Debug, Clone)]
pub struct MdnsRecord {
pub name: String,
pub rtype: u16,
pub rclass: u16,
pub ttl: u32,
pub data: Vec<u8>,
}
impl Default for MdnsConfig {
fn default() -> Self {
Self {
multicast_addr: Ipv4Addr::new(224, 0, 0, 251),
multicast_port: 5353,
announce_interval: Duration::from_secs(60),
query_interval: Duration::from_secs(30),
default_ttl: 120,
discovery_timeout: Duration::from_secs(5),
max_peers: 100,
peer_timeout: Duration::from_secs(300),
}
}
}
impl EnhancedMdnsTransport {
pub async fn new(
entity_id: String,
local_port: u16,
config: Option<MdnsConfig>,
) -> Result<Self> {
let config = config.unwrap_or_default();
let multicast_socket = create_multicast_socket(&config).await?;
let instance_name = format!("{}._synapse._tcp.local.", entity_id);
let service_type = "_synapse._tcp.local.".to_string();
let circuit_config = CircuitBreakerConfig {
failure_threshold: 3, minimum_requests: 2, failure_window: std::time::Duration::from_secs(30), recovery_timeout: std::time::Duration::from_secs(10), half_open_max_calls: 2, success_threshold: 0.7, };
let circuit_breaker = Arc::new(CircuitBreaker::new(circuit_config));
info!("Creating enhanced mDNS transport for {} on port {} with circuit breaker", entity_id, local_port);
Ok(Self {
instance_name,
service_type,
local_port,
entity_id,
multicast_socket: Arc::new(Mutex::new(multicast_socket)),
discovered_peers: Arc::new(RwLock::new(HashMap::new())),
our_announcements: Vec::new(),
config,
metrics: Arc::new(RwLock::new(TransportMetrics::default())),
circuit_breaker,
})
}
pub async fn start(&mut self) -> Result<()> {
info!("Starting enhanced mDNS service for {}", self.entity_id);
self.start_announcements().await?;
self.start_discovery().await?;
self.start_packet_processing().await?;
self.start_cleanup_task().await;
Ok(())
}
pub async fn announce_service(&mut self) -> Result<()> {
let announcement = ServiceAnnouncement {
instance_name: self.instance_name.clone(),
service_type: self.service_type.clone(),
domain: "local.".to_string(),
host_name: format!("{}.local.", self.entity_id),
port: self.local_port,
txt_records: self.build_txt_records(),
ttl: self.config.default_ttl,
announced_at: Instant::now(),
};
self.send_ptr_record(&announcement).await?;
self.send_srv_record(&announcement).await?;
self.send_txt_record(&announcement).await?;
self.send_host_records(&announcement).await?;
self.our_announcements.push(announcement);
info!("Announced Synapse service: {}", self.instance_name);
Ok(())
}
pub async fn discover_services(&self) -> Result<Vec<EnhancedMdnsPeer>> {
info!("Discovering Synapse services on local network");
self.send_ptr_query("_synapse._tcp.local.").await?;
tokio::time::sleep(self.config.discovery_timeout).await;
let peers = self.discovered_peers.read().unwrap();
Ok(peers.values().cloned().collect())
}
pub async fn find_peer(&self, entity_id: &str) -> Option<EnhancedMdnsPeer> {
let peers = self.discovered_peers.read().unwrap();
peers.get(entity_id).cloned()
}
pub async fn get_all_peers(&self) -> Vec<EnhancedMdnsPeer> {
let peers = self.discovered_peers.read().unwrap();
peers.values().cloned().collect()
}
pub async fn send_to_peer(&self, peer: &EnhancedMdnsPeer, message: &SecureMessage) -> Result<String> {
if let Some(addr) = peer.addresses.first() {
let socket_addr = SocketAddr::new(*addr, peer.port);
let tcp_transport = super::tcp::TcpTransport::new(0).await?;
{
let mut metrics = self.metrics.write().unwrap();
metrics.last_updated = Instant::now();
}
match tcp_transport.connect(&addr.to_string(), peer.port).await {
Ok(mut stream) => {
tcp_transport.send_via_stream(&mut stream, message).await?;
{
let mut metrics = self.metrics.write().unwrap();
metrics.reliability_score = (metrics.reliability_score * 0.9 + 0.1).min(1.0);
metrics.last_updated = Instant::now();
}
info!("Sent message to mDNS peer {} at {}", peer.entity_id, socket_addr);
Ok(format!("mdns://{}@{}", peer.entity_id, socket_addr))
}
Err(e) => {
{
let mut metrics = self.metrics.write().unwrap();
metrics.reliability_score = (metrics.reliability_score * 0.9).max(0.0);
metrics.packet_loss = (metrics.packet_loss + 0.1).min(1.0);
metrics.last_updated = Instant::now();
}
error!("Failed to send to mDNS peer {}: {}", peer.entity_id, e);
Err(e)
}
}
} else {
Err(crate::error::SynapseError::TransportError(
format!("No addresses available for peer {}", peer.entity_id)
).into())
}
}
async fn start_announcements(&self) -> Result<()> {
let socket = Arc::clone(&self.multicast_socket);
let config = self.config.clone();
let entity_id = self.entity_id.clone();
tokio::spawn(async move {
let mut announce_interval = interval(config.announce_interval);
loop {
announce_interval.tick().await;
if let Err(e) = Self::send_periodic_announcements(&socket, &entity_id, &config).await {
warn!("Failed to send mDNS announcements: {}", e);
}
}
});
Ok(())
}
async fn start_discovery(&self) -> Result<()> {
let socket = Arc::clone(&self.multicast_socket);
let config = self.config.clone();
tokio::spawn(async move {
let mut query_interval = interval(config.query_interval);
loop {
query_interval.tick().await;
if let Err(e) = Self::send_discovery_queries(&socket, &config).await {
warn!("Failed to send mDNS discovery queries: {}", e);
}
}
});
Ok(())
}
async fn start_packet_processing(&self) -> Result<()> {
let socket = Arc::clone(&self.multicast_socket);
let peers = Arc::clone(&self.discovered_peers);
let metrics = Arc::clone(&self.metrics);
tokio::spawn(async move {
let mut buffer = [0u8; 4096];
loop {
let socket_guard = socket.lock().await;
match socket_guard.recv_from(&mut buffer).await {
Ok((size, src)) => {
drop(socket_guard);
{
let mut m = metrics.write().unwrap();
m.throughput_bps = (m.throughput_bps + size as u64 * 8).max(size as u64 * 8);
m.last_updated = Instant::now();
}
if let Err(e) = Self::process_mdns_packet(&buffer[..size], src, &peers).await {
debug!("Error processing mDNS packet from {}: {}", src, e);
}
}
Err(e) => {
error!("Error receiving mDNS packet: {}", e);
tokio::time::sleep(Duration::from_millis(100)).await;
}
}
}
});
Ok(())
}
async fn start_cleanup_task(&self) {
let peers = Arc::clone(&self.discovered_peers);
let timeout = self.config.peer_timeout;
tokio::spawn(async move {
let mut cleanup_interval = interval(Duration::from_secs(60));
loop {
cleanup_interval.tick().await;
let now = Instant::now();
let mut peers_guard = peers.write().unwrap();
let initial_count = peers_guard.len();
peers_guard.retain(|entity_id, peer| {
if now.duration_since(peer.last_seen) > timeout {
debug!("Removing stale mDNS peer: {}", entity_id);
false
} else {
true
}
});
let removed = initial_count - peers_guard.len();
if removed > 0 {
info!("Cleaned up {} stale mDNS peers", removed);
}
}
});
}
fn build_txt_records(&self) -> HashMap<String, String> {
let mut txt_records = HashMap::new();
txt_records.insert("version".to_string(), "1.0".to_string());
txt_records.insert("protocol".to_string(), "synapse".to_string());
txt_records.insert("entity_id".to_string(), self.entity_id.clone());
txt_records.insert("capabilities".to_string(), "tcp,encrypted,direct".to_string());
txt_records.insert("transport_types".to_string(), "tcp,udp,email".to_string());
txt_records
}
async fn send_ptr_record(&self, announcement: &ServiceAnnouncement) -> Result<()> {
debug!("Sending PTR record for {}", announcement.instance_name);
Ok(())
}
async fn send_srv_record(&self, announcement: &ServiceAnnouncement) -> Result<()> {
debug!("Sending SRV record for {}", announcement.instance_name);
Ok(())
}
async fn send_txt_record(&self, announcement: &ServiceAnnouncement) -> Result<()> {
debug!("Sending TXT record for {}", announcement.instance_name);
Ok(())
}
async fn send_host_records(&self, announcement: &ServiceAnnouncement) -> Result<()> {
debug!("Sending host records for {}", announcement.host_name);
Ok(())
}
async fn send_ptr_query(&self, service_type: &str) -> Result<()> {
debug!("Sending PTR query for {}", service_type);
Ok(())
}
async fn send_periodic_announcements(
_socket: &Arc<Mutex<TokioUdpSocket>>,
_entity_id: &str,
_config: &MdnsConfig,
) -> Result<()> {
Ok(())
}
async fn send_discovery_queries(
_socket: &Arc<Mutex<TokioUdpSocket>>,
_config: &MdnsConfig,
) -> Result<()> {
Ok(())
}
async fn process_mdns_packet(
_packet_data: &[u8],
_src: SocketAddr,
_peers: &Arc<RwLock<HashMap<String, EnhancedMdnsPeer>>>,
) -> Result<()> {
Ok(())
}
}
impl EnhancedMdnsTransport {
pub fn get_circuit_breaker(&self) -> Arc<CircuitBreaker> {
self.circuit_breaker.clone()
}
pub fn get_circuit_stats(&self) -> crate::circuit_breaker::CircuitStats {
self.circuit_breaker.get_stats()
}
pub async fn is_circuit_open(&self) -> bool {
!self.circuit_breaker.can_proceed().await
}
}
#[cfg(feature = "mdns")]
async fn create_multicast_socket(config: &MdnsConfig) -> Result<TokioUdpSocket> {
use std::net::{SocketAddrV4, Ipv4Addr};
use socket2::{Socket, Domain, Type, Protocol};
let socket = Socket::new(Domain::IPV4, Type::DGRAM, Some(Protocol::UDP))?;
socket.set_reuse_address(true)?;
#[cfg(not(windows))]
socket.set_reuse_port(true)?;
let bind_addr = SocketAddrV4::new(Ipv4Addr::UNSPECIFIED, config.multicast_port);
socket.bind(&bind_addr.into())?;
let std_socket: std::net::UdpSocket = socket.into();
std_socket.set_nonblocking(true)?;
let tokio_socket = TokioUdpSocket::from_std(std_socket)?;
Ok(tokio_socket)
}
#[cfg(not(feature = "mdns"))]
async fn create_multicast_socket(config: &MdnsConfig) -> Result<TokioUdpSocket> {
use std::net::{SocketAddrV4, Ipv4Addr};
let bind_addr = SocketAddrV4::new(Ipv4Addr::UNSPECIFIED, config.multicast_port);
let socket = TokioUdpSocket::bind(bind_addr).await?;
socket.join_multicast_v4(config.multicast_addr, Ipv4Addr::UNSPECIFIED)?;
socket.set_multicast_loop_v4(false)?;
socket.set_multicast_ttl_v4(255)?;
info!("Created mDNS multicast socket on port {}", config.multicast_port);
Ok(socket)
}
pub mod utils {
use super::*;
pub fn parse_dns_name(data: &[u8], offset: usize) -> Result<(String, usize)> {
let mut name = String::new();
let mut pos = offset;
let mut jumped = false;
let mut jump_count = 0;
loop {
if pos >= data.len() {
return Err(crate::error::SynapseError::TransportError(
"DNS name parsing: unexpected end of data".to_string()
).into());
}
let len = data[pos] as usize;
if len & 0xC0 == 0xC0 {
if !jumped {
}
let pointer = ((len & 0x3F) << 8) | (data[pos + 1] as usize);
pos = pointer;
jumped = true;
jump_count += 1;
if jump_count > 10 {
return Err(crate::error::SynapseError::TransportError(
"DNS name parsing: too many jumps".to_string()
).into());
}
continue;
}
pos += 1;
if len == 0 {
break;
}
if pos + len > data.len() {
return Err(crate::error::SynapseError::TransportError(
"DNS name parsing: label too long".to_string()
).into());
}
if !name.is_empty() {
name.push('.');
}
name.push_str(&String::from_utf8_lossy(&data[pos..pos + len]));
pos += len;
}
Ok((name, pos))
}
pub fn encode_dns_name(name: &str) -> Vec<u8> {
let mut encoded = Vec::new();
for label in name.split('.') {
if label.is_empty() {
continue;
}
encoded.push(label.len() as u8);
encoded.extend_from_slice(label.as_bytes());
}
encoded.push(0); encoded
}
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn test_enhanced_mdns_creation() {
let transport = EnhancedMdnsTransport::new(
"test_entity".to_string(),
8080,
None,
).await;
assert!(transport.is_ok());
}
#[test]
fn test_dns_name_encoding() {
let name = "test._synapse._tcp.local";
let encoded = utils::encode_dns_name(name);
assert_eq!(encoded[0], 4); assert_eq!(&encoded[1..5], b"test");
assert_eq!(encoded[5], 8); }
}
pub struct EnhancedMdnsServiceBrowser {
service_types: Vec<String>,
service_cache: Arc<RwLock<HashMap<String, ServiceRecord>>>,
config: BrowserConfig,
browse_socket: Arc<Mutex<TokioUdpSocket>>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ServiceRecord {
pub service_name: String,
pub service_type: String,
pub domain: String,
pub host_name: String,
pub addresses: Vec<IpAddr>,
pub port: u16,
pub txt_records: HashMap<String, String>,
pub priority: u16,
pub weight: u16,
pub ttl: u32,
#[serde(skip, default = "Instant::now")]
pub discovered_at: Instant,
#[serde(skip, default = "Instant::now")]
pub last_updated: Instant,
pub service_state: ServiceState,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub enum ServiceState {
Discovering,
Resolving,
Active,
Inactive,
Expired,
}
#[derive(Debug, Clone)]
pub struct BrowserConfig {
pub browse_interval: Duration,
pub cache_ttl: Duration,
pub max_cache_size: usize,
pub continuous_monitoring: bool,
}
impl Default for BrowserConfig {
fn default() -> Self {
Self {
browse_interval: Duration::from_secs(30),
cache_ttl: Duration::from_secs(300),
max_cache_size: 500,
continuous_monitoring: true,
}
}
}
impl EnhancedMdnsServiceBrowser {
pub async fn new(service_types: Vec<String>, config: Option<BrowserConfig>) -> Result<Self> {
let config = config.unwrap_or_default();
let mdns_config = MdnsConfig::default();
let browse_socket = create_multicast_socket(&mdns_config).await?;
Ok(Self {
service_types,
service_cache: Arc::new(RwLock::new(HashMap::new())),
config,
browse_socket: Arc::new(Mutex::new(browse_socket)),
})
}
pub async fn start_browsing(&self) -> Result<()> {
info!("Starting mDNS service browsing for types: {:?}", self.service_types);
for service_type in &self.service_types {
self.start_service_type_browsing(service_type.clone()).await?;
}
self.start_cache_cleanup().await;
Ok(())
}
async fn start_service_type_browsing(&self, service_type: String) -> Result<()> {
let socket = self.browse_socket.clone();
let _cache = self.service_cache.clone();
let interval_duration = self.config.browse_interval;
tokio::spawn(async move {
let mut interval_timer = interval(interval_duration);
loop {
interval_timer.tick().await;
if let Err(e) = Self::send_service_browse_query(&socket, &service_type).await {
error!("Failed to send browse query for {}: {}", service_type, e);
}
}
});
Ok(())
}
async fn send_service_browse_query(
socket: &Arc<Mutex<TokioUdpSocket>>,
service_type: &str
) -> Result<()> {
let query_packet = MdnsPacket::Query {
questions: vec![MdnsQuestion {
name: service_type.to_string(),
qtype: 12, qclass: 1, }],
transaction_id: rand::random(),
};
let packet_bytes = Self::encode_mdns_packet(&query_packet)?;
let socket_guard = socket.lock().await;
socket_guard.send_to(&packet_bytes, "224.0.0.251:5353").await?;
Ok(())
}
pub async fn get_discovered_services(&self) -> Vec<ServiceRecord> {
self.service_cache.read().unwrap().values().cloned().collect()
}
pub async fn get_services_by_type(&self, service_type: &str) -> Vec<ServiceRecord> {
self.service_cache
.read()
.unwrap()
.values()
.filter(|record| record.service_type == service_type)
.cloned()
.collect()
}
pub async fn find_services_by_capability(&self, capability: &str) -> Vec<ServiceRecord> {
self.service_cache
.read()
.unwrap()
.values()
.filter(|record| {
record.txt_records.get("capabilities")
.map(|caps| caps.contains(capability))
.unwrap_or(false)
})
.cloned()
.collect()
}
async fn start_cache_cleanup(&self) {
let cache = self.service_cache.clone();
let ttl = self.config.cache_ttl;
let max_size = self.config.max_cache_size;
tokio::spawn(async move {
let mut cleanup_interval = interval(Duration::from_secs(60));
loop {
cleanup_interval.tick().await;
let mut cache_guard = cache.write().unwrap();
let now = Instant::now();
cache_guard.retain(|_, record| {
now.duration_since(record.last_updated) < ttl
});
if cache_guard.len() > max_size {
let mut entries: Vec<_> = cache_guard.iter()
.map(|(k, v)| (k.clone(), v.last_updated))
.collect();
entries.sort_by_key(|(_, last_updated)| *last_updated);
let to_remove = cache_guard.len() - max_size;
for (key, _) in entries.iter().take(to_remove) {
cache_guard.remove(key);
}
}
}
});
}
fn encode_mdns_packet(packet: &MdnsPacket) -> Result<Vec<u8>> {
let mut bytes = Vec::new();
match packet {
MdnsPacket::Query { questions, transaction_id } => {
bytes.extend_from_slice(&transaction_id.to_be_bytes());
bytes.extend_from_slice(&[0x00, 0x00]); bytes.extend_from_slice(&(questions.len() as u16).to_be_bytes());
bytes.extend_from_slice(&[0x00, 0x00]); bytes.extend_from_slice(&[0x00, 0x00]); bytes.extend_from_slice(&[0x00, 0x00]);
for question in questions {
bytes.extend_from_slice(&utils::encode_dns_name(&question.name));
bytes.extend_from_slice(&question.qtype.to_be_bytes());
bytes.extend_from_slice(&question.qclass.to_be_bytes());
}
}
MdnsPacket::Response { .. } => {
return Err(crate::error::SynapseError::TransportError(
"Response packet encoding not implemented yet".to_string()
).into());
}
}
Ok(bytes)
}
}
pub struct EnhancedMdnsResponder {
our_services: Vec<ServiceRecord>,
config: ResponderConfig,
response_socket: Arc<Mutex<TokioUdpSocket>>,
}
#[derive(Debug, Clone)]
pub struct ResponderConfig {
pub announce_interval: Duration,
pub record_ttl: u32,
pub immediate_response: bool,
}
impl Default for ResponderConfig {
fn default() -> Self {
Self {
announce_interval: Duration::from_secs(120),
record_ttl: 120,
immediate_response: true,
}
}
}
impl EnhancedMdnsResponder {
pub async fn new(config: Option<ResponderConfig>) -> Result<Self> {
let config = config.unwrap_or_default();
let mdns_config = MdnsConfig::default();
let response_socket = create_multicast_socket(&mdns_config).await?;
Ok(Self {
our_services: Vec::new(),
config,
response_socket: Arc::new(Mutex::new(response_socket)),
})
}
pub async fn add_service(&mut self, service: ServiceRecord) -> Result<()> {
info!("Adding service for announcement: {}", service.service_name);
self.our_services.push(service);
Ok(())
}
pub async fn start_responding(&self) -> Result<()> {
info!("Starting mDNS responder for {} services", self.our_services.len());
self.start_periodic_announcements().await?;
self.start_query_response_handler().await?;
Ok(())
}
async fn start_periodic_announcements(&self) -> Result<()> {
let socket = self.response_socket.clone();
let services = self.our_services.clone();
let interval_duration = self.config.announce_interval;
tokio::spawn(async move {
let mut announcement_interval = interval(interval_duration);
loop {
announcement_interval.tick().await;
for service in &services {
if let Err(e) = Self::announce_service(&socket, service).await {
error!("Failed to announce service {}: {}", service.service_name, e);
}
}
}
});
Ok(())
}
async fn announce_service(
_socket: &Arc<Mutex<TokioUdpSocket>>,
service: &ServiceRecord
) -> Result<()> {
debug!("Announcing service: {}", service.service_name);
info!("Announced service {} on port {}", service.service_name, service.port);
Ok(())
}
async fn start_query_response_handler(&self) -> Result<()> {
let socket = self.response_socket.clone();
let services = self.our_services.clone();
tokio::spawn(async move {
let mut buffer = [0u8; 4096];
loop {
let socket_guard = socket.lock().await;
match socket_guard.recv_from(&mut buffer).await {
Ok((size, src)) => {
drop(socket_guard);
if let Err(e) = Self::handle_query(&socket, &buffer[..size], src, &services).await {
debug!("Error handling mDNS query from {}: {}", src, e);
}
}
Err(e) => {
error!("Error receiving mDNS query: {}", e);
tokio::time::sleep(Duration::from_millis(100)).await;
}
}
}
});
Ok(())
}
async fn handle_query(
_socket: &Arc<Mutex<TokioUdpSocket>>,
_packet_data: &[u8],
_src: SocketAddr,
services: &[ServiceRecord]
) -> Result<()> {
debug!("Received mDNS query from client, {} services available", services.len());
Ok(())
}
}
impl EnhancedMdnsTransport {
pub async fn create_service_browser(&self) -> Result<EnhancedMdnsServiceBrowser> {
EnhancedMdnsServiceBrowser::new(
vec![
"_synapse._tcp.local.".to_string(),
"_synapse._tcp.local.".to_string(),
"_synapse-router._tcp.local.".to_string(),
],
None
).await
}
pub async fn create_service_responder(&self) -> Result<EnhancedMdnsResponder> {
let mut responder = EnhancedMdnsResponder::new(None).await?;
let our_service = ServiceRecord {
service_name: self.instance_name.clone(),
service_type: self.service_type.clone(),
domain: "local.".to_string(),
host_name: format!("{}.local.", self.entity_id),
addresses: vec![], port: self.local_port,
txt_records: HashMap::from([
("entity_id".to_string(), self.entity_id.clone()),
("version".to_string(), "1.0".to_string()),
("capabilities".to_string(), "routing,discovery,secure_messaging".to_string()),
]),
priority: 10,
weight: 5,
ttl: 120,
discovered_at: Instant::now(),
last_updated: Instant::now(),
service_state: ServiceState::Active,
};
responder.add_service(our_service).await?;
Ok(responder)
}
pub async fn comprehensive_discovery(&self) -> Result<Vec<EnhancedMdnsPeer>> {
let browser = self.create_service_browser().await?;
browser.start_browsing().await?;
tokio::time::sleep(Duration::from_secs(3)).await;
let services = browser.get_discovered_services().await;
let mut peers = Vec::new();
for service in services {
if service.service_state == ServiceState::Active {
let capabilities = service.txt_records.get("capabilities")
.map(|caps| caps.split(',').map(|s| s.trim().to_string()).collect())
.unwrap_or_default();
let protocol_version = service.txt_records.get("version")
.cloned()
.unwrap_or_else(|| "1.0".to_string());
let entity_id = service.txt_records.get("entity_id")
.cloned()
.unwrap_or_else(|| service.service_name.clone());
let peer = EnhancedMdnsPeer {
entity_id,
instance_name: service.service_name,
service_type: service.service_type,
host_name: service.host_name,
addresses: service.addresses,
port: service.port,
txt_records: service.txt_records,
priority: service.priority,
weight: service.weight,
ttl: service.ttl,
discovered_at: service.discovered_at,
last_seen: service.last_updated,
capabilities,
protocol_version,
};
peers.push(peer);
}
}
Ok(peers)
}
}