feat: add DHT crawler and search application with UDP buffer management
- Implemented DHT types and options for configuration. - Created UDP buffer pool for efficient packet handling. - Developed UDP ingress handling with worker management. - Established a new DHT search application with modular API structure. - Defined domain models for fingerprinting and torrent metadata. - Integrated storage solutions using RocksDB for persistent data management. - Added telemetry for structured logging and metrics.
This commit is contained in:
@@ -18,3 +18,5 @@ torrents/
|
||||
|
||||
# 个人脚本(不提交到仓库)
|
||||
scripts/
|
||||
|
||||
opencodes/
|
||||
@@ -0,0 +1,111 @@
|
||||
# DHT 元数据搜索服务开发约定
|
||||
|
||||
## 项目目标
|
||||
|
||||
本项目用于持续或间歇地从 BitTorrent DHT 网络发现 infohash 获取元数据并提供本地全文搜索和高性能过滤能力
|
||||
|
||||
基础 DHT 协议和抓取能力保留在 `dht-crawler` 中
|
||||
|
||||
面向最终用户运行的服务代码统一放在 `dht-search` 中
|
||||
|
||||
## 技术方案
|
||||
|
||||
- Tokio 负责异步任务调度网络任务和有界队列
|
||||
- RocksDB 负责权威数据持久化精确去重抓取状态和索引状态
|
||||
- Tantivy 负责可重建的全文搜索过滤排序和结果聚合
|
||||
- Axum 负责 HTTP 搜索接口详情接口和运行状态接口
|
||||
- BLAKE3 负责计算规范化内容结构指纹
|
||||
- Serde 负责配置领域对象和接口数据的序列化
|
||||
- Tracing 负责结构化日志和故障定位
|
||||
|
||||
RocksDB 是唯一权威数据源
|
||||
|
||||
Tantivy 索引必须能够从 RocksDB 完整重建
|
||||
|
||||
## 数据处理流程
|
||||
|
||||
1. DHT 采集器发现 infohash
|
||||
2. 内存近期缓存过滤高频重复
|
||||
3. RocksDB 精确判断 infohash 是否已处理
|
||||
4. 未处理的 infohash 进入有界 Metadata 下载队列
|
||||
5. Metadata 完成校验和规范化后计算内容指纹
|
||||
6. 使用 RocksDB WriteBatch 原子保存元数据去重映射和待索引状态
|
||||
7. 后台索引任务批量写入 Tantivy
|
||||
8. Tantivy 提交成功后将记录状态更新为已索引
|
||||
9. Axum 只通过搜索和存储抽象读取数据
|
||||
|
||||
所有任务队列必须有明确容量并在队列满时产生背压
|
||||
|
||||
禁止通过无限队列维持表面吞吐
|
||||
|
||||
## 去重规则
|
||||
|
||||
第一层以 infohash 做精确去重并阻止相同 Metadata 被重复下载
|
||||
|
||||
第二层根据规范化文件路径和文件大小计算 content key 将不同 infohash 的相同内容聚合展示
|
||||
|
||||
模糊名称相似度只用于搜索结果聚合不得直接删除数据
|
||||
|
||||
Bloom Filter 只能作为前置加速结构不得作为最终去重依据
|
||||
|
||||
再次发现已有 infohash 时只更新最后发现时间和发现次数
|
||||
|
||||
## 持久化和恢复
|
||||
|
||||
程序必须支持间歇运行和跨重启恢复
|
||||
|
||||
正常关闭时先停止接收新任务再排空或持久化队列最后提交搜索索引
|
||||
|
||||
异常退出后依靠 RocksDB WAL 恢复已提交数据
|
||||
|
||||
每条搜索文档必须具有待索引和已索引状态以便启动后补建索引
|
||||
|
||||
内存队列不得成为任何权威状态的唯一保存位置
|
||||
|
||||
数据目录必须通过配置指定且不得依赖当前工作目录
|
||||
|
||||
## 代码边界
|
||||
|
||||
- `crawler` 只负责协调 DHT 事件和 Metadata 下载
|
||||
- `domain` 只定义领域模型规范化规则和内容指纹
|
||||
- `storage` 只负责 RocksDB 数据布局原子写入查询和恢复
|
||||
- `search` 只负责 Tantivy schema 文档转换索引和查询
|
||||
- `api` 只负责 HTTP 协议参数校验和响应转换
|
||||
- `config` 只负责读取校验和暴露配置
|
||||
- `telemetry` 只负责日志指标和运行观测
|
||||
- `shutdown` 只负责关闭信号和优雅退出协调
|
||||
|
||||
模块之间通过明确的数据结构和 trait 通信不得跨层直接访问内部实现
|
||||
|
||||
## 资源约束
|
||||
|
||||
- Metadata 下载并发必须可配置
|
||||
- 单个 Metadata 大小文件数量和路径长度必须有限制
|
||||
- RocksDB 写入使用批处理并明确控制 block cache
|
||||
- Tantivy 写入使用批量提交并明确控制 IndexWriter 内存预算
|
||||
- 日志不得输出完整 Metadata 或大文件列表
|
||||
- 长期运行的集合必须有容量上限过期规则或磁盘持久化方案
|
||||
|
||||
## 开发规则
|
||||
|
||||
新业务代码写入 `dht-search`
|
||||
|
||||
`dht-crawler` 只接受可复用的 DHT 基础能力不得包含数据库搜索接口或部署逻辑
|
||||
|
||||
每个 Rust 文件顶部必须使用中文行注释描述该文件的功能边界且注释行尾不添加标点
|
||||
|
||||
新增行为必须包含与风险相称的测试
|
||||
|
||||
修改持久化 key 或编码格式时必须提供兼容或迁移方案
|
||||
|
||||
不得把 `opencodes` 下的参考项目纳入 workspace 或修改其内容
|
||||
|
||||
## 构建准备
|
||||
|
||||
应用骨架默认不启用 RocksDB 原生构建
|
||||
|
||||
实现存储层时通过 `rocksdb-storage` feature 启用 RocksDB
|
||||
|
||||
Windows 构建 RocksDB 前需要安装 LLVM 并确保 `LIBCLANG_PATH` 指向包含 `libclang.dll` 的目录
|
||||
|
||||
启用后的检查命令为 `cargo check -p dht-search --features rocksdb-storage`
|
||||
+13
-53
@@ -1,60 +1,21 @@
|
||||
[package]
|
||||
name = "dht-crawler"
|
||||
version = "0.2.1"
|
||||
[workspace]
|
||||
members = ["dht-crawler", "dht-search"]
|
||||
exclude = ["opencodes"]
|
||||
resolver = "3"
|
||||
|
||||
[workspace.package]
|
||||
edition = "2024"
|
||||
authors = ["桥下红药 <1121744186@qq.com>"]
|
||||
description = "高性能的 Rust DHT (Distributed Hash Table) 爬虫库 | A high-performance Rust DHT crawler library for fetching torrent information from the BitTorrent DHT network"
|
||||
authors = ["chuan <pchuan98@qq.com>"]
|
||||
license = "MIT"
|
||||
documentation = "https://docs.rs/dht-crawler"
|
||||
repository = "https://github.com/0xddy/dht-crawler"
|
||||
keywords = ["dht", "bittorrent", "crawler", "p2p", "torrent"]
|
||||
categories = ["network-programming", "asynchronous"]
|
||||
readme = "README.md"
|
||||
repository = "https://git.pchuan.top/tools/dht"
|
||||
|
||||
[lib]
|
||||
name = "dht_crawler"
|
||||
path = "src/lib.rs"
|
||||
crate-type = ["rlib"]
|
||||
|
||||
[dependencies]
|
||||
tokio = { version = "1.35", features = ["rt", "rt-multi-thread", "net", "sync", "time", "macros"] }
|
||||
tokio-util = { version = "0.7" }
|
||||
[workspace.dependencies]
|
||||
serde = { version = "1.0", features = ["derive"] }
|
||||
serde_bencode = "0.2"
|
||||
sha1 = "0.10"
|
||||
hex = "0.4"
|
||||
rand = "0.8"
|
||||
log = "0.4"
|
||||
thiserror = "1.0"
|
||||
socket2 = { version = "0.5", features = ["all"] }
|
||||
rbit = "0.2"
|
||||
bytes = "1.0"
|
||||
ahash = "0.8"
|
||||
serde_bytes = "0.11.19"
|
||||
metrics = { version = "0.24", optional = true }
|
||||
async-channel = "2.5.0"
|
||||
crossbeam-queue = "0.3"
|
||||
arc-swap = "1.7"
|
||||
|
||||
[dev-dependencies]
|
||||
tokio = { version = "1.35", features = ["signal"] }
|
||||
thiserror = "2.0.19"
|
||||
tokio = "1.53"
|
||||
tokio-util = "0.7"
|
||||
tracing = "0.1"
|
||||
tracing-subscriber = { version = "0.3", features = ["env-filter", "fmt", "tracing-log"] }
|
||||
mimalloc = "0.1"
|
||||
# 禁用默认 features,只启用 http-listener(不包含 hyper-rustls)
|
||||
# 这样可以避免引入 aws-lc-sys
|
||||
metrics-exporter-prometheus = { version = "0.18", default-features = false, features = ["http-listener"] }
|
||||
|
||||
|
||||
[features]
|
||||
default = []
|
||||
metrics = ["dep:metrics"]
|
||||
mimalloc = []
|
||||
|
||||
[[example]]
|
||||
name = "dht_crawler_example"
|
||||
path = "examples/main.rs"
|
||||
|
||||
tracing-subscriber = "0.3"
|
||||
|
||||
[profile.release]
|
||||
opt-level = 3
|
||||
@@ -66,4 +27,3 @@ strip = true
|
||||
[profile.dev]
|
||||
opt-level = 0
|
||||
debug = true
|
||||
|
||||
|
||||
@@ -1,264 +1,15 @@
|
||||
# dht-crawler
|
||||
# DHT 元数据搜索服务
|
||||
|
||||
[](https://crates.io/crates/dht-crawler)
|
||||
[](https://docs.rs/dht-crawler)
|
||||
[](LICENSE)
|
||||
这是一个用于发现持久化索引和搜索 BitTorrent DHT 元数据的 Rust workspace
|
||||
|
||||
基于 Rust 和 Tokio 的 BitTorrent DHT 爬虫库。它参与 BEP-5 DHT 网络,通过 BEP-51
|
||||
`sample_infohashes` 主动发现 InfoHash,也接收 `announce_peer`,并通过 BEP-9
|
||||
`ut_metadata` 获取、校验和解析 torrent 元数据。
|
||||
## 目录
|
||||
|
||||
`dht-crawler` 提供:
|
||||
|
||||
- IPv4、IPv6 和双栈 DHT;
|
||||
- 主动节点发现、BEP-51 InfoHash 采样与 `get_peers` 查询;
|
||||
- 有界、去重的 Metadata 下载队列;
|
||||
- InfoHash 过滤、异步准入、结果交付和完成通知;
|
||||
- 默认可用的运行时统计,以及可选的 `metrics` 集成。
|
||||
|
||||
## 安装
|
||||
|
||||
```bash
|
||||
cargo add dht-crawler
|
||||
cargo add tokio --features rt-multi-thread,macros,signal
|
||||
```text
|
||||
dht-search/ 最终运行的采集存储搜索和接口应用
|
||||
dht-crawler/ 可独立复用的 DHT 协议与 Metadata 获取基础库
|
||||
opencodes/ 不参与构建的参考项目
|
||||
```
|
||||
|
||||
或在 `Cargo.toml` 中添加:
|
||||
当前应用层仅完成目录和依赖准备尚未实现业务逻辑
|
||||
|
||||
```toml
|
||||
[dependencies]
|
||||
dht-crawler = "0.2"
|
||||
tokio = { version = "1", features = ["rt-multi-thread", "macros", "signal"] }
|
||||
```
|
||||
|
||||
## 快速开始
|
||||
|
||||
```rust
|
||||
use dht_crawler::prelude::*;
|
||||
|
||||
#[tokio::main]
|
||||
async fn main() -> Result<()> {
|
||||
let server = DHTServer::new(DHTOptions {
|
||||
port: 6881,
|
||||
netmode: NetMode::Ipv4Only,
|
||||
..Default::default()
|
||||
})
|
||||
.await?;
|
||||
|
||||
server.on_torrent(|torrent| {
|
||||
println!(
|
||||
"{} {} {}",
|
||||
torrent.info_hash,
|
||||
torrent.name,
|
||||
torrent.format_size()
|
||||
);
|
||||
});
|
||||
|
||||
server.on_error(|error| {
|
||||
eprintln!("DHT runtime error: {error}");
|
||||
});
|
||||
|
||||
let shutdown = server.clone();
|
||||
tokio::spawn(async move {
|
||||
if tokio::signal::ctrl_c().await.is_ok() {
|
||||
shutdown.shutdown();
|
||||
}
|
||||
});
|
||||
|
||||
// 一直运行,直到 shutdown() 被调用。
|
||||
server.start().await
|
||||
}
|
||||
```
|
||||
|
||||
运行仓库中的完整示例:
|
||||
|
||||
```bash
|
||||
cargo run --release --example dht_crawler_example
|
||||
```
|
||||
|
||||
## 核心 API
|
||||
|
||||
`DHTServer` 是主要入口:
|
||||
|
||||
| API | 用途 |
|
||||
|---|---|
|
||||
| `DHTServer::new(options)` | 校验配置、绑定 UDP Socket 并创建内部管道 |
|
||||
| `start().await` | 启动 DHT 与爬取任务,等待 `shutdown()` |
|
||||
| `shutdown()` | 停止 UDP、爬取和 Metadata 任务;可重复调用 |
|
||||
| `filter(callback)` | 在 InfoHash 进入队列前执行同步过滤 |
|
||||
| `on_metadata_fetch(callback)` | 在第一次 Peer 下载前执行异步准入 |
|
||||
| `on_torrent(callback)` | 接收已校验的 `TorrentInfo` |
|
||||
| `on_torrent_with_ack(callback)` | 接收结果并显式确认是否接受交付 |
|
||||
| `on_metadata_fetch_complete(callback)` | 接收已准入任务的最终状态 |
|
||||
| `on_error(callback)` | 接收运行期错误 |
|
||||
| `runtime_stats()` | 获取可复制的运行时统计句柄 |
|
||||
|
||||
同类回调重复注册时,新回调会替换旧回调。
|
||||
|
||||
### 过滤与准入
|
||||
|
||||
`filter` 是同步的早期过滤器,适合拦截已处理过的 InfoHash:
|
||||
|
||||
```rust
|
||||
server.filter(|info_hash| !already_exists(info_hash));
|
||||
```
|
||||
|
||||
`on_metadata_fetch` 是异步准入回调,在实际连接 Peer 前调用:
|
||||
|
||||
```rust
|
||||
server.on_metadata_fetch(|info_hash| async move {
|
||||
should_download(&info_hash).await
|
||||
});
|
||||
```
|
||||
|
||||
返回 `false` 会终止任务,不下载 Metadata,也不会触发 torrent 或 completion 回调。
|
||||
未注册准入回调时默认允许下载。
|
||||
|
||||
### 交付确认
|
||||
|
||||
不需要确认下游是否接收时使用 `on_torrent`。需要确认下游是否成功接收时使用
|
||||
`on_torrent_with_ack`:
|
||||
|
||||
```rust
|
||||
server.on_torrent_with_ack(|torrent| {
|
||||
output.try_send(torrent).is_ok()
|
||||
});
|
||||
|
||||
server.on_metadata_fetch_complete(|completion| {
|
||||
println!(
|
||||
"{}: {:?}, attempts={}",
|
||||
completion.info_hash,
|
||||
completion.status,
|
||||
completion.attempts
|
||||
);
|
||||
});
|
||||
```
|
||||
|
||||
完成状态:
|
||||
|
||||
| 状态 | 含义 |
|
||||
|---|---|
|
||||
| `Accepted` | Metadata 下载成功,结果已被回调接受 |
|
||||
| `FetchFailed` | 所有可用 Peer 尝试均失败 |
|
||||
| `DeliveryRejected` | Metadata 下载成功,但结果未被回调接受 |
|
||||
|
||||
`attempts` 只统计实际发起的 Peer 网络请求。异步准入拒绝不会产生 completion 事件。
|
||||
|
||||
## 配置
|
||||
|
||||
大多数调用方可以从 `DHTOptions::default()` 开始,只覆盖监听方式和容量限制:
|
||||
|
||||
```rust
|
||||
let options = DHTOptions {
|
||||
port: 6881,
|
||||
netmode: NetMode::DualStack,
|
||||
hash_queue_capacity: 20_000,
|
||||
metadata: MetadataOptions {
|
||||
timeout_secs: 5,
|
||||
max_queue_size: 20_000,
|
||||
max_worker_count: 256,
|
||||
..Default::default()
|
||||
},
|
||||
crawl: CrawlOptions {
|
||||
rate_limit: RateLimitOptions {
|
||||
max_find_node_rate_per_sec: 300,
|
||||
max_in_flight: 768,
|
||||
..Default::default()
|
||||
},
|
||||
..Default::default()
|
||||
},
|
||||
..Default::default()
|
||||
};
|
||||
```
|
||||
|
||||
配置分组:
|
||||
|
||||
| 类型 | 控制内容 |
|
||||
|---|---|
|
||||
| `DHTOptions` | 监听端口、网络模式和顶层队列 |
|
||||
| `MetadataOptions` | 下载超时、队列、并发和失败 Peer 缓存 |
|
||||
| `PeerLookupOptions` | 主动 `get_peers` 的速率与并发 |
|
||||
| `SampleInfohashesOptions` | BEP-51 采样速率、并发、超时、退避和 Hash 去重容量 |
|
||||
| `RateLimitOptions` | `find_node`、在途请求和 UDP 回复预算 |
|
||||
| `PoolOptions` | 节点池、最近探测记录和响应节点缓存 |
|
||||
| `BootstrapOptions` | Bootstrap 节点与失败退避 |
|
||||
| `TargetOptions` | 主动爬取目标生成策略 |
|
||||
| `SchedulerOptions` | 内部事件队列、批处理与快照限制 |
|
||||
|
||||
完整字段和默认值以 [docs.rs API 文档](https://docs.rs/dht-crawler) 为准。需要注意:
|
||||
|
||||
- `DHTOptions::default()` 使用 `Ipv4Only`;
|
||||
- `NetMode::DualStack` 会分别绑定 IPv4 和 IPv6 Socket;
|
||||
- `DHTServer::new()` 会立即在所有可用接口上绑定配置的 UDP 端口;
|
||||
- 空节点池默认每 30 秒重新尝试 Bootstrap,每轮最多使用 16 个已解析端点;
|
||||
- `MetadataOptions::timeout_secs` 是单个 Peer 尝试的端到端期限;
|
||||
- `PeerLookupOptions::max_lookups_per_second = 0` 会关闭主动 `get_peers`;
|
||||
- `SampleInfohashesOptions::max_queries_per_second = 0` 会关闭 BEP-51 主动采样;
|
||||
- Metadata 和爬取队列都是有界的,容量应与下游处理能力一起调整。
|
||||
|
||||
## 数据与运行语义
|
||||
|
||||
`TorrentInfo` 包含 `info_hash`、`magnet_link`、`name`、`total_size`、`files`、
|
||||
`piece_length`、`peers` 和 `timestamp`。只有通过 SHA1 校验并成功解析的 Metadata
|
||||
才会交付给 torrent 回调。
|
||||
|
||||
库使用有界队列控制内存占用。队列满或速率预算耗尽时,新事件可能被拒绝、淘汰或计入
|
||||
drop 指标。Metadata 队列按 InfoHash 去重,每个任务可尝试多个候选 Peer;连接超时和
|
||||
连接失败的 Peer 会被短期缓存,避免反复占用 worker。
|
||||
|
||||
`start()` 返回后,当前实例不能再次启动。如需重新运行,请创建新的 `DHTServer`。
|
||||
|
||||
## 可观测性
|
||||
|
||||
运行时快照无需启用 Cargo feature:
|
||||
|
||||
```rust
|
||||
let stats = server.runtime_stats();
|
||||
let snapshot = stats.snapshot();
|
||||
|
||||
println!(
|
||||
"nodes={} metadata={}/{} workers={}",
|
||||
snapshot.node_pool_size,
|
||||
snapshot.metadata_queue_depth,
|
||||
snapshot.metadata_queue_max,
|
||||
snapshot.metadata_in_flight,
|
||||
);
|
||||
```
|
||||
|
||||
`observability_snapshot()` 提供 UDP、查询、队列、Metadata 失败原因和固定桶直方图。
|
||||
这些快照面向监控,读取时不是跨字段事务视图。
|
||||
|
||||
启用 `metrics` 后,库通过 [`metrics`](https://crates.io/crates/metrics) facade 记录
|
||||
指标,但不会安装 recorder 或启动 HTTP 服务:
|
||||
|
||||
```toml
|
||||
[dependencies]
|
||||
dht-crawler = { version = "0.2", features = ["metrics"] }
|
||||
```
|
||||
|
||||
指标名称、类型、标签和单位见 [docs/metrics.md](docs/metrics.md)。
|
||||
|
||||
## Cargo features
|
||||
|
||||
默认不启用任何 feature。
|
||||
|
||||
| Feature | 用途 |
|
||||
|---|---|
|
||||
| `metrics` | 通过 `metrics` facade 记录指标 |
|
||||
| `mimalloc` | 将 mimalloc 注册为全局分配器 |
|
||||
|
||||
启用 `mimalloc` 前,请确认最终二进制没有注册其他全局分配器。
|
||||
|
||||
## 开发
|
||||
|
||||
```bash
|
||||
cargo fmt --all --check
|
||||
cargo check --all-targets --all-features
|
||||
cargo clippy --all-targets --all-features -- -D warnings
|
||||
cargo test --all-features
|
||||
cargo doc --no-deps --all-features
|
||||
```
|
||||
|
||||
## 许可证
|
||||
|
||||
[MIT](LICENSE)
|
||||
基础库的使用方式和指标说明见 [`dht-crawler/README.md`](dht-crawler/README.md)
|
||||
|
||||
@@ -0,0 +1,54 @@
|
||||
# 定义可复用 DHT 爬虫基础库的包元数据依赖和功能开关
|
||||
|
||||
[package]
|
||||
name = "dht-crawler"
|
||||
version = "0.2.1"
|
||||
edition.workspace = true
|
||||
authors.workspace = true
|
||||
description = "高性能的 Rust DHT 爬虫基础库"
|
||||
license.workspace = true
|
||||
documentation = "https://docs.rs/dht-crawler"
|
||||
repository.workspace = true
|
||||
keywords = ["dht", "bittorrent", "crawler", "p2p", "torrent"]
|
||||
categories = ["network-programming", "asynchronous"]
|
||||
readme = "README.md"
|
||||
|
||||
[lib]
|
||||
name = "dht_crawler"
|
||||
crate-type = ["rlib"]
|
||||
|
||||
[dependencies]
|
||||
ahash = "0.8"
|
||||
arc-swap = "1.7"
|
||||
async-channel = "2.5.0"
|
||||
bytes = "1.0"
|
||||
crossbeam-queue = "0.3"
|
||||
hex = "0.4"
|
||||
log = "0.4"
|
||||
metrics = { version = "0.24", optional = true }
|
||||
rand = "0.10.2"
|
||||
rbit = "0.2"
|
||||
serde.workspace = true
|
||||
serde_bencode = "0.2"
|
||||
serde_bytes = "0.11.19"
|
||||
sha1 = "0.11.0"
|
||||
socket2 = { version = "0.6.5", features = ["all"] }
|
||||
thiserror.workspace = true
|
||||
tokio = { workspace = true, features = ["rt", "rt-multi-thread", "net", "sync", "time", "macros"] }
|
||||
tokio-util.workspace = true
|
||||
|
||||
[dev-dependencies]
|
||||
metrics-exporter-prometheus = { version = "0.18", default-features = false, features = ["http-listener"] }
|
||||
mimalloc = "0.1"
|
||||
tokio = { workspace = true, features = ["signal"] }
|
||||
tracing.workspace = true
|
||||
tracing-subscriber = { workspace = true, features = ["env-filter", "fmt", "tracing-log"] }
|
||||
|
||||
[features]
|
||||
default = []
|
||||
metrics = ["dep:metrics"]
|
||||
mimalloc = []
|
||||
|
||||
[[example]]
|
||||
name = "dht_crawler_example"
|
||||
path = "examples/main.rs"
|
||||
@@ -0,0 +1,264 @@
|
||||
# dht-crawler
|
||||
|
||||
[](https://crates.io/crates/dht-crawler)
|
||||
[](https://docs.rs/dht-crawler)
|
||||
[](../LICENSE)
|
||||
|
||||
基于 Rust 和 Tokio 的 BitTorrent DHT 爬虫库。它参与 BEP-5 DHT 网络,通过 BEP-51
|
||||
`sample_infohashes` 主动发现 InfoHash,也接收 `announce_peer`,并通过 BEP-9
|
||||
`ut_metadata` 获取、校验和解析 torrent 元数据。
|
||||
|
||||
`dht-crawler` 提供:
|
||||
|
||||
- IPv4、IPv6 和双栈 DHT;
|
||||
- 主动节点发现、BEP-51 InfoHash 采样与 `get_peers` 查询;
|
||||
- 有界、去重的 Metadata 下载队列;
|
||||
- InfoHash 过滤、异步准入、结果交付和完成通知;
|
||||
- 默认可用的运行时统计,以及可选的 `metrics` 集成。
|
||||
|
||||
## 安装
|
||||
|
||||
```bash
|
||||
cargo add dht-crawler
|
||||
cargo add tokio --features rt-multi-thread,macros,signal
|
||||
```
|
||||
|
||||
或在 `Cargo.toml` 中添加:
|
||||
|
||||
```toml
|
||||
[dependencies]
|
||||
dht-crawler = "0.2"
|
||||
tokio = { version = "1", features = ["rt-multi-thread", "macros", "signal"] }
|
||||
```
|
||||
|
||||
## 快速开始
|
||||
|
||||
```rust
|
||||
use dht_crawler::prelude::*;
|
||||
|
||||
#[tokio::main]
|
||||
async fn main() -> Result<()> {
|
||||
let server = DHTServer::new(DHTOptions {
|
||||
port: 6881,
|
||||
netmode: NetMode::Ipv4Only,
|
||||
..Default::default()
|
||||
})
|
||||
.await?;
|
||||
|
||||
server.on_torrent(|torrent| {
|
||||
println!(
|
||||
"{} {} {}",
|
||||
torrent.info_hash,
|
||||
torrent.name,
|
||||
torrent.format_size()
|
||||
);
|
||||
});
|
||||
|
||||
server.on_error(|error| {
|
||||
eprintln!("DHT runtime error: {error}");
|
||||
});
|
||||
|
||||
let shutdown = server.clone();
|
||||
tokio::spawn(async move {
|
||||
if tokio::signal::ctrl_c().await.is_ok() {
|
||||
shutdown.shutdown();
|
||||
}
|
||||
});
|
||||
|
||||
// 一直运行,直到 shutdown() 被调用。
|
||||
server.start().await
|
||||
}
|
||||
```
|
||||
|
||||
运行仓库中的完整示例:
|
||||
|
||||
```bash
|
||||
cargo run --release --example dht_crawler_example
|
||||
```
|
||||
|
||||
## 核心 API
|
||||
|
||||
`DHTServer` 是主要入口:
|
||||
|
||||
| API | 用途 |
|
||||
|---|---|
|
||||
| `DHTServer::new(options)` | 校验配置、绑定 UDP Socket 并创建内部管道 |
|
||||
| `start().await` | 启动 DHT 与爬取任务,等待 `shutdown()` |
|
||||
| `shutdown()` | 停止 UDP、爬取和 Metadata 任务;可重复调用 |
|
||||
| `filter(callback)` | 在 InfoHash 进入队列前执行同步过滤 |
|
||||
| `on_metadata_fetch(callback)` | 在第一次 Peer 下载前执行异步准入 |
|
||||
| `on_torrent(callback)` | 接收已校验的 `TorrentInfo` |
|
||||
| `on_torrent_with_ack(callback)` | 接收结果并显式确认是否接受交付 |
|
||||
| `on_metadata_fetch_complete(callback)` | 接收已准入任务的最终状态 |
|
||||
| `on_error(callback)` | 接收运行期错误 |
|
||||
| `runtime_stats()` | 获取可复制的运行时统计句柄 |
|
||||
|
||||
同类回调重复注册时,新回调会替换旧回调。
|
||||
|
||||
### 过滤与准入
|
||||
|
||||
`filter` 是同步的早期过滤器,适合拦截已处理过的 InfoHash:
|
||||
|
||||
```rust
|
||||
server.filter(|info_hash| !already_exists(info_hash));
|
||||
```
|
||||
|
||||
`on_metadata_fetch` 是异步准入回调,在实际连接 Peer 前调用:
|
||||
|
||||
```rust
|
||||
server.on_metadata_fetch(|info_hash| async move {
|
||||
should_download(&info_hash).await
|
||||
});
|
||||
```
|
||||
|
||||
返回 `false` 会终止任务,不下载 Metadata,也不会触发 torrent 或 completion 回调。
|
||||
未注册准入回调时默认允许下载。
|
||||
|
||||
### 交付确认
|
||||
|
||||
不需要确认下游是否接收时使用 `on_torrent`。需要确认下游是否成功接收时使用
|
||||
`on_torrent_with_ack`:
|
||||
|
||||
```rust
|
||||
server.on_torrent_with_ack(|torrent| {
|
||||
output.try_send(torrent).is_ok()
|
||||
});
|
||||
|
||||
server.on_metadata_fetch_complete(|completion| {
|
||||
println!(
|
||||
"{}: {:?}, attempts={}",
|
||||
completion.info_hash,
|
||||
completion.status,
|
||||
completion.attempts
|
||||
);
|
||||
});
|
||||
```
|
||||
|
||||
完成状态:
|
||||
|
||||
| 状态 | 含义 |
|
||||
|---|---|
|
||||
| `Accepted` | Metadata 下载成功,结果已被回调接受 |
|
||||
| `FetchFailed` | 所有可用 Peer 尝试均失败 |
|
||||
| `DeliveryRejected` | Metadata 下载成功,但结果未被回调接受 |
|
||||
|
||||
`attempts` 只统计实际发起的 Peer 网络请求。异步准入拒绝不会产生 completion 事件。
|
||||
|
||||
## 配置
|
||||
|
||||
大多数调用方可以从 `DHTOptions::default()` 开始,只覆盖监听方式和容量限制:
|
||||
|
||||
```rust
|
||||
let options = DHTOptions {
|
||||
port: 6881,
|
||||
netmode: NetMode::DualStack,
|
||||
hash_queue_capacity: 20_000,
|
||||
metadata: MetadataOptions {
|
||||
timeout_secs: 5,
|
||||
max_queue_size: 20_000,
|
||||
max_worker_count: 256,
|
||||
..Default::default()
|
||||
},
|
||||
crawl: CrawlOptions {
|
||||
rate_limit: RateLimitOptions {
|
||||
max_find_node_rate_per_sec: 300,
|
||||
max_in_flight: 768,
|
||||
..Default::default()
|
||||
},
|
||||
..Default::default()
|
||||
},
|
||||
..Default::default()
|
||||
};
|
||||
```
|
||||
|
||||
配置分组:
|
||||
|
||||
| 类型 | 控制内容 |
|
||||
|---|---|
|
||||
| `DHTOptions` | 监听端口、网络模式和顶层队列 |
|
||||
| `MetadataOptions` | 下载超时、队列、并发和失败 Peer 缓存 |
|
||||
| `PeerLookupOptions` | 主动 `get_peers` 的速率与并发 |
|
||||
| `SampleInfohashesOptions` | BEP-51 采样速率、并发、超时、退避和 Hash 去重容量 |
|
||||
| `RateLimitOptions` | `find_node`、在途请求和 UDP 回复预算 |
|
||||
| `PoolOptions` | 节点池、最近探测记录和响应节点缓存 |
|
||||
| `BootstrapOptions` | Bootstrap 节点与失败退避 |
|
||||
| `TargetOptions` | 主动爬取目标生成策略 |
|
||||
| `SchedulerOptions` | 内部事件队列、批处理与快照限制 |
|
||||
|
||||
完整字段和默认值以 [docs.rs API 文档](https://docs.rs/dht-crawler) 为准。需要注意:
|
||||
|
||||
- `DHTOptions::default()` 使用 `Ipv4Only`;
|
||||
- `NetMode::DualStack` 会分别绑定 IPv4 和 IPv6 Socket;
|
||||
- `DHTServer::new()` 会立即在所有可用接口上绑定配置的 UDP 端口;
|
||||
- 空节点池默认每 30 秒重新尝试 Bootstrap,每轮最多使用 16 个已解析端点;
|
||||
- `MetadataOptions::timeout_secs` 是单个 Peer 尝试的端到端期限;
|
||||
- `PeerLookupOptions::max_lookups_per_second = 0` 会关闭主动 `get_peers`;
|
||||
- `SampleInfohashesOptions::max_queries_per_second = 0` 会关闭 BEP-51 主动采样;
|
||||
- Metadata 和爬取队列都是有界的,容量应与下游处理能力一起调整。
|
||||
|
||||
## 数据与运行语义
|
||||
|
||||
`TorrentInfo` 包含 `info_hash`、`magnet_link`、`name`、`total_size`、`files`、
|
||||
`piece_length`、`peers` 和 `timestamp`。只有通过 SHA1 校验并成功解析的 Metadata
|
||||
才会交付给 torrent 回调。
|
||||
|
||||
库使用有界队列控制内存占用。队列满或速率预算耗尽时,新事件可能被拒绝、淘汰或计入
|
||||
drop 指标。Metadata 队列按 InfoHash 去重,每个任务可尝试多个候选 Peer;连接超时和
|
||||
连接失败的 Peer 会被短期缓存,避免反复占用 worker。
|
||||
|
||||
`start()` 返回后,当前实例不能再次启动。如需重新运行,请创建新的 `DHTServer`。
|
||||
|
||||
## 可观测性
|
||||
|
||||
运行时快照无需启用 Cargo feature:
|
||||
|
||||
```rust
|
||||
let stats = server.runtime_stats();
|
||||
let snapshot = stats.snapshot();
|
||||
|
||||
println!(
|
||||
"nodes={} metadata={}/{} workers={}",
|
||||
snapshot.node_pool_size,
|
||||
snapshot.metadata_queue_depth,
|
||||
snapshot.metadata_queue_max,
|
||||
snapshot.metadata_in_flight,
|
||||
);
|
||||
```
|
||||
|
||||
`observability_snapshot()` 提供 UDP、查询、队列、Metadata 失败原因和固定桶直方图。
|
||||
这些快照面向监控,读取时不是跨字段事务视图。
|
||||
|
||||
启用 `metrics` 后,库通过 [`metrics`](https://crates.io/crates/metrics) facade 记录
|
||||
指标,但不会安装 recorder 或启动 HTTP 服务:
|
||||
|
||||
```toml
|
||||
[dependencies]
|
||||
dht-crawler = { version = "0.2", features = ["metrics"] }
|
||||
```
|
||||
|
||||
指标名称、类型、标签和单位见 [docs/metrics.md](docs/metrics.md)。
|
||||
|
||||
## Cargo features
|
||||
|
||||
默认不启用任何 feature。
|
||||
|
||||
| Feature | 用途 |
|
||||
|---|---|
|
||||
| `metrics` | 通过 `metrics` facade 记录指标 |
|
||||
| `mimalloc` | 将 mimalloc 注册为全局分配器 |
|
||||
|
||||
启用 `mimalloc` 前,请确认最终二进制没有注册其他全局分配器。
|
||||
|
||||
## 开发
|
||||
|
||||
```bash
|
||||
cargo fmt --all --check
|
||||
cargo check --all-targets --all-features
|
||||
cargo clippy --all-targets --all-features -- -D warnings
|
||||
cargo test --all-features
|
||||
cargo doc --no-deps --all-features
|
||||
```
|
||||
|
||||
## 许可证
|
||||
|
||||
[MIT](LICENSE)
|
||||
@@ -17,7 +17,7 @@ use arc_swap::ArcSwap;
|
||||
use bytes::BytesMut;
|
||||
#[cfg(feature = "metrics")]
|
||||
use metrics::{counter, gauge};
|
||||
use rand::Rng;
|
||||
use rand::RngExt;
|
||||
use std::collections::{HashMap, VecDeque};
|
||||
use std::net::SocketAddr;
|
||||
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||
@@ -684,7 +684,7 @@ impl CrawlActor {
|
||||
.random_walk_percent
|
||||
.saturating_add(self.config.sparse_bucket_percent)
|
||||
.max(1);
|
||||
if rand::thread_rng().gen_range(0..total) < self.config.sparse_bucket_percent {
|
||||
if rand::rng().random_range(0..total) < self.config.sparse_bucket_percent {
|
||||
target_for_bucket(&self.local_id, bucket_index(&node.id, &self.local_id))
|
||||
} else {
|
||||
random_node_id()
|
||||
@@ -1,4 +1,4 @@
|
||||
use rand::Rng;
|
||||
use rand::RngExt;
|
||||
|
||||
pub(crate) type TransactionId = [u8; 8];
|
||||
|
||||
@@ -14,7 +14,7 @@ pub(crate) fn transaction_id_from_bytes(bytes: &[u8]) -> Option<TransactionId> {
|
||||
|
||||
pub(crate) fn random_node_id() -> [u8; 20] {
|
||||
let mut id = [0u8; 20];
|
||||
rand::thread_rng().fill(&mut id);
|
||||
rand::rng().fill(&mut id);
|
||||
id
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
use crate::types::NodeTuple;
|
||||
use rand::seq::SliceRandom;
|
||||
use rand::seq::IndexedRandom;
|
||||
|
||||
#[derive(Default)]
|
||||
pub(crate) struct RoutingSnapshot {
|
||||
@@ -22,15 +22,15 @@ impl RoutingSnapshot {
|
||||
}
|
||||
|
||||
pub(crate) fn random_nodes(&self, count: usize, filter_ipv6: Option<bool>) -> Vec<NodeTuple> {
|
||||
let mut rng = rand::thread_rng();
|
||||
let mut rng = rand::rng();
|
||||
match filter_ipv6 {
|
||||
Some(true) => self.v6.choose_multiple(&mut rng, count).cloned().collect(),
|
||||
Some(false) => self.v4.choose_multiple(&mut rng, count).cloned().collect(),
|
||||
Some(true) => self.v6.sample(&mut rng, count).cloned().collect(),
|
||||
Some(false) => self.v4.sample(&mut rng, count).cloned().collect(),
|
||||
None => {
|
||||
let mut all = Vec::with_capacity(self.v4.len() + self.v6.len());
|
||||
all.extend_from_slice(&self.v4);
|
||||
all.extend_from_slice(&self.v6);
|
||||
all.choose_multiple(&mut rng, count).cloned().collect()
|
||||
all.sample(&mut rng, count).cloned().collect()
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -28,7 +28,7 @@ use arc_swap::ArcSwapOption;
|
||||
use bytes::BytesMut;
|
||||
#[cfg(feature = "metrics")]
|
||||
use metrics::counter;
|
||||
use rand::Rng;
|
||||
use rand::RngExt;
|
||||
use socket2::{Domain, Protocol, Socket, Type};
|
||||
use std::collections::{HashMap, VecDeque};
|
||||
use std::future::Future;
|
||||
@@ -347,7 +347,7 @@ impl DHTServer {
|
||||
|
||||
let node_id = random_node_id();
|
||||
let mut token_secret = [0u8; 10];
|
||||
rand::thread_rng().fill(&mut token_secret);
|
||||
rand::rng().fill(&mut token_secret);
|
||||
|
||||
let (hash_events_tx, hash_rx) =
|
||||
mpsc::channel::<HashDiscovered>(options.hash_queue_capacity);
|
||||
@@ -0,0 +1,28 @@
|
||||
# 定义 DHT 元数据搜索应用的独立依赖和构建入口
|
||||
|
||||
[package]
|
||||
name = "dht-search"
|
||||
version = "0.1.0"
|
||||
edition.workspace = true
|
||||
authors.workspace = true
|
||||
license.workspace = true
|
||||
repository.workspace = true
|
||||
publish = false
|
||||
|
||||
[features]
|
||||
default = []
|
||||
rocksdb-storage = ["dep:rocksdb"]
|
||||
|
||||
[dependencies]
|
||||
axum = "0.8.9"
|
||||
blake3 = "1.8.5"
|
||||
dht-crawler = { path = "../dht-crawler", features = ["metrics"] }
|
||||
rocksdb = { version = "0.24.0", optional = true }
|
||||
serde.workspace = true
|
||||
serde_json = "1.0"
|
||||
tantivy = "0.26.1"
|
||||
thiserror.workspace = true
|
||||
tokio = { workspace = true, features = ["macros", "rt-multi-thread", "signal", "sync", "time"] }
|
||||
tokio-util.workspace = true
|
||||
tracing.workspace = true
|
||||
tracing-subscriber = { workspace = true, features = ["env-filter", "fmt"] }
|
||||
@@ -0,0 +1 @@
|
||||
// 负责处理搜索详情统计和健康检查请求
|
||||
@@ -0,0 +1,5 @@
|
||||
// 负责组合 HTTP 路由和共享接口状态但不直接访问数据库实现
|
||||
|
||||
mod handlers;
|
||||
mod request;
|
||||
mod response;
|
||||
@@ -0,0 +1 @@
|
||||
// 负责定义 HTTP 查询参数和输入校验模型
|
||||
@@ -0,0 +1 @@
|
||||
// 负责定义稳定的 HTTP 响应模型和领域对象转换边界
|
||||
@@ -0,0 +1 @@
|
||||
// 负责连接采集存储索引和接口层并定义应用级启动顺序
|
||||
@@ -0,0 +1 @@
|
||||
// 负责加载校验和提供应用配置但不执行任何业务逻辑
|
||||
@@ -0,0 +1,4 @@
|
||||
// 负责组合 DHT 发现 Metadata 下载和持久化提交管线
|
||||
|
||||
mod pipeline;
|
||||
mod worker;
|
||||
@@ -0,0 +1 @@
|
||||
// 负责定义采集阶段之间的有界队列背压和任务流转规则
|
||||
@@ -0,0 +1 @@
|
||||
// 负责消费 infohash 并执行受资源限制的 Metadata 下载任务
|
||||
@@ -0,0 +1 @@
|
||||
// 负责规范化文件结构并生成用于内容聚合的稳定 BLAKE3 指纹
|
||||
@@ -0,0 +1,4 @@
|
||||
// 负责导出不依赖存储搜索和传输实现的核心领域模型
|
||||
|
||||
mod fingerprint;
|
||||
mod torrent;
|
||||
@@ -0,0 +1 @@
|
||||
// 负责定义种子元数据文件条目发现状态和索引状态模型
|
||||
@@ -0,0 +1 @@
|
||||
// 负责定义应用层统一错误类型和跨模块错误转换边界
|
||||
@@ -0,0 +1,14 @@
|
||||
// 负责组装应用依赖启动运行时并协调服务生命周期
|
||||
|
||||
mod api;
|
||||
mod app;
|
||||
mod config;
|
||||
mod crawler;
|
||||
mod domain;
|
||||
mod error;
|
||||
mod search;
|
||||
mod shutdown;
|
||||
mod storage;
|
||||
mod telemetry;
|
||||
|
||||
fn main() {}
|
||||
@@ -0,0 +1 @@
|
||||
// 负责批量写入删除提交和从权威存储重建 Tantivy 索引
|
||||
@@ -0,0 +1,5 @@
|
||||
// 负责暴露全文搜索抽象并隐藏 Tantivy 的具体实现细节
|
||||
|
||||
mod indexer;
|
||||
mod query;
|
||||
mod schema;
|
||||
@@ -0,0 +1 @@
|
||||
// 负责构建全文查询过滤排序分页和内容聚合逻辑
|
||||
@@ -0,0 +1 @@
|
||||
// 负责定义 Tantivy 字段分词索引存储和快速字段策略
|
||||
@@ -0,0 +1 @@
|
||||
// 负责监听退出信号并协调有界队列排空提交和资源关闭
|
||||
@@ -0,0 +1 @@
|
||||
// 负责定义稳定的 RocksDB 键空间编码和版本边界
|
||||
@@ -0,0 +1,5 @@
|
||||
// 负责暴露持久化抽象并隐藏 RocksDB 的具体实现细节
|
||||
|
||||
mod keys;
|
||||
mod repository;
|
||||
mod rocks;
|
||||
@@ -0,0 +1 @@
|
||||
// 负责定义元数据去重状态恢复和索引任务所需的存储接口
|
||||
@@ -0,0 +1 @@
|
||||
// 负责实现 RocksDB 打开配置批量写入精确查询和关闭流程
|
||||
@@ -0,0 +1 @@
|
||||
// 负责初始化结构化日志指标和运行状态观测
|
||||
Reference in New Issue
Block a user