use std::sync::Arc;
use iroh_base::{NodeId, RelayUrl, SecretKey};
use iroh_relay::node_info::{EncodingError, NodeInfo};
use n0_future::{
boxed::BoxStream,
task::{self, AbortOnDropHandle},
time::{self, Duration, Instant},
};
use n0_watcher::{Disconnected, Watchable, Watcher as _};
use pkarr::{
SignedPacket,
errors::{PublicKeyError, SignedPacketVerifyError},
};
use snafu::{ResultExt, Snafu};
use tracing::{Instrument, debug, error_span, warn};
use url::Url;
use super::{DiscoveryContext, DiscoveryError, IntoDiscovery, IntoDiscoveryError};
#[cfg(not(wasm_browser))]
use crate::dns::DnsResolver;
use crate::{
discovery::{Discovery, DiscoveryItem, NodeData},
endpoint::force_staging_infra,
};
#[cfg(feature = "discovery-pkarr-dht")]
pub mod dht;
#[allow(missing_docs)]
#[derive(Debug, Snafu)]
#[non_exhaustive]
pub enum PkarrError {
#[snafu(display("Invalid public key"))]
PublicKey { source: PublicKeyError },
#[snafu(display("Packet failed to verify"))]
Verify { source: SignedPacketVerifyError },
#[snafu(display("Invalid relay URL"))]
InvalidRelayUrl { url: RelayUrl },
#[snafu(display("Error sending http request"))]
HttpSend { source: reqwest::Error },
#[snafu(display("Error resolving http request"))]
HttpRequest { status: reqwest::StatusCode },
#[snafu(display("Http payload error"))]
HttpPayload { source: reqwest::Error },
#[snafu(display("EncodingError"))]
Encoding { source: EncodingError },
}
impl From<PkarrError> for DiscoveryError {
fn from(err: PkarrError) -> Self {
DiscoveryError::from_err("pkarr", err)
}
}
pub const N0_DNS_PKARR_RELAY_PROD: &str = "https://dns.iroh.link/pkarr";
pub const N0_DNS_PKARR_RELAY_STAGING: &str = "https://staging-dns.iroh.link/pkarr";
pub const DEFAULT_PKARR_TTL: u32 = 30;
pub const DEFAULT_REPUBLISH_INTERVAL: Duration = Duration::from_secs(60 * 5);
#[derive(Debug)]
pub struct PkarrPublisherBuilder {
pkarr_relay: Url,
ttl: u32,
republish_interval: Duration,
#[cfg(not(wasm_browser))]
dns_resolver: Option<DnsResolver>,
}
impl PkarrPublisherBuilder {
fn new(pkarr_relay: Url) -> Self {
Self {
pkarr_relay,
ttl: DEFAULT_PKARR_TTL,
republish_interval: DEFAULT_REPUBLISH_INTERVAL,
#[cfg(not(wasm_browser))]
dns_resolver: None,
}
}
fn n0_dns() -> Self {
let pkarr_relay = match force_staging_infra() {
true => N0_DNS_PKARR_RELAY_STAGING,
false => N0_DNS_PKARR_RELAY_PROD,
};
let pkarr_relay: Url = pkarr_relay.parse().expect("url is valid");
Self::new(pkarr_relay)
}
pub fn ttl(mut self, ttl: u32) -> Self {
self.ttl = ttl;
self
}
pub fn republish_interval(mut self, republish_interval: Duration) -> Self {
self.republish_interval = republish_interval;
self
}
#[cfg(not(wasm_browser))]
pub fn dns_resolver(mut self, dns_resolver: DnsResolver) -> Self {
self.dns_resolver = Some(dns_resolver);
self
}
pub fn build(self, secret_key: SecretKey) -> PkarrPublisher {
PkarrPublisher::new(
secret_key,
self.pkarr_relay,
self.ttl,
self.republish_interval,
#[cfg(not(wasm_browser))]
self.dns_resolver,
)
}
}
impl IntoDiscovery for PkarrPublisherBuilder {
fn into_discovery(
mut self,
context: &DiscoveryContext,
) -> Result<impl Discovery, IntoDiscoveryError> {
#[cfg(not(wasm_browser))]
if self.dns_resolver.is_none() {
self.dns_resolver = Some(context.dns_resolver().clone());
}
Ok(self.build(context.secret_key().clone()))
}
}
#[derive(derive_more::Debug, Clone)]
pub struct PkarrPublisher {
node_id: NodeId,
watchable: Watchable<Option<NodeInfo>>,
_drop_guard: Arc<AbortOnDropHandle<()>>,
}
impl PkarrPublisher {
pub fn builder(pkarr_relay: Url) -> PkarrPublisherBuilder {
PkarrPublisherBuilder::new(pkarr_relay)
}
fn new(
secret_key: SecretKey,
pkarr_relay: Url,
ttl: u32,
republish_interval: Duration,
#[cfg(not(wasm_browser))] dns_resolver: Option<DnsResolver>,
) -> Self {
debug!("creating pkarr publisher that publishes to {pkarr_relay}");
let node_id = secret_key.public();
#[cfg(wasm_browser)]
let pkarr_client = PkarrRelayClient::new(pkarr_relay);
#[cfg(not(wasm_browser))]
let pkarr_client = if let Some(dns_resolver) = dns_resolver {
PkarrRelayClient::with_dns_resolver(pkarr_relay, dns_resolver)
} else {
PkarrRelayClient::new(pkarr_relay)
};
let watchable = Watchable::default();
let service = PublisherService {
ttl,
watcher: watchable.watch(),
secret_key,
pkarr_client,
republish_interval,
};
let join_handle = task::spawn(
service
.run()
.instrument(error_span!("pkarr_publish", me=%node_id.fmt_short())),
);
Self {
watchable,
node_id,
_drop_guard: Arc::new(AbortOnDropHandle::new(join_handle)),
}
}
pub fn n0_dns() -> PkarrPublisherBuilder {
PkarrPublisherBuilder::n0_dns()
}
pub fn update_node_data(&self, data: &NodeData) {
let mut data = data.clone();
if data.relay_url().is_some() {
data.clear_direct_addresses();
}
let info = NodeInfo::from_parts(self.node_id, data);
self.watchable.set(Some(info)).ok();
}
}
impl Discovery for PkarrPublisher {
fn publish(&self, data: &NodeData) {
self.update_node_data(data);
}
}
#[derive(derive_more::Debug, Clone)]
struct PublisherService {
#[debug("SecretKey")]
secret_key: SecretKey,
#[debug("PkarrClient")]
pkarr_client: PkarrRelayClient,
watcher: n0_watcher::Direct<Option<NodeInfo>>,
ttl: u32,
republish_interval: Duration,
}
impl PublisherService {
async fn run(mut self) {
let mut failed_attempts = 0;
let republish = time::sleep(Duration::MAX);
tokio::pin!(republish);
loop {
if !self.watcher.is_connected() {
break;
}
if let Some(info) = self.watcher.get() {
match self.publish_current(info).await {
Err(err) => {
failed_attempts += 1;
let retry_after = Duration::from_secs(failed_attempts);
republish.as_mut().reset(Instant::now() + retry_after);
warn!(
err = %format!("{err:#}"),
url = %self.pkarr_client.pkarr_relay_url ,
?retry_after,
%failed_attempts,
"Failed to publish to pkarr",
);
}
_ => {
failed_attempts = 0;
republish
.as_mut()
.reset(Instant::now() + self.republish_interval);
}
}
}
tokio::select! {
res = self.watcher.updated() => match res {
Ok(_) => debug!("Publish node info to pkarr (info changed)"),
Err(Disconnected { .. }) => break,
},
_ = &mut republish => debug!("Publish node info to pkarr (interval elapsed)"),
}
}
}
async fn publish_current(&self, info: NodeInfo) -> Result<(), PkarrError> {
debug!(
data = ?info.data,
pkarr_relay = %self.pkarr_client.pkarr_relay_url,
"Publish node info to pkarr"
);
let signed_packet = info
.to_pkarr_signed_packet(&self.secret_key, self.ttl)
.context(EncodingSnafu)?;
self.pkarr_client.publish(&signed_packet).await?;
Ok(())
}
}
#[derive(Debug)]
pub struct PkarrResolverBuilder {
pkarr_relay: Url,
#[cfg(not(wasm_browser))]
dns_resolver: Option<DnsResolver>,
}
impl PkarrResolverBuilder {
#[cfg(not(wasm_browser))]
pub fn dns_resolver(mut self, dns_resolver: DnsResolver) -> Self {
self.dns_resolver = Some(dns_resolver);
self
}
pub fn build(self) -> PkarrResolver {
#[cfg(wasm_browser)]
let pkarr_client = PkarrRelayClient::new(self.pkarr_relay);
#[cfg(not(wasm_browser))]
let pkarr_client = if let Some(dns_resolver) = self.dns_resolver {
PkarrRelayClient::with_dns_resolver(self.pkarr_relay, dns_resolver)
} else {
PkarrRelayClient::new(self.pkarr_relay)
};
PkarrResolver { pkarr_client }
}
}
impl IntoDiscovery for PkarrResolverBuilder {
fn into_discovery(
mut self,
context: &DiscoveryContext,
) -> Result<impl Discovery, IntoDiscoveryError> {
#[cfg(not(wasm_browser))]
if self.dns_resolver.is_none() {
self.dns_resolver = Some(context.dns_resolver().clone());
}
Ok(self.build())
}
}
#[derive(derive_more::Debug, Clone)]
pub struct PkarrResolver {
pkarr_client: PkarrRelayClient,
}
impl PkarrResolver {
pub fn builder(pkarr_relay: Url) -> PkarrResolverBuilder {
PkarrResolverBuilder {
pkarr_relay,
#[cfg(not(wasm_browser))]
dns_resolver: None,
}
}
pub fn n0_dns() -> PkarrResolverBuilder {
let pkarr_relay = match force_staging_infra() {
true => N0_DNS_PKARR_RELAY_STAGING,
false => N0_DNS_PKARR_RELAY_PROD,
};
let pkarr_relay: Url = pkarr_relay.parse().expect("url is valid");
Self::builder(pkarr_relay)
}
}
impl Discovery for PkarrResolver {
fn resolve(&self, node_id: NodeId) -> Option<BoxStream<Result<DiscoveryItem, DiscoveryError>>> {
let pkarr_client = self.pkarr_client.clone();
let fut = async move {
let signed_packet = pkarr_client.resolve(node_id).await?;
let info = NodeInfo::from_pkarr_signed_packet(&signed_packet)
.map_err(|err| DiscoveryError::from_err("pkarr", err))?;
let item = DiscoveryItem::new(info, "pkarr", None);
Ok(item)
};
let stream = n0_future::stream::once_future(fut);
Some(Box::pin(stream))
}
}
#[derive(Debug, Clone)]
pub struct PkarrRelayClient {
http_client: reqwest::Client,
pkarr_relay_url: Url,
}
impl PkarrRelayClient {
pub fn new(pkarr_relay_url: Url) -> Self {
Self {
http_client: reqwest::Client::new(),
pkarr_relay_url,
}
}
#[cfg(not(wasm_browser))]
pub fn with_dns_resolver(pkarr_relay_url: Url, dns_resolver: crate::dns::DnsResolver) -> Self {
let http_client = reqwest::Client::builder()
.dns_resolver(Arc::new(dns_resolver))
.build()
.expect("failed to create request client");
Self {
http_client,
pkarr_relay_url,
}
}
pub async fn resolve(&self, node_id: NodeId) -> Result<SignedPacket, DiscoveryError> {
let public_key = pkarr::PublicKey::try_from(node_id.as_bytes()).context(PublicKeySnafu)?;
let mut url = self.pkarr_relay_url.clone();
url.path_segments_mut()
.map_err(|_| {
InvalidRelayUrlSnafu {
url: self.pkarr_relay_url.clone(),
}
.build()
})?
.push(&public_key.to_z32());
let response = self
.http_client
.get(url)
.send()
.await
.context(HttpSendSnafu)?;
if !response.status().is_success() {
return Err(HttpRequestSnafu {
status: response.status(),
}
.build()
.into());
}
let payload = response.bytes().await.context(HttpPayloadSnafu)?;
let packet =
SignedPacket::from_relay_payload(&public_key, &payload).context(VerifySnafu)?;
Ok(packet)
}
pub async fn publish(&self, signed_packet: &SignedPacket) -> Result<(), PkarrError> {
let mut url = self.pkarr_relay_url.clone();
url.path_segments_mut()
.map_err(|_| {
InvalidRelayUrlSnafu {
url: self.pkarr_relay_url.clone(),
}
.build()
})?
.push(&signed_packet.public_key().to_z32());
let response = self
.http_client
.put(url)
.body(signed_packet.to_relay_payload())
.send()
.await
.context(HttpSendSnafu)?;
if !response.status().is_success() {
return Err(HttpRequestSnafu {
status: response.status(),
}
.build());
}
Ok(())
}
}