From a1399ef7427b7e44285f2f0a35b2fc14ff819561 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=A1=A5=E4=B8=8B=E7=BA=A2=E8=8D=AF?= Date: Thu, 30 Jul 2026 13:55:25 +0800 Subject: [PATCH] docs: rewrite developer README --- README.md | 456 ++++++++++++++++++++---------------------------------- 1 file changed, 166 insertions(+), 290 deletions(-) diff --git a/README.md b/README.md index 5cdafd7..6d0aa2c 100644 --- a/README.md +++ b/README.md @@ -4,63 +4,32 @@ [![Documentation](https://docs.rs/dht-crawler/badge.svg)](https://docs.rs/dht-crawler) [![License](https://img.shields.io/crates/l/dht-crawler.svg)](LICENSE) -基于 Rust、Tokio 的 BitTorrent DHT 爬虫库。它加入 BEP-5 网络,接收有效的 -`announce_peer`,并通过 BEP-9 `ut_metadata` 下载和校验种子元数据。 +基于 Rust 和 Tokio 的 BitTorrent DHT 爬虫库。它参与 BEP-5 DHT 网络,接收 +`announce_peer`,并通过 BEP-9 `ut_metadata` 获取、校验和解析 torrent 元数据。 -当前版本:`0.2.0`。0.2 重做了节点池、主动爬取、Metadata 调度和运行时观测接口, -从 0.1 升级时请先阅读[迁移说明](#从-01-迁移到-02)和 [CHANGELOG](CHANGELOG.md)。 +`dht-crawler` 提供: -## 文档导航 - -- [快速开始](#快速开始):最小可运行示例和优雅停机。 -- [架构与背压](#架构与背压):UDP、主动爬取和 Metadata 管道。 -- [配置默认值](#配置默认值):所有公开 `DHTOptions` 字段。 -- [回调与生命周期](#回调与生命周期):抓取准入、交付确认和完成状态。 -- [运行时观测](#运行时观测):无 exporter 快照、Prometheus 和完整[指标表](docs/metrics.md)。 -- [JNI](#jni):Java 集成入口;[0.1 → 0.2 迁移](#从-01-迁移到-02)。 - -## 主要能力 - -- IPv4、IPv6 和双栈 DHT Socket。 -- 单所有者 crawl actor:严格 FIFO 节点池、最近探测状态、在途请求和所有速率预算 - 由一个 actor 管理,UDP worker 不锁节点池。 -- 查询 QPS、新目标/分钟、节点替换/分钟、总在途、子网在途、回复包、回复字节和 - 单来源回复分别限流。 -- 有界 Metadata 队列按 InfoHash 去重,并保留最多三个新鲜 Peer。 -- Metadata 总超时覆盖 TCP 连接、BitTorrent/扩展握手、分片下载、SHA1 校验和解析。 -- 按 `SocketAddr` 缓存 Peer 的超时/连接失败,避免坏 Peer 反复占用 worker。 -- 传输无关的原子运行时快照和固定桶直方图;可选 `metrics` feature。 -- 可选 JNI 接口和 Java 示例。 +- IPv4、IPv6 和双栈 DHT; +- 主动节点发现与 `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"] } ``` -如果应用需要通过 `metrics` facade 输出指标: - -```toml -[dependencies] -dht-crawler = { version = "0.2", features = ["metrics"] } -metrics-exporter-prometheus = { version = "0.18", default-features = false, features = ["http-listener"] } -``` - -`metrics` feature 只负责记录指标,不会在库内启动 HTTP 服务。应用必须自行安装 -recorder/exporter;完整指标清单见 [docs/metrics.md](docs/metrics.md)。 - -### Cargo features - -| Feature | 默认启用 | 作用 | -|---|---|---| -| `metrics` | 否 | 通过 `metrics` facade 记录低基数指标 | -| `jni` | 否 | 构建 Java JNI 接口和 `cdylib` | -| `mimalloc` | 否 | 将 mimalloc 注册为全局分配器 | - -库的默认 feature 集为空。启用 `mimalloc` 前,请确认最终二进制没有注册其他全局分配器。 - ## 快速开始 ```rust @@ -68,303 +37,217 @@ use dht_crawler::prelude::*; #[tokio::main] async fn main() -> Result<()> { - let options = DHTOptions { - port: 12313, + let server = DHTServer::new(DHTOptions { + port: 6881, netmode: NetMode::Ipv4Only, - metadata: MetadataOptions { - timeout_secs: 4, - max_queue_size: 10_000, - max_worker_count: 256, - ..Default::default() - }, - crawl: CrawlOptions { - rate_limit: RateLimitOptions { - max_find_node_rate_per_sec: 200, - burst: 40, - max_in_flight: 512, - ..Default::default() - }, - ..Default::default() - }, ..Default::default() - }; + }) + .await?; - let server = DHTServer::new(options).await?; - - // 返回 true 才允许该 InfoHash 进入实际 Peer 下载阶段。 - server.on_metadata_fetch(|_info_hash| async move { true }); - - // 简单回调总是接受交付。 server.on_torrent(|torrent| { - println!("{}: {}", torrent.info_hash, torrent.name); + println!( + "{} {} {}", + torrent.info_hash, + torrent.name, + torrent.format_size() + ); }); - server.on_error(|error| eprintln!("DHT runtime error: {error}")); + server.on_error(|error| { + eprintln!("DHT runtime error: {error}"); + }); - let shutdown_server = server.clone(); + let shutdown = server.clone(); tokio::spawn(async move { if tokio::signal::ctrl_c().await.is_ok() { - shutdown_server.shutdown(); + shutdown.shutdown(); } }); - // start() 阻塞到 shutdown() 被调用。 + // 一直运行,直到 shutdown() 被调用。 server.start().await } ``` -可运行版本见 [examples/main.rs](examples/main.rs)。如果自己的 Tokio 依赖没有启用 -`signal` feature,可以使用其他取消源调用 `shutdown()`。 +运行仓库中的完整示例: -## 架构与背压 - -```text -UDP sockets - ├─ bounded UDP worker queues ──→ KRPC workers ──→ bounded crawl events - │ │ - │ └─ announce_peer → bounded hash ingress - │ - └─ crawl egress ← single crawl actor ← priority/discovery events - │ - ├─ strict FIFO node pool + recent-probe set - ├─ pending transaction map + subnet counters - └─ ArcSwap responsive-node snapshot - -hash ingress → deduplicating Metadata queue → bounded workers → torrent callback +```bash +cargo run --release --example dht_crawler_example ``` -所有跨任务入口都是有界队列。达到容量时,事件会被拒绝、淘汰或计入 drop 指标, -不会依靠无限增长的缓冲区掩盖下游过载。主动爬取预算还会随 Metadata 队列压力下降。 +## 核心 API -### 主动爬取 +`DHTServer` 是主要入口: -- 新地址通常只发送一次 `find_node`;默认等待回复 `2s`,不做同目标重试。 -- 超时不会自动降低配置 QPS。Metadata 队列压力达到 80% 后才开始自动降速,95% 时 - 降到 `metadata_pressure_floor_percent` 指定的比例;默认下限为配置 QPS 的 25%。 -- 节点池是严格 FIFO。重复地址、无效公网地址和超出 replacement budget 的替换会被拒绝。 -- 响应成功的节点进入一个独立、有界、带 TTL 的 responsive ring,用于回复其他 DHT - 节点和 revisit 查询;它不是第二个爬取池。 -- 节点池低于 `low_watermark` 时触发 bootstrap。失败的 bootstrap 来源按配置退避。 -- UDP 回复总包数、总字节数和单来源包数分别限流,其中 10% 包/字节预算保留给 - `ping` 和 `get_peers` 的保底回复,但不会突破配置的总上限。 +| 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()` | 获取可复制的运行时统计句柄 | -### Metadata 调度 +同类回调重复注册时,新回调会替换旧回调。 -- Hash ingress 和 Metadata pending queue 都是有界的。 -- Pending queue 按 InfoHash 去重,每个 Hash 最多保留三个不同且新鲜的 Peer。 -- Pending 项固定在 60 秒后过期。队列满时,较新的 Hash 可以淘汰最旧项;比当前 - 最旧项还旧的事件直接视为 stale。 -- worker 优先分派最新的可用 Hash,以提高 Peer 仍在线的概率。 -- `timeout_secs` 是一次 Peer 尝试的端到端期限,不会在连接、握手和下载阶段重复叠加。 -- 单个 metadata payload 上限为 10 MiB;下载完成后必须通过 SHA1 和 bencode 解析。 -- Peer failure cache 只缓存 `timeout` 和 `connect_failed`,键为完整 `SocketAddr` - (IP + port)。缓存命中不会发起网络请求,也不计入三次真实 Peer 尝试。 -- `peer_failure_cache_capacity = 0` 或 `peer_failure_ttl_secs = 0` 会关闭缓存。 +### 过滤与准入 -## 配置默认值 +`filter` 是同步的早期过滤器,适合拦截已处理过的 InfoHash: -库本身不包含 P1/P15 等档位概念。应用如需档位,应将其转换成下列具体选项。 - -### `DHTOptions` 与 Metadata - -| 字段 | 默认值 | 说明 | -|---|---:|---| -| `port` | `6881` | DHT UDP 监听端口 | -| `netmode` | `Ipv4Only` | `DHTOptions::default()` 的网络模式 | -| `hash_queue_capacity` | `10000` | announce 到 Metadata scheduler 的 ingress 容量 | -| `metadata.timeout_secs` | `4` | 单 Peer 端到端超时 | -| `metadata.max_queue_size` | `10000` | 去重 Pending Hash 容量 | -| `metadata.max_worker_count` | `256` | 最大并发 Metadata job 数 | -| `metadata.peer_failure_cache_capacity` | `200000` | 坏 Peer 缓存容量 | -| `metadata.peer_failure_ttl_secs` | `60` | 坏 Peer 缓存 TTL | -| `peer_lookup.max_lookups_per_second` | `32` | 每秒启动的主动 `get_peers` 查询数;`0` 为关闭 | -| `peer_lookup.burst` | `32` | 空闲后可立即消费的查询预算 | -| `peer_lookup.max_active_lookups` | `64` | 同时活跃的 InfoHash 查询上限 | - -### `crawl.rate_limit` - -| 字段 | 默认值 | -|---|---:| -| `max_find_node_rate_per_sec` | `200` | -| `burst` | `40` | -| `max_in_flight` | `512` | -| `request_timeout_secs` | `2` | -| `max_new_destinations_per_minute` | `10000` | -| `max_replacements_per_minute` | `25000` | -| `max_response_rate_per_sec` | `500` | -| `max_response_bytes_per_sec` | `1048576` | -| `max_response_rate_per_source` | `40` | -| `metadata_pressure_floor_percent` | `25` | -| `max_in_flight_per_subnet` | `8` | - -### Pool、Bootstrap、Target 与 Scheduler - -| 字段 | 默认值 | -|---|---:| -| `pool.capacity` | `100000` | -| `pool.recent_probe_ttl_secs` | `600` | -| `pool.responsive_capacity` | `16384` | -| `pool.responsive_ttl_secs` | `900` | -| `pool.low_watermark` | `10000` | -| `bootstrap.interval_secs` | `300` | -| `bootstrap.max_nodes_per_round` | `3` | -| `bootstrap.source_backoff_base_secs` | `300` | -| `bootstrap.source_backoff_max_secs` | `3600` | -| `target.random_walk_percent` | `70` | -| `target.sparse_bucket_percent` | `30` | -| `target.neighbor_sender_id` | `true` | -| `scheduler.priority_event_channel_capacity` | `8192` | -| `scheduler.discovery_event_channel_capacity` | `16384` | -| `scheduler.event_batch_limit` | `256` | -| `scheduler.node_batch_limit` | `4096` | -| `scheduler.routing_snapshot_size` | `4096` | -| `scheduler.snapshot_refresh_millis` | `1000` | - -默认 bootstrap 来源: - -```text -router.bittorrent.com:6881 -dht.transmissionbt.com:6881 -router.utorrent.com:6881 -dht.aelitis.com:6881 +```rust +server.filter(|info_hash| !already_exists(info_hash)); ``` -内部会对不安全的零值和百分比做归一化,例如容量/在途至少为 1、百分比最大为 100、 -`low_watermark` 不超过 pool capacity。建议调用方仍显式传入有效配置,不依赖归一化。 +`on_metadata_fetch` 是异步准入回调,在实际连接 Peer 前调用: -## 回调与生命周期 +```rust +server.on_metadata_fetch(|info_hash| async move { + should_download(&info_hash).await +}); +``` -### `on_metadata_fetch` +返回 `false` 会终止任务,不下载 Metadata,也不会触发 torrent 或 completion 回调。 +未注册准入回调时默认允许下载。 -在 Hash 首次准备进入 Peer 下载前调用。返回 `false` 表示 gate reject:不下载、不触发 -`on_torrent`,也不会触发 `on_metadata_fetch_complete`。 +### 交付确认 -### `on_torrent` 与 `on_torrent_with_ack` +不需要确认下游是否接收时使用 `on_torrent`。需要确认下游是否成功接收时使用 +`on_torrent_with_ack`: -- `on_torrent` 适合无需确认交付的调用方;回调返回后视为 `Accepted`。 -- `on_torrent_with_ack` 返回 `true` 表示应用接受交付,返回 `false` 表示 - `DeliveryRejected`。后者代表 Metadata 已成功下载,但业务层没有接收,不等同于 - `FetchFailed`。 -- 如果没有注册任何 torrent callback,成功下载的 Metadata 同样按 `DeliveryRejected` - 结束。 -- Torrent 回调发生 panic 时会被捕获并按拒绝交付处理。 +```rust +server.on_torrent_with_ack(|torrent| { + output.try_send(torrent).is_ok() +}); -### `on_metadata_fetch_complete` +server.on_metadata_fetch_complete(|completion| { + println!( + "{}: {:?}, attempts={}", + completion.info_hash, + completion.status, + completion.attempts + ); +}); +``` -一个通过 gate 的 Hash 最终只发出一次完成通知: +完成状态: | 状态 | 含义 | |---|---| -| `Accepted` | 下载成功,torrent callback 接受交付 | -| `FetchFailed` | 所有可用 Peer 尝试失败 | -| `DeliveryRejected` | 下载成功,torrent callback 拒绝交付 | +| `Accepted` | Metadata 下载成功,结果已被回调接受 | +| `FetchFailed` | 所有可用 Peer 尝试均失败 | +| `DeliveryRejected` | Metadata 下载成功,但结果未被回调接受 | -`attempts` 只统计实际发起的 Peer 网络尝试;failure cache 命中不计入。 +`attempts` 只统计实际发起的 Peer 网络请求。异步准入拒绝不会产生 completion 事件。 -### 启动与停止 +## 配置 -- `DHTServer::new()` 创建并绑定所需 Socket,失败直接返回 `Err`。 -- `start()` 启动后台任务并等待取消,因此通常应在应用主任务中 await。 -- `shutdown()` 可从 clone handle 调用,取消 DHT、crawl 和 Metadata 后台任务。 -- 已 shutdown 的同一实例不能重新 start;需要重新构造 `DHTServer`。 -- `on_error` 用于接收运行期协议/worker 错误;初始化错误仍通过 `Result` 返回。 +大多数调用方可以从 `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() +}; +``` -不启用 `metrics` feature 也可以读取运行时快照: +配置分组: + +| 类型 | 控制内容 | +|---|---| +| `DHTOptions` | 监听端口、网络模式和顶层队列 | +| `MetadataOptions` | 下载超时、队列、并发和失败 Peer 缓存 | +| `PeerLookupOptions` | 主动 `get_peers` 的速率与并发 | +| `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 端口; +- `MetadataOptions::timeout_secs` 是单个 Peer 尝试的端到端期限; +- `PeerLookupOptions::max_lookups_per_second = 0` 会关闭主动 `get_peers`; +- 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(); -let runtime = stats.snapshot(); println!( - "nodes={} metadata={}/{} in_flight={}", - runtime.node_pool_size, - runtime.metadata_queue_depth, - runtime.metadata_queue_max, - runtime.metadata_in_flight, -); - -let observability = stats.observability_snapshot(); -println!( - "udp rx={}B tx={}B fetch_p95={:?}ms", - observability.udp_rx_bytes, - observability.udp_tx_bytes, - observability.fetch_duration_ms.percentile(0.95), + "nodes={} metadata={}/{} workers={}", + snapshot.node_pool_size, + snapshot.metadata_queue_depth, + snapshot.metadata_queue_max, + snapshot.metadata_in_flight, ); ``` -快照使用 relaxed atomic load,适合监控,不是跨字段事务视图。计数器是进程生命周期 -累计值,调用方通过相邻快照差值计算 rate,并应处理进程重启导致的 counter reset。 +`observability_snapshot()` 提供 UDP、查询、队列、Metadata 失败原因和固定桶直方图。 +这些快照面向监控,读取时不是跨字段事务视图。 -固定桶: +启用 `metrics` 后,库通过 [`metrics`](https://crates.io/crates/metrics) facade 记录 +指标,但不会安装 recorder 或启动 HTTP 服务: -| 直方图 | 边界/单位 | -|---|---| -| Metadata queue wait | `10/50/100/250/500/1000/2000/5000 ms` | -| Metadata fetch | `250/500/1000/2000/4000/6000/10000 ms` | -| Metadata payload size | `16/32/64/128/256/512/1024 KiB, 10 MiB` | - -`counts[i]` 表示小于等于 `bounds[i]` 的非累计桶计数,`overflow` 表示超过最后边界的 -数量。`percentile()` 返回桶上界;落入 overflow 时只能返回最后一个上界,因此它是 -有界近似值,不是精确分位数。 - -高级快照类型从 crate 根导出;当前 prelude 只重导出 `DhtRuntimeStats` 和 -`DhtRuntimeSnapshot`。 - -## Prometheus / metrics - -启用 `metrics` 后,库通过 `metrics` facade 记录 counter、gauge 和 histogram。 -应用必须在创建/启动 server 前安装全局 recorder。示例: - -```rust -use metrics_exporter_prometheus::PrometheusBuilder; - -PrometheusBuilder::new() - .with_http_listener("127.0.0.1:9000".parse().unwrap()) - .install() - .unwrap(); +```toml +[dependencies] +dht-crawler = { version = "0.2", features = ["metrics"] } ``` -完整名称、标签和单位见 [docs/metrics.md](docs/metrics.md)。 +指标名称、类型、标签和单位见 [docs/metrics.md](docs/metrics.md)。 -## 应用层集成 +## Cargo features -本 crate 只提供 DHT、BEP-9 Metadata、回调和观测能力,不包含 Redis、Manticore、 -HTTP 看板或 `P1`~`P17` 性能档位。同级的 `dht-crawler-node` 项目负责这些应用层策略, -并把档位转换成具体的 `DHTOptions`。开发两个项目时应保持下面的目录关系: +默认不启用任何 feature。 -```text -workspace-parent/ -├── dht-crawler/ -└── dht-crawler-node/ -``` - -## JNI - -启用 `jni` feature 可构建 `cdylib`。Java 示例、线程模型和 JNI 配置字段见 -[examples-jni/README.md](examples-jni/README.md)。Java `DHTOptions` 是 Rust 配置的 -扁平化子集,未暴露的 Bootstrap、Target、Scheduler 和 Peer failure cache 字段使用 -Rust 默认值。 - -## 从 0.1 迁移到 0.2 - -0.2 是 breaking release: - -| 0.1 | 0.2 | +| Feature | 用途 | |---|---| -| `metadata_timeout` | `metadata.timeout_secs` | -| `max_metadata_queue_size` | `metadata.max_queue_size` | -| `max_metadata_worker_count` | `metadata.max_worker_count` | -| `node_queue_capacity` | `crawl.pool.capacity` | -| 旧 active/candidate frontier | 单所有者严格 FIFO pool + responsive ring | -| 无交付确认 | `on_torrent_with_ack` + `DeliveryRejected` | -| 粗粒度统计 | `runtime_stats()` 的两类原子快照和固定桶 | +| `metrics` | 通过 `metrics` facade 记录指标 | +| `jni` | 构建 Java JNI 接口和 `cdylib` | +| `mimalloc` | 将 mimalloc 注册为全局分配器 | -`DHTOptions` 不提供旧字段兼容层,升级时必须修改构造代码。配置档位属于应用策略, -不在库内实现。 +JNI 的构建方式和 Java 示例见 [examples-jni/README.md](examples-jni/README.md)。 +启用 `mimalloc` 前,请确认最终二进制没有注册其他全局分配器。 -## 构建与验证 +## 开发 ```bash cargo fmt --all --check @@ -372,13 +255,6 @@ cargo check --all-targets --all-features cargo clippy --all-targets --all-features -- -D warnings cargo test --all-features cargo doc --no-deps --all-features -cargo run --release --example dht_crawler_example -``` - -JNI: - -```bash -cargo build --release --features jni ``` ## 许可证