feat(cluster): harden collector transport and diagnostics

This commit is contained in:
chuan
2026-08-12 00:53:24 +08:00
parent d24ac2afe9
commit 787bcfc877
16 changed files with 752 additions and 67 deletions
+1
View File
@@ -1,6 +1,7 @@
# Rust
/target/
/.tools/
/.run-*
/.run-data/
/.remote-data/
/data/
Generated
+1
View File
@@ -681,6 +681,7 @@ dependencies = [
"tracing-subscriber",
"unicode-normalization",
"windows-sys 0.61.2",
"zstd",
]
[[package]]
+11 -1
View File
@@ -33,10 +33,14 @@ Tantivy 使用代际影子索引处理 Schema 文档格式损坏和缺失等全
| `coordinator` | 唯一存储搜索协调器 不启动 DHT |
| `collector` | 只运行 DHT 采集有效性验证和持久待发送箱 |
协调器和采集器必须使用相同的 `DHT_SEARCH_CLUSTER_TOKEN` 环境变量 值至少三十二个字符 内部接口使用独立端口且不应通过公开 Caddy 站点暴露
协调器和采集器必须使用相同的 `DHT_SEARCH_CLUSTER_TOKEN` 环境变量 值至少三十二个字符 内部接口使用独立端口且不得无保护暴露 公网中继必须同时使用防火墙或反向代理来源 IP 白名单限制采集器来源
采集器第一次连接时自动生成并持久化节点 ID 协调器保存每个节点的期望配置和修订号 Metadata 在采集器本地 RocksDB 待发送箱落盘后才确认下载成功 协调器断开或积压达到高水位时采集器会暂停并在恢复后继续 采集器心跳同时上报 DHT Peer 失败分类队列和当前进程资源 协调器把集群汇总与逐节点快照写入可删除的诊断历史
Metadata 待发送箱以最多 32 条和 1 MiB 原始编码软上限组成批次 超过软上限的单条 Metadata 仍会独立发送 采集器最多并发 8 个批次并使用 Zstandard 压缩 协调器限制请求体和解压后数据为 16 MiB 只有收到数量一致的完整批次结果后才从本地 RocksDB 原子确认删除 失败批次继续保留并按有上限的指数退避重试
诊断页的本次运行接收和发送表示所选采集器容器从本次启动开始的非回环网络接口累计字节 包含 DHT UDP Peer Metadata 下载 协调器上传 心跳和配置同步 全部采集器视图显示各采集器最近上报值的合计 不包含宿主机其他服务并在采集器容器重启后重新计数
本地启动一个协调器和两个采集器
```shell
@@ -143,6 +147,12 @@ cargo test --workspace --all-targets --all-features
cargo clippy --workspace --all-targets --all-features -- -D warnings
```
手动验证高延迟限速和请求失败条件下的持久待发送箱批量清空
```powershell
cargo test -p dht-search --all-features compressed_batches_drain_a_mixed_outbox_over_a_limited_link -- --ignored --nocapture
```
Web 验证命令
```shell
+2 -1
View File
@@ -85,6 +85,7 @@
- [ ] 使用同一版本完成二十四小时连续运行
- [ ] 使用同一版本完成七天连续运行
- [ ] 在真实独立设备上部署一个协调器和至少两个采集器完成二十四小时小流量验收
- [ ] 验证采样节点冷却表修复后两个远端采集器的内存进入稳定区间
- [ ] 记录连续运行期间的私有内存队列深度 Metadata 成功率磁盘增长和索引提交状态
- [ ]`README.md` 补充简洁的备份恢复和常见故障排查步骤
@@ -97,4 +98,4 @@
## 当前下一步
在真实独立设备上小流量部署分布式采集 验证隧道断开恢复每节点资源占用待发送箱增长和有效性验证任务转移
继续观察两个远端采集器的长期内存和待发送箱状态并完成二十四小时分布式运行验收
+129 -19
View File
@@ -15,7 +15,8 @@ use arc_swap::{ArcSwap, ArcSwapOption};
use bytes::BytesMut;
#[cfg(feature = "metrics")]
use metrics::{counter, gauge};
use std::collections::VecDeque;
use std::cmp::Reverse;
use std::collections::{BinaryHeap, VecDeque};
use std::future::Future;
use std::net::SocketAddr;
use std::sync::Arc;
@@ -43,6 +44,63 @@ struct PendingRequest {
source: DiscoverySource,
}
struct NodeCooldowns {
entries: AHashMap<SocketAddr, Instant>,
expiry: BinaryHeap<Reverse<(Instant, SocketAddr)>>,
capacity: usize,
}
impl NodeCooldowns {
fn new(capacity: usize) -> Self {
Self {
entries: AHashMap::with_capacity(capacity.min(16_384)),
expiry: BinaryHeap::with_capacity(capacity.min(16_384)),
capacity,
}
}
fn is_blocked(&mut self, addr: &SocketAddr, now: Instant) -> bool {
self.expire(now);
self.entries.contains_key(addr)
}
fn insert(&mut self, addr: SocketAddr, deadline: Instant, now: Instant) {
if self.capacity == 0 {
return;
}
self.expire(now);
if self.entries.contains_key(&addr) {
return;
}
while self.entries.len() >= self.capacity {
let Some(Reverse((old_deadline, old_addr))) = self.expiry.pop() else {
break;
};
if self.entries.get(&old_addr).copied() == Some(old_deadline) {
self.entries.remove(&old_addr);
}
}
self.entries.insert(addr, deadline);
self.expiry.push(Reverse((deadline, addr)));
}
fn expire(&mut self, now: Instant) {
while let Some(Reverse((deadline, addr))) = self.expiry.peek().copied() {
if deadline > now {
break;
}
self.expiry.pop();
if self.entries.get(&addr).copied() == Some(deadline) {
self.entries.remove(&addr);
}
}
}
fn len(&self) -> usize {
self.entries.len()
}
}
struct SampleResponse {
remote_addr: SocketAddr,
tid: TransactionId,
@@ -239,7 +297,7 @@ pub(crate) fn spawn_sample_infohashes(
dedup_capacity: options.dedup_capacity,
seen_hashes: AHashSet::new(),
seen_order: VecDeque::new(),
next_allowed: AHashMap::new(),
cooldowns: NodeCooldowns::new(options.cooldown_capacity),
pending_addrs: AHashSet::new(),
pending: AHashMap::new(),
pending_expiry: VecDeque::new(),
@@ -273,7 +331,7 @@ struct SampleInfohashesActor {
dedup_capacity: usize,
seen_hashes: AHashSet<[u8; 20]>,
seen_order: VecDeque<[u8; 20]>,
next_allowed: AHashMap<SocketAddr, Instant>,
cooldowns: NodeCooldowns,
pending_addrs: AHashSet<SocketAddr>,
pending: AHashMap<PendingKey, PendingRequest>,
pending_expiry: VecDeque<(Instant, PendingKey)>,
@@ -392,11 +450,7 @@ impl SampleInfohashesActor {
if sent_count >= budget {
break;
}
if self.pending_addrs.contains(&node.addr)
|| self
.next_allowed
.get(&node.addr)
.is_some_and(|deadline| *deadline > now)
if self.pending_addrs.contains(&node.addr) || self.cooldowns.is_blocked(&node.addr, now)
{
continue;
}
@@ -417,11 +471,7 @@ impl SampleInfohashesActor {
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)
if self.pending_addrs.contains(&node.addr) || self.cooldowns.is_blocked(&node.addr, now)
{
continue;
}
@@ -448,8 +498,8 @@ impl SampleInfohashesActor {
if socket.send_to(&buffer, node.addr).await.is_err() {
self.outbound_query_budget.refund_one();
self.runtime_stats.sample_send_failed();
self.next_allowed
.insert(node.addr, now + self.unsupported_backoff);
self.cooldowns
.insert(node.addr, now + self.unsupported_backoff, now);
return false;
}
let key = PendingKey {
@@ -552,7 +602,7 @@ impl SampleInfohashesActor {
} else {
protocol_interval.saturating_add(self.unsupported_backoff)
};
self.next_allowed.insert(event.remote_addr, now + delay);
self.cooldowns.insert(event.remote_addr, now + delay, now);
#[cfg(feature = "metrics")]
{
counter!("dht_sample_infohashes_responses_total").increment(1);
@@ -576,6 +626,7 @@ impl SampleInfohashesActor {
}
fn expire(&mut self, now: Instant) {
self.cooldowns.expire(now);
while let Some((deadline, key)) = self.pending_expiry.front().copied() {
if deadline > now {
break;
@@ -590,14 +641,17 @@ impl SampleInfohashesActor {
}
self.pending.remove(&key);
self.pending_addrs.remove(&key.addr);
self.next_allowed
.insert(key.addr, now + self.unsupported_backoff);
self.cooldowns
.insert(key.addr, now + self.unsupported_backoff, now);
self.runtime_stats.sample_timeout();
#[cfg(feature = "metrics")]
counter!("dht_sample_infohashes_timeouts_total").increment(1);
}
#[cfg(feature = "metrics")]
gauge!("dht_sample_infohashes_in_flight").set(self.pending.len() as f64);
{
gauge!("dht_sample_infohashes_in_flight").set(self.pending.len() as f64);
gauge!("dht_sample_infohashes_cooldown_entries").set(self.cooldowns.len() as f64);
}
}
fn next_transaction_id(&mut self) -> TransactionId {
@@ -613,6 +667,62 @@ mod tests {
use super::*;
use crate::protocol::DhtMessage;
fn test_addr(port: u16) -> SocketAddr {
SocketAddr::from(([8, 8, 8, 8], port))
}
#[test]
fn node_cooldowns_expire_without_retaining_addresses() {
let start = Instant::now();
let mut cooldowns = NodeCooldowns::new(4);
let addr = test_addr(1);
cooldowns.insert(addr, start + Duration::from_secs(5), start);
assert!(cooldowns.is_blocked(&addr, start));
assert!(!cooldowns.is_blocked(&addr, start + Duration::from_secs(5)));
assert_eq!(cooldowns.len(), 0);
}
#[test]
fn node_cooldowns_evict_the_oldest_entry_at_capacity() {
let start = Instant::now();
let mut cooldowns = NodeCooldowns::new(2);
cooldowns.insert(test_addr(1), start + Duration::from_secs(10), start);
cooldowns.insert(test_addr(2), start + Duration::from_secs(20), start);
cooldowns.insert(test_addr(3), start + Duration::from_secs(30), start);
assert_eq!(cooldowns.len(), 2);
assert!(!cooldowns.is_blocked(&test_addr(1), start));
assert!(cooldowns.is_blocked(&test_addr(2), start));
assert!(cooldowns.is_blocked(&test_addr(3), start));
}
#[test]
fn repeated_node_cooldown_does_not_duplicate_expiry_entries() {
let start = Instant::now();
let mut cooldowns = NodeCooldowns::new(2);
let addr = test_addr(1);
cooldowns.insert(addr, start + Duration::from_secs(10), start);
cooldowns.insert(addr, start + Duration::from_secs(20), start);
assert_eq!(cooldowns.len(), 1);
assert_eq!(cooldowns.expiry.len(), 1);
}
#[test]
fn node_cooldowns_expire_by_deadline_instead_of_insertion_order() {
let start = Instant::now();
let mut cooldowns = NodeCooldowns::new(2);
cooldowns.insert(test_addr(1), start + Duration::from_secs(20), start);
cooldowns.insert(test_addr(2), start + Duration::from_secs(5), start);
cooldowns.expire(start + Duration::from_secs(5));
assert_eq!(cooldowns.len(), 1);
assert!(cooldowns.is_blocked(&test_addr(1), start));
assert!(!cooldowns.is_blocked(&test_addr(2), start));
}
fn node_for_lane(want_direct: bool, percent: u8) -> NodeTuple {
(1..=u16::MAX)
.map(|port| NodeTuple {
+3
View File
@@ -216,6 +216,8 @@ pub struct SampleInfohashesOptions {
pub request_timeout_millis: u64,
/// Retry delay for timeouts or nodes that do not return samples, in seconds.
pub unsupported_backoff_secs: u64,
/// Maximum node cooldowns retained to prevent unbounded growth on long-running crawlers.
pub cooldown_capacity: usize,
/// Maximum sampled InfoHashes retained for bounded in-memory deduplication.
pub dedup_capacity: usize,
}
@@ -371,6 +373,7 @@ impl Default for SampleInfohashesOptions {
fallback_to_iterative: true,
request_timeout_millis: 1_500,
unsupported_backoff_secs: 300,
cooldown_capacity: 100_000,
dedup_capacity: 1_000_000,
}
}
+1
View File
@@ -39,6 +39,7 @@ tracing-appender = "0.2"
tracing-subscriber = { workspace = true, features = ["env-filter", "fmt"] }
tower-http = { version = "0.6", features = ["fs"] }
unicode-normalization = "0.1"
zstd = "0.13"
[dev-dependencies]
tempfile = "3.27"
+2
View File
@@ -188,6 +188,8 @@ pub(crate) async fn run(config: AppConfig, config_service: ConfigService) -> Res
cpu_time_millis: process.cpu_time_millis,
thread_count: process.thread_count,
handle_count: process.handle_count,
network_received_bytes: process.network_received_bytes,
network_transmitted_bytes: process.network_transmitted_bytes,
outbox_items: persistence.queue_depth as u64,
outbox_bytes: 0,
last_error: None,
+280 -23
View File
@@ -27,12 +27,18 @@ use super::{
outbox::{Outbox, OutboxError},
protocol::{
AdmissionRequest, AdmissionResponse, AdmissionStatus, CollectorDesiredConfig,
CollectorMetrics, ControlResponse, HeartbeatRequest, MetadataEnvelope,
MetadataSubmitResponse, PROTOCOL_VERSION, RegisterRequest, VerificationClaimRequest,
VerificationClaimResponse, VerificationJob, VerificationResultRequest,
CollectorMetrics, ControlResponse, HeartbeatRequest, MetadataBatchRequest,
MetadataBatchResponse, MetadataEnvelope, PROTOCOL_VERSION, RegisterRequest,
VerificationClaimRequest, VerificationClaimResponse, VerificationJob,
VerificationResultRequest,
},
};
const METADATA_UPLOAD_CONCURRENCY: usize = 8;
const METADATA_UPLOAD_BATCH_ITEMS: usize = 32;
const METADATA_UPLOAD_BATCH_BYTES: u64 = 1024 * 1024;
const METADATA_UPLOAD_TIMEOUT: Duration = Duration::from_secs(120);
#[derive(Clone)]
struct CoordinatorClient {
http: reqwest::Client,
@@ -179,18 +185,42 @@ impl CoordinatorClient {
.await
}
async fn submit(&self, envelope: &MetadataEnvelope) -> Result<MetadataSubmitResponse, String> {
let bytes = rmp_serde::to_vec_named(envelope).map_err(|error| error.to_string())?;
async fn submit_batch(
&self,
node_id: &str,
envelopes: Vec<MetadataEnvelope>,
) -> Result<(), String> {
let item_count = envelopes.len();
let request = MetadataBatchRequest {
protocol_version: PROTOCOL_VERSION,
node_id: node_id.to_owned(),
envelopes,
};
let bytes = tokio::task::spawn_blocking(move || {
let bytes = rmp_serde::to_vec_named(&request).map_err(|error| error.to_string())?;
zstd::stream::encode_all(bytes.as_slice(), 1).map_err(|error| error.to_string())
})
.await
.map_err(|error| error.to_string())??;
let response = self
.http
.post(format!("{}/internal/metadata", self.base_url))
.post(format!("{}/internal/metadata/batch", self.base_url))
.timeout(METADATA_UPLOAD_TIMEOUT)
.bearer_auth(self.token.as_ref())
.header(reqwest::header::CONTENT_TYPE, "application/msgpack")
.header(reqwest::header::CONTENT_ENCODING, "zstd")
.body(bytes)
.send()
.await
.map_err(|error| error.to_string())?;
decode_response(response).await
let response: MetadataBatchResponse = decode_response(response).await?;
if response.results.len() != item_count {
return Err(format!(
"协调器返回的 Metadata 批次结果数量不一致: 预期 {item_count} 实际 {}",
response.results.len()
));
}
Ok(())
}
async fn claim_verification(
@@ -290,42 +320,99 @@ async fn heartbeat_loop(
}
async fn sender_loop(shared: Arc<CollectorShared>, cancel: CancellationToken) {
let mut ticker = tokio::time::interval(Duration::from_millis(100));
let mut ticker = tokio::time::interval(Duration::from_millis(50));
ticker.set_missed_tick_behavior(MissedTickBehavior::Skip);
let mut uploads: tokio::task::JoinSet<(Vec<super::outbox::OutboxItem>, Result<(), String>)> =
tokio::task::JoinSet::new();
let mut in_flight = HashSet::new();
loop {
tokio::select! {
_ = cancel.cancelled() => break,
_ = ticker.tick() => {
let outbox = shared.outbox.clone();
let now = unix_timestamp();
let item = match tokio::task::spawn_blocking(move || outbox.next_ready(now)).await {
Ok(Ok(item)) => item,
Ok(Err(error)) => {
set_error(&shared, error.to_string());
continue;
}
_ = cancel.cancelled() => {
uploads.abort_all();
break;
}
completed = uploads.join_next(), if !uploads.is_empty() => {
let Some(completed) = completed else { continue };
let (items, result) = match completed {
Ok(completed) => completed,
Err(error) => {
set_error(&shared, error.to_string());
uploads.abort_all();
while uploads.join_next().await.is_some() {}
in_flight.clear();
set_error(&shared, format!("Metadata 上传任务异常结束: {error}"));
continue;
}
};
let Some(item) = item else { continue };
match shared.client.submit(&item.envelope).await {
for item in &items {
in_flight.remove(&item.sequence);
}
match result {
Ok(_) => {
let outbox = shared.outbox.clone();
match tokio::task::spawn_blocking(move || outbox.acknowledge(&item)).await {
match tokio::task::spawn_blocking(move || outbox.acknowledge_batch(&items)).await {
Ok(Ok(())) => {}
Ok(Err(error)) => set_error(&shared, error.to_string()),
Err(error) => set_error(&shared, error.to_string()),
}
}
Err(error) => {
let first = items.first();
if let Some(first) = first
&& first.sequence % 128 == 0
{
tracing::warn!(
sequence = first.sequence,
attempts = first.attempts,
batch_items = items.len(),
%error,
"Metadata 批量上传协调器失败"
);
}
set_error(&shared, error);
let outbox = shared.outbox.clone();
let _ = tokio::task::spawn_blocking(move || outbox.retry(item, now)).await;
let now = unix_timestamp();
let _ = tokio::task::spawn_blocking(move || outbox.retry_batch(items, now)).await;
}
}
}
_ = ticker.tick() => {}
}
while uploads.len() < METADATA_UPLOAD_CONCURRENCY {
let outbox = shared.outbox.clone();
let excluded = in_flight.clone();
let now = unix_timestamp();
let items = match tokio::task::spawn_blocking(move || {
outbox.ready_batch_excluding(
now,
&excluded,
METADATA_UPLOAD_BATCH_ITEMS,
METADATA_UPLOAD_BATCH_BYTES,
)
})
.await
{
Ok(Ok(item)) => item,
Ok(Err(error)) => {
set_error(&shared, error.to_string());
break;
}
Err(error) => {
set_error(&shared, error.to_string());
break;
}
};
if items.is_empty() {
break;
}
in_flight.extend(items.iter().map(|item| item.sequence));
let client = shared.client.clone();
let node_id = shared.node_id.clone();
uploads.spawn(async move {
let envelopes = items.iter().map(|item| item.envelope.clone()).collect();
let result = client.submit_batch(&node_id, envelopes).await;
(items, result)
});
}
}
}
@@ -650,6 +737,8 @@ fn metrics(shared: &CollectorShared) -> CollectorMetrics {
cpu_time_millis: process.cpu_time_millis,
thread_count: process.thread_count,
handle_count: process.handle_count,
network_received_bytes: process.network_received_bytes,
network_transmitted_bytes: process.network_transmitted_bytes,
outbox_items: outbox.items,
outbox_bytes: outbox.bytes,
last_error: shared
@@ -680,3 +769,171 @@ fn unix_timestamp() -> u64 {
.unwrap_or_default()
.as_secs()
}
#[cfg(test)]
mod tests {
use std::sync::atomic::{AtomicUsize, Ordering};
use axum::{
Json, Router,
body::Bytes,
extract::State,
http::StatusCode,
response::{IntoResponse, Response},
routing::post,
};
use tempfile::TempDir;
use tokio::time::Instant;
use super::*;
use crate::cluster::protocol::{MetadataSubmitResponse, MetadataSubmitStatus};
#[derive(Clone)]
struct LimitedLink {
available_at: Arc<tokio::sync::Mutex<Instant>>,
requests: Arc<AtomicUsize>,
compressed_bytes: Arc<AtomicU64>,
}
async fn accept_batch(State(link): State<LimitedLink>, body: Bytes) -> Response {
let request_number = link.requests.fetch_add(1, Ordering::Relaxed) + 1;
link.compressed_bytes
.fetch_add(body.len() as u64, Ordering::Relaxed);
let decoded = zstd::stream::decode_all(body.as_ref()).unwrap();
let request: MetadataBatchRequest = rmp_serde::from_slice(&decoded).unwrap();
let transmission = Duration::from_secs_f64(body.len() as f64 / 162_500.0);
let now = Instant::now();
let mut available_at = link.available_at.lock().await;
let completed_at = (*available_at).max(now) + transmission;
*available_at = completed_at;
drop(available_at);
let loss_penalty = if request_number.is_multiple_of(5) {
Duration::from_millis(440)
} else {
Duration::ZERO
};
tokio::time::sleep_until(completed_at + Duration::from_millis(220) + loss_penalty).await;
if request_number.is_multiple_of(7) {
return (StatusCode::SERVICE_UNAVAILABLE, "simulated request loss").into_response();
}
Json(MetadataBatchResponse {
results: request
.envelopes
.iter()
.map(|_| MetadataSubmitResponse {
status: MetadataSubmitStatus::Inserted,
reason: None,
})
.collect(),
})
.into_response()
}
fn test_envelope(index: usize, file_count: usize) -> MetadataEnvelope {
let hash = format!("{index:040x}");
let files = (0..file_count)
.map(|file_index| {
let digest = blake3::hash(format!("{index}:{file_index}").as_bytes()).to_hex();
dht_crawler::FileInfo {
path: format!("release/{digest}/{digest}-{file_index:05}.bin"),
size: file_index as u64 + 1,
}
})
.collect();
MetadataEnvelope {
protocol_version: PROTOCOL_VERSION,
node_id: "limited-link-node".to_owned(),
lease_token: None,
torrent: dht_crawler::TorrentInfo {
info_hash: hash.clone(),
magnet_link: format!("magnet:?xt=urn:btih:{hash}"),
name: format!("limited-link-{index}"),
total_size: file_count as u64,
files,
piece_length: 16_384,
peers: Vec::new(),
timestamp: 1,
},
}
}
#[tokio::test]
#[ignore = "手动验证高延迟限速链路上的批量上传吞吐"]
async fn compressed_batches_drain_a_mixed_outbox_over_a_limited_link() {
let link = LimitedLink {
available_at: Arc::new(tokio::sync::Mutex::new(Instant::now())),
requests: Arc::new(AtomicUsize::new(0)),
compressed_bytes: Arc::new(AtomicU64::new(0)),
};
let app = Router::new()
.route("/internal/metadata/batch", post(accept_batch))
.with_state(link.clone());
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let address = listener.local_addr().unwrap();
let server_cancel = CancellationToken::new();
let server_shutdown = server_cancel.clone();
let server = tokio::spawn(async move {
axum::serve(listener, app)
.with_graceful_shutdown(server_shutdown.cancelled_owned())
.await
});
let directory = TempDir::new().unwrap();
let outbox = Arc::new(Outbox::open(directory.path(), 1_000, 64 * 1024 * 1024).unwrap());
for index in 1..=200 {
outbox.append(test_envelope(index, 12)).unwrap();
}
for index in 201..=208 {
outbox.append(test_envelope(index, 4_000)).unwrap();
}
for index in 209..=210 {
outbox.append(test_envelope(index, 20_000)).unwrap();
}
let initial = outbox.snapshot();
let shared = Arc::new(CollectorShared {
node_id: "limited-link-node".to_owned(),
client: CoordinatorClient::new(&format!("http://{address}"), "token".to_owned())
.unwrap(),
outbox: outbox.clone(),
leases: Mutex::new(HashMap::new()),
stats: RwLock::new(DhtRuntimeStats::default()),
server: tokio::sync::RwLock::new(None),
applied_revision: AtomicU64::new(0),
connected: AtomicBool::new(true),
capacity_paused: AtomicBool::new(false),
last_error: RwLock::new(None),
});
let cancel = CancellationToken::new();
let started = Instant::now();
let sender = tokio::spawn(sender_loop(shared, cancel.clone()));
tokio::time::timeout(Duration::from_secs(60), async {
while outbox.snapshot().items != 0 {
tokio::time::sleep(Duration::from_millis(50)).await;
}
})
.await
.unwrap();
let elapsed = started.elapsed();
cancel.cancel();
sender.await.unwrap();
server_cancel.cancel();
server.await.unwrap().unwrap();
let compressed = link.compressed_bytes.load(Ordering::Relaxed);
eprintln!(
"limited-link items={} raw_bytes={} compressed_bytes={} requests={} elapsed_ms={}",
initial.items,
initial.bytes,
compressed,
link.requests.load(Ordering::Relaxed),
elapsed.as_millis()
);
assert_eq!(outbox.snapshot().items, 0);
assert!(compressed < initial.bytes / 2);
assert!(elapsed < Duration::from_secs(45));
}
}
+132 -3
View File
@@ -2,6 +2,7 @@
use std::{
collections::HashMap,
io::Read,
net::SocketAddr,
str::FromStr,
sync::{Arc, Mutex, RwLock},
@@ -31,9 +32,9 @@ use super::{
protocol::{
AdmissionDecision, AdmissionRequest, AdmissionResponse, AdmissionStatus,
CollectorConfigUpdate, CollectorMetrics, CollectorView, ControlResponse, HeartbeatRequest,
MetadataEnvelope, MetadataSubmitResponse, MetadataSubmitStatus, PROTOCOL_VERSION,
RegisterRequest, VerificationClaimRequest, VerificationClaimResponse, VerificationJob,
VerificationResultRequest,
MetadataBatchRequest, MetadataBatchResponse, MetadataEnvelope, MetadataSubmitResponse,
MetadataSubmitStatus, PROTOCOL_VERSION, RegisterRequest, VerificationClaimRequest,
VerificationClaimResponse, VerificationJob, VerificationResultRequest,
},
registry::{CollectorRecord, CollectorRegistry},
};
@@ -180,6 +181,14 @@ impl CoordinatorRuntime {
sum_optional(aggregate.cpu_time_millis, metrics.cpu_time_millis);
aggregate.thread_count = sum_optional(aggregate.thread_count, metrics.thread_count);
aggregate.handle_count = sum_optional(aggregate.handle_count, metrics.handle_count);
aggregate.network_received_bytes = sum_optional(
aggregate.network_received_bytes,
metrics.network_received_bytes,
);
aggregate.network_transmitted_bytes = sum_optional(
aggregate.network_transmitted_bytes,
metrics.network_transmitted_bytes,
);
}
aggregate.metadata_succeeded = aggregate
.metadata_succeeded
@@ -325,6 +334,7 @@ fn internal_router(runtime: CoordinatorRuntime) -> Router {
.route("/internal/heartbeat", post(heartbeat))
.route("/internal/admissions", post(admit))
.route("/internal/metadata", post(submit_metadata))
.route("/internal/metadata/batch", post(submit_metadata_batch))
.route("/internal/verification/claim", post(claim_verification))
.route("/internal/verification/result", post(finish_verification))
.layer(DefaultBodyLimit::max(16 * 1024 * 1024))
@@ -486,6 +496,75 @@ async fn submit_metadata(
Ok(Json(response))
}
async fn submit_metadata_batch(
State(runtime): State<CoordinatorRuntime>,
headers: HeaderMap,
body: Bytes,
) -> HandlerResult<MetadataBatchResponse> {
const MAX_BATCH_ITEMS: usize = 32;
const MAX_DECOMPRESSED_BYTES: u64 = 16 * 1024 * 1024;
authorize(&runtime, &headers)?;
let content_encoding = headers
.get(axum::http::header::CONTENT_ENCODING)
.and_then(|value| value.to_str().ok());
if content_encoding.is_some_and(|value| value != "zstd") {
return Err(bad_request("Metadata 批次只支持 zstd 内容编码"));
}
let compressed = content_encoding == Some("zstd");
let body = tokio::task::spawn_blocking(move || {
if compressed {
decode_zstd_limited(body.as_ref(), MAX_DECOMPRESSED_BYTES)
} else {
Ok(body.to_vec())
}
})
.await
.map_err(|error| internal(error.to_string()))?
.map_err(bad_request)?;
let request: MetadataBatchRequest =
rmp_serde::from_slice(&body).map_err(|error| bad_request(error.to_string()))?;
validate_version(request.protocol_version)?;
ensure_registered(&runtime, &request.node_id)?;
if request.envelopes.is_empty() || request.envelopes.len() > MAX_BATCH_ITEMS {
return Err(bad_request(format!(
"Metadata 批次条目数量必须在 1 到 {MAX_BATCH_ITEMS} 之间"
)));
}
for envelope in &request.envelopes {
validate_version(envelope.protocol_version)?;
if envelope.node_id != request.node_id {
return Err(bad_request("Metadata 批次包含其他采集器的数据"));
}
}
let state = runtime.inner.clone();
let results = tokio::task::spawn_blocking(move || {
request
.envelopes
.into_iter()
.map(|envelope| ingest_metadata(&state, envelope))
.collect::<Result<Vec<_>, _>>()
})
.await
.map_err(|error| internal(error.to_string()))?
.map_err(|error| internal(error.to_string()))?;
Ok(Json(MetadataBatchResponse { results }))
}
fn decode_zstd_limited(body: &[u8], max_bytes: u64) -> Result<Vec<u8>, String> {
let decoder = zstd::stream::read::Decoder::new(body).map_err(|error| error.to_string())?;
let mut decoded = Vec::new();
decoder
.take(max_bytes.saturating_add(1))
.read_to_end(&mut decoded)
.map_err(|error| error.to_string())?;
if decoded.len() as u64 > max_bytes {
return Err("Metadata 批次解压后超过 16 MiB".to_owned());
}
Ok(decoded)
}
fn ingest_metadata(
state: &CoordinatorState,
mut envelope: MetadataEnvelope,
@@ -902,6 +981,52 @@ mod tests {
let info_hash = InfoHash::from_str(&hash).unwrap();
let record = repository.get(info_hash).unwrap().unwrap();
assert_eq!(record.seen_count, 1);
let second_hash = "02".repeat(20);
let mut second_envelope = envelope.clone();
second_envelope.torrent.info_hash = second_hash.clone();
second_envelope.torrent.magnet_link = format!("magnet:?xt=urn:btih:{second_hash}");
let batch = MetadataBatchRequest {
protocol_version: PROTOCOL_VERSION,
node_id: NODE_ID.to_owned(),
envelopes: vec![envelope, second_envelope],
};
let batch = rmp_serde::to_vec_named(&batch).unwrap();
let batch = zstd::stream::encode_all(batch.as_slice(), 1).unwrap();
let response = app
.oneshot(
Request::post("/internal/metadata/batch")
.header("authorization", format!("Bearer {TOKEN}"))
.header("content-type", "application/msgpack")
.header("content-encoding", "zstd")
.body(Body::from(batch))
.unwrap(),
)
.await
.unwrap();
let body = to_bytes(response.into_body(), usize::MAX).await.unwrap();
let submitted: MetadataBatchResponse = serde_json::from_slice(&body).unwrap();
assert_eq!(submitted.results.len(), 2);
assert_eq!(
submitted.results[0].status,
MetadataSubmitStatus::AlreadyKnown
);
assert_eq!(submitted.results[1].status, MetadataSubmitStatus::Inserted);
assert!(
repository
.get(InfoHash::from_str(&second_hash).unwrap())
.unwrap()
.is_some()
);
}
#[test]
fn compressed_batch_decoder_enforces_the_expanded_size_limit() {
let source = vec![7_u8; 1_025];
let compressed = zstd::stream::encode_all(source.as_slice(), 1).unwrap();
assert_eq!(decode_zstd_limited(&compressed, 1_025).unwrap(), source);
assert!(decode_zstd_limited(&compressed, 1_024).is_err());
}
#[tokio::test]
@@ -945,6 +1070,8 @@ mod tests {
metadata_succeeded: succeeded,
metadata_failure_timeout: timeout,
resident_memory_bytes: Some(100),
network_received_bytes: Some(1_000),
network_transmitted_bytes: Some(500),
..CollectorMetrics::default()
},
);
@@ -955,5 +1082,7 @@ mod tests {
assert_eq!(aggregate.metadata_succeeded, 8);
assert_eq!(aggregate.metadata_failure_timeout, 16);
assert_eq!(aggregate.resident_memory_bytes, Some(200));
assert_eq!(aggregate.network_received_bytes, Some(2_000));
assert_eq!(aggregate.network_transmitted_bytes, Some(1_000));
}
}
+110 -18
View File
@@ -1,6 +1,6 @@
// 负责持久保存采集器尚未送达的 Metadata 并实施字节和条目双重上限
use std::{path::Path, sync::Mutex};
use std::{collections::HashSet, path::Path, sync::Mutex};
use rocksdb::{DB, Direction, IteratorMode, WriteBatch};
use serde::{Deserialize, Serialize};
@@ -155,7 +155,15 @@ impl Outbox {
Ok(())
}
pub(crate) fn next_ready(&self, now: u64) -> Result<Option<OutboxItem>, OutboxError> {
pub(crate) fn ready_batch_excluding(
&self,
now: u64,
excluded: &HashSet<u64>,
max_items: usize,
max_bytes: u64,
) -> Result<Vec<OutboxItem>, OutboxError> {
let mut ready = Vec::with_capacity(max_items);
let mut bytes = 0_u64;
for entry in self
.db
.iterator(IteratorMode::From(&[ITEM_PREFIX], Direction::Forward))
@@ -165,31 +173,54 @@ impl Outbox {
break;
}
let item: OutboxItem = rmp_serde::from_slice(&value)?;
if item.next_attempt_at <= now {
return Ok(Some(item));
if item.next_attempt_at <= now && !excluded.contains(&item.sequence) {
if !ready.is_empty() && bytes.saturating_add(item.encoded_bytes) > max_bytes {
break;
}
bytes = bytes.saturating_add(item.encoded_bytes);
ready.push(item);
if ready.len() >= max_items || bytes >= max_bytes {
break;
}
}
}
Ok(None)
Ok(ready)
}
pub(crate) fn acknowledge(&self, item: &OutboxItem) -> Result<(), OutboxError> {
pub(crate) fn acknowledge_batch(&self, items: &[OutboxItem]) -> Result<(), OutboxError> {
if items.is_empty() {
return Ok(());
}
let mut state = self
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
self.db.delete(item_key(item.sequence))?;
state.items = state.items.saturating_sub(1);
state.bytes = state.bytes.saturating_sub(item.encoded_bytes);
let mut batch = WriteBatch::default();
let mut bytes = 0_u64;
for item in items {
batch.delete(item_key(item.sequence));
bytes = bytes.saturating_add(item.encoded_bytes);
}
self.db.write(batch)?;
state.items = state.items.saturating_sub(items.len() as u64);
state.bytes = state.bytes.saturating_sub(bytes);
Ok(())
}
pub(crate) fn retry(&self, mut item: OutboxItem, now: u64) -> Result<(), OutboxError> {
item.attempts = item.attempts.saturating_add(1);
let exponent = item.attempts.min(6);
let delay = 1_u64 << exponent;
item.next_attempt_at = now.saturating_add(delay.min(60));
self.db
.put(item_key(item.sequence), rmp_serde::to_vec_named(&item)?)?;
pub(crate) fn retry_batch(
&self,
mut items: Vec<OutboxItem>,
now: u64,
) -> Result<(), OutboxError> {
let mut batch = WriteBatch::default();
for item in &mut items {
item.attempts = item.attempts.saturating_add(1);
let exponent = item.attempts.min(6);
let delay = 1_u64 << exponent;
item.next_attempt_at = now.saturating_add(delay.min(60));
batch.put(item_key(item.sequence), rmp_serde::to_vec_named(&*item)?);
}
self.db.write(batch)?;
Ok(())
}
@@ -265,11 +296,72 @@ mod tests {
let outbox = Outbox::open(directory.path(), 10, 1_000_000).unwrap();
assert_eq!(outbox.identity().unwrap(), identity);
assert_eq!(outbox.snapshot().items, 1);
let item = outbox.next_ready(1).unwrap().unwrap();
outbox.acknowledge(&item).unwrap();
let items = outbox
.ready_batch_excluding(1, &HashSet::new(), 10, 1_000_000)
.unwrap();
outbox.acknowledge_batch(&items).unwrap();
assert_eq!(outbox.snapshot().items, 0);
}
#[test]
fn ready_item_can_exclude_an_upload_in_progress() {
let directory = TempDir::new().unwrap();
let outbox = Outbox::open(directory.path(), 10, 1_000_000).unwrap();
let first = outbox.append(envelope(&"01".repeat(20))).unwrap();
let second = outbox.append(envelope(&"02".repeat(20))).unwrap();
let items = outbox
.ready_batch_excluding(1, &HashSet::from([first]), 10, 1_000_000)
.unwrap();
assert_eq!(items.len(), 1);
assert_eq!(items[0].sequence, second);
}
#[test]
fn ready_batch_respects_item_and_byte_limits() {
let directory = TempDir::new().unwrap();
let outbox = Outbox::open(directory.path(), 10, 1_000_000).unwrap();
outbox.append(envelope(&"01".repeat(20))).unwrap();
outbox.append(envelope(&"02".repeat(20))).unwrap();
outbox.append(envelope(&"03".repeat(20))).unwrap();
let by_items = outbox
.ready_batch_excluding(1, &HashSet::new(), 2, 1_000_000)
.unwrap();
assert_eq!(by_items.len(), 2);
let first_bytes = by_items[0].encoded_bytes;
let by_bytes = outbox
.ready_batch_excluding(1, &HashSet::new(), 10, first_bytes)
.unwrap();
assert_eq!(by_bytes.len(), 1);
}
#[test]
fn retry_updates_a_whole_batch_atomically() {
let directory = TempDir::new().unwrap();
let outbox = Outbox::open(directory.path(), 10, 1_000_000).unwrap();
outbox.append(envelope(&"01".repeat(20))).unwrap();
outbox.append(envelope(&"02".repeat(20))).unwrap();
let items = outbox
.ready_batch_excluding(1, &HashSet::new(), 10, 1_000_000)
.unwrap();
outbox.retry_batch(items, 10).unwrap();
assert!(
outbox
.ready_batch_excluding(11, &HashSet::new(), 10, 1_000_000)
.unwrap()
.is_empty()
);
let retried = outbox
.ready_batch_excluding(12, &HashSet::new(), 10, 1_000_000)
.unwrap();
assert_eq!(retried.len(), 2);
assert!(retried.iter().all(|item| item.attempts == 1));
}
#[test]
fn capacity_is_enforced_before_write() {
let directory = TempDir::new().unwrap();
+14
View File
@@ -55,6 +55,8 @@ pub(crate) struct CollectorMetrics {
pub(crate) cpu_time_millis: Option<u64>,
pub(crate) thread_count: Option<u64>,
pub(crate) handle_count: Option<u64>,
pub(crate) network_received_bytes: Option<u64>,
pub(crate) network_transmitted_bytes: Option<u64>,
pub(crate) outbox_items: u64,
pub(crate) outbox_bytes: u64,
pub(crate) last_error: Option<String>,
@@ -119,6 +121,18 @@ pub(crate) struct MetadataSubmitResponse {
pub(crate) reason: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) struct MetadataBatchRequest {
pub(crate) protocol_version: u16,
pub(crate) node_id: String,
pub(crate) envelopes: Vec<MetadataEnvelope>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub(crate) struct MetadataBatchResponse {
pub(crate) results: Vec<MetadataSubmitResponse>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub(crate) struct VerificationClaimRequest {
pub(crate) protocol_version: u16,
+2
View File
@@ -33,6 +33,8 @@ pub(crate) struct ProcessDiagnostics {
pub(crate) cpu_time_millis: Option<u64>,
pub(crate) thread_count: Option<u64>,
pub(crate) handle_count: Option<u64>,
pub(crate) network_received_bytes: Option<u64>,
pub(crate) network_transmitted_bytes: Option<u64>,
}
#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, PartialEq, Eq)]
+44 -1
View File
@@ -1,4 +1,4 @@
// 负责采集当前进程内存 CPU 线程和句柄资源快照
// 负责采集当前进程资源以及容器网络命名空间的累计流量快照
use super::model::ProcessDiagnostics;
@@ -59,6 +59,8 @@ pub(crate) fn snapshot() -> ProcessDiagnostics {
.then(|| filetime_ticks(kernel).saturating_add(filetime_ticks(user)) / 10_000),
thread_count: windows_thread_count(),
handle_count: handles_ok.then_some(u64::from(handles)),
network_received_bytes: None,
network_transmitted_bytes: None,
};
fn filetime_ticks(value: FILETIME) -> u64 {
@@ -95,6 +97,7 @@ pub(crate) fn snapshot() -> ProcessDiagnostics {
#[cfg(target_os = "linux")]
pub(crate) fn snapshot() -> ProcessDiagnostics {
let status = std::fs::read_to_string("/proc/self/status").unwrap_or_default();
let network = linux_container_network_bytes();
ProcessDiagnostics {
resident_memory_bytes: status_kib(&status, "VmRSS:"),
private_memory_bytes: status_kib(&status, "RssAnon:")
@@ -104,9 +107,42 @@ pub(crate) fn snapshot() -> ProcessDiagnostics {
handle_count: std::fs::read_dir("/proc/self/fd")
.ok()
.map(|entries| entries.count().min(u64::MAX as usize) as u64),
network_received_bytes: network.map(|value| value.0),
network_transmitted_bytes: network.map(|value| value.1),
}
}
#[cfg(target_os = "linux")]
fn linux_container_network_bytes() -> Option<(u64, u64)> {
if !std::path::Path::new("/.dockerenv").exists() {
return None;
}
parse_linux_network_bytes(&std::fs::read_to_string("/proc/self/net/dev").ok()?)
}
#[cfg(any(target_os = "linux", test))]
fn parse_linux_network_bytes(value: &str) -> Option<(u64, u64)> {
let mut received = 0_u64;
let mut transmitted = 0_u64;
let mut interfaces = 0_u64;
for line in value.lines() {
let Some((name, counters)) = line.split_once(':') else {
continue;
};
if name.trim() == "lo" {
continue;
}
let fields: Vec<_> = counters.split_whitespace().collect();
let (Some(rx), Some(tx)) = (fields.first(), fields.get(8)) else {
continue;
};
received = received.saturating_add(rx.parse::<u64>().ok()?);
transmitted = transmitted.saturating_add(tx.parse::<u64>().ok()?);
interfaces = interfaces.saturating_add(1);
}
(interfaces > 0).then_some((received, transmitted))
}
#[cfg(target_os = "linux")]
fn status_kib(status: &str, key: &str) -> Option<u64> {
status_value(status, key).map(|value| value.saturating_mul(1_024))
@@ -164,4 +200,11 @@ mod tests {
assert!(snapshot.handle_count.is_some_and(|value| value > 0));
}
}
#[test]
fn container_network_parser_ignores_loopback_and_sums_interfaces() {
let value = "Inter-| Receive | Transmit\n lo: 10 0 0 0 0 0 0 0 20 0 0 0 0 0 0 0\n eth0: 100 0 0 0 0 0 0 0 200 0 0 0 0 0 0 0\n eth1: 3 0 0 0 0 0 0 0 4 0 0 0 0 0 0 0\n";
assert_eq!(super::parse_linux_network_bytes(value), Some((103, 204)));
}
}
+16 -1
View File
@@ -73,6 +73,8 @@ function selectedProcess(item: DiagnosticSample) {
cpu_time_millis: collector.metrics.cpu_time_millis,
thread_count: collector.metrics.thread_count,
handle_count: collector.metrics.handle_count,
network_received_bytes: collector.metrics.network_received_bytes,
network_transmitted_bytes: collector.metrics.network_transmitted_bytes,
}
}
@@ -164,6 +166,19 @@ const metadataPeerSummary = computed(() => {
return `成功率 ${successRate.toFixed(2)}%`
})
function networkMetric(field: 'network_received_bytes' | 'network_transmitted_bytes'): number | null {
if (selectedCollector.value) return selectedCollector.value.metrics[field]
const values = collectors.value.map((collector) => collector.metrics[field]).filter((value): value is number => typeof value === 'number')
return values.length ? values.reduce((sum, value) => sum + value, 0) : null
}
const networkSummary = computed(() => {
const received = networkMetric('network_received_bytes')
const transmitted = networkMetric('network_transmitted_bytes')
if (received === null || transmitted === null) return '本次运行流量未知'
return `本次运行接收 ${formatBytes(received)} · 发送 ${formatBytes(transmitted)}`
})
const torrentChart = computed(() => chartOption([
{ name: '已保存种子', values: samples.value.map((item) => item.storage.stored_torrents) },
], '个', 0, true))
@@ -352,7 +367,7 @@ onBeforeUnmount(() => {
<div class="grid gap-4 lg:grid-cols-2">
<article class="chart-card"><div class="chart-card-header"><h2>进程内存</h2><p>{{ visibleProcess?.thread_count ?? '未知' }} 个线程 · {{ visibleProcess?.handle_count ?? '未知' }} 个句柄</p></div><VChart class="diagnostic-chart" :option="memoryChart" autoresize /></article>
<article v-if="!selectedCollector" class="chart-card"><div class="chart-card-header"><h2>RocksDB 内存与压缩</h2><p>SST {{ sample.storage.live_sst_bytes === null ? '未知' : formatBytes(sample.storage.live_sst_bytes) }}</p></div><VChart class="diagnostic-chart" :option="storageChart" autoresize /></article>
<article class="chart-card"><div class="chart-card-header"><h2>网络与 Peer 下载吞吐</h2></div><VChart class="diagnostic-chart" :option="throughputChart" autoresize /></article>
<article class="chart-card"><div class="chart-card-header"><h2>网络与 Peer 下载吞吐</h2><p>{{ networkSummary }}</p></div><VChart class="diagnostic-chart" :option="throughputChart" autoresize /></article>
<article class="chart-card" :class="selectedCollector ? 'lg:col-span-2' : ''"><div class="chart-card-header"><h2>Peer 失败原因构成</h2><p>{{ metadataPeerSummary }}</p></div><VChart class="diagnostic-chart" :option="metadataFailureChart" autoresize /></article>
<article v-if="!selectedCollector" class="chart-card lg:col-span-2"><div class="chart-card-header"><h2>种子数量趋势</h2><p>{{ sample.storage.stored_torrents === null ? '等待新采样' : `当前 ${sample.storage.stored_torrents.toLocaleString()} 个` }}</p></div><VChart class="diagnostic-chart" :option="torrentChart" autoresize /></article>
</div>
+4
View File
@@ -179,6 +179,8 @@ export interface ProcessDiagnostics {
cpu_time_millis: number | null
thread_count: number | null
handle_count: number | null
network_received_bytes: number | null
network_transmitted_bytes: number | null
}
export interface StorageDiagnostics {
@@ -408,6 +410,8 @@ export interface CollectorMetrics {
cpu_time_millis: number | null
thread_count: number | null
handle_count: number | null
network_received_bytes: number | null
network_transmitted_bytes: number | null
outbox_items: number
outbox_bytes: number
last_error: string | null