refactor: 拆分核心运行和存储模块

This commit is contained in:
chuan
2026-08-10 15:48:24 +08:00
parent 9f0911c596
commit c70d7e5c7b
13 changed files with 2459 additions and 2394 deletions
+8 -848
View File
@@ -1,107 +1,16 @@
// 负责收集 DHT 抓取调度网络和 Metadata 运行指标
use crate::types::DiscoverySource;
use std::sync::Arc;
use std::sync::atomic::{AtomicU32, AtomicU64, AtomicUsize, Ordering};
pub const METADATA_QUEUE_WAIT_BUCKETS_MS: [u64; 8] = [10, 50, 100, 250, 500, 1_000, 2_000, 5_000];
pub const METADATA_FETCH_BUCKETS_MS: [u64; 7] = [250, 500, 1_000, 2_000, 4_000, 6_000, 10_000];
pub const METADATA_SIZE_BUCKETS_BYTES: [u64; 8] = [
16 * 1024,
32 * 1024,
64 * 1024,
128 * 1024,
256 * 1024,
512 * 1024,
1024 * 1024,
10 * 1024 * 1024,
];
mod histogram;
mod recorders;
#[derive(Debug, Clone, PartialEq, Eq)]
/// Snapshot of a non-cumulative fixed-bucket histogram.
pub struct FixedHistogramSnapshot {
/// Inclusive upper bound for each bucket.
pub bounds: Vec<u64>,
/// Per-bucket counts; values are not cumulative.
pub counts: Vec<u64>,
/// Values larger than the final bound.
pub overflow: u64,
/// Total observations including overflow.
pub count: u64,
/// Wrapping sum of all observed values.
pub sum: u64,
}
impl FixedHistogramSnapshot {
/// Returns the inclusive bucket bound containing the requested approximate percentile.
///
/// `percentile` is clamped to `0.0..=1.0`. Overflow observations return the final finite
/// bound because the histogram intentionally stores no dynamic maximum.
pub fn percentile(&self, percentile: f64) -> Option<u64> {
if self.count == 0 {
return None;
}
let rank = ((self.count as f64 * percentile.clamp(0.0, 1.0)).ceil() as u64).max(1);
let mut cumulative = 0u64;
for (bound, count) in self.bounds.iter().zip(&self.counts) {
cumulative = cumulative.saturating_add(*count);
if cumulative >= rank {
return Some(*bound);
}
}
self.bounds.last().copied()
}
}
struct AtomicFixedHistogram<const N: usize> {
bounds: [u64; N],
counts: [AtomicU64; N],
overflow: AtomicU64,
count: AtomicU64,
sum: AtomicU64,
}
impl<const N: usize> AtomicFixedHistogram<N> {
fn new(bounds: [u64; N]) -> Self {
Self {
bounds,
counts: std::array::from_fn(|_| AtomicU64::new(0)),
overflow: AtomicU64::new(0),
count: AtomicU64::new(0),
sum: AtomicU64::new(0),
}
}
fn record(&self, value: u64) {
if let Some(index) = self.bounds.iter().position(|bound| value <= *bound) {
self.counts[index].fetch_add(1, Ordering::Relaxed);
} else {
self.overflow.fetch_add(1, Ordering::Relaxed);
}
self.count.fetch_add(1, Ordering::Relaxed);
self.sum.fetch_add(value, Ordering::Relaxed);
}
fn snapshot(&self) -> FixedHistogramSnapshot {
FixedHistogramSnapshot {
bounds: self.bounds.to_vec(),
counts: self
.counts
.iter()
.map(|count| count.load(Ordering::Relaxed))
.collect(),
overflow: self.overflow.load(Ordering::Relaxed),
count: self.count.load(Ordering::Relaxed),
sum: self.sum.load(Ordering::Relaxed),
}
}
}
impl<const N: usize> Default for AtomicFixedHistogram<N> {
fn default() -> Self {
Self::new([0; N])
}
}
use histogram::AtomicFixedHistogram;
pub use histogram::{
FixedHistogramSnapshot, METADATA_FETCH_BUCKETS_MS, METADATA_QUEUE_WAIT_BUCKETS_MS,
METADATA_SIZE_BUCKETS_BYTES,
};
/// Transport-neutral counters and fixed-bucket histograms used by dashboards.
#[derive(Debug, Clone, PartialEq, Eq)]
@@ -809,756 +718,7 @@ impl DhtRuntimeStats {
metadata_size_bytes: inner.metadata_size_bytes.snapshot(),
}
}
pub(crate) fn hash_received(&self) {
self.inner.hashes_received.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn hash_ingress_dropped(&self) {
self.inner
.hash_ingress_dropped
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn set_hash_ingress_queue_depth(&self, depth: usize) {
self.inner
.hash_ingress_queue_depth
.store(depth, Ordering::Relaxed);
}
pub(crate) fn set_crawl_priority_queue_depth(&self, depth: usize) {
self.inner
.crawl_priority_queue_depth
.store(depth, Ordering::Relaxed);
}
pub(crate) fn set_crawl_discovery_queue_depth(&self, depth: usize) {
self.inner
.crawl_discovery_queue_depth
.store(depth, Ordering::Relaxed);
}
pub(crate) fn set_metadata_queue(&self, depth: usize, in_flight: usize) {
self.inner
.metadata_queue_depth
.store(depth, Ordering::Relaxed);
self.inner
.metadata_in_flight
.store(in_flight, Ordering::Relaxed);
}
pub(crate) fn metadata_queue_deduplicated(&self) {
self.inner
.metadata_queue_deduplicated
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_queue_inserted(&self) {
self.inner
.metadata_queue_inserted
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_queue_evicted(&self) {
self.inner
.metadata_queue_evicted
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_queue_stale(&self, count: usize) {
self.inner
.metadata_queue_stale
.fetch_add(count as u64, Ordering::Relaxed);
}
pub(crate) fn metadata_race_started(&self) {
self.inner
.metadata_races_started
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_peer_candidate(&self) {
self.inner
.metadata_peer_candidates
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_live_peer_join(&self) {
self.inner
.metadata_live_peer_joins
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_peer_canceled(&self, count: usize) {
self.inner
.metadata_peer_canceled
.fetch_add(count as u64, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_requested(&self) {
self.inner
.peer_lookup_requested
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_started(&self) {
self.inner
.peer_lookup_started
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_rate_limited(&self) {
self.inner
.peer_lookup_rate_limited
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_empty(&self) {
self.inner.peer_lookup_empty.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_query(&self) {
self.inner
.peer_lookup_queries
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_response(&self) {
self.inner
.peer_lookup_responses
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_timeout(&self) {
self.inner
.peer_lookup_timeouts
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_send_failed(&self) {
self.inner
.peer_lookup_send_failures
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_peer_found(&self) {
self.inner
.peer_lookup_peers_found
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_response_dropped(&self) {
self.inner
.peer_lookup_response_dropped
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_output_dropped(&self) {
self.inner
.peer_lookup_output_dropped
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_preferred_succeeded(&self) {
self.inner
.peer_lookup_preferred_succeeded
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_fallback(&self) {
self.inner
.peer_lookup_fallbacks
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn configure_sample_candidate_queue(&self, capacity: usize) {
self.inner
.sample_candidate_queue_capacity
.store(capacity, Ordering::Relaxed);
}
pub(crate) fn set_sample_candidate_queue_depth(&self, depth: usize) {
self.inner
.sample_candidate_queue_depth
.store(depth, Ordering::Relaxed);
}
pub(crate) fn sample_candidate_routed(&self) {
self.inner
.sample_candidates_routed
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn sample_candidate_fallback(&self) {
self.inner
.sample_candidates_fallback
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn sample_query(&self, direct: bool) {
self.inner
.sample_infohashes_queries
.fetch_add(1, Ordering::Relaxed);
let source = if direct {
&self.inner.sample_infohashes_direct_queries
} else {
&self.inner.sample_infohashes_snapshot_queries
};
source.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn sample_response(&self, direct: bool) {
self.inner
.sample_infohashes_responses
.fetch_add(1, Ordering::Relaxed);
let source = if direct {
&self.inner.sample_infohashes_direct_responses
} else {
&self.inner.sample_infohashes_snapshot_responses
};
source.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn sample_timeout(&self) {
self.inner
.sample_infohashes_timeouts
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn sample_send_failed(&self) {
self.inner
.sample_infohashes_send_failures
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn sample_response_dropped(&self) {
self.inner
.sample_infohashes_response_dropped
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn sample_hash_discovered(&self) {
self.inner
.sample_infohashes_hashes_discovered
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn sample_hash_filtered(&self, count: usize) {
self.inner
.sample_infohashes_hashes_filtered
.fetch_add(count as u64, Ordering::Relaxed);
}
pub(crate) fn sample_hash_duplicate(&self) {
self.inner
.sample_infohashes_hashes_duplicate
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn sample_hash_dropped(&self) {
self.inner
.sample_infohashes_hashes_dropped
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_peer_attempt(&self) {
self.inner
.metadata_peer_attempts
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_peer_succeeded(&self) {
self.inner
.metadata_peer_succeeded
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_success_from(&self, source: DiscoverySource) {
let counter = match source {
DiscoverySource::AnnouncePeer => &self.inner.metadata_success_from_announce,
DiscoverySource::SampleDirect => &self.inner.metadata_success_from_sample_direct,
DiscoverySource::SampleSnapshot => &self.inner.metadata_success_from_sample_snapshot,
DiscoverySource::ActiveLookup => &self.inner.metadata_success_from_active_lookup,
};
counter.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_peer_failed(&self) {
self.inner
.metadata_peer_failed
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_peer_timeout(&self) {
self.inner
.metadata_peer_timeouts
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_connect_failed(&self) {
self.inner
.metadata_connect_failed
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_no_extension(&self) {
self.inner
.metadata_no_extension
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_peer_failure_cache_hit(&self) {
self.inner
.metadata_peer_failure_cache_hits
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn set_metadata_peer_failure_cache_entries(&self, count: usize) {
self.inner
.metadata_peer_failure_cache_entries
.store(count, Ordering::Relaxed);
}
pub(crate) fn set_node_pool_size(&self, size: usize) {
self.inner.node_pool_size.store(size, Ordering::Relaxed);
}
pub(crate) fn set_find_node_in_flight(&self, count: usize) {
self.inner
.find_node_in_flight
.store(count, Ordering::Relaxed);
}
pub(crate) fn set_find_node_effective_rate(&self, rate: u32) {
self.inner
.find_node_effective_rate_per_sec
.store(rate, Ordering::Relaxed);
}
pub(crate) fn query_new(&self) {
self.inner.queries_new.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn query_revisit(&self) {
self.inner.queries_revisit.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn query_bootstrap(&self) {
self.inner.queries_bootstrap.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn response(&self) {
self.inner.responses.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn unmatched_response(&self) {
self.inner
.unmatched_responses
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn timeout(&self) {
self.inner.timeouts.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn send_failure(&self) {
self.inner.send_failures.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn crawl_event_dropped_response(&self) {
self.inner
.crawl_events_dropped_response
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn crawl_event_dropped_discovered(&self) {
self.inner
.crawl_events_dropped_discovered
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn udp_received(&self) {
self.inner.udp_received.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn udp_received_bytes(&self, bytes: usize) {
self.inner
.udp_rx_bytes
.fetch_add(bytes as u64, Ordering::Relaxed);
}
pub(crate) fn udp_sent(&self, bytes: usize) {
self.inner.udp_tx_packets.fetch_add(1, Ordering::Relaxed);
self.inner
.udp_tx_bytes
.fetch_add(bytes as u64, Ordering::Relaxed);
}
pub(crate) fn inbound_query(&self, query: &str) {
let counter = match query {
"ping" => &self.inner.inbound_ping,
"find_node" => &self.inner.inbound_find_node,
"get_peers" => &self.inner.inbound_get_peers,
"announce_peer" => &self.inner.inbound_announce_peer,
_ => &self.inner.inbound_other,
};
counter.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn response_normal(&self) {
self.inner.response_normal.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn response_send_failed(&self) {
self.inner
.response_send_failed
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn announce_accepted(&self) {
self.inner.announce_accepted.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn announce_invalid_token(&self) {
self.inner
.announce_invalid_token
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn announce_filtered(&self) {
self.inner.announce_filtered.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn node_admitted(&self) {
self.inner.node_admitted.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn node_replaced(&self) {
self.inner.node_replaced.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn node_dropped_duplicate(&self) {
self.inner
.node_dropped_duplicate
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn node_dropped_rate_limited(&self) {
self.inner
.node_dropped_rate_limited
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn node_dropped_invalid(&self) {
self.inner
.node_dropped_invalid
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_bytes_downloaded(&self, bytes: usize) {
self.inner
.metadata_bytes_downloaded
.fetch_add(bytes as u64, Ordering::Relaxed);
}
pub(crate) fn metadata_failure_timeout(&self) {
self.inner
.metadata_failure_timeout
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_failure_connect(&self) {
self.inner
.metadata_failure_connect
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_failure_no_extension(&self) {
self.inner
.metadata_failure_no_extension
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_failure_send(&self) {
self.inner
.metadata_failure_send
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_failure_size_limit(&self) {
self.inner
.metadata_failure_size_limit
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_failure_sha1(&self) {
self.inner
.metadata_failure_sha1
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_failure_parse(&self) {
self.inner
.metadata_failure_parse
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_failure_other(&self) {
self.inner
.metadata_failure_other
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_cache_hit_timeout(&self) {
self.inner
.peer_cache_timeout_hits
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_cache_hit_connect(&self) {
self.inner
.peer_cache_connect_hits
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn observe_metadata_queue_wait(&self, millis: u64) {
self.inner.queue_wait_ms.record(millis);
}
pub(crate) fn observe_metadata_fetch_duration(&self, millis: u64) {
self.inner.fetch_duration_ms.record(millis);
}
pub(crate) fn observe_metadata_size(&self, bytes: usize) {
self.inner.metadata_size_bytes.record(bytes as u64);
}
pub(crate) fn udp_queue_full(&self) {
self.inner.udp_queue_full.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn udp_invalid(&self) {
self.inner.udp_invalid.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn udp_response_rate_limited(&self) {
self.inner
.udp_responses_rate_limited
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn udp_response_priority_reserved(&self) {
self.inner
.udp_responses_priority_reserved
.fetch_add(1, Ordering::Relaxed);
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn cloned_handle_updates_one_snapshot() {
let stats = DhtRuntimeStats::with_limits(DhtRuntimeLimits {
metadata_queue: 100,
node_pool: 1_000,
node_pool_low_watermark: 10,
find_node_in_flight: 512,
initial_find_node_rate: 200,
hash_ingress_queue: 20,
crawl_priority_queue: 30,
crawl_discovery_queue: 40,
});
let writer = stats.clone();
writer.hash_received();
writer.hash_ingress_dropped();
writer.set_hash_ingress_queue_depth(11);
writer.set_crawl_priority_queue_depth(12);
writer.set_crawl_discovery_queue_depth(13);
writer.set_metadata_queue(42, 7);
writer.metadata_queue_inserted();
writer.metadata_queue_deduplicated();
writer.metadata_queue_evicted();
writer.metadata_queue_stale(3);
writer.metadata_race_started();
writer.metadata_peer_candidate();
writer.metadata_live_peer_join();
writer.metadata_peer_canceled(2);
writer.peer_lookup_requested();
writer.peer_lookup_started();
writer.peer_lookup_rate_limited();
writer.peer_lookup_empty();
writer.peer_lookup_query();
writer.peer_lookup_response();
writer.peer_lookup_timeout();
writer.peer_lookup_send_failed();
writer.peer_lookup_peer_found();
writer.peer_lookup_response_dropped();
writer.peer_lookup_output_dropped();
writer.peer_lookup_preferred_succeeded();
writer.peer_lookup_fallback();
writer.configure_sample_candidate_queue(20);
writer.set_sample_candidate_queue_depth(14);
writer.sample_candidate_routed();
writer.sample_candidate_fallback();
writer.sample_query(true);
writer.sample_response(true);
writer.sample_timeout();
writer.sample_send_failed();
writer.sample_response_dropped();
writer.sample_hash_discovered();
writer.sample_hash_filtered(1);
writer.sample_hash_duplicate();
writer.sample_hash_dropped();
writer.metadata_peer_attempt();
writer.metadata_peer_succeeded();
writer.metadata_success_from(DiscoverySource::AnnouncePeer);
writer.metadata_success_from(DiscoverySource::SampleDirect);
writer.metadata_success_from(DiscoverySource::SampleSnapshot);
writer.metadata_success_from(DiscoverySource::ActiveLookup);
writer.metadata_peer_failed();
writer.metadata_peer_timeout();
writer.metadata_connect_failed();
writer.metadata_no_extension();
writer.metadata_peer_failure_cache_hit();
writer.set_metadata_peer_failure_cache_entries(9);
writer.set_node_pool_size(321);
writer.set_find_node_in_flight(12);
writer.set_find_node_effective_rate(150);
writer.query_new();
writer.query_revisit();
writer.query_bootstrap();
writer.response();
writer.unmatched_response();
writer.timeout();
writer.send_failure();
writer.crawl_event_dropped_response();
writer.crawl_event_dropped_discovered();
writer.udp_received();
writer.udp_queue_full();
writer.udp_invalid();
writer.udp_response_rate_limited();
writer.udp_response_priority_reserved();
assert_eq!(
stats.snapshot(),
DhtRuntimeSnapshot {
hashes_received: 1,
hash_ingress_dropped: 1,
hash_ingress_queue_depth: 11,
hash_ingress_queue_capacity: 20,
crawl_priority_queue_depth: 12,
crawl_priority_queue_capacity: 30,
crawl_discovery_queue_depth: 13,
crawl_discovery_queue_capacity: 40,
metadata_queue_depth: 42,
metadata_queue_max: 100,
metadata_in_flight: 7,
metadata_queue_inserted: 1,
metadata_queue_deduplicated: 1,
metadata_queue_evicted: 1,
metadata_queue_stale: 3,
metadata_races_started: 1,
metadata_peer_candidates: 1,
metadata_live_peer_joins: 1,
metadata_peer_canceled: 2,
peer_lookup_requested: 1,
peer_lookup_started: 1,
peer_lookup_rate_limited: 1,
peer_lookup_empty: 1,
peer_lookup_queries: 1,
peer_lookup_responses: 1,
peer_lookup_timeouts: 1,
peer_lookup_send_failures: 1,
peer_lookup_peers_found: 1,
peer_lookup_response_dropped: 1,
peer_lookup_output_dropped: 1,
peer_lookup_preferred_succeeded: 1,
peer_lookup_fallbacks: 1,
sample_candidate_queue_depth: 14,
sample_candidate_queue_capacity: 20,
sample_candidates_routed: 1,
sample_candidates_fallback: 1,
sample_infohashes_queries: 1,
sample_infohashes_direct_queries: 1,
sample_infohashes_snapshot_queries: 0,
sample_infohashes_responses: 1,
sample_infohashes_direct_responses: 1,
sample_infohashes_snapshot_responses: 0,
sample_infohashes_timeouts: 1,
sample_infohashes_send_failures: 1,
sample_infohashes_response_dropped: 1,
sample_infohashes_hashes_discovered: 1,
sample_infohashes_hashes_filtered: 1,
sample_infohashes_hashes_duplicate: 1,
sample_infohashes_hashes_dropped: 1,
metadata_peer_attempts: 1,
metadata_peer_succeeded: 1,
metadata_success_from_announce: 1,
metadata_success_from_sample_direct: 1,
metadata_success_from_sample_snapshot: 1,
metadata_success_from_active_lookup: 1,
metadata_peer_failed: 1,
metadata_peer_timeouts: 1,
metadata_connect_failed: 1,
metadata_no_extension: 1,
metadata_peer_failure_cache_hits: 1,
metadata_peer_failure_cache_entries: 9,
node_pool_size: 321,
node_pool_capacity: 1_000,
node_pool_low_watermark: 10,
find_node_in_flight: 12,
find_node_in_flight_max: 512,
find_node_effective_rate_per_sec: 150,
queries_new: 1,
queries_revisit: 1,
queries_bootstrap: 1,
responses: 1,
unmatched_responses: 1,
timeouts: 1,
send_failures: 1,
crawl_events_dropped_response: 1,
crawl_events_dropped_discovered: 1,
udp_received: 1,
udp_queue_full: 1,
udp_invalid: 1,
udp_responses_rate_limited: 1,
udp_responses_priority_reserved: 1,
}
);
}
#[test]
fn observability_snapshot_tracks_fixed_histograms_and_categories() {
let stats = DhtRuntimeStats::default();
stats.udp_received();
stats.udp_received_bytes(128);
stats.udp_sent(64);
stats.inbound_query("ping");
stats.inbound_query("unknown");
stats.node_admitted();
stats.metadata_bytes_downloaded(1_024);
stats.metadata_failure_parse();
stats.observe_metadata_queue_wait(75);
stats.observe_metadata_queue_wait(600);
stats.observe_metadata_fetch_duration(1_500);
stats.observe_metadata_size(70_000);
let snapshot = stats.observability_snapshot();
assert_eq!(snapshot.udp_rx_packets, 1);
assert_eq!(snapshot.udp_rx_bytes, 128);
assert_eq!(snapshot.udp_tx_packets, 1);
assert_eq!(snapshot.udp_tx_bytes, 64);
assert_eq!(snapshot.inbound_ping, 1);
assert_eq!(snapshot.inbound_other, 1);
assert_eq!(snapshot.node_admitted, 1);
assert_eq!(snapshot.metadata_bytes_downloaded, 1_024);
assert_eq!(snapshot.metadata_failure_parse, 1);
assert_eq!(snapshot.queue_wait_ms.percentile(0.50), Some(100));
assert_eq!(snapshot.queue_wait_ms.percentile(0.95), Some(1_000));
assert_eq!(snapshot.fetch_duration_ms.percentile(0.95), Some(2_000));
assert_eq!(snapshot.metadata_size_bytes.percentile(0.50), Some(131_072));
}
}
mod tests;
+102
View File
@@ -0,0 +1,102 @@
// 负责固定分桶直方图的记录快照和近似分位数计算
use std::sync::atomic::{AtomicU64, Ordering};
pub const METADATA_QUEUE_WAIT_BUCKETS_MS: [u64; 8] = [10, 50, 100, 250, 500, 1_000, 2_000, 5_000];
pub const METADATA_FETCH_BUCKETS_MS: [u64; 7] = [250, 500, 1_000, 2_000, 4_000, 6_000, 10_000];
pub const METADATA_SIZE_BUCKETS_BYTES: [u64; 8] = [
16 * 1024,
32 * 1024,
64 * 1024,
128 * 1024,
256 * 1024,
512 * 1024,
1024 * 1024,
10 * 1024 * 1024,
];
#[derive(Debug, Clone, PartialEq, Eq)]
/// Snapshot of a non-cumulative fixed-bucket histogram.
pub struct FixedHistogramSnapshot {
/// Inclusive upper bound for each bucket.
pub bounds: Vec<u64>,
/// Per-bucket counts; values are not cumulative.
pub counts: Vec<u64>,
/// Values larger than the final bound.
pub overflow: u64,
/// Total observations including overflow.
pub count: u64,
/// Wrapping sum of all observed values.
pub sum: u64,
}
impl FixedHistogramSnapshot {
/// Returns the inclusive bucket bound containing the requested approximate percentile.
///
/// `percentile` is clamped to `0.0..=1.0`. Overflow observations return the final finite
/// bound because the histogram intentionally stores no dynamic maximum.
pub fn percentile(&self, percentile: f64) -> Option<u64> {
if self.count == 0 {
return None;
}
let rank = ((self.count as f64 * percentile.clamp(0.0, 1.0)).ceil() as u64).max(1);
let mut cumulative = 0u64;
for (bound, count) in self.bounds.iter().zip(&self.counts) {
cumulative = cumulative.saturating_add(*count);
if cumulative >= rank {
return Some(*bound);
}
}
self.bounds.last().copied()
}
}
pub(super) struct AtomicFixedHistogram<const N: usize> {
bounds: [u64; N],
counts: [AtomicU64; N],
overflow: AtomicU64,
count: AtomicU64,
sum: AtomicU64,
}
impl<const N: usize> AtomicFixedHistogram<N> {
pub(super) fn new(bounds: [u64; N]) -> Self {
Self {
bounds,
counts: std::array::from_fn(|_| AtomicU64::new(0)),
overflow: AtomicU64::new(0),
count: AtomicU64::new(0),
sum: AtomicU64::new(0),
}
}
pub(super) fn record(&self, value: u64) {
if let Some(index) = self.bounds.iter().position(|bound| value <= *bound) {
self.counts[index].fetch_add(1, Ordering::Relaxed);
} else {
self.overflow.fetch_add(1, Ordering::Relaxed);
}
self.count.fetch_add(1, Ordering::Relaxed);
self.sum.fetch_add(value, Ordering::Relaxed);
}
pub(super) fn snapshot(&self) -> FixedHistogramSnapshot {
FixedHistogramSnapshot {
bounds: self.bounds.to_vec(),
counts: self
.counts
.iter()
.map(|count| count.load(Ordering::Relaxed))
.collect(),
overflow: self.overflow.load(Ordering::Relaxed),
count: self.count.load(Ordering::Relaxed),
sum: self.sum.load(Ordering::Relaxed),
}
}
}
impl<const N: usize> Default for AtomicFixedHistogram<N> {
fn default() -> Self {
Self::new([0; N])
}
}
+549
View File
@@ -0,0 +1,549 @@
// 负责写入 DHT 运行计数器状态和观测直方图
use super::{DhtRuntimeStats, Ordering};
use crate::types::DiscoverySource;
impl DhtRuntimeStats {
pub(crate) fn hash_received(&self) {
self.inner.hashes_received.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn hash_ingress_dropped(&self) {
self.inner
.hash_ingress_dropped
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn set_hash_ingress_queue_depth(&self, depth: usize) {
self.inner
.hash_ingress_queue_depth
.store(depth, Ordering::Relaxed);
}
pub(crate) fn set_crawl_priority_queue_depth(&self, depth: usize) {
self.inner
.crawl_priority_queue_depth
.store(depth, Ordering::Relaxed);
}
pub(crate) fn set_crawl_discovery_queue_depth(&self, depth: usize) {
self.inner
.crawl_discovery_queue_depth
.store(depth, Ordering::Relaxed);
}
pub(crate) fn set_metadata_queue(&self, depth: usize, in_flight: usize) {
self.inner
.metadata_queue_depth
.store(depth, Ordering::Relaxed);
self.inner
.metadata_in_flight
.store(in_flight, Ordering::Relaxed);
}
pub(crate) fn metadata_queue_deduplicated(&self) {
self.inner
.metadata_queue_deduplicated
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_queue_inserted(&self) {
self.inner
.metadata_queue_inserted
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_queue_evicted(&self) {
self.inner
.metadata_queue_evicted
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_queue_stale(&self, count: usize) {
self.inner
.metadata_queue_stale
.fetch_add(count as u64, Ordering::Relaxed);
}
pub(crate) fn metadata_race_started(&self) {
self.inner
.metadata_races_started
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_peer_candidate(&self) {
self.inner
.metadata_peer_candidates
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_live_peer_join(&self) {
self.inner
.metadata_live_peer_joins
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_peer_canceled(&self, count: usize) {
self.inner
.metadata_peer_canceled
.fetch_add(count as u64, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_requested(&self) {
self.inner
.peer_lookup_requested
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_started(&self) {
self.inner
.peer_lookup_started
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_rate_limited(&self) {
self.inner
.peer_lookup_rate_limited
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_empty(&self) {
self.inner.peer_lookup_empty.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_query(&self) {
self.inner
.peer_lookup_queries
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_response(&self) {
self.inner
.peer_lookup_responses
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_timeout(&self) {
self.inner
.peer_lookup_timeouts
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_send_failed(&self) {
self.inner
.peer_lookup_send_failures
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_peer_found(&self) {
self.inner
.peer_lookup_peers_found
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_response_dropped(&self) {
self.inner
.peer_lookup_response_dropped
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_output_dropped(&self) {
self.inner
.peer_lookup_output_dropped
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_preferred_succeeded(&self) {
self.inner
.peer_lookup_preferred_succeeded
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_lookup_fallback(&self) {
self.inner
.peer_lookup_fallbacks
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn configure_sample_candidate_queue(&self, capacity: usize) {
self.inner
.sample_candidate_queue_capacity
.store(capacity, Ordering::Relaxed);
}
pub(crate) fn set_sample_candidate_queue_depth(&self, depth: usize) {
self.inner
.sample_candidate_queue_depth
.store(depth, Ordering::Relaxed);
}
pub(crate) fn sample_candidate_routed(&self) {
self.inner
.sample_candidates_routed
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn sample_candidate_fallback(&self) {
self.inner
.sample_candidates_fallback
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn sample_query(&self, direct: bool) {
self.inner
.sample_infohashes_queries
.fetch_add(1, Ordering::Relaxed);
let source = if direct {
&self.inner.sample_infohashes_direct_queries
} else {
&self.inner.sample_infohashes_snapshot_queries
};
source.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn sample_response(&self, direct: bool) {
self.inner
.sample_infohashes_responses
.fetch_add(1, Ordering::Relaxed);
let source = if direct {
&self.inner.sample_infohashes_direct_responses
} else {
&self.inner.sample_infohashes_snapshot_responses
};
source.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn sample_timeout(&self) {
self.inner
.sample_infohashes_timeouts
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn sample_send_failed(&self) {
self.inner
.sample_infohashes_send_failures
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn sample_response_dropped(&self) {
self.inner
.sample_infohashes_response_dropped
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn sample_hash_discovered(&self) {
self.inner
.sample_infohashes_hashes_discovered
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn sample_hash_filtered(&self, count: usize) {
self.inner
.sample_infohashes_hashes_filtered
.fetch_add(count as u64, Ordering::Relaxed);
}
pub(crate) fn sample_hash_duplicate(&self) {
self.inner
.sample_infohashes_hashes_duplicate
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn sample_hash_dropped(&self) {
self.inner
.sample_infohashes_hashes_dropped
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_peer_attempt(&self) {
self.inner
.metadata_peer_attempts
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_peer_succeeded(&self) {
self.inner
.metadata_peer_succeeded
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_success_from(&self, source: DiscoverySource) {
let counter = match source {
DiscoverySource::AnnouncePeer => &self.inner.metadata_success_from_announce,
DiscoverySource::SampleDirect => &self.inner.metadata_success_from_sample_direct,
DiscoverySource::SampleSnapshot => &self.inner.metadata_success_from_sample_snapshot,
DiscoverySource::ActiveLookup => &self.inner.metadata_success_from_active_lookup,
};
counter.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_peer_failed(&self) {
self.inner
.metadata_peer_failed
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_peer_timeout(&self) {
self.inner
.metadata_peer_timeouts
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_connect_failed(&self) {
self.inner
.metadata_connect_failed
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_no_extension(&self) {
self.inner
.metadata_no_extension
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_peer_failure_cache_hit(&self) {
self.inner
.metadata_peer_failure_cache_hits
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn set_metadata_peer_failure_cache_entries(&self, count: usize) {
self.inner
.metadata_peer_failure_cache_entries
.store(count, Ordering::Relaxed);
}
pub(crate) fn set_node_pool_size(&self, size: usize) {
self.inner.node_pool_size.store(size, Ordering::Relaxed);
}
pub(crate) fn set_find_node_in_flight(&self, count: usize) {
self.inner
.find_node_in_flight
.store(count, Ordering::Relaxed);
}
pub(crate) fn set_find_node_effective_rate(&self, rate: u32) {
self.inner
.find_node_effective_rate_per_sec
.store(rate, Ordering::Relaxed);
}
pub(crate) fn query_new(&self) {
self.inner.queries_new.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn query_revisit(&self) {
self.inner.queries_revisit.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn query_bootstrap(&self) {
self.inner.queries_bootstrap.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn response(&self) {
self.inner.responses.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn unmatched_response(&self) {
self.inner
.unmatched_responses
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn timeout(&self) {
self.inner.timeouts.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn send_failure(&self) {
self.inner.send_failures.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn crawl_event_dropped_response(&self) {
self.inner
.crawl_events_dropped_response
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn crawl_event_dropped_discovered(&self) {
self.inner
.crawl_events_dropped_discovered
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn udp_received(&self) {
self.inner.udp_received.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn udp_received_bytes(&self, bytes: usize) {
self.inner
.udp_rx_bytes
.fetch_add(bytes as u64, Ordering::Relaxed);
}
pub(crate) fn udp_sent(&self, bytes: usize) {
self.inner.udp_tx_packets.fetch_add(1, Ordering::Relaxed);
self.inner
.udp_tx_bytes
.fetch_add(bytes as u64, Ordering::Relaxed);
}
pub(crate) fn inbound_query(&self, query: &str) {
let counter = match query {
"ping" => &self.inner.inbound_ping,
"find_node" => &self.inner.inbound_find_node,
"get_peers" => &self.inner.inbound_get_peers,
"announce_peer" => &self.inner.inbound_announce_peer,
_ => &self.inner.inbound_other,
};
counter.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn response_normal(&self) {
self.inner.response_normal.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn response_send_failed(&self) {
self.inner
.response_send_failed
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn announce_accepted(&self) {
self.inner.announce_accepted.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn announce_invalid_token(&self) {
self.inner
.announce_invalid_token
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn announce_filtered(&self) {
self.inner.announce_filtered.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn node_admitted(&self) {
self.inner.node_admitted.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn node_replaced(&self) {
self.inner.node_replaced.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn node_dropped_duplicate(&self) {
self.inner
.node_dropped_duplicate
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn node_dropped_rate_limited(&self) {
self.inner
.node_dropped_rate_limited
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn node_dropped_invalid(&self) {
self.inner
.node_dropped_invalid
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_bytes_downloaded(&self, bytes: usize) {
self.inner
.metadata_bytes_downloaded
.fetch_add(bytes as u64, Ordering::Relaxed);
}
pub(crate) fn metadata_failure_timeout(&self) {
self.inner
.metadata_failure_timeout
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_failure_connect(&self) {
self.inner
.metadata_failure_connect
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_failure_no_extension(&self) {
self.inner
.metadata_failure_no_extension
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_failure_send(&self) {
self.inner
.metadata_failure_send
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_failure_size_limit(&self) {
self.inner
.metadata_failure_size_limit
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_failure_sha1(&self) {
self.inner
.metadata_failure_sha1
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_failure_parse(&self) {
self.inner
.metadata_failure_parse
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn metadata_failure_other(&self) {
self.inner
.metadata_failure_other
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_cache_hit_timeout(&self) {
self.inner
.peer_cache_timeout_hits
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn peer_cache_hit_connect(&self) {
self.inner
.peer_cache_connect_hits
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn observe_metadata_queue_wait(&self, millis: u64) {
self.inner.queue_wait_ms.record(millis);
}
pub(crate) fn observe_metadata_fetch_duration(&self, millis: u64) {
self.inner.fetch_duration_ms.record(millis);
}
pub(crate) fn observe_metadata_size(&self, bytes: usize) {
self.inner.metadata_size_bytes.record(bytes as u64);
}
pub(crate) fn udp_queue_full(&self) {
self.inner.udp_queue_full.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn udp_invalid(&self) {
self.inner.udp_invalid.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn udp_response_rate_limited(&self) {
self.inner
.udp_responses_rate_limited
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn udp_response_priority_reserved(&self) {
self.inner
.udp_responses_priority_reserved
.fetch_add(1, Ordering::Relaxed);
}
}
+208
View File
@@ -0,0 +1,208 @@
// 负责验证 DHT 运行指标快照和直方图统计行为
use super::*;
use crate::types::DiscoverySource;
#[test]
fn cloned_handle_updates_one_snapshot() {
let stats = DhtRuntimeStats::with_limits(DhtRuntimeLimits {
metadata_queue: 100,
node_pool: 1_000,
node_pool_low_watermark: 10,
find_node_in_flight: 512,
initial_find_node_rate: 200,
hash_ingress_queue: 20,
crawl_priority_queue: 30,
crawl_discovery_queue: 40,
});
let writer = stats.clone();
writer.hash_received();
writer.hash_ingress_dropped();
writer.set_hash_ingress_queue_depth(11);
writer.set_crawl_priority_queue_depth(12);
writer.set_crawl_discovery_queue_depth(13);
writer.set_metadata_queue(42, 7);
writer.metadata_queue_inserted();
writer.metadata_queue_deduplicated();
writer.metadata_queue_evicted();
writer.metadata_queue_stale(3);
writer.metadata_race_started();
writer.metadata_peer_candidate();
writer.metadata_live_peer_join();
writer.metadata_peer_canceled(2);
writer.peer_lookup_requested();
writer.peer_lookup_started();
writer.peer_lookup_rate_limited();
writer.peer_lookup_empty();
writer.peer_lookup_query();
writer.peer_lookup_response();
writer.peer_lookup_timeout();
writer.peer_lookup_send_failed();
writer.peer_lookup_peer_found();
writer.peer_lookup_response_dropped();
writer.peer_lookup_output_dropped();
writer.peer_lookup_preferred_succeeded();
writer.peer_lookup_fallback();
writer.configure_sample_candidate_queue(20);
writer.set_sample_candidate_queue_depth(14);
writer.sample_candidate_routed();
writer.sample_candidate_fallback();
writer.sample_query(true);
writer.sample_response(true);
writer.sample_timeout();
writer.sample_send_failed();
writer.sample_response_dropped();
writer.sample_hash_discovered();
writer.sample_hash_filtered(1);
writer.sample_hash_duplicate();
writer.sample_hash_dropped();
writer.metadata_peer_attempt();
writer.metadata_peer_succeeded();
writer.metadata_success_from(DiscoverySource::AnnouncePeer);
writer.metadata_success_from(DiscoverySource::SampleDirect);
writer.metadata_success_from(DiscoverySource::SampleSnapshot);
writer.metadata_success_from(DiscoverySource::ActiveLookup);
writer.metadata_peer_failed();
writer.metadata_peer_timeout();
writer.metadata_connect_failed();
writer.metadata_no_extension();
writer.metadata_peer_failure_cache_hit();
writer.set_metadata_peer_failure_cache_entries(9);
writer.set_node_pool_size(321);
writer.set_find_node_in_flight(12);
writer.set_find_node_effective_rate(150);
writer.query_new();
writer.query_revisit();
writer.query_bootstrap();
writer.response();
writer.unmatched_response();
writer.timeout();
writer.send_failure();
writer.crawl_event_dropped_response();
writer.crawl_event_dropped_discovered();
writer.udp_received();
writer.udp_queue_full();
writer.udp_invalid();
writer.udp_response_rate_limited();
writer.udp_response_priority_reserved();
assert_eq!(
stats.snapshot(),
DhtRuntimeSnapshot {
hashes_received: 1,
hash_ingress_dropped: 1,
hash_ingress_queue_depth: 11,
hash_ingress_queue_capacity: 20,
crawl_priority_queue_depth: 12,
crawl_priority_queue_capacity: 30,
crawl_discovery_queue_depth: 13,
crawl_discovery_queue_capacity: 40,
metadata_queue_depth: 42,
metadata_queue_max: 100,
metadata_in_flight: 7,
metadata_queue_inserted: 1,
metadata_queue_deduplicated: 1,
metadata_queue_evicted: 1,
metadata_queue_stale: 3,
metadata_races_started: 1,
metadata_peer_candidates: 1,
metadata_live_peer_joins: 1,
metadata_peer_canceled: 2,
peer_lookup_requested: 1,
peer_lookup_started: 1,
peer_lookup_rate_limited: 1,
peer_lookup_empty: 1,
peer_lookup_queries: 1,
peer_lookup_responses: 1,
peer_lookup_timeouts: 1,
peer_lookup_send_failures: 1,
peer_lookup_peers_found: 1,
peer_lookup_response_dropped: 1,
peer_lookup_output_dropped: 1,
peer_lookup_preferred_succeeded: 1,
peer_lookup_fallbacks: 1,
sample_candidate_queue_depth: 14,
sample_candidate_queue_capacity: 20,
sample_candidates_routed: 1,
sample_candidates_fallback: 1,
sample_infohashes_queries: 1,
sample_infohashes_direct_queries: 1,
sample_infohashes_snapshot_queries: 0,
sample_infohashes_responses: 1,
sample_infohashes_direct_responses: 1,
sample_infohashes_snapshot_responses: 0,
sample_infohashes_timeouts: 1,
sample_infohashes_send_failures: 1,
sample_infohashes_response_dropped: 1,
sample_infohashes_hashes_discovered: 1,
sample_infohashes_hashes_filtered: 1,
sample_infohashes_hashes_duplicate: 1,
sample_infohashes_hashes_dropped: 1,
metadata_peer_attempts: 1,
metadata_peer_succeeded: 1,
metadata_success_from_announce: 1,
metadata_success_from_sample_direct: 1,
metadata_success_from_sample_snapshot: 1,
metadata_success_from_active_lookup: 1,
metadata_peer_failed: 1,
metadata_peer_timeouts: 1,
metadata_connect_failed: 1,
metadata_no_extension: 1,
metadata_peer_failure_cache_hits: 1,
metadata_peer_failure_cache_entries: 9,
node_pool_size: 321,
node_pool_capacity: 1_000,
node_pool_low_watermark: 10,
find_node_in_flight: 12,
find_node_in_flight_max: 512,
find_node_effective_rate_per_sec: 150,
queries_new: 1,
queries_revisit: 1,
queries_bootstrap: 1,
responses: 1,
unmatched_responses: 1,
timeouts: 1,
send_failures: 1,
crawl_events_dropped_response: 1,
crawl_events_dropped_discovered: 1,
udp_received: 1,
udp_queue_full: 1,
udp_invalid: 1,
udp_responses_rate_limited: 1,
udp_responses_priority_reserved: 1,
}
);
}
#[test]
fn observability_snapshot_tracks_fixed_histograms_and_categories() {
let stats = DhtRuntimeStats::default();
stats.udp_received();
stats.udp_received_bytes(128);
stats.udp_sent(64);
stats.inbound_query("ping");
stats.inbound_query("unknown");
stats.node_admitted();
stats.metadata_bytes_downloaded(1_024);
stats.metadata_failure_parse();
stats.observe_metadata_queue_wait(75);
stats.observe_metadata_queue_wait(600);
stats.observe_metadata_fetch_duration(1_500);
stats.observe_metadata_size(70_000);
let snapshot = stats.observability_snapshot();
assert_eq!(snapshot.udp_rx_packets, 1);
assert_eq!(snapshot.udp_rx_bytes, 128);
assert_eq!(snapshot.udp_tx_packets, 1);
assert_eq!(snapshot.udp_tx_bytes, 64);
assert_eq!(snapshot.inbound_ping, 1);
assert_eq!(snapshot.inbound_other, 1);
assert_eq!(snapshot.node_admitted, 1);
assert_eq!(snapshot.metadata_bytes_downloaded, 1_024);
assert_eq!(snapshot.metadata_failure_parse, 1);
assert_eq!(snapshot.queue_wait_ms.percentile(0.50), Some(100));
assert_eq!(snapshot.queue_wait_ms.percentile(0.95), Some(1_000));
assert_eq!(snapshot.fetch_duration_ms.percentile(0.95), Some(2_000));
assert_eq!(snapshot.metadata_size_bytes.percentile(0.50), Some(131_072));
}
+10 -764
View File
@@ -1,19 +1,17 @@
// 负责管理有界 Metadata 任务队列 Peer 竞速和完成回调
use crate::metadata::{FetchedMetadata, MetadataFetchOutcome, RbitFetcher};
use crate::metadata::RbitFetcher;
use crate::peer_lookup::PeerLookupRequest;
#[cfg(test)]
use crate::runtime_stats::DhtRuntimeLimits;
use crate::runtime_stats::DhtRuntimeStats;
use crate::types::{
DiscoverySource, HashDiscovered, MetadataFetchCompletion, MetadataFetchCompletionStatus,
TorrentInfo,
HashDiscovered, MetadataFetchCompletion, MetadataFetchCompletionStatus, TorrentInfo,
};
use arc_swap::ArcSwapOption;
#[cfg(feature = "metrics")]
use metrics::{counter, gauge, histogram};
use std::collections::{BTreeMap, HashMap, HashSet};
use std::future::Future;
use std::collections::{HashMap, HashSet};
use std::net::SocketAddr;
use std::panic::{AssertUnwindSafe, catch_unwind};
use std::sync::Arc;
@@ -31,6 +29,12 @@ const LIVE_PEER_IDLE_GRACE: Duration = Duration::from_millis(300);
const ACTIVE_PEER_LOOKUP_DELAY: Duration = Duration::from_millis(100);
const DISPATCH_TICK: Duration = Duration::from_millis(25);
mod peer_race;
mod pending_queue;
use peer_race::race_peer_fetches;
use pending_queue::{PeerCandidate, PendingHashQueue, QueuePushKind, QueuedHash};
/// Callback returning whether the application accepted a downloaded torrent.
pub type TorrentAckCallback = Box<dyn Fn(TorrentInfo) -> bool + Send + Sync + 'static>;
/// Callback invoked once when an admitted Metadata job reaches a terminal state.
@@ -68,171 +72,6 @@ pub(crate) struct MetadataSchedulerRuntime {
pub(crate) peer_lookup_tx: Option<mpsc::Sender<PeerLookupRequest>>,
}
#[derive(Debug, Clone, Copy)]
struct PeerCandidate {
addr: SocketAddr,
source: DiscoverySource,
discovered_at: Instant,
}
#[derive(Debug)]
struct QueuedHash {
info_hash: String,
peers: Vec<PeerCandidate>,
queued_at: Instant,
ready_at: Instant,
latest_at: Instant,
order_key: (Instant, u64),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum QueuePushKind {
Inserted,
Updated,
EvictedOldest,
Stale,
}
#[derive(Debug)]
struct PendingHashQueue {
capacity: usize,
ttl: Duration,
entries: HashMap<String, QueuedHash>,
order: BTreeMap<(Instant, u64), String>,
sequence: u64,
}
impl PendingHashQueue {
fn new(capacity: usize, ttl: Duration) -> Self {
Self {
capacity: capacity.max(1),
ttl,
entries: HashMap::new(),
order: BTreeMap::new(),
sequence: 0,
}
}
fn len(&self) -> usize {
self.entries.len()
}
fn is_empty(&self) -> bool {
self.entries.is_empty()
}
#[cfg(test)]
fn contains(&self, info_hash: &str) -> bool {
self.entries.contains_key(info_hash)
}
fn next_order_key(&mut self, at: Instant) -> (Instant, u64) {
self.sequence = self.sequence.wrapping_add(1);
(at, self.sequence)
}
fn push(&mut self, event: HashDiscovered, now: Instant) -> QueuePushKind {
if now
.checked_duration_since(event.discovered_at)
.unwrap_or_default()
> self.ttl
{
return QueuePushKind::Stale;
}
if let Some(mut entry) = self.entries.remove(&event.info_hash) {
self.order.remove(&entry.order_key);
let discovered_at = entry
.peers
.iter()
.find(|peer| peer.addr == event.peer_addr)
.map(|peer| peer.discovered_at.max(event.discovered_at))
.unwrap_or(event.discovered_at);
entry.peers.retain(|peer| peer.addr != event.peer_addr);
entry.peers.push(PeerCandidate {
addr: event.peer_addr,
source: event.source,
discovered_at,
});
entry
.peers
.sort_unstable_by_key(|peer| std::cmp::Reverse(peer.discovered_at));
entry.peers.truncate(MAX_METADATA_PEERS_PER_HASH);
entry.latest_at = entry.latest_at.max(event.discovered_at);
entry.order_key = self.next_order_key(entry.latest_at);
self.order.insert(entry.order_key, entry.info_hash.clone());
self.entries.insert(entry.info_hash.clone(), entry);
return QueuePushKind::Updated;
}
let mut result = QueuePushKind::Inserted;
if self.entries.len() >= self.capacity {
let Some((&oldest_key, oldest_hash)) = self.order.first_key_value() else {
return QueuePushKind::Stale;
};
if event.discovered_at <= oldest_key.0 {
return QueuePushKind::Stale;
}
let oldest_hash = oldest_hash.clone();
self.order.remove(&oldest_key);
self.entries.remove(&oldest_hash);
result = QueuePushKind::EvictedOldest;
}
let order_key = self.next_order_key(event.discovered_at);
let info_hash = event.info_hash;
self.order.insert(order_key, info_hash.clone());
self.entries.insert(
info_hash.clone(),
QueuedHash {
info_hash,
peers: vec![PeerCandidate {
addr: event.peer_addr,
source: event.source,
discovered_at: event.discovered_at,
}],
queued_at: now,
ready_at: now + PEER_COALESCE_WINDOW,
latest_at: event.discovered_at,
order_key,
},
);
result
}
fn remove(&mut self, info_hash: &str) -> Option<QueuedHash> {
let entry = self.entries.remove(info_hash)?;
self.order.remove(&entry.order_key);
Some(entry)
}
fn pop_newest_ready(
&mut self,
now: Instant,
in_flight: &HashMap<String, InFlightState>,
) -> Option<QueuedHash> {
let info_hash = self.order.iter().rev().find_map(|(_, hash)| {
let entry = self.entries.get(hash)?;
(!in_flight.contains_key(hash) && entry.ready_at <= now).then(|| hash.clone())
})?;
self.remove(&info_hash)
}
fn expire(&mut self, now: Instant) -> usize {
let mut expired = 0;
while let Some((&oldest_key, oldest_hash)) = self.order.first_key_value() {
if now.checked_duration_since(oldest_key.0).unwrap_or_default() <= self.ttl {
break;
}
let oldest_hash = oldest_hash.clone();
self.order.remove(&oldest_key);
self.entries.remove(&oldest_hash);
expired += 1;
}
expired
}
}
#[derive(Debug)]
struct MetadataJob {
info_hash: String,
@@ -767,598 +606,5 @@ impl MetadataScheduler {
}
}
async fn race_peer_fetches<F, Fut>(
peers: Vec<PeerCandidate>,
peer_rx: &mut mpsc::Receiver<PeerCandidate>,
fetch: F,
runtime_stats: &DhtRuntimeStats,
) -> Option<(PeerCandidate, FetchedMetadata)>
where
F: Fn(PeerCandidate) -> Fut + Clone + Send + Sync + 'static,
Fut: Future<Output = (PeerCandidate, MetadataFetchOutcome)> + Send + 'static,
{
runtime_stats.metadata_race_started();
#[cfg(feature = "metrics")]
counter!("dht_metadata_peer_races_total").increment(1);
let mut tasks = JoinSet::new();
let mut candidates = 0usize;
for peer in peers {
let fetch = fetch.clone();
tasks.spawn(fetch(peer));
candidates += 1;
runtime_stats.metadata_peer_candidate();
#[cfg(feature = "metrics")]
counter!("dht_metadata_peer_candidates_total", "source" => "initial").increment(1);
}
let mut accepting_live_peers = candidates < MAX_METADATA_PEERS_PER_HASH;
loop {
if tasks.is_empty() {
if !accepting_live_peers {
return None;
}
match tokio::time::timeout(LIVE_PEER_IDLE_GRACE, peer_rx.recv()).await {
Ok(Some(peer)) => {
let fetch = fetch.clone();
tasks.spawn(fetch(peer));
candidates += 1;
runtime_stats.metadata_peer_candidate();
#[cfg(feature = "metrics")]
counter!("dht_metadata_peer_candidates_total", "source" => "live").increment(1);
accepting_live_peers = candidates < MAX_METADATA_PEERS_PER_HASH;
}
Ok(None) | Err(_) => return None,
}
continue;
}
tokio::select! {
result = tasks.join_next() => {
let Some(Ok((peer, outcome))) = result else {
continue;
};
if let MetadataFetchOutcome::Fetched(metadata) = outcome {
let canceled = tasks.len();
runtime_stats.metadata_peer_canceled(canceled);
#[cfg(feature = "metrics")]
counter!("dht_metadata_peer_canceled_total").increment(canceled as u64);
peer_rx.close();
tasks.abort_all();
while tasks.join_next().await.is_some() {}
return Some((peer, metadata));
}
}
peer = peer_rx.recv(), if accepting_live_peers => {
match peer {
Some(peer) => {
let fetch = fetch.clone();
tasks.spawn(fetch(peer));
candidates += 1;
runtime_stats.metadata_peer_candidate();
#[cfg(feature = "metrics")]
counter!("dht_metadata_peer_candidates_total", "source" => "live")
.increment(1);
accepting_live_peers = candidates < MAX_METADATA_PEERS_PER_HASH;
}
None => accepting_live_peers = false,
}
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::AtomicBool;
struct DropSignal(Arc<AtomicBool>);
impl Drop for DropSignal {
fn drop(&mut self) {
self.0.store(true, Ordering::Relaxed);
}
}
fn empty_torrent_callback() -> Arc<ArcSwapOption<TorrentAckCallback>> {
Arc::new(ArcSwapOption::empty())
}
fn fetch_gate<F, Fut>(callback: F) -> Arc<ArcSwapOption<MetadataFetchCallback>>
where
F: Fn(String) -> Fut + Send + Sync + 'static,
Fut: std::future::Future<Output = bool> + Send + 'static,
{
let holder = Arc::new(ArcSwapOption::empty());
let callback: Arc<MetadataFetchCallback> =
Arc::new(Box::new(move |hash| Box::pin(callback(hash))));
holder.store(Some(callback));
holder
}
fn completion_callback<F>(callback: F) -> Arc<ArcSwapOption<MetadataCompletionCallback>>
where
F: Fn(MetadataFetchCompletion) + Send + Sync + 'static,
{
let holder = Arc::new(ArcSwapOption::empty());
let callback: Arc<MetadataCompletionCallback> = Arc::new(Box::new(callback));
holder.store(Some(callback));
holder
}
fn event(hash: &str, port: u16, discovered_at: Instant) -> HashDiscovered {
HashDiscovered {
info_hash: hash.to_string(),
peer_addr: SocketAddr::from(([127, 0, 0, 1], port)),
source: DiscoverySource::AnnouncePeer,
discovered_at,
}
}
#[test]
fn runtime_stats_track_queue_outcomes_and_depth() {
let (_hash_tx, hash_rx) = mpsc::channel(4);
let stats = DhtRuntimeStats::with_limits(DhtRuntimeLimits {
metadata_queue: 2,
..DhtRuntimeLimits::default()
});
let scheduler = MetadataScheduler::new_with_runtime_stats(
hash_rx,
Arc::new(RbitFetcher::new(1)),
MetadataSchedulerLimits {
queue_size: 2,
concurrency: 1,
},
MetadataSchedulerCallbacks {
torrent: empty_torrent_callback(),
fetch_gate: Arc::new(ArcSwapOption::empty()),
completion: Arc::new(ArcSwapOption::empty()),
},
Arc::new(AtomicUsize::new(0)),
CancellationToken::new(),
stats.clone(),
);
let start = Instant::now();
let mut queue = PendingHashQueue::new(2, HASH_QUEUE_TTL);
let mut in_flight = HashMap::new();
scheduler.enqueue(&mut queue, &mut in_flight, event("old", 1000, start));
scheduler.enqueue(
&mut queue,
&mut in_flight,
event("old", 1001, start + Duration::from_secs(1)),
);
scheduler.enqueue(
&mut queue,
&mut in_flight,
event("middle", 1002, start + Duration::from_secs(2)),
);
scheduler.enqueue(
&mut queue,
&mut in_flight,
event("new", 1003, start + Duration::from_secs(3)),
);
scheduler.enqueue(
&mut queue,
&mut in_flight,
event(
"stale",
1004,
start
.checked_sub(HASH_QUEUE_TTL + Duration::from_secs(1))
.unwrap(),
),
);
scheduler.sync_queue_len(queue.len(), 1);
let snapshot = stats.snapshot();
assert_eq!(snapshot.metadata_queue_inserted, 3);
assert_eq!(snapshot.metadata_queue_deduplicated, 1);
assert_eq!(snapshot.metadata_queue_evicted, 1);
assert_eq!(snapshot.metadata_queue_stale, 1);
assert_eq!(snapshot.metadata_queue_depth, 3);
assert_eq!(snapshot.metadata_queue_max, 2);
assert_eq!(snapshot.metadata_in_flight, 1);
}
#[test]
fn queue_deduplicates_hash_and_keeps_twelve_newest_unique_peers() {
let start = Instant::now();
let mut queue = PendingHashQueue::new(10, Duration::from_secs(60));
for offset in 0..13 {
assert_ne!(
queue.push(
event(
"hash",
1000 + offset,
start + Duration::from_secs(offset as u64)
),
start + Duration::from_secs(offset as u64),
),
QueuePushKind::Stale
);
}
queue.push(
event("hash", 1012, start + Duration::from_secs(20)),
start + Duration::from_secs(20),
);
assert_eq!(queue.len(), 1);
let entry = queue.remove("hash").unwrap();
let ports: Vec<_> = entry.peers.iter().map(|peer| peer.addr.port()).collect();
assert_eq!(ports, (1001..=1012).rev().collect::<Vec<_>>());
}
#[test]
fn coalescing_window_is_fixed_from_first_enqueue() {
let start = Instant::now();
let mut queue = PendingHashQueue::new(10, Duration::from_secs(60));
let in_flight = HashMap::new();
queue.push(event("hash", 1000, start), start);
queue.push(
event("hash", 1001, start + Duration::from_millis(15)),
start + Duration::from_millis(15),
);
assert!(
queue
.pop_newest_ready(start + Duration::from_millis(24), &in_flight)
.is_none()
);
let entry = queue
.pop_newest_ready(start + Duration::from_millis(25), &in_flight)
.expect("duplicate arrivals must not extend the coalescing deadline");
assert_eq!(entry.peers.len(), 2);
}
#[test]
fn in_flight_hash_routes_at_most_twelve_unique_peers_to_live_race() {
let (_hash_tx, hash_rx) = mpsc::channel(4);
let stats = DhtRuntimeStats::default();
let scheduler = MetadataScheduler::new_with_runtime_stats(
hash_rx,
Arc::new(RbitFetcher::new(1)),
MetadataSchedulerLimits {
queue_size: 4,
concurrency: 1,
},
MetadataSchedulerCallbacks {
torrent: empty_torrent_callback(),
fetch_gate: Arc::new(ArcSwapOption::empty()),
completion: Arc::new(ArcSwapOption::empty()),
},
Arc::new(AtomicUsize::new(0)),
CancellationToken::new(),
stats.clone(),
);
let (peer_tx, mut peer_rx) = mpsc::channel(MAX_METADATA_PEERS_PER_HASH);
let mut in_flight = HashMap::from([(
"hash".to_string(),
InFlightState {
peer_tx,
scheduled_peers: HashSet::from([SocketAddr::from(([127, 0, 0, 1], 1000))]),
},
)]);
let mut queue = PendingHashQueue::new(4, HASH_QUEUE_TTL);
let now = Instant::now();
for port in 1001..=1013 {
scheduler.enqueue(&mut queue, &mut in_flight, event("hash", port, now));
}
scheduler.enqueue(&mut queue, &mut in_flight, event("hash", 1001, now));
let mut received = Vec::new();
while let Ok(peer) = peer_rx.try_recv() {
received.push(peer.addr.port());
}
assert_eq!(received, (1001..=1011).collect::<Vec<_>>());
assert_eq!(
in_flight["hash"].scheduled_peers.len(),
MAX_METADATA_PEERS_PER_HASH
);
assert!(queue.is_empty());
let snapshot = stats.snapshot();
assert_eq!(snapshot.metadata_live_peer_joins, 11);
assert_eq!(snapshot.metadata_queue_deduplicated, 14);
}
#[tokio::test]
async fn peer_race_returns_first_success_and_cancels_remaining_fetches() {
let now = Instant::now();
let peers = vec![
PeerCandidate {
addr: SocketAddr::from(([127, 0, 0, 1], 1000)),
source: DiscoverySource::AnnouncePeer,
discovered_at: now,
},
PeerCandidate {
addr: SocketAddr::from(([127, 0, 0, 1], 1001)),
source: DiscoverySource::SampleDirect,
discovered_at: now,
},
PeerCandidate {
addr: SocketAddr::from(([127, 0, 0, 1], 1002)),
source: DiscoverySource::ActiveLookup,
discovered_at: now,
},
];
let started = Arc::new(AtomicUsize::new(0));
let cancelled = Arc::new(AtomicBool::new(false));
let barrier = Arc::new(tokio::sync::Barrier::new(peers.len()));
let started_for_race = started.clone();
let cancelled_for_race = cancelled.clone();
let (_peer_tx, mut peer_rx) = mpsc::channel(MAX_METADATA_PEERS_PER_HASH);
let stats = DhtRuntimeStats::default();
let result = race_peer_fetches(
peers,
&mut peer_rx,
move |peer| {
let started = started_for_race.clone();
let cancelled = cancelled_for_race.clone();
let barrier = barrier.clone();
async move {
started.fetch_add(1, Ordering::Relaxed);
barrier.wait().await;
match peer.addr.port() {
1000 => {
tokio::time::sleep(Duration::from_millis(10)).await;
(
peer,
MetadataFetchOutcome::Fetched((
"winner".to_string(),
1,
Vec::new(),
0,
)),
)
}
1001 => {
tokio::time::sleep(Duration::from_secs(1)).await;
(peer, MetadataFetchOutcome::Failed)
}
_ => {
let _signal = DropSignal(cancelled);
std::future::pending::<()>().await;
unreachable!()
}
}
}
},
&stats,
)
.await
.expect("one peer should win the race");
assert_eq!(result.0.addr.port(), 1000);
assert_eq!(result.1.0, "winner");
assert!(cancelled.load(Ordering::Relaxed));
let snapshot = stats.snapshot();
assert_eq!(snapshot.metadata_races_started, 1);
assert_eq!(snapshot.metadata_peer_candidates, 3);
assert_eq!(snapshot.metadata_peer_canceled, 2);
}
#[tokio::test]
async fn live_peer_joins_running_race_and_cancels_slow_initial_peer() {
let now = Instant::now();
let initial = vec![PeerCandidate {
addr: SocketAddr::from(([127, 0, 0, 1], 1000)),
source: DiscoverySource::AnnouncePeer,
discovered_at: now,
}];
let live = PeerCandidate {
addr: SocketAddr::from(([127, 0, 0, 1], 1001)),
source: DiscoverySource::ActiveLookup,
discovered_at: now,
};
let (peer_tx, mut peer_rx) = mpsc::channel(MAX_METADATA_PEERS_PER_HASH);
let initial_started = Arc::new(tokio::sync::Notify::new());
let initial_started_for_fetch = initial_started.clone();
let initial_cancelled = Arc::new(AtomicBool::new(false));
let initial_cancelled_for_fetch = initial_cancelled.clone();
let stats = DhtRuntimeStats::default();
let sender = tokio::spawn(async move {
initial_started.notified().await;
peer_tx.send(live).await.unwrap();
});
let result = tokio::time::timeout(
Duration::from_secs(1),
race_peer_fetches(
initial,
&mut peer_rx,
move |peer| {
let initial_started = initial_started_for_fetch.clone();
let initial_cancelled = initial_cancelled_for_fetch.clone();
async move {
if peer.addr.port() == 1000 {
let _signal = DropSignal(initial_cancelled);
initial_started.notify_one();
std::future::pending::<()>().await;
unreachable!()
}
(
peer,
MetadataFetchOutcome::Fetched((
"live-winner".to_string(),
1,
Vec::new(),
0,
)),
)
}
},
&stats,
),
)
.await
.unwrap()
.expect("live peer should win");
sender.await.unwrap();
assert_eq!(result.0.addr.port(), 1001);
assert!(initial_cancelled.load(Ordering::Relaxed));
let snapshot = stats.snapshot();
assert_eq!(snapshot.metadata_peer_candidates, 2);
assert_eq!(snapshot.metadata_peer_canceled, 1);
}
#[test]
fn full_queue_evicts_oldest_for_newer_hash() {
let start = Instant::now();
let mut queue = PendingHashQueue::new(2, Duration::from_secs(60));
queue.push(event("old", 1000, start), start);
queue.push(
event("middle", 1001, start + Duration::from_secs(1)),
start + Duration::from_secs(1),
);
let result = queue.push(
event("new", 1002, start + Duration::from_secs(2)),
start + Duration::from_secs(2),
);
assert_eq!(result, QueuePushKind::EvictedOldest);
assert!(!queue.contains("old"));
assert!(queue.contains("middle"));
assert!(queue.contains("new"));
}
#[test]
fn queue_preserves_peer_discovery_source() {
let start = Instant::now();
let mut queue = PendingHashQueue::new(2, Duration::from_secs(60));
let mut direct = event("direct", 1000, start);
direct.source = DiscoverySource::SampleDirect;
queue.push(direct, start);
let entry = queue
.pop_newest_ready(start + Duration::from_secs(1), &HashMap::new())
.unwrap();
assert_eq!(entry.peers[0].source, DiscoverySource::SampleDirect);
}
#[test]
fn queue_expires_stale_hashes_and_pops_newest_first() {
let start = Instant::now();
let mut queue = PendingHashQueue::new(4, Duration::from_secs(10));
queue.push(event("older", 1000, start), start);
queue.push(
event("newer", 1001, start + Duration::from_secs(1)),
start + Duration::from_secs(1),
);
let in_flight = HashMap::new();
assert!(
queue
.pop_newest_ready(start + Duration::from_millis(24), &in_flight)
.is_none()
);
let newest = queue
.pop_newest_ready(start + Duration::from_secs(2), &in_flight)
.unwrap();
assert_eq!(newest.info_hash, "newer");
assert_eq!(queue.expire(start + Duration::from_secs(11)), 1);
assert!(queue.is_empty());
}
#[tokio::test]
async fn gate_rejection_does_not_emit_completion() {
let (hash_tx, hash_rx) = mpsc::channel(4);
let gate_calls = Arc::new(AtomicUsize::new(0));
let gate_calls_for_callback = gate_calls.clone();
let gate = fetch_gate(move |_| {
gate_calls_for_callback.fetch_add(1, Ordering::Relaxed);
async { false }
});
let (completion_tx, mut completion_rx) = mpsc::unbounded_channel();
let completion = completion_callback(move |result| {
let _ = completion_tx.send(result);
});
let shutdown = CancellationToken::new();
let scheduler = MetadataScheduler::new(
hash_rx,
Arc::new(RbitFetcher::new(1)),
MetadataSchedulerLimits {
queue_size: 4,
concurrency: 1,
},
MetadataSchedulerCallbacks {
torrent: empty_torrent_callback(),
fetch_gate: gate,
completion,
},
Arc::new(AtomicUsize::new(0)),
shutdown.clone(),
);
let scheduler_task = tokio::spawn(scheduler.run());
hash_tx
.send(event("bad", 1000, Instant::now()))
.await
.unwrap();
tokio::time::timeout(Duration::from_secs(1), async {
while gate_calls.load(Ordering::Relaxed) == 0 {
tokio::task::yield_now().await;
}
})
.await
.unwrap();
assert!(
tokio::time::timeout(Duration::from_millis(50), completion_rx.recv())
.await
.is_err()
);
shutdown.cancel();
scheduler_task.await.unwrap();
}
#[tokio::test]
async fn admitted_failure_emits_one_completion() {
let (hash_tx, hash_rx) = mpsc::channel(4);
let gate = fetch_gate(|_| async { true });
let (completion_tx, mut completion_rx) = mpsc::unbounded_channel();
let completion = completion_callback(move |result| {
let _ = completion_tx.send(result);
});
let shutdown = CancellationToken::new();
let scheduler = MetadataScheduler::new(
hash_rx,
Arc::new(RbitFetcher::new(1)),
MetadataSchedulerLimits {
queue_size: 4,
concurrency: 1,
},
MetadataSchedulerCallbacks {
torrent: empty_torrent_callback(),
fetch_gate: gate,
completion,
},
Arc::new(AtomicUsize::new(0)),
shutdown.clone(),
);
let scheduler_task = tokio::spawn(scheduler.run());
hash_tx
.send(event("bad", 1000, Instant::now()))
.await
.unwrap();
let result = tokio::time::timeout(Duration::from_secs(1), completion_rx.recv())
.await
.unwrap()
.unwrap();
assert_eq!(result.info_hash, "bad");
assert_eq!(result.status, MetadataFetchCompletionStatus::FetchFailed);
assert_eq!(result.attempts, 0);
assert!(
tokio::time::timeout(Duration::from_millis(50), completion_rx.recv())
.await
.is_err()
);
shutdown.cancel();
scheduler_task.await.unwrap();
}
}
mod tests;
+91
View File
@@ -0,0 +1,91 @@
// 负责并发竞速 Metadata Peer 并在首个成功结果后取消剩余任务
use super::{LIVE_PEER_IDLE_GRACE, MAX_METADATA_PEERS_PER_HASH, PeerCandidate};
use crate::metadata::{FetchedMetadata, MetadataFetchOutcome};
use crate::runtime_stats::DhtRuntimeStats;
#[cfg(feature = "metrics")]
use metrics::counter;
use std::future::Future;
use tokio::sync::mpsc;
use tokio::task::JoinSet;
pub(super) async fn race_peer_fetches<F, Fut>(
peers: Vec<PeerCandidate>,
peer_rx: &mut mpsc::Receiver<PeerCandidate>,
fetch: F,
runtime_stats: &DhtRuntimeStats,
) -> Option<(PeerCandidate, FetchedMetadata)>
where
F: Fn(PeerCandidate) -> Fut + Clone + Send + Sync + 'static,
Fut: Future<Output = (PeerCandidate, MetadataFetchOutcome)> + Send + 'static,
{
runtime_stats.metadata_race_started();
#[cfg(feature = "metrics")]
counter!("dht_metadata_peer_races_total").increment(1);
let mut tasks = JoinSet::new();
let mut candidates = 0usize;
for peer in peers {
let fetch = fetch.clone();
tasks.spawn(fetch(peer));
candidates += 1;
runtime_stats.metadata_peer_candidate();
#[cfg(feature = "metrics")]
counter!("dht_metadata_peer_candidates_total", "source" => "initial").increment(1);
}
let mut accepting_live_peers = candidates < MAX_METADATA_PEERS_PER_HASH;
loop {
if tasks.is_empty() {
if !accepting_live_peers {
return None;
}
match tokio::time::timeout(LIVE_PEER_IDLE_GRACE, peer_rx.recv()).await {
Ok(Some(peer)) => {
let fetch = fetch.clone();
tasks.spawn(fetch(peer));
candidates += 1;
runtime_stats.metadata_peer_candidate();
#[cfg(feature = "metrics")]
counter!("dht_metadata_peer_candidates_total", "source" => "live").increment(1);
accepting_live_peers = candidates < MAX_METADATA_PEERS_PER_HASH;
}
Ok(None) | Err(_) => return None,
}
continue;
}
tokio::select! {
result = tasks.join_next() => {
let Some(Ok((peer, outcome))) = result else {
continue;
};
if let MetadataFetchOutcome::Fetched(metadata) = outcome {
let canceled = tasks.len();
runtime_stats.metadata_peer_canceled(canceled);
#[cfg(feature = "metrics")]
counter!("dht_metadata_peer_canceled_total").increment(canceled as u64);
peer_rx.close();
tasks.abort_all();
while tasks.join_next().await.is_some() {}
return Some((peer, metadata));
}
}
peer = peer_rx.recv(), if accepting_live_peers => {
match peer {
Some(peer) => {
let fetch = fetch.clone();
tasks.spawn(fetch(peer));
candidates += 1;
runtime_stats.metadata_peer_candidate();
#[cfg(feature = "metrics")]
counter!("dht_metadata_peer_candidates_total", "source" => "live")
.increment(1);
accepting_live_peers = candidates < MAX_METADATA_PEERS_PER_HASH;
}
None => accepting_live_peers = false,
}
}
}
}
}
+172
View File
@@ -0,0 +1,172 @@
// 负责维护按新鲜度排序并按 infohash 去重的有界待处理队列
use super::{InFlightState, MAX_METADATA_PEERS_PER_HASH, PEER_COALESCE_WINDOW};
use crate::types::{DiscoverySource, HashDiscovered};
use std::collections::{BTreeMap, HashMap};
use std::net::SocketAddr;
use std::time::{Duration, Instant};
#[derive(Debug, Clone, Copy)]
pub(super) struct PeerCandidate {
pub(super) addr: SocketAddr,
pub(super) source: DiscoverySource,
pub(super) discovered_at: Instant,
}
#[derive(Debug)]
pub(super) struct QueuedHash {
pub(super) info_hash: String,
pub(super) peers: Vec<PeerCandidate>,
pub(super) queued_at: Instant,
pub(super) ready_at: Instant,
pub(super) latest_at: Instant,
pub(super) order_key: (Instant, u64),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(super) enum QueuePushKind {
Inserted,
Updated,
EvictedOldest,
Stale,
}
#[derive(Debug)]
pub(super) struct PendingHashQueue {
capacity: usize,
ttl: Duration,
entries: HashMap<String, QueuedHash>,
order: BTreeMap<(Instant, u64), String>,
sequence: u64,
}
impl PendingHashQueue {
pub(super) fn new(capacity: usize, ttl: Duration) -> Self {
Self {
capacity: capacity.max(1),
ttl,
entries: HashMap::new(),
order: BTreeMap::new(),
sequence: 0,
}
}
pub(super) fn len(&self) -> usize {
self.entries.len()
}
pub(super) fn is_empty(&self) -> bool {
self.entries.is_empty()
}
#[cfg(test)]
pub(super) fn contains(&self, info_hash: &str) -> bool {
self.entries.contains_key(info_hash)
}
pub(super) fn next_order_key(&mut self, at: Instant) -> (Instant, u64) {
self.sequence = self.sequence.wrapping_add(1);
(at, self.sequence)
}
pub(super) fn push(&mut self, event: HashDiscovered, now: Instant) -> QueuePushKind {
if now
.checked_duration_since(event.discovered_at)
.unwrap_or_default()
> self.ttl
{
return QueuePushKind::Stale;
}
if let Some(mut entry) = self.entries.remove(&event.info_hash) {
self.order.remove(&entry.order_key);
let discovered_at = entry
.peers
.iter()
.find(|peer| peer.addr == event.peer_addr)
.map(|peer| peer.discovered_at.max(event.discovered_at))
.unwrap_or(event.discovered_at);
entry.peers.retain(|peer| peer.addr != event.peer_addr);
entry.peers.push(PeerCandidate {
addr: event.peer_addr,
source: event.source,
discovered_at,
});
entry
.peers
.sort_unstable_by_key(|peer| std::cmp::Reverse(peer.discovered_at));
entry.peers.truncate(MAX_METADATA_PEERS_PER_HASH);
entry.latest_at = entry.latest_at.max(event.discovered_at);
entry.order_key = self.next_order_key(entry.latest_at);
self.order.insert(entry.order_key, entry.info_hash.clone());
self.entries.insert(entry.info_hash.clone(), entry);
return QueuePushKind::Updated;
}
let mut result = QueuePushKind::Inserted;
if self.entries.len() >= self.capacity {
let Some((&oldest_key, oldest_hash)) = self.order.first_key_value() else {
return QueuePushKind::Stale;
};
if event.discovered_at <= oldest_key.0 {
return QueuePushKind::Stale;
}
let oldest_hash = oldest_hash.clone();
self.order.remove(&oldest_key);
self.entries.remove(&oldest_hash);
result = QueuePushKind::EvictedOldest;
}
let order_key = self.next_order_key(event.discovered_at);
let info_hash = event.info_hash;
self.order.insert(order_key, info_hash.clone());
self.entries.insert(
info_hash.clone(),
QueuedHash {
info_hash,
peers: vec![PeerCandidate {
addr: event.peer_addr,
source: event.source,
discovered_at: event.discovered_at,
}],
queued_at: now,
ready_at: now + PEER_COALESCE_WINDOW,
latest_at: event.discovered_at,
order_key,
},
);
result
}
pub(super) fn remove(&mut self, info_hash: &str) -> Option<QueuedHash> {
let entry = self.entries.remove(info_hash)?;
self.order.remove(&entry.order_key);
Some(entry)
}
pub(super) fn pop_newest_ready(
&mut self,
now: Instant,
in_flight: &HashMap<String, InFlightState>,
) -> Option<QueuedHash> {
let info_hash = self.order.iter().rev().find_map(|(_, hash)| {
let entry = self.entries.get(hash)?;
(!in_flight.contains_key(hash) && entry.ready_at <= now).then(|| hash.clone())
})?;
self.remove(&info_hash)
}
pub(super) fn expire(&mut self, now: Instant) -> usize {
let mut expired = 0;
while let Some((&oldest_key, oldest_hash)) = self.order.first_key_value() {
if now.checked_duration_since(oldest_key.0).unwrap_or_default() <= self.ttl {
break;
}
let oldest_hash = oldest_hash.clone();
self.order.remove(&oldest_key);
self.entries.remove(&oldest_hash);
expired += 1;
}
expired
}
}
+510
View File
@@ -0,0 +1,510 @@
// 负责验证 Metadata 队列调度 Peer 竞速和完成回调行为
use super::*;
use crate::metadata::MetadataFetchOutcome;
use crate::types::DiscoverySource;
use std::sync::atomic::AtomicBool;
struct DropSignal(Arc<AtomicBool>);
impl Drop for DropSignal {
fn drop(&mut self) {
self.0.store(true, Ordering::Relaxed);
}
}
fn empty_torrent_callback() -> Arc<ArcSwapOption<TorrentAckCallback>> {
Arc::new(ArcSwapOption::empty())
}
fn fetch_gate<F, Fut>(callback: F) -> Arc<ArcSwapOption<MetadataFetchCallback>>
where
F: Fn(String) -> Fut + Send + Sync + 'static,
Fut: std::future::Future<Output = bool> + Send + 'static,
{
let holder = Arc::new(ArcSwapOption::empty());
let callback: Arc<MetadataFetchCallback> =
Arc::new(Box::new(move |hash| Box::pin(callback(hash))));
holder.store(Some(callback));
holder
}
fn completion_callback<F>(callback: F) -> Arc<ArcSwapOption<MetadataCompletionCallback>>
where
F: Fn(MetadataFetchCompletion) + Send + Sync + 'static,
{
let holder = Arc::new(ArcSwapOption::empty());
let callback: Arc<MetadataCompletionCallback> = Arc::new(Box::new(callback));
holder.store(Some(callback));
holder
}
fn event(hash: &str, port: u16, discovered_at: Instant) -> HashDiscovered {
HashDiscovered {
info_hash: hash.to_string(),
peer_addr: SocketAddr::from(([127, 0, 0, 1], port)),
source: DiscoverySource::AnnouncePeer,
discovered_at,
}
}
#[test]
fn runtime_stats_track_queue_outcomes_and_depth() {
let (_hash_tx, hash_rx) = mpsc::channel(4);
let stats = DhtRuntimeStats::with_limits(DhtRuntimeLimits {
metadata_queue: 2,
..DhtRuntimeLimits::default()
});
let scheduler = MetadataScheduler::new_with_runtime_stats(
hash_rx,
Arc::new(RbitFetcher::new(1)),
MetadataSchedulerLimits {
queue_size: 2,
concurrency: 1,
},
MetadataSchedulerCallbacks {
torrent: empty_torrent_callback(),
fetch_gate: Arc::new(ArcSwapOption::empty()),
completion: Arc::new(ArcSwapOption::empty()),
},
Arc::new(AtomicUsize::new(0)),
CancellationToken::new(),
stats.clone(),
);
let start = Instant::now();
let mut queue = PendingHashQueue::new(2, HASH_QUEUE_TTL);
let mut in_flight = HashMap::new();
scheduler.enqueue(&mut queue, &mut in_flight, event("old", 1000, start));
scheduler.enqueue(
&mut queue,
&mut in_flight,
event("old", 1001, start + Duration::from_secs(1)),
);
scheduler.enqueue(
&mut queue,
&mut in_flight,
event("middle", 1002, start + Duration::from_secs(2)),
);
scheduler.enqueue(
&mut queue,
&mut in_flight,
event("new", 1003, start + Duration::from_secs(3)),
);
scheduler.enqueue(
&mut queue,
&mut in_flight,
event(
"stale",
1004,
start
.checked_sub(HASH_QUEUE_TTL + Duration::from_secs(1))
.unwrap(),
),
);
scheduler.sync_queue_len(queue.len(), 1);
let snapshot = stats.snapshot();
assert_eq!(snapshot.metadata_queue_inserted, 3);
assert_eq!(snapshot.metadata_queue_deduplicated, 1);
assert_eq!(snapshot.metadata_queue_evicted, 1);
assert_eq!(snapshot.metadata_queue_stale, 1);
assert_eq!(snapshot.metadata_queue_depth, 3);
assert_eq!(snapshot.metadata_queue_max, 2);
assert_eq!(snapshot.metadata_in_flight, 1);
}
#[test]
fn queue_deduplicates_hash_and_keeps_twelve_newest_unique_peers() {
let start = Instant::now();
let mut queue = PendingHashQueue::new(10, Duration::from_secs(60));
for offset in 0..13 {
assert_ne!(
queue.push(
event(
"hash",
1000 + offset,
start + Duration::from_secs(offset as u64)
),
start + Duration::from_secs(offset as u64),
),
QueuePushKind::Stale
);
}
queue.push(
event("hash", 1012, start + Duration::from_secs(20)),
start + Duration::from_secs(20),
);
assert_eq!(queue.len(), 1);
let entry = queue.remove("hash").unwrap();
let ports: Vec<_> = entry.peers.iter().map(|peer| peer.addr.port()).collect();
assert_eq!(ports, (1001..=1012).rev().collect::<Vec<_>>());
}
#[test]
fn coalescing_window_is_fixed_from_first_enqueue() {
let start = Instant::now();
let mut queue = PendingHashQueue::new(10, Duration::from_secs(60));
let in_flight = HashMap::new();
queue.push(event("hash", 1000, start), start);
queue.push(
event("hash", 1001, start + Duration::from_millis(15)),
start + Duration::from_millis(15),
);
assert!(
queue
.pop_newest_ready(start + Duration::from_millis(24), &in_flight)
.is_none()
);
let entry = queue
.pop_newest_ready(start + Duration::from_millis(25), &in_flight)
.expect("duplicate arrivals must not extend the coalescing deadline");
assert_eq!(entry.peers.len(), 2);
}
#[test]
fn in_flight_hash_routes_at_most_twelve_unique_peers_to_live_race() {
let (_hash_tx, hash_rx) = mpsc::channel(4);
let stats = DhtRuntimeStats::default();
let scheduler = MetadataScheduler::new_with_runtime_stats(
hash_rx,
Arc::new(RbitFetcher::new(1)),
MetadataSchedulerLimits {
queue_size: 4,
concurrency: 1,
},
MetadataSchedulerCallbacks {
torrent: empty_torrent_callback(),
fetch_gate: Arc::new(ArcSwapOption::empty()),
completion: Arc::new(ArcSwapOption::empty()),
},
Arc::new(AtomicUsize::new(0)),
CancellationToken::new(),
stats.clone(),
);
let (peer_tx, mut peer_rx) = mpsc::channel(MAX_METADATA_PEERS_PER_HASH);
let mut in_flight = HashMap::from([(
"hash".to_string(),
InFlightState {
peer_tx,
scheduled_peers: HashSet::from([SocketAddr::from(([127, 0, 0, 1], 1000))]),
},
)]);
let mut queue = PendingHashQueue::new(4, HASH_QUEUE_TTL);
let now = Instant::now();
for port in 1001..=1013 {
scheduler.enqueue(&mut queue, &mut in_flight, event("hash", port, now));
}
scheduler.enqueue(&mut queue, &mut in_flight, event("hash", 1001, now));
let mut received = Vec::new();
while let Ok(peer) = peer_rx.try_recv() {
received.push(peer.addr.port());
}
assert_eq!(received, (1001..=1011).collect::<Vec<_>>());
assert_eq!(
in_flight["hash"].scheduled_peers.len(),
MAX_METADATA_PEERS_PER_HASH
);
assert!(queue.is_empty());
let snapshot = stats.snapshot();
assert_eq!(snapshot.metadata_live_peer_joins, 11);
assert_eq!(snapshot.metadata_queue_deduplicated, 14);
}
#[tokio::test]
async fn peer_race_returns_first_success_and_cancels_remaining_fetches() {
let now = Instant::now();
let peers = vec![
PeerCandidate {
addr: SocketAddr::from(([127, 0, 0, 1], 1000)),
source: DiscoverySource::AnnouncePeer,
discovered_at: now,
},
PeerCandidate {
addr: SocketAddr::from(([127, 0, 0, 1], 1001)),
source: DiscoverySource::SampleDirect,
discovered_at: now,
},
PeerCandidate {
addr: SocketAddr::from(([127, 0, 0, 1], 1002)),
source: DiscoverySource::ActiveLookup,
discovered_at: now,
},
];
let started = Arc::new(AtomicUsize::new(0));
let cancelled = Arc::new(AtomicBool::new(false));
let barrier = Arc::new(tokio::sync::Barrier::new(peers.len()));
let started_for_race = started.clone();
let cancelled_for_race = cancelled.clone();
let (_peer_tx, mut peer_rx) = mpsc::channel(MAX_METADATA_PEERS_PER_HASH);
let stats = DhtRuntimeStats::default();
let result = race_peer_fetches(
peers,
&mut peer_rx,
move |peer| {
let started = started_for_race.clone();
let cancelled = cancelled_for_race.clone();
let barrier = barrier.clone();
async move {
started.fetch_add(1, Ordering::Relaxed);
barrier.wait().await;
match peer.addr.port() {
1000 => {
tokio::time::sleep(Duration::from_millis(10)).await;
(
peer,
MetadataFetchOutcome::Fetched(("winner".to_string(), 1, Vec::new(), 0)),
)
}
1001 => {
tokio::time::sleep(Duration::from_secs(1)).await;
(peer, MetadataFetchOutcome::Failed)
}
_ => {
let _signal = DropSignal(cancelled);
std::future::pending::<()>().await;
unreachable!()
}
}
}
},
&stats,
)
.await
.expect("one peer should win the race");
assert_eq!(result.0.addr.port(), 1000);
assert_eq!(result.1.0, "winner");
assert!(cancelled.load(Ordering::Relaxed));
let snapshot = stats.snapshot();
assert_eq!(snapshot.metadata_races_started, 1);
assert_eq!(snapshot.metadata_peer_candidates, 3);
assert_eq!(snapshot.metadata_peer_canceled, 2);
}
#[tokio::test]
async fn live_peer_joins_running_race_and_cancels_slow_initial_peer() {
let now = Instant::now();
let initial = vec![PeerCandidate {
addr: SocketAddr::from(([127, 0, 0, 1], 1000)),
source: DiscoverySource::AnnouncePeer,
discovered_at: now,
}];
let live = PeerCandidate {
addr: SocketAddr::from(([127, 0, 0, 1], 1001)),
source: DiscoverySource::ActiveLookup,
discovered_at: now,
};
let (peer_tx, mut peer_rx) = mpsc::channel(MAX_METADATA_PEERS_PER_HASH);
let initial_started = Arc::new(tokio::sync::Notify::new());
let initial_started_for_fetch = initial_started.clone();
let initial_cancelled = Arc::new(AtomicBool::new(false));
let initial_cancelled_for_fetch = initial_cancelled.clone();
let stats = DhtRuntimeStats::default();
let sender = tokio::spawn(async move {
initial_started.notified().await;
peer_tx.send(live).await.unwrap();
});
let result = tokio::time::timeout(
Duration::from_secs(1),
race_peer_fetches(
initial,
&mut peer_rx,
move |peer| {
let initial_started = initial_started_for_fetch.clone();
let initial_cancelled = initial_cancelled_for_fetch.clone();
async move {
if peer.addr.port() == 1000 {
let _signal = DropSignal(initial_cancelled);
initial_started.notify_one();
std::future::pending::<()>().await;
unreachable!()
}
(
peer,
MetadataFetchOutcome::Fetched((
"live-winner".to_string(),
1,
Vec::new(),
0,
)),
)
}
},
&stats,
),
)
.await
.unwrap()
.expect("live peer should win");
sender.await.unwrap();
assert_eq!(result.0.addr.port(), 1001);
assert!(initial_cancelled.load(Ordering::Relaxed));
let snapshot = stats.snapshot();
assert_eq!(snapshot.metadata_peer_candidates, 2);
assert_eq!(snapshot.metadata_peer_canceled, 1);
}
#[test]
fn full_queue_evicts_oldest_for_newer_hash() {
let start = Instant::now();
let mut queue = PendingHashQueue::new(2, Duration::from_secs(60));
queue.push(event("old", 1000, start), start);
queue.push(
event("middle", 1001, start + Duration::from_secs(1)),
start + Duration::from_secs(1),
);
let result = queue.push(
event("new", 1002, start + Duration::from_secs(2)),
start + Duration::from_secs(2),
);
assert_eq!(result, QueuePushKind::EvictedOldest);
assert!(!queue.contains("old"));
assert!(queue.contains("middle"));
assert!(queue.contains("new"));
}
#[test]
fn queue_preserves_peer_discovery_source() {
let start = Instant::now();
let mut queue = PendingHashQueue::new(2, Duration::from_secs(60));
let mut direct = event("direct", 1000, start);
direct.source = DiscoverySource::SampleDirect;
queue.push(direct, start);
let entry = queue
.pop_newest_ready(start + Duration::from_secs(1), &HashMap::new())
.unwrap();
assert_eq!(entry.peers[0].source, DiscoverySource::SampleDirect);
}
#[test]
fn queue_expires_stale_hashes_and_pops_newest_first() {
let start = Instant::now();
let mut queue = PendingHashQueue::new(4, Duration::from_secs(10));
queue.push(event("older", 1000, start), start);
queue.push(
event("newer", 1001, start + Duration::from_secs(1)),
start + Duration::from_secs(1),
);
let in_flight = HashMap::new();
assert!(
queue
.pop_newest_ready(start + Duration::from_millis(24), &in_flight)
.is_none()
);
let newest = queue
.pop_newest_ready(start + Duration::from_secs(2), &in_flight)
.unwrap();
assert_eq!(newest.info_hash, "newer");
assert_eq!(queue.expire(start + Duration::from_secs(11)), 1);
assert!(queue.is_empty());
}
#[tokio::test]
async fn gate_rejection_does_not_emit_completion() {
let (hash_tx, hash_rx) = mpsc::channel(4);
let gate_calls = Arc::new(AtomicUsize::new(0));
let gate_calls_for_callback = gate_calls.clone();
let gate = fetch_gate(move |_| {
gate_calls_for_callback.fetch_add(1, Ordering::Relaxed);
async { false }
});
let (completion_tx, mut completion_rx) = mpsc::unbounded_channel();
let completion = completion_callback(move |result| {
let _ = completion_tx.send(result);
});
let shutdown = CancellationToken::new();
let scheduler = MetadataScheduler::new(
hash_rx,
Arc::new(RbitFetcher::new(1)),
MetadataSchedulerLimits {
queue_size: 4,
concurrency: 1,
},
MetadataSchedulerCallbacks {
torrent: empty_torrent_callback(),
fetch_gate: gate,
completion,
},
Arc::new(AtomicUsize::new(0)),
shutdown.clone(),
);
let scheduler_task = tokio::spawn(scheduler.run());
hash_tx
.send(event("bad", 1000, Instant::now()))
.await
.unwrap();
tokio::time::timeout(Duration::from_secs(1), async {
while gate_calls.load(Ordering::Relaxed) == 0 {
tokio::task::yield_now().await;
}
})
.await
.unwrap();
assert!(
tokio::time::timeout(Duration::from_millis(50), completion_rx.recv())
.await
.is_err()
);
shutdown.cancel();
scheduler_task.await.unwrap();
}
#[tokio::test]
async fn admitted_failure_emits_one_completion() {
let (hash_tx, hash_rx) = mpsc::channel(4);
let gate = fetch_gate(|_| async { true });
let (completion_tx, mut completion_rx) = mpsc::unbounded_channel();
let completion = completion_callback(move |result| {
let _ = completion_tx.send(result);
});
let shutdown = CancellationToken::new();
let scheduler = MetadataScheduler::new(
hash_rx,
Arc::new(RbitFetcher::new(1)),
MetadataSchedulerLimits {
queue_size: 4,
concurrency: 1,
},
MetadataSchedulerCallbacks {
torrent: empty_torrent_callback(),
fetch_gate: gate,
completion,
},
Arc::new(AtomicUsize::new(0)),
shutdown.clone(),
);
let scheduler_task = tokio::spawn(scheduler.run());
hash_tx
.send(event("bad", 1000, Instant::now()))
.await
.unwrap();
let result = tokio::time::timeout(Duration::from_secs(1), completion_rx.recv())
.await
.unwrap()
.unwrap();
assert_eq!(result.info_hash, "bad");
assert_eq!(result.status, MetadataFetchCompletionStatus::FetchFailed);
assert_eq!(result.attempts, 0);
assert!(
tokio::time::timeout(Duration::from_millis(50), completion_rx.recv())
.await
.is_err()
);
shutdown.cancel();
scheduler_task.await.unwrap();
}