perf: 增加新节点直接采样策略
This commit is contained in:
@@ -255,7 +255,10 @@
|
||||
- [x] 实现样本来源节点单点 `get_peers` 优先和失败后有限递归降级
|
||||
- [x] 将 BEP-51 最大在途请求和采样失败后的迭代回退暴露为应用配置
|
||||
- [x] 使用 Bitmagnet 等效并发完成一分钟资源测试并确认主网卡无丢包无错误且持久化无积压
|
||||
- [ ] 评估将新发现但尚未验证的节点直接作为 BEP-51 候选并避免与 `find_node` 重复探测
|
||||
- [x] 将新发现但尚未验证的节点按地址稳定分流到有界 BEP-51 通道并避免与 `find_node` 重复探测
|
||||
- [x] 增加直接采样候选队列请求响应重复过滤丢弃和 Metadata 成功来源转化指标
|
||||
- [x] 使用相同网络预算对旧快照采样和新节点直接采样完成十分钟对比
|
||||
- [x] 确认直接采样在相近 UDP 流量下 Metadata 成功数提高约百分之六点六且单条成功 UDP 成本降低约百分之六点六
|
||||
- [x] 完成首轮三分钟对比并验证 Peer Lookup UDP 从 `278` 降至 `254` 且网络稳定
|
||||
- [ ] 通过多轮或更长时间运行评估随机 DHT 样本下的 Metadata 成功率
|
||||
- [ ] 统计按需验证的 Peer 发现率握手成功率和平均验证耗时
|
||||
@@ -322,6 +325,6 @@
|
||||
- [x] 将规模基准拆分为参数数据集工作负载采样报告和编排模块
|
||||
- [x] 明确单元组件集成端到端和性能测试层级并增加公开 API 集成测试
|
||||
|
||||
对比 Bitmagnet 新节点直接采样策略并测量有效样本数每次成功 Metadata 的 UDP 和 TCP 成本
|
||||
完成二十四小时持续运行并继续观察私有内存 Metadata 成功率候选队列深度和每条成功 Metadata 的网络成本
|
||||
|
||||
随后完成二十四小时持续运行并根据数据决定资源参数和正则查询优化优先级
|
||||
随后增加磁盘剩余空间保护日志轮转和 RocksDB 检查点恢复
|
||||
|
||||
@@ -180,7 +180,7 @@ let options = DHTOptions {
|
||||
| `DHTOptions` | 监听端口、网络模式、顶层队列和主动 UDP 查询总预算 |
|
||||
| `MetadataOptions` | 下载超时、队列、并发、每秒 TCP 建连和失败 Peer 缓存 |
|
||||
| `PeerLookupOptions` | 主动 `get_peers` 的速率与并发 |
|
||||
| `SampleInfohashesOptions` | BEP-51 采样速率、并发、超时、退避和 Hash 去重容量 |
|
||||
| `SampleInfohashesOptions` | BEP-51 采样速率、并发、新节点稳定分流、有界候选队列、超时、退避和 Hash 去重容量 |
|
||||
| `RateLimitOptions` | `find_node`、在途请求和 UDP 回复预算 |
|
||||
| `PoolOptions` | 节点池、最近探测记录和响应节点缓存 |
|
||||
| `BootstrapOptions` | Bootstrap 节点与失败退避 |
|
||||
@@ -191,6 +191,8 @@ let options = DHTOptions {
|
||||
|
||||
- BEP-51 采样 hash 可以通过 `DHTServer::on_sampled_hashes` 批量异步准入
|
||||
- 采样准入队列有固定容量并在压力升高时暂停新的 BEP-51 查询
|
||||
- `new_node_sample_percent` 大于零时按节点地址稳定分流并优先直接采样 避免同一次发现同时执行 `find_node`
|
||||
- 新节点采样通道满时回退到抓取池 不会依靠无限队列维持吞吐
|
||||
- 带首选节点的 Peer Lookup 先执行单点查询只有失败后才进入有限迭代查找
|
||||
|
||||
- `DHTOptions::default()` 使用 `Ipv4Only`;
|
||||
|
||||
@@ -13,6 +13,7 @@ use crate::routing_snapshot::RoutingSnapshot;
|
||||
#[cfg(test)]
|
||||
use crate::runtime_stats::DhtRuntimeLimits;
|
||||
use crate::runtime_stats::DhtRuntimeStats;
|
||||
use crate::sample_infohashes::SampleCandidateRouter;
|
||||
use crate::types::{NetMode, NodeTuple};
|
||||
use ahash::AHashMap;
|
||||
use arc_swap::ArcSwap;
|
||||
@@ -80,6 +81,7 @@ pub(crate) struct CrawlEngine {
|
||||
pub(crate) node_count: Arc<AtomicUsize>,
|
||||
runtime_stats: DhtRuntimeStats,
|
||||
outbound_query_budget: SharedRateBudget,
|
||||
sample_candidate_router: Option<SampleCandidateRouter>,
|
||||
}
|
||||
|
||||
impl CrawlEngine {
|
||||
@@ -87,6 +89,7 @@ impl CrawlEngine {
|
||||
config: ResolvedCrawlConfig,
|
||||
runtime_stats: DhtRuntimeStats,
|
||||
outbound_query_budget: SharedRateBudget,
|
||||
sample_candidate_router: Option<SampleCandidateRouter>,
|
||||
) -> Self {
|
||||
let (priority_tx, priority_rx) = mpsc::channel(config.priority_event_channel_capacity);
|
||||
let (discovery_tx, discovery_rx) = mpsc::channel(config.discovery_event_channel_capacity);
|
||||
@@ -99,10 +102,18 @@ impl CrawlEngine {
|
||||
node_count: Arc::new(AtomicUsize::new(0)),
|
||||
runtime_stats,
|
||||
outbound_query_budget,
|
||||
sample_candidate_router,
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn route_discovered(&self, node: NodeTuple) {
|
||||
if self
|
||||
.sample_candidate_router
|
||||
.as_ref()
|
||||
.is_some_and(|router| router.route(node))
|
||||
{
|
||||
return;
|
||||
}
|
||||
let enqueue_result = self.discovery_tx.try_send(node);
|
||||
self.runtime_stats.set_crawl_discovery_queue_depth(
|
||||
self.discovery_tx
|
||||
@@ -935,6 +946,7 @@ mod tests {
|
||||
config,
|
||||
stats.clone(),
|
||||
SharedRateBudget::per_second(10_000, 10_000, true),
|
||||
None,
|
||||
);
|
||||
|
||||
engine.route_discovered(node(1, "8.8.8.8:1"));
|
||||
|
||||
@@ -40,7 +40,7 @@ pub use runtime_stats::{
|
||||
pub use scheduler::{MetadataScheduler, MetadataSchedulerCallbacks, MetadataSchedulerLimits};
|
||||
pub use server::{DHTServer, HashDiscovered};
|
||||
pub use types::{
|
||||
BootstrapOptions, CrawlOptions, DHTOptions, FileInfo, MetadataFetchCompletion,
|
||||
BootstrapOptions, CrawlOptions, DHTOptions, DiscoverySource, FileInfo, MetadataFetchCompletion,
|
||||
MetadataFetchCompletionStatus, MetadataOptions, NetMode, NodeTuple, PeerLookupOptions,
|
||||
PeerLookupResult, PoolOptions, RateLimitOptions, SampleInfohashesOptions, SchedulerOptions,
|
||||
TargetOptions, TorrentInfo,
|
||||
@@ -55,9 +55,9 @@ pub mod prelude {
|
||||
};
|
||||
pub use crate::server::DHTServer;
|
||||
pub use crate::types::{
|
||||
BootstrapOptions, CrawlOptions, DHTOptions, FileInfo, MetadataFetchCompletion,
|
||||
MetadataFetchCompletionStatus, MetadataOptions, NetMode, NodeTuple, PeerLookupOptions,
|
||||
PeerLookupResult, PoolOptions, RateLimitOptions, SampleInfohashesOptions, SchedulerOptions,
|
||||
TargetOptions, TorrentInfo,
|
||||
BootstrapOptions, CrawlOptions, DHTOptions, DiscoverySource, FileInfo,
|
||||
MetadataFetchCompletion, MetadataFetchCompletionStatus, MetadataOptions, NetMode,
|
||||
NodeTuple, PeerLookupOptions, PeerLookupResult, PoolOptions, RateLimitOptions,
|
||||
SampleInfohashesOptions, SchedulerOptions, TargetOptions, TorrentInfo,
|
||||
};
|
||||
}
|
||||
|
||||
@@ -7,7 +7,7 @@ use crate::protocol::DhtResponse;
|
||||
use crate::routing_snapshot::{RoutingSnapshot, xor_distance_cmp};
|
||||
use crate::runtime_stats::DhtRuntimeStats;
|
||||
use crate::server::HashDiscovered;
|
||||
use crate::types::{NetMode, NodeTuple, PeerLookupOptions, PeerLookupResult};
|
||||
use crate::types::{DiscoverySource, NetMode, NodeTuple, PeerLookupOptions, PeerLookupResult};
|
||||
use ahash::{AHashMap, AHashSet};
|
||||
use arc_swap::ArcSwap;
|
||||
use bytes::BytesMut;
|
||||
@@ -57,6 +57,7 @@ struct LookupState {
|
||||
deadline: Instant,
|
||||
preferred_phase: bool,
|
||||
allow_iterative_fallback: bool,
|
||||
source: DiscoverySource,
|
||||
completion: Option<oneshot::Sender<PeerLookupResult>>,
|
||||
}
|
||||
|
||||
@@ -92,6 +93,7 @@ pub(crate) struct PeerLookupRequest {
|
||||
pub(crate) info_hash: [u8; 20],
|
||||
pub(crate) preferred_node: Option<NodeTuple>,
|
||||
pub(crate) allow_iterative_fallback: bool,
|
||||
pub(crate) source: DiscoverySource,
|
||||
pub(crate) completion: Option<oneshot::Sender<PeerLookupResult>>,
|
||||
}
|
||||
|
||||
@@ -101,6 +103,7 @@ impl PeerLookupRequest {
|
||||
info_hash,
|
||||
preferred_node: None,
|
||||
allow_iterative_fallback: true,
|
||||
source: DiscoverySource::ActiveLookup,
|
||||
completion: None,
|
||||
}
|
||||
}
|
||||
@@ -135,6 +138,7 @@ impl PeerLookupHandle {
|
||||
info_hash,
|
||||
preferred_node: None,
|
||||
allow_iterative_fallback: true,
|
||||
source: DiscoverySource::ActiveLookup,
|
||||
completion: Some(completion),
|
||||
})
|
||||
.await
|
||||
@@ -369,6 +373,7 @@ impl PeerLookupActor {
|
||||
deadline: now + LOOKUP_TIMEOUT,
|
||||
preferred_phase: preferred.is_some(),
|
||||
allow_iterative_fallback: request.allow_iterative_fallback,
|
||||
source: request.source,
|
||||
completion: request.completion,
|
||||
},
|
||||
);
|
||||
@@ -495,12 +500,14 @@ impl PeerLookupActor {
|
||||
let lookup_id = pending.lookup_id;
|
||||
let preferred_succeeded = pending.preferred_phase && !discovered.is_empty();
|
||||
let preferred_without_fallback = pending.preferred_phase && !state.allow_iterative_fallback;
|
||||
let source = state.source;
|
||||
let _ = state;
|
||||
|
||||
for peer in discovered {
|
||||
let event = HashDiscovered {
|
||||
info_hash: hash.clone(),
|
||||
peer_addr: peer,
|
||||
source,
|
||||
discovered_at: now,
|
||||
};
|
||||
if self.hash_tx.try_send(event).is_ok() {
|
||||
@@ -666,6 +673,7 @@ mod tests {
|
||||
addr: remote_addr,
|
||||
}),
|
||||
allow_iterative_fallback: true,
|
||||
source: DiscoverySource::SampleDirect,
|
||||
completion: Some(completion),
|
||||
})
|
||||
.await
|
||||
@@ -706,6 +714,7 @@ mod tests {
|
||||
.unwrap();
|
||||
assert_eq!(event.info_hash, hex::encode([3; 20]));
|
||||
assert_eq!(event.peer_addr, "8.8.4.4:6881".parse().unwrap());
|
||||
assert_eq!(event.source, DiscoverySource::SampleDirect);
|
||||
let result = tokio::time::timeout(Duration::from_secs(1), result_rx)
|
||||
.await
|
||||
.unwrap()
|
||||
@@ -763,6 +772,7 @@ mod tests {
|
||||
addr: preferred_addr,
|
||||
}),
|
||||
allow_iterative_fallback: true,
|
||||
source: DiscoverySource::SampleDirect,
|
||||
completion: None,
|
||||
})
|
||||
.await
|
||||
@@ -848,6 +858,7 @@ mod tests {
|
||||
addr: preferred_addr,
|
||||
}),
|
||||
allow_iterative_fallback: false,
|
||||
source: DiscoverySource::SampleDirect,
|
||||
completion: Some(completion),
|
||||
})
|
||||
.await
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
// 负责收集 DHT 抓取调度网络和 Metadata 运行指标
|
||||
|
||||
use crate::types::DiscoverySource;
|
||||
use std::sync::Arc;
|
||||
use std::sync::atomic::{AtomicU32, AtomicU64, AtomicUsize, Ordering};
|
||||
|
||||
@@ -252,10 +253,26 @@ pub struct DhtRuntimeSnapshot {
|
||||
pub peer_lookup_preferred_succeeded: u64,
|
||||
/// Preferred-node attempts that fell back to iterative lookup.
|
||||
pub peer_lookup_fallbacks: u64,
|
||||
/// Current newly discovered node sampling-lane depth.
|
||||
pub sample_candidate_queue_depth: usize,
|
||||
/// Configured newly discovered node sampling-lane capacity.
|
||||
pub sample_candidate_queue_capacity: usize,
|
||||
/// Newly discovered nodes routed directly to BEP-51 sampling.
|
||||
pub sample_candidates_routed: u64,
|
||||
/// Sampling candidates sent to the crawl pool because the bounded lane was full.
|
||||
pub sample_candidates_fallback: u64,
|
||||
/// BEP-51 requests sent.
|
||||
pub sample_infohashes_queries: u64,
|
||||
/// BEP-51 requests sent directly to newly discovered nodes.
|
||||
pub sample_infohashes_direct_queries: u64,
|
||||
/// BEP-51 requests sent to the responsive-node snapshot.
|
||||
pub sample_infohashes_snapshot_queries: u64,
|
||||
/// Matched BEP-51 responses.
|
||||
pub sample_infohashes_responses: u64,
|
||||
/// Matched responses from directly sampled new nodes.
|
||||
pub sample_infohashes_direct_responses: u64,
|
||||
/// Matched responses from responsive-node snapshot sampling.
|
||||
pub sample_infohashes_snapshot_responses: u64,
|
||||
/// BEP-51 requests that timed out.
|
||||
pub sample_infohashes_timeouts: u64,
|
||||
/// BEP-51 UDP sends that failed immediately.
|
||||
@@ -274,6 +291,14 @@ pub struct DhtRuntimeSnapshot {
|
||||
pub metadata_peer_attempts: u64,
|
||||
/// Successful Peer downloads and parses.
|
||||
pub metadata_peer_succeeded: u64,
|
||||
/// Successful Metadata whose winning Peer came from `announce_peer`.
|
||||
pub metadata_success_from_announce: u64,
|
||||
/// Successful Metadata whose winning Peer came from direct new-node sampling.
|
||||
pub metadata_success_from_sample_direct: u64,
|
||||
/// Successful Metadata whose winning Peer came from responsive snapshot sampling.
|
||||
pub metadata_success_from_sample_snapshot: u64,
|
||||
/// Successful Metadata whose winning Peer came from delayed active lookup.
|
||||
pub metadata_success_from_active_lookup: u64,
|
||||
/// Failed real Peer attempts.
|
||||
pub metadata_peer_failed: u64,
|
||||
/// End-to-end Peer timeouts.
|
||||
@@ -366,8 +391,16 @@ struct DhtRuntimeStatsInner {
|
||||
peer_lookup_output_dropped: AtomicU64,
|
||||
peer_lookup_preferred_succeeded: AtomicU64,
|
||||
peer_lookup_fallbacks: AtomicU64,
|
||||
sample_candidate_queue_depth: AtomicUsize,
|
||||
sample_candidate_queue_capacity: AtomicUsize,
|
||||
sample_candidates_routed: AtomicU64,
|
||||
sample_candidates_fallback: AtomicU64,
|
||||
sample_infohashes_queries: AtomicU64,
|
||||
sample_infohashes_direct_queries: AtomicU64,
|
||||
sample_infohashes_snapshot_queries: AtomicU64,
|
||||
sample_infohashes_responses: AtomicU64,
|
||||
sample_infohashes_direct_responses: AtomicU64,
|
||||
sample_infohashes_snapshot_responses: AtomicU64,
|
||||
sample_infohashes_timeouts: AtomicU64,
|
||||
sample_infohashes_send_failures: AtomicU64,
|
||||
sample_infohashes_response_dropped: AtomicU64,
|
||||
@@ -377,6 +410,10 @@ struct DhtRuntimeStatsInner {
|
||||
sample_infohashes_hashes_dropped: AtomicU64,
|
||||
metadata_peer_attempts: AtomicU64,
|
||||
metadata_peer_succeeded: AtomicU64,
|
||||
metadata_success_from_announce: AtomicU64,
|
||||
metadata_success_from_sample_direct: AtomicU64,
|
||||
metadata_success_from_sample_snapshot: AtomicU64,
|
||||
metadata_success_from_active_lookup: AtomicU64,
|
||||
metadata_peer_failed: AtomicU64,
|
||||
metadata_peer_timeouts: AtomicU64,
|
||||
metadata_connect_failed: AtomicU64,
|
||||
@@ -479,8 +516,16 @@ impl Default for DhtRuntimeStatsInner {
|
||||
peer_lookup_output_dropped: AtomicU64::new(0),
|
||||
peer_lookup_preferred_succeeded: AtomicU64::new(0),
|
||||
peer_lookup_fallbacks: AtomicU64::new(0),
|
||||
sample_candidate_queue_depth: AtomicUsize::new(0),
|
||||
sample_candidate_queue_capacity: AtomicUsize::new(0),
|
||||
sample_candidates_routed: AtomicU64::new(0),
|
||||
sample_candidates_fallback: AtomicU64::new(0),
|
||||
sample_infohashes_queries: AtomicU64::new(0),
|
||||
sample_infohashes_direct_queries: AtomicU64::new(0),
|
||||
sample_infohashes_snapshot_queries: AtomicU64::new(0),
|
||||
sample_infohashes_responses: AtomicU64::new(0),
|
||||
sample_infohashes_direct_responses: AtomicU64::new(0),
|
||||
sample_infohashes_snapshot_responses: AtomicU64::new(0),
|
||||
sample_infohashes_timeouts: AtomicU64::new(0),
|
||||
sample_infohashes_send_failures: AtomicU64::new(0),
|
||||
sample_infohashes_response_dropped: AtomicU64::new(0),
|
||||
@@ -490,6 +535,10 @@ impl Default for DhtRuntimeStatsInner {
|
||||
sample_infohashes_hashes_dropped: AtomicU64::new(0),
|
||||
metadata_peer_attempts: AtomicU64::new(0),
|
||||
metadata_peer_succeeded: AtomicU64::new(0),
|
||||
metadata_success_from_announce: AtomicU64::new(0),
|
||||
metadata_success_from_sample_direct: AtomicU64::new(0),
|
||||
metadata_success_from_sample_snapshot: AtomicU64::new(0),
|
||||
metadata_success_from_active_lookup: AtomicU64::new(0),
|
||||
metadata_peer_failed: AtomicU64::new(0),
|
||||
metadata_peer_timeouts: AtomicU64::new(0),
|
||||
metadata_connect_failed: AtomicU64::new(0),
|
||||
@@ -619,8 +668,28 @@ impl DhtRuntimeStats {
|
||||
.peer_lookup_preferred_succeeded
|
||||
.load(Ordering::Relaxed),
|
||||
peer_lookup_fallbacks: inner.peer_lookup_fallbacks.load(Ordering::Relaxed),
|
||||
sample_candidate_queue_depth: inner
|
||||
.sample_candidate_queue_depth
|
||||
.load(Ordering::Relaxed),
|
||||
sample_candidate_queue_capacity: inner
|
||||
.sample_candidate_queue_capacity
|
||||
.load(Ordering::Relaxed),
|
||||
sample_candidates_routed: inner.sample_candidates_routed.load(Ordering::Relaxed),
|
||||
sample_candidates_fallback: inner.sample_candidates_fallback.load(Ordering::Relaxed),
|
||||
sample_infohashes_queries: inner.sample_infohashes_queries.load(Ordering::Relaxed),
|
||||
sample_infohashes_direct_queries: inner
|
||||
.sample_infohashes_direct_queries
|
||||
.load(Ordering::Relaxed),
|
||||
sample_infohashes_snapshot_queries: inner
|
||||
.sample_infohashes_snapshot_queries
|
||||
.load(Ordering::Relaxed),
|
||||
sample_infohashes_responses: inner.sample_infohashes_responses.load(Ordering::Relaxed),
|
||||
sample_infohashes_direct_responses: inner
|
||||
.sample_infohashes_direct_responses
|
||||
.load(Ordering::Relaxed),
|
||||
sample_infohashes_snapshot_responses: inner
|
||||
.sample_infohashes_snapshot_responses
|
||||
.load(Ordering::Relaxed),
|
||||
sample_infohashes_timeouts: inner.sample_infohashes_timeouts.load(Ordering::Relaxed),
|
||||
sample_infohashes_send_failures: inner
|
||||
.sample_infohashes_send_failures
|
||||
@@ -642,6 +711,18 @@ impl DhtRuntimeStats {
|
||||
.load(Ordering::Relaxed),
|
||||
metadata_peer_attempts: inner.metadata_peer_attempts.load(Ordering::Relaxed),
|
||||
metadata_peer_succeeded: inner.metadata_peer_succeeded.load(Ordering::Relaxed),
|
||||
metadata_success_from_announce: inner
|
||||
.metadata_success_from_announce
|
||||
.load(Ordering::Relaxed),
|
||||
metadata_success_from_sample_direct: inner
|
||||
.metadata_success_from_sample_direct
|
||||
.load(Ordering::Relaxed),
|
||||
metadata_success_from_sample_snapshot: inner
|
||||
.metadata_success_from_sample_snapshot
|
||||
.load(Ordering::Relaxed),
|
||||
metadata_success_from_active_lookup: inner
|
||||
.metadata_success_from_active_lookup
|
||||
.load(Ordering::Relaxed),
|
||||
metadata_peer_failed: inner.metadata_peer_failed.load(Ordering::Relaxed),
|
||||
metadata_peer_timeouts: inner.metadata_peer_timeouts.load(Ordering::Relaxed),
|
||||
metadata_connect_failed: inner.metadata_connect_failed.load(Ordering::Relaxed),
|
||||
@@ -890,16 +971,52 @@ impl DhtRuntimeStats {
|
||||
.fetch_add(1, Ordering::Relaxed);
|
||||
}
|
||||
|
||||
pub(crate) fn sample_query(&self) {
|
||||
pub(crate) fn configure_sample_candidate_queue(&self, capacity: usize) {
|
||||
self.inner
|
||||
.sample_infohashes_queries
|
||||
.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_response(&self) {
|
||||
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) {
|
||||
@@ -956,6 +1073,16 @@ impl DhtRuntimeStats {
|
||||
.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
|
||||
@@ -1272,8 +1399,12 @@ mod tests {
|
||||
writer.peer_lookup_output_dropped();
|
||||
writer.peer_lookup_preferred_succeeded();
|
||||
writer.peer_lookup_fallback();
|
||||
writer.sample_query();
|
||||
writer.sample_response();
|
||||
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();
|
||||
@@ -1283,6 +1414,10 @@ mod tests {
|
||||
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();
|
||||
@@ -1342,8 +1477,16 @@ mod tests {
|
||||
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,
|
||||
@@ -1353,6 +1496,10 @@ mod tests {
|
||||
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,
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
// 负责执行有界 BEP-51 采样并将准入后的 infohash 交给 Peer 查找
|
||||
|
||||
use crate::addr::is_valid_node_addr;
|
||||
use crate::budget::{RateBucket, SharedRateBudget};
|
||||
use crate::crawl_engine::CrawlEngine;
|
||||
use crate::krpc::{encode_sample_infohashes_query, for_each_response_node};
|
||||
@@ -8,7 +9,7 @@ use crate::peer_lookup::{PeerLookupHandle, PeerLookupRequest};
|
||||
use crate::protocol::DhtResponse;
|
||||
use crate::routing_snapshot::RoutingSnapshot;
|
||||
use crate::runtime_stats::DhtRuntimeStats;
|
||||
use crate::types::{NetMode, NodeTuple, SampleInfohashesOptions};
|
||||
use crate::types::{DiscoverySource, NetMode, NodeTuple, SampleInfohashesOptions};
|
||||
use ahash::{AHashMap, AHashSet};
|
||||
use arc_swap::{ArcSwap, ArcSwapOption};
|
||||
use bytes::BytesMut;
|
||||
@@ -39,6 +40,7 @@ struct PendingKey {
|
||||
#[derive(Debug, Clone, Copy)]
|
||||
struct PendingRequest {
|
||||
deadline: Instant,
|
||||
source: DiscoverySource,
|
||||
}
|
||||
|
||||
struct SampleResponse {
|
||||
@@ -49,6 +51,7 @@ struct SampleResponse {
|
||||
|
||||
struct SampleAdmissionBatch {
|
||||
preferred_node: NodeTuple,
|
||||
source: DiscoverySource,
|
||||
hashes: Vec<[u8; 20]>,
|
||||
}
|
||||
|
||||
@@ -65,6 +68,73 @@ pub(crate) struct SampleInfohashesHandle {
|
||||
runtime_stats: DhtRuntimeStats,
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
pub(crate) struct SampleCandidateRouter {
|
||||
sender: mpsc::Sender<NodeTuple>,
|
||||
sample_percent: u8,
|
||||
runtime_stats: DhtRuntimeStats,
|
||||
}
|
||||
|
||||
impl SampleCandidateRouter {
|
||||
pub(crate) fn route(&self, node: NodeTuple) -> bool {
|
||||
if self.sample_percent == 0
|
||||
|| !is_valid_node_addr(&node.addr)
|
||||
|| stable_address_bucket(node.addr) >= self.sample_percent
|
||||
{
|
||||
return false;
|
||||
}
|
||||
match self.sender.try_send(node) {
|
||||
Ok(()) => {
|
||||
self.runtime_stats.sample_candidate_routed();
|
||||
self.runtime_stats.set_sample_candidate_queue_depth(
|
||||
self.sender
|
||||
.max_capacity()
|
||||
.saturating_sub(self.sender.capacity()),
|
||||
);
|
||||
true
|
||||
}
|
||||
Err(mpsc::error::TrySendError::Full(_)) => {
|
||||
self.runtime_stats.sample_candidate_fallback();
|
||||
self.runtime_stats
|
||||
.set_sample_candidate_queue_depth(self.sender.max_capacity());
|
||||
false
|
||||
}
|
||||
Err(mpsc::error::TrySendError::Closed(_)) => false,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn sample_candidate_lane(
|
||||
options: &SampleInfohashesOptions,
|
||||
runtime_stats: DhtRuntimeStats,
|
||||
) -> (SampleCandidateRouter, mpsc::Receiver<NodeTuple>) {
|
||||
let capacity = options.candidate_queue_capacity.max(1);
|
||||
let (sender, receiver) = mpsc::channel(capacity);
|
||||
runtime_stats.configure_sample_candidate_queue(capacity);
|
||||
(
|
||||
SampleCandidateRouter {
|
||||
sender,
|
||||
sample_percent: options.new_node_sample_percent.min(100),
|
||||
runtime_stats,
|
||||
},
|
||||
receiver,
|
||||
)
|
||||
}
|
||||
|
||||
fn stable_address_bucket(addr: SocketAddr) -> u8 {
|
||||
let mut hash = 0xcbf2_9ce4_8422_2325u64;
|
||||
let mut mix = |byte: u8| {
|
||||
hash ^= u64::from(byte);
|
||||
hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
|
||||
};
|
||||
match addr.ip() {
|
||||
std::net::IpAddr::V4(ip) => ip.octets().into_iter().for_each(&mut mix),
|
||||
std::net::IpAddr::V6(ip) => ip.octets().into_iter().for_each(&mut mix),
|
||||
}
|
||||
addr.port().to_be_bytes().into_iter().for_each(&mut mix);
|
||||
(hash % 100) as u8
|
||||
}
|
||||
|
||||
impl SampleInfohashesHandle {
|
||||
pub(crate) fn route_response(
|
||||
&self,
|
||||
@@ -101,15 +171,27 @@ pub(crate) struct SampleInfohashesRuntime {
|
||||
pub(crate) shutdown: CancellationToken,
|
||||
}
|
||||
|
||||
pub(crate) struct SampleInfohashesInputs<'a> {
|
||||
pub(crate) sockets: &'a std::collections::HashMap<SocketAddr, Arc<UdpSocket>>,
|
||||
pub(crate) snapshot: Arc<ArcSwap<RoutingSnapshot>>,
|
||||
pub(crate) candidate_rx: mpsc::Receiver<NodeTuple>,
|
||||
pub(crate) crawl_engine: Arc<CrawlEngine>,
|
||||
pub(crate) peer_lookup: PeerLookupHandle,
|
||||
}
|
||||
|
||||
pub(crate) fn spawn_sample_infohashes(
|
||||
netmode: NetMode,
|
||||
local_id: [u8; 20],
|
||||
sockets: &std::collections::HashMap<SocketAddr, Arc<UdpSocket>>,
|
||||
snapshot: Arc<ArcSwap<RoutingSnapshot>>,
|
||||
crawl_engine: Arc<CrawlEngine>,
|
||||
peer_lookup: PeerLookupHandle,
|
||||
inputs: SampleInfohashesInputs<'_>,
|
||||
runtime: SampleInfohashesRuntime,
|
||||
) -> SampleInfohashesHandle {
|
||||
let SampleInfohashesInputs {
|
||||
sockets,
|
||||
snapshot,
|
||||
candidate_rx,
|
||||
crawl_engine,
|
||||
peer_lookup,
|
||||
} = inputs;
|
||||
let SampleInfohashesRuntime {
|
||||
options,
|
||||
stats,
|
||||
@@ -140,6 +222,8 @@ pub(crate) fn spawn_sample_infohashes(
|
||||
socket_v4,
|
||||
socket_v6,
|
||||
snapshot,
|
||||
candidate_rx,
|
||||
direct_only: options.new_node_sample_percent > 0,
|
||||
crawl_engine,
|
||||
admission_tx,
|
||||
response_rx,
|
||||
@@ -177,6 +261,8 @@ struct SampleInfohashesActor {
|
||||
socket_v4: Option<Arc<UdpSocket>>,
|
||||
socket_v6: Option<Arc<UdpSocket>>,
|
||||
snapshot: Arc<ArcSwap<RoutingSnapshot>>,
|
||||
candidate_rx: mpsc::Receiver<NodeTuple>,
|
||||
direct_only: bool,
|
||||
crawl_engine: Arc<CrawlEngine>,
|
||||
admission_tx: mpsc::Sender<SampleAdmissionBatch>,
|
||||
response_rx: mpsc::Receiver<SampleResponse>,
|
||||
@@ -225,6 +311,7 @@ fn spawn_sample_admission(
|
||||
info_hash,
|
||||
preferred_node: Some(batch.preferred_node),
|
||||
allow_iterative_fallback: fallback_to_iterative,
|
||||
source: batch.source,
|
||||
completion: None,
|
||||
};
|
||||
let sent = tokio::select! {
|
||||
@@ -259,6 +346,7 @@ impl SampleInfohashesActor {
|
||||
}
|
||||
}
|
||||
}
|
||||
self.runtime_stats.set_sample_candidate_queue_depth(0);
|
||||
}
|
||||
|
||||
async fn dispatch(&mut self, now: Instant) {
|
||||
@@ -283,6 +371,23 @@ impl SampleInfohashesActor {
|
||||
.load()
|
||||
.random_nodes((budget * 8).max(64), filter_ipv6);
|
||||
let mut sent_count = 0usize;
|
||||
while sent_count < budget {
|
||||
let Some(node) = self.next_direct_candidate(now) else {
|
||||
break;
|
||||
};
|
||||
if self
|
||||
.send_query(node, DiscoverySource::SampleDirect, now)
|
||||
.await
|
||||
{
|
||||
sent_count += 1;
|
||||
}
|
||||
}
|
||||
if self.direct_only {
|
||||
if sent_count < budget {
|
||||
self.query_budget.refund(budget - sent_count);
|
||||
}
|
||||
return;
|
||||
}
|
||||
for node in candidates {
|
||||
if sent_count >= budget {
|
||||
break;
|
||||
@@ -295,7 +400,10 @@ impl SampleInfohashesActor {
|
||||
{
|
||||
continue;
|
||||
}
|
||||
if self.send_query(node, now).await {
|
||||
if self
|
||||
.send_query(node, DiscoverySource::SampleSnapshot, now)
|
||||
.await
|
||||
{
|
||||
sent_count += 1;
|
||||
}
|
||||
}
|
||||
@@ -304,7 +412,25 @@ impl SampleInfohashesActor {
|
||||
}
|
||||
}
|
||||
|
||||
async fn send_query(&mut self, node: NodeTuple, now: Instant) -> bool {
|
||||
fn next_direct_candidate(&mut self, now: Instant) -> Option<NodeTuple> {
|
||||
for _ in 0..64 {
|
||||
let node = self.candidate_rx.try_recv().ok()?;
|
||||
self.runtime_stats
|
||||
.set_sample_candidate_queue_depth(self.candidate_rx.len());
|
||||
if self.pending_addrs.contains(&node.addr)
|
||||
|| self
|
||||
.next_allowed
|
||||
.get(&node.addr)
|
||||
.is_some_and(|deadline| *deadline > now)
|
||||
{
|
||||
continue;
|
||||
}
|
||||
return Some(node);
|
||||
}
|
||||
None
|
||||
}
|
||||
|
||||
async fn send_query(&mut self, node: NodeTuple, source: DiscoverySource, now: Instant) -> bool {
|
||||
let socket = if node.addr.is_ipv4() {
|
||||
self.socket_v4.clone()
|
||||
} else {
|
||||
@@ -331,11 +457,13 @@ impl SampleInfohashesActor {
|
||||
tid,
|
||||
};
|
||||
let deadline = now + self.request_timeout;
|
||||
self.pending.insert(key, PendingRequest { deadline });
|
||||
self.pending
|
||||
.insert(key, PendingRequest { deadline, source });
|
||||
self.pending_expiry.push_back((deadline, key));
|
||||
self.pending_addrs.insert(node.addr);
|
||||
self.runtime_stats.udp_sent(buffer.len());
|
||||
self.runtime_stats.sample_query();
|
||||
self.runtime_stats
|
||||
.sample_query(source == DiscoverySource::SampleDirect);
|
||||
#[cfg(feature = "metrics")]
|
||||
{
|
||||
counter!("dht_sample_infohashes_queries_total").increment(1);
|
||||
@@ -349,11 +477,12 @@ impl SampleInfohashesActor {
|
||||
addr: event.remote_addr,
|
||||
tid: event.tid,
|
||||
};
|
||||
if self.pending.remove(&key).is_none() {
|
||||
let Some(pending) = self.pending.remove(&key) else {
|
||||
return;
|
||||
}
|
||||
};
|
||||
self.pending_addrs.remove(&event.remote_addr);
|
||||
self.runtime_stats.sample_response();
|
||||
self.runtime_stats
|
||||
.sample_response(pending.source == DiscoverySource::SampleDirect);
|
||||
|
||||
for_each_response_node(&event.response, self.netmode, |node| {
|
||||
self.crawl_engine.route_discovered(node)
|
||||
@@ -395,6 +524,7 @@ impl SampleInfohashesActor {
|
||||
.admission_tx
|
||||
.try_send(SampleAdmissionBatch {
|
||||
preferred_node,
|
||||
source: pending.source,
|
||||
hashes,
|
||||
})
|
||||
.is_ok()
|
||||
@@ -483,6 +613,16 @@ mod tests {
|
||||
use super::*;
|
||||
use crate::protocol::DhtMessage;
|
||||
|
||||
fn node_for_lane(want_direct: bool, percent: u8) -> NodeTuple {
|
||||
(1..=u16::MAX)
|
||||
.map(|port| NodeTuple {
|
||||
id: [port as u8; 20],
|
||||
addr: SocketAddr::from(([8, 8, 8, 8], port)),
|
||||
})
|
||||
.find(|node| (stable_address_bucket(node.addr) < percent) == want_direct)
|
||||
.expect("both routing lanes have an address")
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn sample_transaction_ids_have_a_reserved_tag() {
|
||||
assert!(is_sample_infohashes_tid(&[
|
||||
@@ -508,6 +648,40 @@ mod tests {
|
||||
assert_eq!(response.samples.unwrap().len(), 40);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn new_node_routing_is_stable_and_bounded() {
|
||||
let options = SampleInfohashesOptions {
|
||||
new_node_sample_percent: 50,
|
||||
candidate_queue_capacity: 1,
|
||||
..SampleInfohashesOptions::default()
|
||||
};
|
||||
let stats = DhtRuntimeStats::default();
|
||||
let (router, mut receiver) = sample_candidate_lane(&options, stats.clone());
|
||||
let direct = node_for_lane(true, 50);
|
||||
let crawl = node_for_lane(false, 50);
|
||||
|
||||
assert!(router.route(direct));
|
||||
assert!(!router.route(crawl));
|
||||
assert!(!router.route(direct));
|
||||
assert_eq!(receiver.try_recv().unwrap(), direct);
|
||||
|
||||
let snapshot = stats.snapshot();
|
||||
assert_eq!(snapshot.sample_candidate_queue_capacity, 1);
|
||||
assert_eq!(snapshot.sample_candidates_routed, 1);
|
||||
assert_eq!(snapshot.sample_candidates_fallback, 1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn zero_percent_preserves_snapshot_only_sampling() {
|
||||
let options = SampleInfohashesOptions {
|
||||
new_node_sample_percent: 0,
|
||||
..SampleInfohashesOptions::default()
|
||||
};
|
||||
let (router, mut receiver) = sample_candidate_lane(&options, DhtRuntimeStats::default());
|
||||
assert!(!router.route(node_for_lane(false, 0)));
|
||||
assert!(receiver.try_recv().is_err());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn batch_admission_only_forwards_application_approved_hashes() {
|
||||
let (batch_tx, batch_rx) = mpsc::channel(1);
|
||||
@@ -533,6 +707,7 @@ mod tests {
|
||||
id: [9; 20],
|
||||
addr: "127.0.0.1:6881".parse().unwrap(),
|
||||
},
|
||||
source: DiscoverySource::SampleDirect,
|
||||
hashes: vec![[1; 20], [2; 20]],
|
||||
})
|
||||
.await
|
||||
@@ -544,6 +719,7 @@ mod tests {
|
||||
.unwrap();
|
||||
assert_eq!(request.info_hash, [2; 20]);
|
||||
assert!(!request.allow_iterative_fallback);
|
||||
assert_eq!(request.source, DiscoverySource::SampleDirect);
|
||||
let snapshot = stats.snapshot();
|
||||
assert_eq!(snapshot.sample_infohashes_hashes_filtered, 1);
|
||||
assert_eq!(snapshot.sample_infohashes_hashes_discovered, 1);
|
||||
|
||||
@@ -6,7 +6,9 @@ use crate::peer_lookup::PeerLookupRequest;
|
||||
use crate::runtime_stats::DhtRuntimeLimits;
|
||||
use crate::runtime_stats::DhtRuntimeStats;
|
||||
use crate::server::HashDiscovered;
|
||||
use crate::types::{MetadataFetchCompletion, MetadataFetchCompletionStatus, TorrentInfo};
|
||||
use crate::types::{
|
||||
DiscoverySource, MetadataFetchCompletion, MetadataFetchCompletionStatus, TorrentInfo,
|
||||
};
|
||||
use arc_swap::ArcSwapOption;
|
||||
#[cfg(feature = "metrics")]
|
||||
use metrics::{counter, gauge, histogram};
|
||||
@@ -69,6 +71,7 @@ pub(crate) struct MetadataSchedulerRuntime {
|
||||
#[derive(Debug, Clone, Copy)]
|
||||
struct PeerCandidate {
|
||||
addr: SocketAddr,
|
||||
source: DiscoverySource,
|
||||
discovered_at: Instant,
|
||||
}
|
||||
|
||||
@@ -148,6 +151,7 @@ impl PendingHashQueue {
|
||||
entry.peers.retain(|peer| peer.addr != event.peer_addr);
|
||||
entry.peers.push(PeerCandidate {
|
||||
addr: event.peer_addr,
|
||||
source: event.source,
|
||||
discovered_at,
|
||||
});
|
||||
entry
|
||||
@@ -184,6 +188,7 @@ impl PendingHashQueue {
|
||||
info_hash,
|
||||
peers: vec![PeerCandidate {
|
||||
addr: event.peer_addr,
|
||||
source: event.source,
|
||||
discovered_at: event.discovered_at,
|
||||
}],
|
||||
queued_at: now,
|
||||
@@ -450,6 +455,7 @@ impl MetadataScheduler {
|
||||
} else {
|
||||
let peer = PeerCandidate {
|
||||
addr: hash.peer_addr,
|
||||
source: hash.source,
|
||||
discovered_at: hash.discovered_at,
|
||||
};
|
||||
match state.peer_tx.try_send(peer) {
|
||||
@@ -690,6 +696,7 @@ impl MetadataScheduler {
|
||||
let attempts = attempts.load(Ordering::Relaxed);
|
||||
|
||||
if let Some((peer, (name, total_size, files, piece_length))) = fetched {
|
||||
runtime_stats.metadata_success_from(peer.source);
|
||||
let metadata = TorrentInfo {
|
||||
info_hash: job.info_hash.clone(),
|
||||
name,
|
||||
@@ -884,6 +891,7 @@ mod tests {
|
||||
HashDiscovered {
|
||||
info_hash: hash.to_string(),
|
||||
peer_addr: SocketAddr::from(([127, 0, 0, 1], port)),
|
||||
source: DiscoverySource::AnnouncePeer,
|
||||
discovered_at,
|
||||
}
|
||||
}
|
||||
@@ -1061,14 +1069,17 @@ mod tests {
|
||||
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,
|
||||
},
|
||||
];
|
||||
@@ -1134,10 +1145,12 @@ mod tests {
|
||||
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);
|
||||
@@ -1212,6 +1225,20 @@ mod tests {
|
||||
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();
|
||||
|
||||
@@ -17,15 +17,18 @@ use crate::response_limiter::{
|
||||
};
|
||||
use crate::runtime_stats::{DhtRuntimeLimits, DhtRuntimeStats};
|
||||
use crate::sample_infohashes::{
|
||||
SampleHashAdmissionCallback, SampleInfohashesHandle, SampleInfohashesRuntime,
|
||||
is_sample_infohashes_tid, spawn_sample_infohashes,
|
||||
SampleHashAdmissionCallback, SampleInfohashesHandle, SampleInfohashesInputs,
|
||||
SampleInfohashesRuntime, is_sample_infohashes_tid, sample_candidate_lane,
|
||||
spawn_sample_infohashes,
|
||||
};
|
||||
use crate::scheduler::{
|
||||
MetadataCompletionCallback, MetadataFetchCallback, MetadataScheduler,
|
||||
MetadataSchedulerCallbacks, MetadataSchedulerLimits, MetadataSchedulerRuntime,
|
||||
TorrentAckCallback,
|
||||
};
|
||||
use crate::types::{DHTOptions, MetadataFetchCompletion, NetMode, NodeTuple, TorrentInfo};
|
||||
use crate::types::{
|
||||
DHTOptions, DiscoverySource, MetadataFetchCompletion, NetMode, NodeTuple, TorrentInfo,
|
||||
};
|
||||
use crate::udp_buffer::UdpBufferPool;
|
||||
use crate::udp_ingress::{WorkerHandle, spawn_udp_listener};
|
||||
use arc_swap::ArcSwapOption;
|
||||
@@ -64,6 +67,8 @@ pub struct HashDiscovered {
|
||||
pub info_hash: String,
|
||||
/// Peer endpoint derived from announce `port`/`implied_port`.
|
||||
pub peer_addr: SocketAddr,
|
||||
/// Network path that supplied this Peer candidate.
|
||||
pub source: DiscoverySource,
|
||||
/// Monotonic discovery time used for freshness and queue ordering.
|
||||
pub discovered_at: std::time::Instant,
|
||||
}
|
||||
@@ -184,10 +189,13 @@ impl DHTServer {
|
||||
options.outbound_query_burst,
|
||||
true,
|
||||
);
|
||||
let (sample_candidate_router, sample_candidate_rx) =
|
||||
sample_candidate_lane(&options.sample_infohashes, runtime_stats.clone());
|
||||
let crawl_engine = Arc::new(CrawlEngine::new(
|
||||
crawl_config.clone(),
|
||||
runtime_stats.clone(),
|
||||
outbound_query_budget.clone(),
|
||||
Some(sample_candidate_router),
|
||||
));
|
||||
let peer_lookup = spawn_peer_lookup(
|
||||
options.netmode,
|
||||
@@ -205,10 +213,13 @@ impl DHTServer {
|
||||
let sample_infohashes = spawn_sample_infohashes(
|
||||
options.netmode,
|
||||
node_id,
|
||||
&sockets_by_bind_addr,
|
||||
crawl_engine.snapshot.clone(),
|
||||
crawl_engine.clone(),
|
||||
peer_lookup.clone(),
|
||||
SampleInfohashesInputs {
|
||||
sockets: &sockets_by_bind_addr,
|
||||
snapshot: crawl_engine.snapshot.clone(),
|
||||
candidate_rx: sample_candidate_rx,
|
||||
crawl_engine: crawl_engine.clone(),
|
||||
peer_lookup: peer_lookup.clone(),
|
||||
},
|
||||
SampleInfohashesRuntime {
|
||||
options: options.sample_infohashes.clone(),
|
||||
stats: runtime_stats.clone(),
|
||||
@@ -687,6 +698,7 @@ impl DHTServer {
|
||||
let event = HashDiscovered {
|
||||
info_hash: hash_hex,
|
||||
peer_addr: SocketAddr::new(addr.ip(), port),
|
||||
source: DiscoverySource::AnnouncePeer,
|
||||
discovered_at: std::time::Instant::now(),
|
||||
};
|
||||
|
||||
|
||||
@@ -161,6 +161,19 @@ pub struct PeerLookupOptions {
|
||||
pub max_active_lookups: usize,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
/// Network path that supplied a Peer candidate for Metadata download.
|
||||
pub enum DiscoverySource {
|
||||
/// A remote node announced itself as a Peer.
|
||||
AnnouncePeer,
|
||||
/// A Peer was found for a Hash sampled directly from a newly discovered node.
|
||||
SampleDirect,
|
||||
/// A Peer was found for a Hash sampled from the responsive-node snapshot.
|
||||
SampleSnapshot,
|
||||
/// A Peer was found by the Metadata scheduler's delayed active lookup.
|
||||
ActiveLookup,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
/// Terminal result of one bounded active `get_peers` lookup.
|
||||
pub struct PeerLookupResult {
|
||||
@@ -179,6 +192,10 @@ pub struct SampleInfohashesOptions {
|
||||
pub burst: u32,
|
||||
/// Maximum outstanding BEP-51 requests.
|
||||
pub max_in_flight: usize,
|
||||
/// Percentage of newly discovered node addresses routed directly to BEP-51 sampling.
|
||||
pub new_node_sample_percent: u8,
|
||||
/// Capacity of the bounded newly discovered node sampling lane.
|
||||
pub candidate_queue_capacity: usize,
|
||||
/// Whether a sampled Hash should fall back to iterative get_peers after the sampling node
|
||||
/// returns no Peer.
|
||||
pub fallback_to_iterative: bool,
|
||||
@@ -336,6 +353,8 @@ impl Default for SampleInfohashesOptions {
|
||||
max_queries_per_second: 1,
|
||||
burst: 1,
|
||||
max_in_flight: 4,
|
||||
new_node_sample_percent: 0,
|
||||
candidate_queue_capacity: 4_096,
|
||||
fallback_to_iterative: true,
|
||||
request_timeout_millis: 1_500,
|
||||
unsupported_backoff_secs: 300,
|
||||
|
||||
@@ -25,6 +25,8 @@ metadata_workers = 400
|
||||
metadata_connects_per_second = 400
|
||||
sample_queries_per_second = 60
|
||||
sample_max_in_flight = 100
|
||||
sample_new_node_percent = 50
|
||||
sample_candidate_queue_capacity = 8192
|
||||
sample_fallback_to_iterative = false
|
||||
peer_lookups_per_second = 200
|
||||
peer_lookup_max_active = 200
|
||||
|
||||
@@ -51,6 +51,8 @@ cargo run -p dht-search -- --data-dir D:\data\dht-search --run-duration-secs 360
|
||||
| `peer_lookup_max_active` | `200` | 同时运行的 Peer 查找数量 |
|
||||
| `sample_queries_per_second` | `60` | BEP-51 采样查询速率 |
|
||||
| `sample_max_in_flight` | `100` | 同时等待响应的 BEP-51 采样请求数量 |
|
||||
| `sample_new_node_percent` | `50` | 按节点地址稳定分流到直接采样通道的比例 |
|
||||
| `sample_candidate_queue_capacity` | `8192` | 新发现节点直接采样通道的有界容量 |
|
||||
| `sample_fallback_to_iterative` | `false` | 采样节点没有返回 Peer 时快速结束以扩大覆盖面 |
|
||||
| `metadata_workers` | `400` | 同时处理的 Metadata 任务数量 |
|
||||
| `metadata_connects_per_second` | `400` | 每秒真正开始的 Peer TCP 连接硬上限 |
|
||||
@@ -113,6 +115,8 @@ Invoke-RestMethod http://127.0.0.1:8080/stats | ConvertTo-Json -Depth 5
|
||||
|
||||
BEP-51 返回的 infohash 会先进入有界批量准入队列并由 RocksDB 精确判断
|
||||
|
||||
新发现节点按地址稳定分流 同一地址只进入直接采样或 `find_node` 通道 直接采样队列满时会安全回退到抓取池 `/stats` 会分别报告候选队列直接采样请求响应 Hash 去重丢弃以及成功 Metadata 的网络来源
|
||||
|
||||
已有 infohash 只更新最后发现时间和发现次数不会再次执行 Peer Lookup
|
||||
|
||||
未知 infohash 首先只向返回样本的 DHT 节点查询一次 `get_peers` 只有单点查询没有返回 Peer 时才降级为有限迭代查找
|
||||
|
||||
@@ -50,12 +50,28 @@ pub(crate) async fn stats(State(state): State<ApiState>) -> Json<StatsResponse>
|
||||
peer_lookup_queries: dht.peer_lookup_queries,
|
||||
peer_lookup_preferred_succeeded: dht.peer_lookup_preferred_succeeded,
|
||||
peer_lookup_fallbacks: dht.peer_lookup_fallbacks,
|
||||
sample_candidate_queue: dht.sample_candidate_queue_depth,
|
||||
sample_candidate_queue_capacity: dht.sample_candidate_queue_capacity,
|
||||
sample_candidates_routed: dht.sample_candidates_routed,
|
||||
sample_candidates_fallback: dht.sample_candidates_fallback,
|
||||
sample_queries: dht.sample_infohashes_queries,
|
||||
sample_direct_queries: dht.sample_infohashes_direct_queries,
|
||||
sample_snapshot_queries: dht.sample_infohashes_snapshot_queries,
|
||||
sample_responses: dht.sample_infohashes_responses,
|
||||
sample_direct_responses: dht.sample_infohashes_direct_responses,
|
||||
sample_snapshot_responses: dht.sample_infohashes_snapshot_responses,
|
||||
sample_timeouts: dht.sample_infohashes_timeouts,
|
||||
sampled_hashes: dht.sample_infohashes_hashes_discovered,
|
||||
sampled_hashes_filtered: dht.sample_infohashes_hashes_filtered,
|
||||
sampled_hashes_duplicate: dht.sample_infohashes_hashes_duplicate,
|
||||
sampled_hashes_dropped: dht.sample_infohashes_hashes_dropped,
|
||||
metadata_peer_attempts: dht.metadata_peer_attempts,
|
||||
metadata_in_flight: dht.metadata_in_flight,
|
||||
metadata_ok: dht.metadata_peer_succeeded,
|
||||
metadata_ok_from_announce: dht.metadata_success_from_announce,
|
||||
metadata_ok_from_sample_direct: dht.metadata_success_from_sample_direct,
|
||||
metadata_ok_from_sample_snapshot: dht.metadata_success_from_sample_snapshot,
|
||||
metadata_ok_from_active_lookup: dht.metadata_success_from_active_lookup,
|
||||
metadata_failed: dht.metadata_peer_failed,
|
||||
metadata_filtered: observability
|
||||
.metadata_failure_size_limit
|
||||
|
||||
@@ -22,12 +22,28 @@ pub(crate) struct StatsResponse {
|
||||
pub(crate) peer_lookup_queries: u64,
|
||||
pub(crate) peer_lookup_preferred_succeeded: u64,
|
||||
pub(crate) peer_lookup_fallbacks: u64,
|
||||
pub(crate) sample_candidate_queue: usize,
|
||||
pub(crate) sample_candidate_queue_capacity: usize,
|
||||
pub(crate) sample_candidates_routed: u64,
|
||||
pub(crate) sample_candidates_fallback: u64,
|
||||
pub(crate) sample_queries: u64,
|
||||
pub(crate) sample_direct_queries: u64,
|
||||
pub(crate) sample_snapshot_queries: u64,
|
||||
pub(crate) sample_responses: u64,
|
||||
pub(crate) sample_direct_responses: u64,
|
||||
pub(crate) sample_snapshot_responses: u64,
|
||||
pub(crate) sample_timeouts: u64,
|
||||
pub(crate) sampled_hashes: u64,
|
||||
pub(crate) sampled_hashes_filtered: u64,
|
||||
pub(crate) sampled_hashes_duplicate: u64,
|
||||
pub(crate) sampled_hashes_dropped: u64,
|
||||
pub(crate) metadata_peer_attempts: u64,
|
||||
pub(crate) metadata_in_flight: usize,
|
||||
pub(crate) metadata_ok: u64,
|
||||
pub(crate) metadata_ok_from_announce: u64,
|
||||
pub(crate) metadata_ok_from_sample_direct: u64,
|
||||
pub(crate) metadata_ok_from_sample_snapshot: u64,
|
||||
pub(crate) metadata_ok_from_active_lookup: u64,
|
||||
pub(crate) metadata_failed: u64,
|
||||
pub(crate) metadata_filtered: u64,
|
||||
pub(crate) metadata_filtered_too_large: u64,
|
||||
|
||||
@@ -52,6 +52,8 @@ pub(crate) struct DhtConfig {
|
||||
pub(crate) metadata_connects_per_second: u32,
|
||||
pub(crate) sample_queries_per_second: u32,
|
||||
pub(crate) sample_max_in_flight: usize,
|
||||
pub(crate) sample_new_node_percent: u8,
|
||||
pub(crate) sample_candidate_queue_capacity: usize,
|
||||
pub(crate) sample_fallback_to_iterative: bool,
|
||||
pub(crate) peer_lookups_per_second: u32,
|
||||
pub(crate) peer_lookup_max_active: usize,
|
||||
@@ -155,6 +157,8 @@ impl AppConfig {
|
||||
max_queries_per_second: self.dht.sample_queries_per_second,
|
||||
burst: self.dht.sample_queries_per_second.max(1),
|
||||
max_in_flight: self.dht.sample_max_in_flight,
|
||||
new_node_sample_percent: self.dht.sample_new_node_percent,
|
||||
candidate_queue_capacity: self.dht.sample_candidate_queue_capacity,
|
||||
fallback_to_iterative: self.dht.sample_fallback_to_iterative,
|
||||
..defaults.sample_infohashes
|
||||
},
|
||||
@@ -241,11 +245,17 @@ impl AppConfig {
|
||||
|| self.dht.outbound_query_burst == 0
|
||||
|| self.dht.metadata_connects_per_second == 0
|
||||
|| self.dht.sample_max_in_flight == 0
|
||||
|| self.dht.sample_candidate_queue_capacity == 0
|
||||
|| self.dht.peer_lookup_max_active == 0
|
||||
|| self.dht.find_node_max_in_flight == 0
|
||||
{
|
||||
return Err(AppError::Config("网络速率和并发上限必须大于零".to_owned()));
|
||||
}
|
||||
if self.dht.sample_new_node_percent > 100 {
|
||||
return Err(AppError::Config(
|
||||
"sample_new_node_percent 必须在 0 到 100 之间".to_owned(),
|
||||
));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
@@ -293,6 +303,8 @@ impl Default for DhtConfig {
|
||||
metadata_connects_per_second: 400,
|
||||
sample_queries_per_second: 60,
|
||||
sample_max_in_flight: 100,
|
||||
sample_new_node_percent: 50,
|
||||
sample_candidate_queue_capacity: 8_192,
|
||||
sample_fallback_to_iterative: false,
|
||||
peer_lookups_per_second: 200,
|
||||
peer_lookup_max_active: 200,
|
||||
@@ -368,6 +380,8 @@ mod tests {
|
||||
assert_eq!(options.max_outbound_queries_per_second, 1_000);
|
||||
assert_eq!(options.sample_infohashes.max_queries_per_second, 60);
|
||||
assert_eq!(options.sample_infohashes.max_in_flight, 100);
|
||||
assert_eq!(options.sample_infohashes.new_node_sample_percent, 50);
|
||||
assert_eq!(options.sample_infohashes.candidate_queue_capacity, 8_192);
|
||||
assert!(!options.sample_infohashes.fallback_to_iterative);
|
||||
assert_eq!(options.peer_lookup.max_lookups_per_second, 200);
|
||||
assert_eq!(options.peer_lookup.max_active_lookups, 200);
|
||||
@@ -415,4 +429,11 @@ mod tests {
|
||||
config.metadata_limits.max_files = 0;
|
||||
assert!(matches!(config.validate(), Err(AppError::Config(_))));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn invalid_new_node_sample_percent_is_rejected() {
|
||||
let mut config = AppConfig::default();
|
||||
config.dht.sample_new_node_percent = 101;
|
||||
assert!(matches!(config.validate(), Err(AppError::Config(_))));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -44,12 +44,23 @@ pub(crate) async fn run(
|
||||
peer_lookup_preferred_succeeded = dht.peer_lookup_preferred_succeeded,
|
||||
peer_lookup_fallbacks = dht.peer_lookup_fallbacks,
|
||||
sample_queries = dht.sample_infohashes_queries,
|
||||
sample_direct_queries = dht.sample_infohashes_direct_queries,
|
||||
sample_direct_responses = dht.sample_infohashes_direct_responses,
|
||||
sample_candidate_queue = dht.sample_candidate_queue_depth,
|
||||
sample_candidates_routed = dht.sample_candidates_routed,
|
||||
sample_candidates_fallback = dht.sample_candidates_fallback,
|
||||
sampled_hashes = dht.sample_infohashes_hashes_discovered,
|
||||
sampled_hashes_filtered = dht.sample_infohashes_hashes_filtered,
|
||||
sampled_hashes_duplicate = dht.sample_infohashes_hashes_duplicate,
|
||||
sampled_hashes_dropped = dht.sample_infohashes_hashes_dropped,
|
||||
peers = dht.peer_lookup_peers_found,
|
||||
metadata_connects_per_second,
|
||||
metadata_in_flight = dht.metadata_in_flight,
|
||||
metadata_ok = dht.metadata_peer_succeeded,
|
||||
metadata_ok_from_announce = dht.metadata_success_from_announce,
|
||||
metadata_ok_from_sample_direct = dht.metadata_success_from_sample_direct,
|
||||
metadata_ok_from_sample_snapshot = dht.metadata_success_from_sample_snapshot,
|
||||
metadata_ok_from_active_lookup = dht.metadata_success_from_active_lookup,
|
||||
metadata_failed = dht.metadata_peer_failed,
|
||||
metadata_filtered = observability
|
||||
.metadata_failure_size_limit
|
||||
|
||||
Reference in New Issue
Block a user