Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion iroh-dns-server/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,7 @@ http = "1.0.0"
humantime = "2.2.0"
humantime-serde = "1.1.1"
iroh-base = { version = "1.2.0", path = "../iroh-base", features = ["key"] }
iroh-metrics = { version = "1.0.1", features = ["service"] }
iroh-metrics = { version = "1.0.2", features = ["service"] }
iroh-dns = { version = "1.3.0", path = "../iroh-dns" }
lru = "0.18.0"
n0-mainline = "0.6"
Expand Down
2 changes: 1 addition & 1 deletion iroh-relay/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@ hyper = { version = "1", features = ["server", "client", "http1"] }
hyper-util = "0.1.1"
iroh-base = { version = "1.2.0", path = "../iroh-base", default-features = false, features = ["key", "relay"] }
iroh-dns = { version = "1.3.0", path = "../iroh-dns" }
iroh-metrics = { version = "1.0.1", default-features = false }
iroh-metrics = { version = "1.0.2", default-features = false }
lru = "0.18.0"
n0-error = "1.0.0"
n0-future = "0.3"
Expand Down
2 changes: 1 addition & 1 deletion iroh/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -65,7 +65,7 @@ url = { version = "2.5", features = ["serde"] }
webpki_types = { package = "rustls-pki-types", version = "1.12" }

# metrics
iroh-metrics = { version = "1.0.1", default-features = false }
iroh-metrics = { version = "1.0.2", default-features = false }

futures-util = "0.3"

Expand Down
2 changes: 1 addition & 1 deletion iroh/bench/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ publish = false
bytes = "1.11"
hdrhistogram = { version = "7.2", default-features = false }
iroh = { version = "1.2.0", path = "..", default-features = false }
iroh-metrics = { version = "1.0.1", optional = true }
iroh-metrics = { version = "1.0.2", optional = true }
n0-future = "0.3"
n0-error = "1.0.0"
noq = "1.3.0"
Expand Down
90 changes: 86 additions & 4 deletions iroh/src/address_lookup.rs
Original file line number Diff line number Diff line change
Expand Up @@ -120,11 +120,13 @@ use crate::{Endpoint, endpoint::EndpointError};
#[cfg(not(wasm_browser))]
pub mod dns;
pub mod memory;
mod metrics;
pub mod pkarr;

#[cfg(not(wasm_browser))]
pub use dns::*;
pub use memory::*;
pub use metrics::{Metrics, ServiceLabels};
pub use pkarr::*;

/// Trait for structs that can be converted into [`AddressLookup`]s.
Expand Down Expand Up @@ -465,9 +467,19 @@ pub struct AddressLookupServices {
last_data: Arc<RwLock<Option<EndpointData>>>,
/// Optional filter applied to all data before publishing to any service.
addr_filter: Arc<RwLock<Option<AddrFilter>>>,
/// Metrics for lookup outcomes.
metrics: Arc<Metrics>,
}

impl AddressLookupServices {
/// Creates a registry recording lookup outcomes in `metrics`.
pub(crate) fn with_metrics(metrics: Arc<Metrics>) -> Self {
Self {
metrics,
..Default::default()
}
}

/// Sets the address filter applied before publishing to any service.
///
/// When set, all address data is filtered once before being distributed
Expand Down Expand Up @@ -554,14 +566,15 @@ impl AddressLookupServices {
&self,
endpoint_id: EndpointId,
) -> impl Stream<Item = Result<Result<Item, Error>, AddressLookupFailed>> + use<> {
self.metrics.lookups.inc();
let services = self.services.read().expect("poisoned");
if services.is_empty() {
AddressLookupStream::empty()
AddressLookupStream::empty(self.metrics.clone())
} else {
let streams = services
.iter()
.filter_map(|service| service.resolve(endpoint_id));
AddressLookupStream::new(streams)
AddressLookupStream::new(streams, self.metrics.clone())
}
}
}
Expand All @@ -584,24 +597,30 @@ struct AddressLookupStream {
errors: Vec<Error>,
did_emit: bool,
closed: bool,
metrics: Arc<Metrics>,
}

impl AddressLookupStream {
fn empty() -> Self {
fn empty(metrics: Arc<Metrics>) -> Self {
Self {
streams: None,
errors: Vec::new(),
did_emit: false,
closed: false,
metrics,
}
}

fn new(streams: impl Iterator<Item = BoxStream<Result<Item, Error>>>) -> Self {
fn new(
streams: impl Iterator<Item = BoxStream<Result<Item, Error>>>,
metrics: Arc<Metrics>,
) -> Self {
Self {
streams: Some(MergeBounded::from_iter(streams)),
errors: Vec::new(),
did_emit: false,
closed: false,
metrics,
}
}
}
Expand All @@ -621,22 +640,32 @@ impl Stream for AddressLookupStream {
Some(inner) => inner,
None => {
this.closed = true;
this.metrics.lookups_failed.inc();
return Poll::Ready(Some(Err(e!(AddressLookupFailed::NoServiceConfigured))));
}
};
let item = match ready!(Pin::new(&mut inner).poll_next(cx)) {
Some(Ok(item)) => {
this.did_emit = true;
this.metrics
.service_results
.get_or_create(&ServiceLabels::new(item.provenance()))
.inc();
Some(Ok(Ok(item)))
}
Some(Err(error)) => {
debug!("address lookup error: {error:#}");
this.metrics
.service_errors
.get_or_create(&ServiceLabels::new(error.provenance))
.inc();
this.errors.push(error.clone());
Some(Ok(Err(error)))
}
None => {
this.closed = true;
if !this.did_emit {
this.metrics.lookups_failed.inc();
let errors = std::mem::take(&mut this.errors);
Some(Err(e!(AddressLookupFailed::NoResults { errors })))
} else {
Expand Down Expand Up @@ -1140,6 +1169,59 @@ mod tests {
Ok(())
}

/// Lookup outcomes are counted per service.
#[tokio::test]
#[traced_test]
async fn address_lookup_metrics() -> Result {
let mut rng = rand_chacha::ChaCha8Rng::seed_from_u64(0u64);
let endpoint_id = SecretKey::from_bytes(&rng.random()).public();

// One succeeding and one failing service.
let succeeding = MemoryLookup::with_provenance("static-test");
let data = EndpointData::from_iter([TransportAddr::Ip("127.0.0.1:1".parse().unwrap())]);
succeeding.add_endpoint_info(EndpointInfo::from_parts(endpoint_id, data));
let services = AddressLookupServices::default();
services.add(succeeding);
services.add(FailingAddressLookup {
delay: Duration::from_millis(10),
});
let _results: Vec<_> = services.resolve(endpoint_id).collect().await;

let metrics = &services.metrics;
assert_eq!(metrics.lookups.get(), 1);
assert_eq!(metrics.lookups_failed.get(), 0);
assert_eq!(
metrics
.service_results
.get(&ServiceLabels::new("static-test"))
.map(|counter| counter.get()),
Some(1)
);
assert_eq!(
metrics
.service_errors
.get(&ServiceLabels::new("failing-test"))
.map(|counter| counter.get()),
Some(1)
);

// Only failing services: the lookup itself is counted as failed.
let services = AddressLookupServices::default();
services.add(FailingAddressLookup {
delay: Duration::from_millis(10),
});
let _results: Vec<_> = services.resolve(endpoint_id).collect().await;
assert_eq!(services.metrics.lookups_failed.get(), 1);

// No services configured: also counted as failed.
let services = AddressLookupServices::default();
let _results: Vec<_> = services.resolve(endpoint_id).collect().await;
assert_eq!(services.metrics.lookups.get(), 1);
assert_eq!(services.metrics.lookups_failed.get(), 1);

Ok(())
}

#[test]
fn concurrent_address_lookup_addr_filter() {
use iroh_base::RelayUrl;
Expand Down
48 changes: 48 additions & 0 deletions iroh/src/address_lookup/metrics.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
//! Metrics for address lookup.

use iroh_metrics::{Counter, EncodeLabelSet, Family, MetricsGroup};
use serde::{Deserialize, Serialize};

/// Labels identifying an address lookup service by its provenance string,
/// see [`crate::address_lookup::Item::provenance`].
#[derive(
Debug, Clone, Hash, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, EncodeLabelSet,
)]
pub struct ServiceLabels {
/// Provenance string of the service.
pub service: String,
}

impl ServiceLabels {
/// Creates labels for the given service provenance.
pub fn new(service: impl Into<String>) -> Self {
Self {
service: service.into(),
}
}
}

/// Metrics collected by address lookup.
///
/// A lookup is one call to [`AddressLookupServices::resolve`] and queries all
/// configured services at once. Each service can yield several results and
/// errors per lookup; those are counted in the `service_*` counters, labeled
/// by the service's provenance (e.g. `dns`, `pkarr`).
///
/// [`AddressLookupServices::resolve`]: crate::address_lookup::AddressLookupServices::resolve
#[derive(Debug, Serialize, Deserialize, MetricsGroup)]
#[non_exhaustive]
#[metrics(name = "address_lookup", default)]
pub struct Metrics {
/// Lookups started.
pub lookups: Counter,
/// Lookups that ended without a single result.
///
/// Includes lookups with no services configured. Lookups abandoned early
/// (e.g. once a connection is established) are not counted.
pub lookups_failed: Counter,
/// Results yielded per service.
pub service_results: Family<ServiceLabels, Counter>,
/// Errors yielded per service.
pub service_errors: Family<ServiceLabels, Counter>,
}
22 changes: 21 additions & 1 deletion iroh/src/metrics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,10 @@ use iroh_metrics::MetricsGroupSet;
pub use iroh_relay::server::Metrics as RelayMetrics;
use serde::{Deserialize, Serialize};

pub use crate::{net_report::Metrics as NetReportMetrics, socket::Metrics as SocketMetrics};
pub use crate::{
address_lookup::Metrics as AddressLookupMetrics, net_report::Metrics as NetReportMetrics,
socket::Metrics as SocketMetrics,
};

/// Metrics collected by an [`crate::endpoint::Endpoint`].
///
Expand All @@ -19,19 +22,36 @@ pub struct EndpointMetrics {
pub socket: Arc<SocketMetrics>,
/// Metrics collected by net reports.
pub net_report: Arc<NetReportMetrics>,
/// Metrics collected by address lookup.
pub address_lookup: Arc<AddressLookupMetrics>,
}

#[cfg(test)]
mod tests {
use super::EndpointMetrics;
#[test]
fn test_serde() {
use crate::address_lookup::ServiceLabels;

let metrics = EndpointMetrics::default();
metrics.socket.actor_link_change.inc();
metrics.net_report.reports.inc_by(10);
metrics
.address_lookup
.service_results
.get_or_create(&ServiceLabels::new("dns"))
.inc_by(3);
let encoded = postcard::to_stdvec(&metrics).unwrap();
let decoded: EndpointMetrics = postcard::from_bytes(&encoded).unwrap();
assert_eq!(decoded.socket.actor_link_change.get(), 1);
assert_eq!(decoded.net_report.reports.get(), 10);
assert_eq!(
decoded
.address_lookup
.service_results
.get(&ServiceLabels::new("dns"))
.map(|counter| counter.get()),
Some(3)
);
}
}
3 changes: 2 additions & 1 deletion iroh/src/socket.rs
Original file line number Diff line number Diff line change
Expand Up @@ -889,7 +889,8 @@ impl EndpointInner {
configured_addrs,
} = opts;

let address_lookup = address_lookup::AddressLookupServices::default();
let address_lookup =
address_lookup::AddressLookupServices::with_metrics(metrics.address_lookup.clone());
let port_mapper = portmapper::create_client(&portmapper_config);

let relay_transport_configs: Vec<_> = transport_configs
Expand Down
Loading