use std::{marker::PhantomData, net::IpAddr};
use futures_core::Stream;
use futures_util::StreamExt;
use mdns_sd::{ServiceDaemon, ServiceEvent, ServiceInfo};
use tracing::debug;
use wot_td::{
extend::{Extend, ExtendablePiece, ExtendableThing},
hlist::Nil,
thing::Thing,
};
#[derive(thiserror::Error, Debug)]
#[non_exhaustive]
pub enum Error {
#[error("mdns cannot be accessed {0}")]
Mdns(#[from] mdns_sd::Error),
#[error("reqwest error {0}")]
Reqwest(#[from] reqwest::Error),
#[error("Missing address")]
NoAddress,
}
pub type Result<T> = std::result::Result<T, Error>;
const WELL_KNOWN: &str = "/.well-known/wot";
pub struct Discoverer<Other: ExtendableThing + ExtendablePiece = Nil> {
mdns: ServiceDaemon,
service_type: String,
_other: PhantomData<Other>,
}
pub struct Discovered<Other: ExtendableThing + ExtendablePiece> {
pub thing: Thing<Other>,
info: ServiceInfo,
scheme: String,
}
impl<Other: ExtendableThing + ExtendablePiece> Discovered<Other> {
pub fn get_addresses(&self) -> Vec<IpAddr> {
self.info
.get_addresses()
.iter()
.map(|ip| ip.to_owned().into())
.collect()
}
pub fn get_port(&self) -> u16 {
self.info.get_port()
}
pub fn get_hostname(&self) -> &str {
self.info.get_hostname()
}
pub fn get_scheme(&self) -> &str {
&self.scheme
}
}
async fn get_thing<Other: ExtendableThing + ExtendablePiece>(
info: ServiceInfo,
) -> Result<Discovered<Other>> {
let host = info.get_addresses().iter().next().ok_or(Error::NoAddress)?;
let port = info.get_port();
let props = info.get_properties();
let path = props.get_property_val_str("td").unwrap_or(WELL_KNOWN);
let proto = props
.get_property_val_str("scheme")
.or_else(|| {
props
.get_property_val_str("tls")
.map(|tls| if tls == "1" { "https" } else { "http" })
})
.unwrap_or("http");
debug!("Got {proto} {host} {port} {path}");
let r = reqwest::get(format!("{proto}://{host}:{port}{path}")).await?;
let thing = r.json().await?;
let scheme = proto.to_owned();
let d = Discovered {
thing,
info,
scheme,
};
Ok(d)
}
impl Discoverer {
pub fn new() -> Result<Self> {
let mdns = ServiceDaemon::new()?;
let service_type = "_wot._tcp.local.".to_owned();
Ok(Self {
mdns,
service_type,
_other: PhantomData,
})
}
}
impl<Other: ExtendableThing + ExtendablePiece> Discoverer<Other> {
pub fn ext<T>(self) -> Discoverer<Other::Target>
where
Other: Extend<T>,
Other::Target: ExtendableThing + ExtendablePiece,
{
let Discoverer {
mdns,
service_type,
_other,
} = self;
Discoverer {
mdns,
service_type,
_other: PhantomData,
}
}
pub fn stream(&self) -> Result<impl Stream<Item = Result<Discovered<Other>>>> {
let receiver = self.mdns.browse(&self.service_type)?;
let s = receiver.into_stream().filter_map(|v| async move {
tracing::info!("{:?}", v);
if let ServiceEvent::ServiceResolved(info) = v {
let t = get_thing(info).await;
Some(t)
} else {
None
}
});
Ok(s)
}
}