Skip to main content

dht_crawler/
runtime_stats.rs

1use std::sync::Arc;
2use std::sync::atomic::{AtomicU32, AtomicU64, AtomicUsize, Ordering};
3
4pub const METADATA_QUEUE_WAIT_BUCKETS_MS: [u64; 8] = [10, 50, 100, 250, 500, 1_000, 2_000, 5_000];
5pub const METADATA_FETCH_BUCKETS_MS: [u64; 7] = [250, 500, 1_000, 2_000, 4_000, 6_000, 10_000];
6pub const METADATA_SIZE_BUCKETS_BYTES: [u64; 8] = [
7    16 * 1024,
8    32 * 1024,
9    64 * 1024,
10    128 * 1024,
11    256 * 1024,
12    512 * 1024,
13    1024 * 1024,
14    10 * 1024 * 1024,
15];
16
17#[derive(Debug, Clone, PartialEq, Eq)]
18/// Snapshot of a non-cumulative fixed-bucket histogram.
19pub struct FixedHistogramSnapshot {
20    /// Inclusive upper bound for each bucket.
21    pub bounds: Vec<u64>,
22    /// Per-bucket counts; values are not cumulative.
23    pub counts: Vec<u64>,
24    /// Values larger than the final bound.
25    pub overflow: u64,
26    /// Total observations including overflow.
27    pub count: u64,
28    /// Wrapping sum of all observed values.
29    pub sum: u64,
30}
31
32impl FixedHistogramSnapshot {
33    /// Returns the inclusive bucket bound containing the requested approximate percentile.
34    ///
35    /// `percentile` is clamped to `0.0..=1.0`. Overflow observations return the final finite
36    /// bound because the histogram intentionally stores no dynamic maximum.
37    pub fn percentile(&self, percentile: f64) -> Option<u64> {
38        if self.count == 0 {
39            return None;
40        }
41        let rank = ((self.count as f64 * percentile.clamp(0.0, 1.0)).ceil() as u64).max(1);
42        let mut cumulative = 0u64;
43        for (bound, count) in self.bounds.iter().zip(&self.counts) {
44            cumulative = cumulative.saturating_add(*count);
45            if cumulative >= rank {
46                return Some(*bound);
47            }
48        }
49        self.bounds.last().copied()
50    }
51}
52
53struct AtomicFixedHistogram<const N: usize> {
54    bounds: [u64; N],
55    counts: [AtomicU64; N],
56    overflow: AtomicU64,
57    count: AtomicU64,
58    sum: AtomicU64,
59}
60
61impl<const N: usize> AtomicFixedHistogram<N> {
62    fn new(bounds: [u64; N]) -> Self {
63        Self {
64            bounds,
65            counts: std::array::from_fn(|_| AtomicU64::new(0)),
66            overflow: AtomicU64::new(0),
67            count: AtomicU64::new(0),
68            sum: AtomicU64::new(0),
69        }
70    }
71
72    fn record(&self, value: u64) {
73        if let Some(index) = self.bounds.iter().position(|bound| value <= *bound) {
74            self.counts[index].fetch_add(1, Ordering::Relaxed);
75        } else {
76            self.overflow.fetch_add(1, Ordering::Relaxed);
77        }
78        self.count.fetch_add(1, Ordering::Relaxed);
79        self.sum.fetch_add(value, Ordering::Relaxed);
80    }
81
82    fn snapshot(&self) -> FixedHistogramSnapshot {
83        FixedHistogramSnapshot {
84            bounds: self.bounds.to_vec(),
85            counts: self
86                .counts
87                .iter()
88                .map(|count| count.load(Ordering::Relaxed))
89                .collect(),
90            overflow: self.overflow.load(Ordering::Relaxed),
91            count: self.count.load(Ordering::Relaxed),
92            sum: self.sum.load(Ordering::Relaxed),
93        }
94    }
95}
96
97impl<const N: usize> Default for AtomicFixedHistogram<N> {
98    fn default() -> Self {
99        Self::new([0; N])
100    }
101}
102
103/// Transport-neutral counters and fixed-bucket histograms used by dashboards.
104#[derive(Debug, Clone, PartialEq, Eq)]
105pub struct DhtObservabilitySnapshot {
106    /// UDP datagrams received before validation.
107    pub udp_rx_packets: u64,
108    /// UDP bytes received before validation.
109    pub udp_rx_bytes: u64,
110    /// Successfully sent UDP query and response datagrams.
111    pub udp_tx_packets: u64,
112    /// Successfully sent UDP query and response bytes.
113    pub udp_tx_bytes: u64,
114    /// Inbound ping queries.
115    pub inbound_ping: u64,
116    /// Inbound find_node queries.
117    pub inbound_find_node: u64,
118    /// Inbound get_peers queries.
119    pub inbound_get_peers: u64,
120    /// Inbound announce_peer queries.
121    pub inbound_announce_peer: u64,
122    /// Other or invalid inbound query names.
123    pub inbound_other: u64,
124    /// Replies admitted by the regular response budget.
125    pub response_normal: u64,
126    /// Priority ping/get_peers replies admitted by the reserve.
127    pub response_fallback: u64,
128    /// Replies rejected by final response limits.
129    pub response_rate_limited: u64,
130    /// Admitted replies whose `send_to` failed.
131    pub response_send_failed: u64,
132    /// Valid, unfiltered announces with a usable Peer port.
133    pub announce_accepted: u64,
134    /// Announces with a missing or invalid token.
135    pub announce_invalid_token: u64,
136    /// Announces rejected by the application Hash filter.
137    pub announce_filtered: u64,
138    /// New nodes admitted to the FIFO pool.
139    pub node_admitted: u64,
140    /// Full-pool replacements.
141    pub node_replaced: u64,
142    /// Nodes rejected as queued or recently probed duplicates.
143    pub node_dropped_duplicate: u64,
144    /// Nodes rejected by the replacement budget.
145    pub node_dropped_rate_limited: u64,
146    /// Nodes rejected because their endpoint is not usable.
147    pub node_dropped_invalid: u64,
148    /// BEP-9 Metadata piece payload bytes received.
149    pub metadata_bytes_downloaded: u64,
150    /// Metadata Peer attempts reaching the end-to-end timeout.
151    pub metadata_failure_timeout: u64,
152    /// Metadata Peer connection failures.
153    pub metadata_failure_connect: u64,
154    /// Peers without extension-protocol support.
155    pub metadata_failure_no_extension: u64,
156    /// Extension handshake or piece-request send failures.
157    pub metadata_failure_send: u64,
158    /// Metadata payloads exceeding the 10 MiB limit.
159    pub metadata_failure_size_limit: u64,
160    /// Metadata payloads failing InfoHash SHA1 validation.
161    pub metadata_failure_sha1: u64,
162    /// Validated payloads that could not be parsed as torrent info.
163    pub metadata_failure_parse: u64,
164    /// Other receive, piece or worker failures.
165    pub metadata_failure_other: u64,
166    /// Failure-cache hits for previously timed-out Peers.
167    pub peer_cache_timeout_hits: u64,
168    /// Failure-cache hits for previous connection failures.
169    pub peer_cache_connect_hits: u64,
170    /// Queue wait observations in milliseconds.
171    pub queue_wait_ms: FixedHistogramSnapshot,
172    /// End-to-end Peer attempt observations in milliseconds.
173    pub fetch_duration_ms: FixedHistogramSnapshot,
174    /// Complete bencoded info payload observations in bytes.
175    pub metadata_size_bytes: FixedHistogramSnapshot,
176}
177
178/// Cheap, cloneable handle for reading live DHT runtime statistics.
179#[derive(Clone, Default)]
180pub struct DhtRuntimeStats {
181    inner: Arc<DhtRuntimeStatsInner>,
182}
183
184/// Point-in-time view returned by [`DhtRuntimeStats::snapshot`].
185#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
186pub struct DhtRuntimeSnapshot {
187    /// Valid announce hashes observed before ingress publication.
188    pub hashes_received: u64,
189    /// Hashes dropped because bounded ingress was full.
190    pub hash_ingress_dropped: u64,
191    /// Current hash-ingress queue depth.
192    pub hash_ingress_queue_depth: usize,
193    /// Configured hash-ingress queue capacity.
194    pub hash_ingress_queue_capacity: usize,
195
196    /// Current priority crawl-event depth.
197    pub crawl_priority_queue_depth: usize,
198    /// Configured priority crawl-event capacity.
199    pub crawl_priority_queue_capacity: usize,
200    /// Current discovered-node event depth.
201    pub crawl_discovery_queue_depth: usize,
202    /// Configured discovered-node event capacity.
203    pub crawl_discovery_queue_capacity: usize,
204
205    /// Deduplicated pending Metadata Hashes.
206    pub metadata_queue_depth: usize,
207    /// Configured pending Metadata capacity.
208    pub metadata_queue_max: usize,
209    /// Current Metadata jobs.
210    pub metadata_in_flight: usize,
211    /// Unique Hash insertions, including insertions that evicted an older Hash.
212    pub metadata_queue_inserted: u64,
213    /// Duplicate Hash updates merged into existing entries.
214    pub metadata_queue_deduplicated: u64,
215    /// Oldest Hashes evicted by newer events.
216    pub metadata_queue_evicted: u64,
217    /// Events rejected or expired because they were stale.
218    pub metadata_queue_stale: u64,
219    /// Metadata Peer races started after admission.
220    pub metadata_races_started: u64,
221    /// Distinct Peer candidates launched into Metadata races.
222    pub metadata_peer_candidates: u64,
223    /// Peer candidates added while a Metadata race was already running.
224    pub metadata_live_peer_joins: u64,
225    /// Losing Peer attempts canceled after another Peer succeeded.
226    pub metadata_peer_canceled: u64,
227    /// Active `get_peers` lookup requests received from Metadata jobs.
228    pub peer_lookup_requested: u64,
229    /// Active `get_peers` lookups admitted by the bounded rate limiter.
230    pub peer_lookup_started: u64,
231    /// Active lookup requests rejected by concurrency or rate limits.
232    pub peer_lookup_rate_limited: u64,
233    /// Active lookup requests skipped because the routing snapshot was empty.
234    pub peer_lookup_empty: u64,
235    /// Active `get_peers` UDP queries sent.
236    pub peer_lookup_queries: u64,
237    /// Matched active `get_peers` responses.
238    pub peer_lookup_responses: u64,
239    /// Active `get_peers` queries that timed out.
240    pub peer_lookup_timeouts: u64,
241    /// Active `get_peers` UDP sends that failed immediately.
242    pub peer_lookup_send_failures: u64,
243    /// Unique Peer endpoints fed back into Metadata scheduling.
244    pub peer_lookup_peers_found: u64,
245    /// Tagged lookup responses dropped before reaching the lookup actor.
246    pub peer_lookup_response_dropped: u64,
247    /// Discovered Peer endpoints dropped because Hash ingress was full.
248    pub peer_lookup_output_dropped: u64,
249    /// Real Peer network attempts.
250    pub metadata_peer_attempts: u64,
251    /// Successful Peer downloads and parses.
252    pub metadata_peer_succeeded: u64,
253    /// Failed real Peer attempts.
254    pub metadata_peer_failed: u64,
255    /// End-to-end Peer timeouts.
256    pub metadata_peer_timeouts: u64,
257    /// Peer connection failures.
258    pub metadata_connect_failed: u64,
259    /// Peers without extension-protocol support.
260    pub metadata_no_extension: u64,
261    /// Peer failure-cache skips.
262    pub metadata_peer_failure_cache_hits: u64,
263    /// Current Peer failure-cache entries.
264    pub metadata_peer_failure_cache_entries: usize,
265
266    /// Current strict FIFO crawl-pool size.
267    pub node_pool_size: usize,
268    /// Configured FIFO capacity.
269    pub node_pool_capacity: usize,
270    /// Bootstrap low-water mark.
271    pub node_pool_low_watermark: usize,
272
273    /// Current pending find_node transactions.
274    pub find_node_in_flight: usize,
275    /// Configured total find_node in-flight limit.
276    pub find_node_in_flight_max: usize,
277    /// Metadata-pressure-adjusted find_node budget per second.
278    pub find_node_effective_rate_per_sec: u32,
279    /// Queries sent to never-before-probed destinations.
280    pub queries_new: u64,
281    /// Queries sent to responsive revisit nodes.
282    pub queries_revisit: u64,
283    /// Queries sent to bootstrap endpoints.
284    pub queries_bootstrap: u64,
285    /// Replies matched to pending transactions.
286    pub responses: u64,
287    /// Replies without a matching pending transaction.
288    pub unmatched_responses: u64,
289    /// Pending transactions that expired.
290    pub timeouts: u64,
291    /// Outbound find_node send failures.
292    pub send_failures: u64,
293    /// Response events dropped before reaching the actor.
294    pub crawl_events_dropped_response: u64,
295    /// Discovered-node events dropped before reaching the actor.
296    pub crawl_events_dropped_discovered: u64,
297
298    /// UDP datagrams received before validation.
299    pub udp_received: u64,
300    /// Valid datagrams dropped because every worker queue was full.
301    pub udp_queue_full: u64,
302    /// Empty, oversized or non-bencoded UDP datagrams.
303    pub udp_invalid: u64,
304    /// DHT replies rejected by final rate limits.
305    pub udp_responses_rate_limited: u64,
306    /// Priority replies admitted by the ping/get_peers reserve.
307    pub udp_responses_priority_reserved: u64,
308}
309
310struct DhtRuntimeStatsInner {
311    hashes_received: AtomicU64,
312    hash_ingress_dropped: AtomicU64,
313    hash_ingress_queue_depth: AtomicUsize,
314    hash_ingress_queue_capacity: usize,
315
316    crawl_priority_queue_depth: AtomicUsize,
317    crawl_priority_queue_capacity: usize,
318    crawl_discovery_queue_depth: AtomicUsize,
319    crawl_discovery_queue_capacity: usize,
320
321    metadata_queue_depth: AtomicUsize,
322    metadata_queue_max: usize,
323    metadata_in_flight: AtomicUsize,
324    metadata_queue_inserted: AtomicU64,
325    metadata_queue_deduplicated: AtomicU64,
326    metadata_queue_evicted: AtomicU64,
327    metadata_queue_stale: AtomicU64,
328    metadata_races_started: AtomicU64,
329    metadata_peer_candidates: AtomicU64,
330    metadata_live_peer_joins: AtomicU64,
331    metadata_peer_canceled: AtomicU64,
332    peer_lookup_requested: AtomicU64,
333    peer_lookup_started: AtomicU64,
334    peer_lookup_rate_limited: AtomicU64,
335    peer_lookup_empty: AtomicU64,
336    peer_lookup_queries: AtomicU64,
337    peer_lookup_responses: AtomicU64,
338    peer_lookup_timeouts: AtomicU64,
339    peer_lookup_send_failures: AtomicU64,
340    peer_lookup_peers_found: AtomicU64,
341    peer_lookup_response_dropped: AtomicU64,
342    peer_lookup_output_dropped: AtomicU64,
343    metadata_peer_attempts: AtomicU64,
344    metadata_peer_succeeded: AtomicU64,
345    metadata_peer_failed: AtomicU64,
346    metadata_peer_timeouts: AtomicU64,
347    metadata_connect_failed: AtomicU64,
348    metadata_no_extension: AtomicU64,
349    metadata_peer_failure_cache_hits: AtomicU64,
350    metadata_peer_failure_cache_entries: AtomicUsize,
351
352    node_pool_size: AtomicUsize,
353    node_pool_capacity: usize,
354    node_pool_low_watermark: usize,
355
356    find_node_in_flight: AtomicUsize,
357    find_node_in_flight_max: usize,
358    find_node_effective_rate_per_sec: AtomicU32,
359    queries_new: AtomicU64,
360    queries_revisit: AtomicU64,
361    queries_bootstrap: AtomicU64,
362    responses: AtomicU64,
363    unmatched_responses: AtomicU64,
364    timeouts: AtomicU64,
365    send_failures: AtomicU64,
366    crawl_events_dropped_response: AtomicU64,
367    crawl_events_dropped_discovered: AtomicU64,
368
369    udp_received: AtomicU64,
370    udp_queue_full: AtomicU64,
371    udp_invalid: AtomicU64,
372    udp_responses_rate_limited: AtomicU64,
373    udp_responses_priority_reserved: AtomicU64,
374
375    udp_rx_bytes: AtomicU64,
376    udp_tx_packets: AtomicU64,
377    udp_tx_bytes: AtomicU64,
378    inbound_ping: AtomicU64,
379    inbound_find_node: AtomicU64,
380    inbound_get_peers: AtomicU64,
381    inbound_announce_peer: AtomicU64,
382    inbound_other: AtomicU64,
383    response_normal: AtomicU64,
384    response_send_failed: AtomicU64,
385    announce_accepted: AtomicU64,
386    announce_invalid_token: AtomicU64,
387    announce_filtered: AtomicU64,
388    node_admitted: AtomicU64,
389    node_replaced: AtomicU64,
390    node_dropped_duplicate: AtomicU64,
391    node_dropped_rate_limited: AtomicU64,
392    node_dropped_invalid: AtomicU64,
393    metadata_bytes_downloaded: AtomicU64,
394    metadata_failure_timeout: AtomicU64,
395    metadata_failure_connect: AtomicU64,
396    metadata_failure_no_extension: AtomicU64,
397    metadata_failure_send: AtomicU64,
398    metadata_failure_size_limit: AtomicU64,
399    metadata_failure_sha1: AtomicU64,
400    metadata_failure_parse: AtomicU64,
401    metadata_failure_other: AtomicU64,
402    peer_cache_timeout_hits: AtomicU64,
403    peer_cache_connect_hits: AtomicU64,
404    queue_wait_ms: AtomicFixedHistogram<8>,
405    fetch_duration_ms: AtomicFixedHistogram<7>,
406    metadata_size_bytes: AtomicFixedHistogram<8>,
407}
408
409impl Default for DhtRuntimeStatsInner {
410    fn default() -> Self {
411        Self {
412            queue_wait_ms: AtomicFixedHistogram::new(METADATA_QUEUE_WAIT_BUCKETS_MS),
413            fetch_duration_ms: AtomicFixedHistogram::new(METADATA_FETCH_BUCKETS_MS),
414            metadata_size_bytes: AtomicFixedHistogram::new(METADATA_SIZE_BUCKETS_BYTES),
415            hashes_received: AtomicU64::new(0),
416            hash_ingress_dropped: AtomicU64::new(0),
417            hash_ingress_queue_depth: AtomicUsize::new(0),
418            hash_ingress_queue_capacity: 0,
419            crawl_priority_queue_depth: AtomicUsize::new(0),
420            crawl_priority_queue_capacity: 0,
421            crawl_discovery_queue_depth: AtomicUsize::new(0),
422            crawl_discovery_queue_capacity: 0,
423            metadata_queue_depth: AtomicUsize::new(0),
424            metadata_queue_max: 0,
425            metadata_in_flight: AtomicUsize::new(0),
426            metadata_queue_inserted: AtomicU64::new(0),
427            metadata_queue_deduplicated: AtomicU64::new(0),
428            metadata_queue_evicted: AtomicU64::new(0),
429            metadata_queue_stale: AtomicU64::new(0),
430            metadata_races_started: AtomicU64::new(0),
431            metadata_peer_candidates: AtomicU64::new(0),
432            metadata_live_peer_joins: AtomicU64::new(0),
433            metadata_peer_canceled: AtomicU64::new(0),
434            peer_lookup_requested: AtomicU64::new(0),
435            peer_lookup_started: AtomicU64::new(0),
436            peer_lookup_rate_limited: AtomicU64::new(0),
437            peer_lookup_empty: AtomicU64::new(0),
438            peer_lookup_queries: AtomicU64::new(0),
439            peer_lookup_responses: AtomicU64::new(0),
440            peer_lookup_timeouts: AtomicU64::new(0),
441            peer_lookup_send_failures: AtomicU64::new(0),
442            peer_lookup_peers_found: AtomicU64::new(0),
443            peer_lookup_response_dropped: AtomicU64::new(0),
444            peer_lookup_output_dropped: AtomicU64::new(0),
445            metadata_peer_attempts: AtomicU64::new(0),
446            metadata_peer_succeeded: AtomicU64::new(0),
447            metadata_peer_failed: AtomicU64::new(0),
448            metadata_peer_timeouts: AtomicU64::new(0),
449            metadata_connect_failed: AtomicU64::new(0),
450            metadata_no_extension: AtomicU64::new(0),
451            metadata_peer_failure_cache_hits: AtomicU64::new(0),
452            metadata_peer_failure_cache_entries: AtomicUsize::new(0),
453            node_pool_size: AtomicUsize::new(0),
454            node_pool_capacity: 0,
455            node_pool_low_watermark: 0,
456            find_node_in_flight: AtomicUsize::new(0),
457            find_node_in_flight_max: 0,
458            find_node_effective_rate_per_sec: AtomicU32::new(0),
459            queries_new: AtomicU64::new(0),
460            queries_revisit: AtomicU64::new(0),
461            queries_bootstrap: AtomicU64::new(0),
462            responses: AtomicU64::new(0),
463            unmatched_responses: AtomicU64::new(0),
464            timeouts: AtomicU64::new(0),
465            send_failures: AtomicU64::new(0),
466            crawl_events_dropped_response: AtomicU64::new(0),
467            crawl_events_dropped_discovered: AtomicU64::new(0),
468            udp_received: AtomicU64::new(0),
469            udp_queue_full: AtomicU64::new(0),
470            udp_invalid: AtomicU64::new(0),
471            udp_responses_rate_limited: AtomicU64::new(0),
472            udp_responses_priority_reserved: AtomicU64::new(0),
473            udp_rx_bytes: AtomicU64::new(0),
474            udp_tx_packets: AtomicU64::new(0),
475            udp_tx_bytes: AtomicU64::new(0),
476            inbound_ping: AtomicU64::new(0),
477            inbound_find_node: AtomicU64::new(0),
478            inbound_get_peers: AtomicU64::new(0),
479            inbound_announce_peer: AtomicU64::new(0),
480            inbound_other: AtomicU64::new(0),
481            response_normal: AtomicU64::new(0),
482            response_send_failed: AtomicU64::new(0),
483            announce_accepted: AtomicU64::new(0),
484            announce_invalid_token: AtomicU64::new(0),
485            announce_filtered: AtomicU64::new(0),
486            node_admitted: AtomicU64::new(0),
487            node_replaced: AtomicU64::new(0),
488            node_dropped_duplicate: AtomicU64::new(0),
489            node_dropped_rate_limited: AtomicU64::new(0),
490            node_dropped_invalid: AtomicU64::new(0),
491            metadata_bytes_downloaded: AtomicU64::new(0),
492            metadata_failure_timeout: AtomicU64::new(0),
493            metadata_failure_connect: AtomicU64::new(0),
494            metadata_failure_no_extension: AtomicU64::new(0),
495            metadata_failure_send: AtomicU64::new(0),
496            metadata_failure_size_limit: AtomicU64::new(0),
497            metadata_failure_sha1: AtomicU64::new(0),
498            metadata_failure_parse: AtomicU64::new(0),
499            metadata_failure_other: AtomicU64::new(0),
500            peer_cache_timeout_hits: AtomicU64::new(0),
501            peer_cache_connect_hits: AtomicU64::new(0),
502        }
503    }
504}
505
506#[derive(Debug, Clone, Copy, Default)]
507pub(crate) struct DhtRuntimeLimits {
508    pub metadata_queue: usize,
509    pub node_pool: usize,
510    pub node_pool_low_watermark: usize,
511    pub find_node_in_flight: usize,
512    pub initial_find_node_rate: u32,
513    pub hash_ingress_queue: usize,
514    pub crawl_priority_queue: usize,
515    pub crawl_discovery_queue: usize,
516}
517
518impl DhtRuntimeStats {
519    pub(crate) fn with_limits(limits: DhtRuntimeLimits) -> Self {
520        Self {
521            inner: Arc::new(DhtRuntimeStatsInner {
522                metadata_queue_max: limits.metadata_queue,
523                node_pool_capacity: limits.node_pool,
524                node_pool_low_watermark: limits.node_pool_low_watermark,
525                find_node_in_flight_max: limits.find_node_in_flight,
526                find_node_effective_rate_per_sec: AtomicU32::new(limits.initial_find_node_rate),
527                hash_ingress_queue_capacity: limits.hash_ingress_queue,
528                crawl_priority_queue_capacity: limits.crawl_priority_queue,
529                crawl_discovery_queue_capacity: limits.crawl_discovery_queue,
530                ..DhtRuntimeStatsInner::default()
531            }),
532        }
533    }
534
535    /// Load every exposed value once using relaxed atomics.
536    /// Loads the core runtime counters and gauges using relaxed atomics.
537    pub fn snapshot(&self) -> DhtRuntimeSnapshot {
538        let inner = &self.inner;
539        DhtRuntimeSnapshot {
540            hashes_received: inner.hashes_received.load(Ordering::Relaxed),
541            hash_ingress_dropped: inner.hash_ingress_dropped.load(Ordering::Relaxed),
542            hash_ingress_queue_depth: inner.hash_ingress_queue_depth.load(Ordering::Relaxed),
543            hash_ingress_queue_capacity: inner.hash_ingress_queue_capacity,
544            crawl_priority_queue_depth: inner.crawl_priority_queue_depth.load(Ordering::Relaxed),
545            crawl_priority_queue_capacity: inner.crawl_priority_queue_capacity,
546            crawl_discovery_queue_depth: inner.crawl_discovery_queue_depth.load(Ordering::Relaxed),
547            crawl_discovery_queue_capacity: inner.crawl_discovery_queue_capacity,
548            metadata_queue_depth: inner.metadata_queue_depth.load(Ordering::Relaxed),
549            metadata_queue_max: inner.metadata_queue_max,
550            metadata_in_flight: inner.metadata_in_flight.load(Ordering::Relaxed),
551            metadata_queue_inserted: inner.metadata_queue_inserted.load(Ordering::Relaxed),
552            metadata_queue_deduplicated: inner.metadata_queue_deduplicated.load(Ordering::Relaxed),
553            metadata_queue_evicted: inner.metadata_queue_evicted.load(Ordering::Relaxed),
554            metadata_queue_stale: inner.metadata_queue_stale.load(Ordering::Relaxed),
555            metadata_races_started: inner.metadata_races_started.load(Ordering::Relaxed),
556            metadata_peer_candidates: inner.metadata_peer_candidates.load(Ordering::Relaxed),
557            metadata_live_peer_joins: inner.metadata_live_peer_joins.load(Ordering::Relaxed),
558            metadata_peer_canceled: inner.metadata_peer_canceled.load(Ordering::Relaxed),
559            peer_lookup_requested: inner.peer_lookup_requested.load(Ordering::Relaxed),
560            peer_lookup_started: inner.peer_lookup_started.load(Ordering::Relaxed),
561            peer_lookup_rate_limited: inner.peer_lookup_rate_limited.load(Ordering::Relaxed),
562            peer_lookup_empty: inner.peer_lookup_empty.load(Ordering::Relaxed),
563            peer_lookup_queries: inner.peer_lookup_queries.load(Ordering::Relaxed),
564            peer_lookup_responses: inner.peer_lookup_responses.load(Ordering::Relaxed),
565            peer_lookup_timeouts: inner.peer_lookup_timeouts.load(Ordering::Relaxed),
566            peer_lookup_send_failures: inner.peer_lookup_send_failures.load(Ordering::Relaxed),
567            peer_lookup_peers_found: inner.peer_lookup_peers_found.load(Ordering::Relaxed),
568            peer_lookup_response_dropped: inner
569                .peer_lookup_response_dropped
570                .load(Ordering::Relaxed),
571            peer_lookup_output_dropped: inner.peer_lookup_output_dropped.load(Ordering::Relaxed),
572            metadata_peer_attempts: inner.metadata_peer_attempts.load(Ordering::Relaxed),
573            metadata_peer_succeeded: inner.metadata_peer_succeeded.load(Ordering::Relaxed),
574            metadata_peer_failed: inner.metadata_peer_failed.load(Ordering::Relaxed),
575            metadata_peer_timeouts: inner.metadata_peer_timeouts.load(Ordering::Relaxed),
576            metadata_connect_failed: inner.metadata_connect_failed.load(Ordering::Relaxed),
577            metadata_no_extension: inner.metadata_no_extension.load(Ordering::Relaxed),
578            metadata_peer_failure_cache_hits: inner
579                .metadata_peer_failure_cache_hits
580                .load(Ordering::Relaxed),
581            metadata_peer_failure_cache_entries: inner
582                .metadata_peer_failure_cache_entries
583                .load(Ordering::Relaxed),
584            node_pool_size: inner.node_pool_size.load(Ordering::Relaxed),
585            node_pool_capacity: inner.node_pool_capacity,
586            node_pool_low_watermark: inner.node_pool_low_watermark,
587            find_node_in_flight: inner.find_node_in_flight.load(Ordering::Relaxed),
588            find_node_in_flight_max: inner.find_node_in_flight_max,
589            find_node_effective_rate_per_sec: inner
590                .find_node_effective_rate_per_sec
591                .load(Ordering::Relaxed),
592            queries_new: inner.queries_new.load(Ordering::Relaxed),
593            queries_revisit: inner.queries_revisit.load(Ordering::Relaxed),
594            queries_bootstrap: inner.queries_bootstrap.load(Ordering::Relaxed),
595            responses: inner.responses.load(Ordering::Relaxed),
596            unmatched_responses: inner.unmatched_responses.load(Ordering::Relaxed),
597            timeouts: inner.timeouts.load(Ordering::Relaxed),
598            send_failures: inner.send_failures.load(Ordering::Relaxed),
599            crawl_events_dropped_response: inner
600                .crawl_events_dropped_response
601                .load(Ordering::Relaxed),
602            crawl_events_dropped_discovered: inner
603                .crawl_events_dropped_discovered
604                .load(Ordering::Relaxed),
605            udp_received: inner.udp_received.load(Ordering::Relaxed),
606            udp_queue_full: inner.udp_queue_full.load(Ordering::Relaxed),
607            udp_invalid: inner.udp_invalid.load(Ordering::Relaxed),
608            udp_responses_rate_limited: inner.udp_responses_rate_limited.load(Ordering::Relaxed),
609            udp_responses_priority_reserved: inner
610                .udp_responses_priority_reserved
611                .load(Ordering::Relaxed),
612        }
613    }
614
615    /// Loads detailed transport counters, failure categories and fixed histograms.
616    pub fn observability_snapshot(&self) -> DhtObservabilitySnapshot {
617        let inner = &self.inner;
618        DhtObservabilitySnapshot {
619            udp_rx_packets: inner.udp_received.load(Ordering::Relaxed),
620            udp_rx_bytes: inner.udp_rx_bytes.load(Ordering::Relaxed),
621            udp_tx_packets: inner.udp_tx_packets.load(Ordering::Relaxed),
622            udp_tx_bytes: inner.udp_tx_bytes.load(Ordering::Relaxed),
623            inbound_ping: inner.inbound_ping.load(Ordering::Relaxed),
624            inbound_find_node: inner.inbound_find_node.load(Ordering::Relaxed),
625            inbound_get_peers: inner.inbound_get_peers.load(Ordering::Relaxed),
626            inbound_announce_peer: inner.inbound_announce_peer.load(Ordering::Relaxed),
627            inbound_other: inner.inbound_other.load(Ordering::Relaxed),
628            response_normal: inner.response_normal.load(Ordering::Relaxed),
629            response_fallback: inner
630                .udp_responses_priority_reserved
631                .load(Ordering::Relaxed),
632            response_rate_limited: inner.udp_responses_rate_limited.load(Ordering::Relaxed),
633            response_send_failed: inner.response_send_failed.load(Ordering::Relaxed),
634            announce_accepted: inner.announce_accepted.load(Ordering::Relaxed),
635            announce_invalid_token: inner.announce_invalid_token.load(Ordering::Relaxed),
636            announce_filtered: inner.announce_filtered.load(Ordering::Relaxed),
637            node_admitted: inner.node_admitted.load(Ordering::Relaxed),
638            node_replaced: inner.node_replaced.load(Ordering::Relaxed),
639            node_dropped_duplicate: inner.node_dropped_duplicate.load(Ordering::Relaxed),
640            node_dropped_rate_limited: inner.node_dropped_rate_limited.load(Ordering::Relaxed),
641            node_dropped_invalid: inner.node_dropped_invalid.load(Ordering::Relaxed),
642            metadata_bytes_downloaded: inner.metadata_bytes_downloaded.load(Ordering::Relaxed),
643            metadata_failure_timeout: inner.metadata_failure_timeout.load(Ordering::Relaxed),
644            metadata_failure_connect: inner.metadata_failure_connect.load(Ordering::Relaxed),
645            metadata_failure_no_extension: inner
646                .metadata_failure_no_extension
647                .load(Ordering::Relaxed),
648            metadata_failure_send: inner.metadata_failure_send.load(Ordering::Relaxed),
649            metadata_failure_size_limit: inner.metadata_failure_size_limit.load(Ordering::Relaxed),
650            metadata_failure_sha1: inner.metadata_failure_sha1.load(Ordering::Relaxed),
651            metadata_failure_parse: inner.metadata_failure_parse.load(Ordering::Relaxed),
652            metadata_failure_other: inner.metadata_failure_other.load(Ordering::Relaxed),
653            peer_cache_timeout_hits: inner.peer_cache_timeout_hits.load(Ordering::Relaxed),
654            peer_cache_connect_hits: inner.peer_cache_connect_hits.load(Ordering::Relaxed),
655            queue_wait_ms: inner.queue_wait_ms.snapshot(),
656            fetch_duration_ms: inner.fetch_duration_ms.snapshot(),
657            metadata_size_bytes: inner.metadata_size_bytes.snapshot(),
658        }
659    }
660
661    pub(crate) fn hash_received(&self) {
662        self.inner.hashes_received.fetch_add(1, Ordering::Relaxed);
663    }
664
665    pub(crate) fn hash_ingress_dropped(&self) {
666        self.inner
667            .hash_ingress_dropped
668            .fetch_add(1, Ordering::Relaxed);
669    }
670
671    pub(crate) fn set_hash_ingress_queue_depth(&self, depth: usize) {
672        self.inner
673            .hash_ingress_queue_depth
674            .store(depth, Ordering::Relaxed);
675    }
676
677    pub(crate) fn set_crawl_priority_queue_depth(&self, depth: usize) {
678        self.inner
679            .crawl_priority_queue_depth
680            .store(depth, Ordering::Relaxed);
681    }
682
683    pub(crate) fn set_crawl_discovery_queue_depth(&self, depth: usize) {
684        self.inner
685            .crawl_discovery_queue_depth
686            .store(depth, Ordering::Relaxed);
687    }
688
689    pub(crate) fn set_metadata_queue(&self, depth: usize, in_flight: usize) {
690        self.inner
691            .metadata_queue_depth
692            .store(depth, Ordering::Relaxed);
693        self.inner
694            .metadata_in_flight
695            .store(in_flight, Ordering::Relaxed);
696    }
697
698    pub(crate) fn metadata_queue_deduplicated(&self) {
699        self.inner
700            .metadata_queue_deduplicated
701            .fetch_add(1, Ordering::Relaxed);
702    }
703
704    pub(crate) fn metadata_queue_inserted(&self) {
705        self.inner
706            .metadata_queue_inserted
707            .fetch_add(1, Ordering::Relaxed);
708    }
709
710    pub(crate) fn metadata_queue_evicted(&self) {
711        self.inner
712            .metadata_queue_evicted
713            .fetch_add(1, Ordering::Relaxed);
714    }
715
716    pub(crate) fn metadata_queue_stale(&self, count: usize) {
717        self.inner
718            .metadata_queue_stale
719            .fetch_add(count as u64, Ordering::Relaxed);
720    }
721
722    pub(crate) fn metadata_race_started(&self) {
723        self.inner
724            .metadata_races_started
725            .fetch_add(1, Ordering::Relaxed);
726    }
727
728    pub(crate) fn metadata_peer_candidate(&self) {
729        self.inner
730            .metadata_peer_candidates
731            .fetch_add(1, Ordering::Relaxed);
732    }
733
734    pub(crate) fn metadata_live_peer_join(&self) {
735        self.inner
736            .metadata_live_peer_joins
737            .fetch_add(1, Ordering::Relaxed);
738    }
739
740    pub(crate) fn metadata_peer_canceled(&self, count: usize) {
741        self.inner
742            .metadata_peer_canceled
743            .fetch_add(count as u64, Ordering::Relaxed);
744    }
745
746    pub(crate) fn peer_lookup_requested(&self) {
747        self.inner
748            .peer_lookup_requested
749            .fetch_add(1, Ordering::Relaxed);
750    }
751
752    pub(crate) fn peer_lookup_started(&self) {
753        self.inner
754            .peer_lookup_started
755            .fetch_add(1, Ordering::Relaxed);
756    }
757
758    pub(crate) fn peer_lookup_rate_limited(&self) {
759        self.inner
760            .peer_lookup_rate_limited
761            .fetch_add(1, Ordering::Relaxed);
762    }
763
764    pub(crate) fn peer_lookup_empty(&self) {
765        self.inner.peer_lookup_empty.fetch_add(1, Ordering::Relaxed);
766    }
767
768    pub(crate) fn peer_lookup_query(&self) {
769        self.inner
770            .peer_lookup_queries
771            .fetch_add(1, Ordering::Relaxed);
772    }
773
774    pub(crate) fn peer_lookup_response(&self) {
775        self.inner
776            .peer_lookup_responses
777            .fetch_add(1, Ordering::Relaxed);
778    }
779
780    pub(crate) fn peer_lookup_timeout(&self) {
781        self.inner
782            .peer_lookup_timeouts
783            .fetch_add(1, Ordering::Relaxed);
784    }
785
786    pub(crate) fn peer_lookup_send_failed(&self) {
787        self.inner
788            .peer_lookup_send_failures
789            .fetch_add(1, Ordering::Relaxed);
790    }
791
792    pub(crate) fn peer_lookup_peer_found(&self) {
793        self.inner
794            .peer_lookup_peers_found
795            .fetch_add(1, Ordering::Relaxed);
796    }
797
798    pub(crate) fn peer_lookup_response_dropped(&self) {
799        self.inner
800            .peer_lookup_response_dropped
801            .fetch_add(1, Ordering::Relaxed);
802    }
803
804    pub(crate) fn peer_lookup_output_dropped(&self) {
805        self.inner
806            .peer_lookup_output_dropped
807            .fetch_add(1, Ordering::Relaxed);
808    }
809
810    pub(crate) fn metadata_peer_attempt(&self) {
811        self.inner
812            .metadata_peer_attempts
813            .fetch_add(1, Ordering::Relaxed);
814    }
815
816    pub(crate) fn metadata_peer_succeeded(&self) {
817        self.inner
818            .metadata_peer_succeeded
819            .fetch_add(1, Ordering::Relaxed);
820    }
821
822    pub(crate) fn metadata_peer_failed(&self) {
823        self.inner
824            .metadata_peer_failed
825            .fetch_add(1, Ordering::Relaxed);
826    }
827
828    pub(crate) fn metadata_peer_timeout(&self) {
829        self.inner
830            .metadata_peer_timeouts
831            .fetch_add(1, Ordering::Relaxed);
832    }
833
834    pub(crate) fn metadata_connect_failed(&self) {
835        self.inner
836            .metadata_connect_failed
837            .fetch_add(1, Ordering::Relaxed);
838    }
839
840    pub(crate) fn metadata_no_extension(&self) {
841        self.inner
842            .metadata_no_extension
843            .fetch_add(1, Ordering::Relaxed);
844    }
845
846    pub(crate) fn metadata_peer_failure_cache_hit(&self) {
847        self.inner
848            .metadata_peer_failure_cache_hits
849            .fetch_add(1, Ordering::Relaxed);
850    }
851
852    pub(crate) fn set_metadata_peer_failure_cache_entries(&self, count: usize) {
853        self.inner
854            .metadata_peer_failure_cache_entries
855            .store(count, Ordering::Relaxed);
856    }
857
858    pub(crate) fn set_node_pool_size(&self, size: usize) {
859        self.inner.node_pool_size.store(size, Ordering::Relaxed);
860    }
861
862    pub(crate) fn set_find_node_in_flight(&self, count: usize) {
863        self.inner
864            .find_node_in_flight
865            .store(count, Ordering::Relaxed);
866    }
867
868    pub(crate) fn set_find_node_effective_rate(&self, rate: u32) {
869        self.inner
870            .find_node_effective_rate_per_sec
871            .store(rate, Ordering::Relaxed);
872    }
873
874    pub(crate) fn query_new(&self) {
875        self.inner.queries_new.fetch_add(1, Ordering::Relaxed);
876    }
877
878    pub(crate) fn query_revisit(&self) {
879        self.inner.queries_revisit.fetch_add(1, Ordering::Relaxed);
880    }
881
882    pub(crate) fn query_bootstrap(&self) {
883        self.inner.queries_bootstrap.fetch_add(1, Ordering::Relaxed);
884    }
885
886    pub(crate) fn response(&self) {
887        self.inner.responses.fetch_add(1, Ordering::Relaxed);
888    }
889
890    pub(crate) fn unmatched_response(&self) {
891        self.inner
892            .unmatched_responses
893            .fetch_add(1, Ordering::Relaxed);
894    }
895
896    pub(crate) fn timeout(&self) {
897        self.inner.timeouts.fetch_add(1, Ordering::Relaxed);
898    }
899
900    pub(crate) fn send_failure(&self) {
901        self.inner.send_failures.fetch_add(1, Ordering::Relaxed);
902    }
903
904    pub(crate) fn crawl_event_dropped_response(&self) {
905        self.inner
906            .crawl_events_dropped_response
907            .fetch_add(1, Ordering::Relaxed);
908    }
909
910    pub(crate) fn crawl_event_dropped_discovered(&self) {
911        self.inner
912            .crawl_events_dropped_discovered
913            .fetch_add(1, Ordering::Relaxed);
914    }
915
916    pub(crate) fn udp_received(&self) {
917        self.inner.udp_received.fetch_add(1, Ordering::Relaxed);
918    }
919
920    pub(crate) fn udp_received_bytes(&self, bytes: usize) {
921        self.inner
922            .udp_rx_bytes
923            .fetch_add(bytes as u64, Ordering::Relaxed);
924    }
925
926    pub(crate) fn udp_sent(&self, bytes: usize) {
927        self.inner.udp_tx_packets.fetch_add(1, Ordering::Relaxed);
928        self.inner
929            .udp_tx_bytes
930            .fetch_add(bytes as u64, Ordering::Relaxed);
931    }
932
933    pub(crate) fn inbound_query(&self, query: &str) {
934        let counter = match query {
935            "ping" => &self.inner.inbound_ping,
936            "find_node" => &self.inner.inbound_find_node,
937            "get_peers" => &self.inner.inbound_get_peers,
938            "announce_peer" => &self.inner.inbound_announce_peer,
939            _ => &self.inner.inbound_other,
940        };
941        counter.fetch_add(1, Ordering::Relaxed);
942    }
943
944    pub(crate) fn response_normal(&self) {
945        self.inner.response_normal.fetch_add(1, Ordering::Relaxed);
946    }
947
948    pub(crate) fn response_send_failed(&self) {
949        self.inner
950            .response_send_failed
951            .fetch_add(1, Ordering::Relaxed);
952    }
953
954    pub(crate) fn announce_accepted(&self) {
955        self.inner.announce_accepted.fetch_add(1, Ordering::Relaxed);
956    }
957
958    pub(crate) fn announce_invalid_token(&self) {
959        self.inner
960            .announce_invalid_token
961            .fetch_add(1, Ordering::Relaxed);
962    }
963
964    pub(crate) fn announce_filtered(&self) {
965        self.inner.announce_filtered.fetch_add(1, Ordering::Relaxed);
966    }
967
968    pub(crate) fn node_admitted(&self) {
969        self.inner.node_admitted.fetch_add(1, Ordering::Relaxed);
970    }
971
972    pub(crate) fn node_replaced(&self) {
973        self.inner.node_replaced.fetch_add(1, Ordering::Relaxed);
974    }
975
976    pub(crate) fn node_dropped_duplicate(&self) {
977        self.inner
978            .node_dropped_duplicate
979            .fetch_add(1, Ordering::Relaxed);
980    }
981
982    pub(crate) fn node_dropped_rate_limited(&self) {
983        self.inner
984            .node_dropped_rate_limited
985            .fetch_add(1, Ordering::Relaxed);
986    }
987
988    pub(crate) fn node_dropped_invalid(&self) {
989        self.inner
990            .node_dropped_invalid
991            .fetch_add(1, Ordering::Relaxed);
992    }
993
994    pub(crate) fn metadata_bytes_downloaded(&self, bytes: usize) {
995        self.inner
996            .metadata_bytes_downloaded
997            .fetch_add(bytes as u64, Ordering::Relaxed);
998    }
999
1000    pub(crate) fn metadata_failure_timeout(&self) {
1001        self.inner
1002            .metadata_failure_timeout
1003            .fetch_add(1, Ordering::Relaxed);
1004    }
1005
1006    pub(crate) fn metadata_failure_connect(&self) {
1007        self.inner
1008            .metadata_failure_connect
1009            .fetch_add(1, Ordering::Relaxed);
1010    }
1011
1012    pub(crate) fn metadata_failure_no_extension(&self) {
1013        self.inner
1014            .metadata_failure_no_extension
1015            .fetch_add(1, Ordering::Relaxed);
1016    }
1017
1018    pub(crate) fn metadata_failure_send(&self) {
1019        self.inner
1020            .metadata_failure_send
1021            .fetch_add(1, Ordering::Relaxed);
1022    }
1023
1024    pub(crate) fn metadata_failure_size_limit(&self) {
1025        self.inner
1026            .metadata_failure_size_limit
1027            .fetch_add(1, Ordering::Relaxed);
1028    }
1029
1030    pub(crate) fn metadata_failure_sha1(&self) {
1031        self.inner
1032            .metadata_failure_sha1
1033            .fetch_add(1, Ordering::Relaxed);
1034    }
1035
1036    pub(crate) fn metadata_failure_parse(&self) {
1037        self.inner
1038            .metadata_failure_parse
1039            .fetch_add(1, Ordering::Relaxed);
1040    }
1041
1042    pub(crate) fn metadata_failure_other(&self) {
1043        self.inner
1044            .metadata_failure_other
1045            .fetch_add(1, Ordering::Relaxed);
1046    }
1047
1048    pub(crate) fn peer_cache_hit_timeout(&self) {
1049        self.inner
1050            .peer_cache_timeout_hits
1051            .fetch_add(1, Ordering::Relaxed);
1052    }
1053
1054    pub(crate) fn peer_cache_hit_connect(&self) {
1055        self.inner
1056            .peer_cache_connect_hits
1057            .fetch_add(1, Ordering::Relaxed);
1058    }
1059
1060    pub(crate) fn observe_metadata_queue_wait(&self, millis: u64) {
1061        self.inner.queue_wait_ms.record(millis);
1062    }
1063
1064    pub(crate) fn observe_metadata_fetch_duration(&self, millis: u64) {
1065        self.inner.fetch_duration_ms.record(millis);
1066    }
1067
1068    pub(crate) fn observe_metadata_size(&self, bytes: usize) {
1069        self.inner.metadata_size_bytes.record(bytes as u64);
1070    }
1071
1072    pub(crate) fn udp_queue_full(&self) {
1073        self.inner.udp_queue_full.fetch_add(1, Ordering::Relaxed);
1074    }
1075
1076    pub(crate) fn udp_invalid(&self) {
1077        self.inner.udp_invalid.fetch_add(1, Ordering::Relaxed);
1078    }
1079
1080    pub(crate) fn udp_response_rate_limited(&self) {
1081        self.inner
1082            .udp_responses_rate_limited
1083            .fetch_add(1, Ordering::Relaxed);
1084    }
1085
1086    pub(crate) fn udp_response_priority_reserved(&self) {
1087        self.inner
1088            .udp_responses_priority_reserved
1089            .fetch_add(1, Ordering::Relaxed);
1090    }
1091}
1092
1093#[cfg(test)]
1094mod tests {
1095    use super::*;
1096
1097    #[test]
1098    fn cloned_handle_updates_one_snapshot() {
1099        let stats = DhtRuntimeStats::with_limits(DhtRuntimeLimits {
1100            metadata_queue: 100,
1101            node_pool: 1_000,
1102            node_pool_low_watermark: 10,
1103            find_node_in_flight: 512,
1104            initial_find_node_rate: 200,
1105            hash_ingress_queue: 20,
1106            crawl_priority_queue: 30,
1107            crawl_discovery_queue: 40,
1108        });
1109        let writer = stats.clone();
1110
1111        writer.hash_received();
1112        writer.hash_ingress_dropped();
1113        writer.set_hash_ingress_queue_depth(11);
1114        writer.set_crawl_priority_queue_depth(12);
1115        writer.set_crawl_discovery_queue_depth(13);
1116        writer.set_metadata_queue(42, 7);
1117        writer.metadata_queue_inserted();
1118        writer.metadata_queue_deduplicated();
1119        writer.metadata_queue_evicted();
1120        writer.metadata_queue_stale(3);
1121        writer.metadata_race_started();
1122        writer.metadata_peer_candidate();
1123        writer.metadata_live_peer_join();
1124        writer.metadata_peer_canceled(2);
1125        writer.peer_lookup_requested();
1126        writer.peer_lookup_started();
1127        writer.peer_lookup_rate_limited();
1128        writer.peer_lookup_empty();
1129        writer.peer_lookup_query();
1130        writer.peer_lookup_response();
1131        writer.peer_lookup_timeout();
1132        writer.peer_lookup_send_failed();
1133        writer.peer_lookup_peer_found();
1134        writer.peer_lookup_response_dropped();
1135        writer.peer_lookup_output_dropped();
1136        writer.metadata_peer_attempt();
1137        writer.metadata_peer_succeeded();
1138        writer.metadata_peer_failed();
1139        writer.metadata_peer_timeout();
1140        writer.metadata_connect_failed();
1141        writer.metadata_no_extension();
1142        writer.metadata_peer_failure_cache_hit();
1143        writer.set_metadata_peer_failure_cache_entries(9);
1144        writer.set_node_pool_size(321);
1145        writer.set_find_node_in_flight(12);
1146        writer.set_find_node_effective_rate(150);
1147        writer.query_new();
1148        writer.query_revisit();
1149        writer.query_bootstrap();
1150        writer.response();
1151        writer.unmatched_response();
1152        writer.timeout();
1153        writer.send_failure();
1154        writer.crawl_event_dropped_response();
1155        writer.crawl_event_dropped_discovered();
1156        writer.udp_received();
1157        writer.udp_queue_full();
1158        writer.udp_invalid();
1159        writer.udp_response_rate_limited();
1160        writer.udp_response_priority_reserved();
1161
1162        assert_eq!(
1163            stats.snapshot(),
1164            DhtRuntimeSnapshot {
1165                hashes_received: 1,
1166                hash_ingress_dropped: 1,
1167                hash_ingress_queue_depth: 11,
1168                hash_ingress_queue_capacity: 20,
1169                crawl_priority_queue_depth: 12,
1170                crawl_priority_queue_capacity: 30,
1171                crawl_discovery_queue_depth: 13,
1172                crawl_discovery_queue_capacity: 40,
1173                metadata_queue_depth: 42,
1174                metadata_queue_max: 100,
1175                metadata_in_flight: 7,
1176                metadata_queue_inserted: 1,
1177                metadata_queue_deduplicated: 1,
1178                metadata_queue_evicted: 1,
1179                metadata_queue_stale: 3,
1180                metadata_races_started: 1,
1181                metadata_peer_candidates: 1,
1182                metadata_live_peer_joins: 1,
1183                metadata_peer_canceled: 2,
1184                peer_lookup_requested: 1,
1185                peer_lookup_started: 1,
1186                peer_lookup_rate_limited: 1,
1187                peer_lookup_empty: 1,
1188                peer_lookup_queries: 1,
1189                peer_lookup_responses: 1,
1190                peer_lookup_timeouts: 1,
1191                peer_lookup_send_failures: 1,
1192                peer_lookup_peers_found: 1,
1193                peer_lookup_response_dropped: 1,
1194                peer_lookup_output_dropped: 1,
1195                metadata_peer_attempts: 1,
1196                metadata_peer_succeeded: 1,
1197                metadata_peer_failed: 1,
1198                metadata_peer_timeouts: 1,
1199                metadata_connect_failed: 1,
1200                metadata_no_extension: 1,
1201                metadata_peer_failure_cache_hits: 1,
1202                metadata_peer_failure_cache_entries: 9,
1203                node_pool_size: 321,
1204                node_pool_capacity: 1_000,
1205                node_pool_low_watermark: 10,
1206                find_node_in_flight: 12,
1207                find_node_in_flight_max: 512,
1208                find_node_effective_rate_per_sec: 150,
1209                queries_new: 1,
1210                queries_revisit: 1,
1211                queries_bootstrap: 1,
1212                responses: 1,
1213                unmatched_responses: 1,
1214                timeouts: 1,
1215                send_failures: 1,
1216                crawl_events_dropped_response: 1,
1217                crawl_events_dropped_discovered: 1,
1218                udp_received: 1,
1219                udp_queue_full: 1,
1220                udp_invalid: 1,
1221                udp_responses_rate_limited: 1,
1222                udp_responses_priority_reserved: 1,
1223            }
1224        );
1225    }
1226
1227    #[test]
1228    fn observability_snapshot_tracks_fixed_histograms_and_categories() {
1229        let stats = DhtRuntimeStats::default();
1230        stats.udp_received();
1231        stats.udp_received_bytes(128);
1232        stats.udp_sent(64);
1233        stats.inbound_query("ping");
1234        stats.inbound_query("unknown");
1235        stats.node_admitted();
1236        stats.metadata_bytes_downloaded(1_024);
1237        stats.metadata_failure_parse();
1238        stats.observe_metadata_queue_wait(75);
1239        stats.observe_metadata_queue_wait(600);
1240        stats.observe_metadata_fetch_duration(1_500);
1241        stats.observe_metadata_size(70_000);
1242
1243        let snapshot = stats.observability_snapshot();
1244        assert_eq!(snapshot.udp_rx_packets, 1);
1245        assert_eq!(snapshot.udp_rx_bytes, 128);
1246        assert_eq!(snapshot.udp_tx_packets, 1);
1247        assert_eq!(snapshot.udp_tx_bytes, 64);
1248        assert_eq!(snapshot.inbound_ping, 1);
1249        assert_eq!(snapshot.inbound_other, 1);
1250        assert_eq!(snapshot.node_admitted, 1);
1251        assert_eq!(snapshot.metadata_bytes_downloaded, 1_024);
1252        assert_eq!(snapshot.metadata_failure_parse, 1);
1253        assert_eq!(snapshot.queue_wait_ms.percentile(0.50), Some(100));
1254        assert_eq!(snapshot.queue_wait_ms.percentile(0.95), Some(1_000));
1255        assert_eq!(snapshot.fetch_duration_ms.percentile(0.95), Some(2_000));
1256        assert_eq!(snapshot.metadata_size_bytes.percentile(0.50), Some(131_072));
1257    }
1258}