use std::sync::atomic::Ordering;
use iroh::endpoint::TransportAddrUsage;
use super::{
session_runtime::PathSubscriptions,
stats::{EndpointStats, NodeAddrInfo, PathInfo, PeerStats},
IrohEndpoint,
};
impl IrohEndpoint {
pub fn endpoint_stats(&self) -> EndpointStats {
let (active_readers, active_writers, active_sessions, total_handles) =
self.inner.ffi.handles.count_handles();
let pool_size = self.inner.http.pool.entry_count_approx() as usize;
let active_connections = self.inner.http.active_connections.load(Ordering::Relaxed);
let active_requests = self.inner.http.active_requests.load(Ordering::Relaxed);
let active_path_subscriptions = self
.inner
.session
.path_subs
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.len();
let active_path_watchers = self
.inner
.session
.active_path_watchers
.load(Ordering::Relaxed);
EndpointStats {
active_readers,
active_writers,
active_sessions,
total_handles,
pool_size,
active_connections,
active_requests,
active_path_subscriptions,
active_path_watchers,
}
}
pub fn bound_sockets(&self) -> Vec<std::net::SocketAddr> {
self.inner.transport.ep.bound_sockets()
}
pub fn transport_alive(&self) -> bool {
!self.inner.transport.ep.is_closed() && !self.bound_sockets().is_empty()
}
pub fn direct_socket_addrs(&self) -> Vec<std::net::SocketAddr> {
let candidates: Vec<std::net::SocketAddr> =
self.inner.transport.ep.addr().ip_addrs().copied().collect();
let bound = self.bound_sockets();
super::bind::reconcile_direct_addr_ports(&candidates, &bound)
}
pub fn node_addr(&self) -> NodeAddrInfo {
let addr = self.inner.transport.ep.addr();
let mut addrs = Vec::new();
for relay in addr.relay_urls() {
addrs.push(relay.to_string());
}
for da in self.direct_socket_addrs() {
addrs.push(da.to_string());
}
NodeAddrInfo {
id: self.inner.transport.node_id_str.clone(),
addrs,
}
}
pub fn dialable_direct_address(&self) -> Option<String> {
select_dialable_direct(&self.direct_socket_addrs())
}
pub fn dialable_direct_addresses(&self) -> Vec<String> {
select_dialable_directs(&self.direct_socket_addrs())
}
pub fn home_relay(&self) -> Option<String> {
self.inner
.transport
.ep
.addr()
.relay_urls()
.next()
.map(|u| u.to_string())
}
pub async fn peer_info(&self, node_id_b32: &str) -> Option<NodeAddrInfo> {
let bytes = crate::base32_decode(node_id_b32).ok()?;
let arr: [u8; 32] = bytes.try_into().ok()?;
let pk = iroh::PublicKey::from_bytes(&arr).ok()?;
let info = self.inner.transport.ep.remote_info(pk).await?;
let id = crate::base32_encode(info.id().as_bytes());
let mut addrs = Vec::new();
for a in info.addrs() {
match a.addr() {
iroh::TransportAddr::Ip(sock) => addrs.push(sock.to_string()),
iroh::TransportAddr::Relay(url) => addrs.push(url.to_string()),
other => addrs.push(format!("{:?}", other)),
}
}
Some(NodeAddrInfo { id, addrs })
}
pub async fn peer_stats(&self, node_id_b32: &str) -> Option<PeerStats> {
let bytes = crate::base32_decode(node_id_b32).ok()?;
let arr: [u8; 32] = bytes.try_into().ok()?;
let pk = iroh::PublicKey::from_bytes(&arr).ok()?;
let info = self.inner.transport.ep.remote_info(pk).await?;
let mut paths = Vec::new();
let mut has_active_relay = false;
let mut active_relay_url: Option<String> = None;
for a in info.addrs() {
let is_relay = a.addr().is_relay();
let is_active = matches!(a.usage(), TransportAddrUsage::Active);
let addr_str = match a.addr() {
iroh::TransportAddr::Ip(sock) => sock.to_string(),
iroh::TransportAddr::Relay(url) => {
if is_active {
has_active_relay = true;
active_relay_url = Some(url.to_string());
}
url.to_string()
}
other => format!("{:?}", other),
};
paths.push(PathInfo {
relay: is_relay,
addr: addr_str,
active: is_active,
});
}
let (rtt_ms, bytes_sent, bytes_received, lost_packets, sent_packets, congestion_window) =
if let Some(pooled) = self.inner.http.pool.get_existing(pk, crate::ALPN).await {
let s = pooled.conn.stats();
let rtt = pooled.conn.rtt(iroh::endpoint::PathId::ZERO);
(
rtt.map(|d| d.as_secs_f64() * 1000.0),
Some(s.udp_tx.bytes),
Some(s.udp_rx.bytes),
None,
None,
None,
)
} else {
(None, None, None, None, None, None)
};
Some(PeerStats {
relay: has_active_relay,
relay_url: active_relay_url,
paths,
rtt_ms,
bytes_sent,
bytes_received,
lost_packets,
sent_packets,
congestion_window,
})
}
pub fn subscribe_path_changes(
&self,
node_id_str: &str,
subscription_id: u32,
) -> tokio::sync::mpsc::UnboundedReceiver<PathInfo> {
let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
let (peer_subscriptions, had_existing_watcher) = {
let mut subscriptions = self
.inner
.session
.path_subs
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let (peer_subscriptions, had_existing) = match subscriptions
.entry(node_id_str.to_string())
{
std::collections::hash_map::Entry::Occupied(entry) => (entry.get().clone(), true),
std::collections::hash_map::Entry::Vacant(entry) => {
let peer_subscriptions = std::sync::Arc::new(PathSubscriptions {
senders: std::sync::Mutex::new(std::collections::HashMap::new()),
});
entry.insert(peer_subscriptions.clone());
(peer_subscriptions, false)
}
};
peer_subscriptions
.senders
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(subscription_id, tx);
(peer_subscriptions, had_existing)
};
if had_existing_watcher {
return rx;
}
let ep = self.clone();
let nid = node_id_str.to_string();
let event_tx = self.inner.session.event_tx.clone();
self.inner
.session
.active_path_watchers
.fetch_add(1, Ordering::Relaxed);
tokio::spawn(async move {
let mut last_key: Option<String> = None;
let mut closed_rx = ep.inner.session.closed_rx.clone();
loop {
if *closed_rx.borrow() {
let mut subscriptions = ep
.inner
.session
.path_subs
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if subscriptions
.get(&nid)
.is_some_and(|current| std::sync::Arc::ptr_eq(current, &peer_subscriptions))
{
subscriptions.remove(&nid);
}
break;
}
let is_closed = {
let mut subscriptions = ep
.inner
.session
.path_subs
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let Some(current) = subscriptions.get(&nid) else {
break;
};
if !std::sync::Arc::ptr_eq(current, &peer_subscriptions) {
break;
}
let is_empty = {
let mut senders = peer_subscriptions
.senders
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
senders.retain(|_, sender| !sender.is_closed());
senders.is_empty()
};
if is_empty {
subscriptions.remove(&nid);
}
is_empty
};
if is_closed {
break;
}
if let Some(stats) = ep.peer_stats(&nid).await {
if let Some(active) = stats.paths.iter().find(|p| p.active) {
let key = format!("{}:{}", active.relay, active.addr);
if Some(&key) != last_key.as_ref() {
last_key = Some(key);
let subscriptions = ep
.inner
.session
.path_subs
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if subscriptions.get(&nid).is_some_and(|current| {
std::sync::Arc::ptr_eq(current, &peer_subscriptions)
}) {
let senders = peer_subscriptions
.senders
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
for sender in senders.values() {
let _ = sender.send(active.clone());
}
}
let _ = event_tx.try_send(
crate::http::events::TransportEvent::path_change(
&nid,
&active.addr,
active.relay,
),
);
}
}
}
tokio::select! {
_ = tokio::time::sleep(std::time::Duration::from_millis(200)) => {}
result = closed_rx.wait_for(|v| *v) => {
let _ = result;
let mut subscriptions = ep.inner
.session
.path_subs
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if subscriptions.get(&nid).is_some_and(|current| {
std::sync::Arc::ptr_eq(current, &peer_subscriptions)
}) {
subscriptions.remove(&nid);
}
break;
}
}
}
ep.inner
.session
.active_path_watchers
.fetch_sub(1, Ordering::Relaxed);
});
rx
}
pub fn unsubscribe_path_changes(&self, node_id_str: &str, subscription_id: u32) {
let subscriptions = self
.inner
.session
.path_subs
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if let Some(peer_subscriptions) = subscriptions.get(node_id_str) {
peer_subscriptions
.senders
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(&subscription_id);
}
}
}
fn select_dialable_direct(addrs: &[std::net::SocketAddr]) -> Option<String> {
addrs
.iter()
.copied()
.find(|a| is_routable_ip(&a.ip()) && !super::bind::is_placeholder_port(a.port()))
.map(|a| a.to_string())
}
fn select_dialable_directs(addrs: &[std::net::SocketAddr]) -> Vec<String> {
let mut out: Vec<String> = Vec::new();
for a in addrs {
if is_routable_ip(&a.ip()) && !super::bind::is_placeholder_port(a.port()) {
let s = a.to_string();
if !out.contains(&s) {
out.push(s);
}
}
}
out
}
fn is_routable_ip(ip: &std::net::IpAddr) -> bool {
if ip.is_loopback() || ip.is_unspecified() {
return false;
}
match ip {
std::net::IpAddr::V4(v4) => !v4.is_link_local(),
std::net::IpAddr::V6(v6) => (v6.segments()[0] & 0xffc0) != 0xfe80,
}
}
#[cfg(test)]
mod dialable_tests {
use super::{is_routable_ip, select_dialable_direct, select_dialable_directs};
use std::net::SocketAddr;
#[test]
fn selects_first_routable_ip_with_reconciled_port() {
let addrs: Vec<SocketAddr> = vec![
"127.0.0.1:59234".parse().unwrap(),
"192.168.1.42:59234".parse().unwrap(),
];
assert_eq!(
select_dialable_direct(&addrs),
Some("192.168.1.42:59234".to_string())
);
}
#[test]
fn skips_loopback_and_unspecified() {
let addrs: Vec<SocketAddr> = vec![
"127.0.0.1:59234".parse().unwrap(),
"0.0.0.0:59234".parse().unwrap(),
];
assert_eq!(select_dialable_direct(&addrs), None);
}
#[test]
fn skips_link_local() {
let addrs: Vec<SocketAddr> = vec![
"169.254.10.1:59234".parse().unwrap(),
"10.0.0.5:59234".parse().unwrap(),
];
assert_eq!(
select_dialable_direct(&addrs),
Some("10.0.0.5:59234".to_string())
);
}
#[test]
fn none_when_only_non_routable() {
let addrs: Vec<SocketAddr> = vec!["169.254.10.1:59234".parse().unwrap()];
assert_eq!(select_dialable_direct(&addrs), None);
}
#[test]
fn skips_placeholder_port_selects_real() {
let addrs: Vec<SocketAddr> = vec![
"192.168.1.42:1".parse().unwrap(),
"10.0.0.5:59234".parse().unwrap(),
];
assert_eq!(
select_dialable_direct(&addrs),
Some("10.0.0.5:59234".to_string())
);
}
#[test]
fn none_when_only_placeholder_port() {
let addrs: Vec<SocketAddr> = vec!["192.168.1.42:1".parse().unwrap()];
assert_eq!(select_dialable_direct(&addrs), None);
}
#[test]
fn brackets_ipv6() {
let addrs: Vec<SocketAddr> = vec!["[2001:db8::1]:443".parse().unwrap()];
assert_eq!(
select_dialable_direct(&addrs),
Some("[2001:db8::1]:443".to_string())
);
}
#[test]
fn routable_ip_predicate() {
assert!(is_routable_ip(&"192.168.1.1".parse().unwrap()));
assert!(is_routable_ip(&"2001:db8::1".parse().unwrap()));
assert!(!is_routable_ip(&"127.0.0.1".parse().unwrap()));
assert!(!is_routable_ip(&"0.0.0.0".parse().unwrap()));
assert!(!is_routable_ip(&"169.254.1.1".parse().unwrap()));
assert!(!is_routable_ip(&"fe80::1".parse().unwrap()));
}
#[test]
fn plural_keeps_all_routable_candidates() {
let addrs: Vec<SocketAddr> = vec![
"127.0.0.1:59234".parse().unwrap(),
"10.12.222.17:56604".parse().unwrap(),
"192.168.50.227:56604".parse().unwrap(),
"169.254.10.1:56604".parse().unwrap(),
"192.168.50.227:1".parse().unwrap(),
"10.12.222.17:56604".parse().unwrap(),
];
assert_eq!(
select_dialable_directs(&addrs),
vec![
"10.12.222.17:56604".to_string(),
"192.168.50.227:56604".to_string(),
]
);
}
#[test]
fn plural_empty_when_none_routable() {
let addrs: Vec<SocketAddr> = vec![
"127.0.0.1:59234".parse().unwrap(),
"169.254.10.1:56604".parse().unwrap(),
];
assert!(select_dialable_directs(&addrs).is_empty());
}
}