feat: 支持即时切换 DHT 离线模式

This commit is contained in:
chuan
2026-08-10 20:24:37 +08:00
parent 368444fe2b
commit db0130e606
22 changed files with 585 additions and 190 deletions
+3 -1
View File
@@ -2,7 +2,7 @@
这是一个用于持续发现持久化索引和搜索 BitTorrent DHT 元数据的 Rust 服务
当前已经具备 DHT 采集 Metadata 下载 RocksDB 精确去重 Tantivy 全文搜索 内容聚合 可用性验证 HTTP API Web 搜索界面 运行诊断 配置管理 备份恢复和 Docker 部署能力
当前已经具备可即时启停的 DHT 采集 Metadata 下载 RocksDB 精确去重 Tantivy 全文搜索 内容聚合 可用性验证 HTTP API Web 搜索界面 运行诊断 配置管理 备份恢复和 Docker 部署能力
## 项目结构
@@ -17,6 +17,8 @@ RocksDB 是唯一权威数据源 Tantivy 索引可以从 RocksDB 完整重建
应用全部配置和内容隐藏规则统一位于 [`config.toml`](config.toml)
Web 右上角的无线电图标可以即时停止或恢复 DHT 持续采集 状态会写回 `dht.enabled` 关闭后不会建立 DHT 和 Metadata 网络任务 但现有 RocksDB 数据仍会继续补建索引并提供本地搜索
## 本地运行
Windows 可以在仓库根目录执行
+1
View File
@@ -182,6 +182,7 @@
- [x] 统一搜索结果数字分页并支持用户选择每页数量
- [x] 将采集索引持久化和验证运行状态集中到系统诊断页并仅在页面打开时每秒刷新
- [x] 将诊断时间范围放入历史趋势区域并精简图表卡片的单行摘要信息
- [x] 在 Web 顶栏提供持久化的 DHT 即时启停开关并保留离线索引搜索能力
- [x] 实现浏览器持久化明暗主题
- [x] 为加载空结果接口错误和失败重试提供明确界面状态
- [x] 将生产静态资源交给 Axum 提供并支持单页回退
+1
View File
@@ -77,6 +77,7 @@ action = "hide"
reason = "libtorrent 分片边界填充目录"
[dht]
enabled = false
port = 12313
netmode = "ipv4-only"
hash_queue_capacity = 20000
+3 -1
View File
@@ -167,7 +167,9 @@ SQLite 使用 WAL 和单独 writer 线程 原始采样超过 24 小时自动删
更新前会执行与启动时相同的完整校验 保存时先在同目录写入并同步临时文件再原子替换目标文件 旧修订号返回 `409 Conflict` 防止多个页面互相覆盖
当前版本不在线修改正在运行的 DHT 存储索引和监听器 保存成功后返回 `restart_required = true` 重启服务后统一生效 命令行覆盖字段也会单独返回并继续优先于文件配置
`dht.enabled`当前版本不在线修改正在运行的 DHT 存储索引和监听器 保存其他配置后返回 `restart_required = true` 并在重启服务后统一生效 命令行覆盖字段也会单独返回并继续优先于文件配置
Web 右上角采集开关通过 `/crawler` 即时停止或重新创建 DHT 运行时并把状态持久化到 `dht.enabled` 关闭时同步停止 Metadata 下载和按需有效性验证 RocksDB Tantivy HTTP 和 Web 保持运行
通过 Web 保存会按 DTO 重新生成 TOML 原有手写注释不会保留 管理接口默认随 HTTP 服务提供 因此生产部署不应把 `/config` 暴露到不受信任的公网入口
+51 -9
View File
@@ -19,10 +19,13 @@ use axum::{
use super::{
ApiState,
request::{ContentVariantsRequest, DiagnosticHistoryRequest, SearchRequest, TorrentRequest},
request::{
ContentVariantsRequest, CrawlerUpdateRequest, DiagnosticHistoryRequest, SearchRequest,
TorrentRequest,
},
response::{
ContentVariantsResponse, ErrorResponse, StatsResponse, StatusResponse, TorrentResponse,
TorrentVariantResponse,
ContentVariantsResponse, CrawlerStatusResponse, ErrorResponse, StatsResponse,
StatusResponse, TorrentResponse, TorrentVariantResponse,
},
};
use crate::config::{ConfigServiceError, ConfigSnapshot, ConfigUpdateRequest};
@@ -37,16 +40,17 @@ pub(crate) async fn ready() -> Json<StatusResponse> {
}
pub(crate) async fn stats(State(state): State<ApiState>) -> Json<StatsResponse> {
let dht = state.dht_stats.snapshot();
let observability = state.dht_stats.observability_snapshot();
let dht_stats = state.crawler.stats();
let dht = dht_stats.snapshot();
let observability = dht_stats.observability_snapshot();
let persistence = state.persistence.snapshot();
let disk = state.disk_guard.snapshot();
let backup = state.backup_stats.snapshot();
let http = state.http_stats.snapshot();
let filtered = persistence.filtered;
let verification = state
.verification
.as_ref()
.crawler
.verification()
.map(|ingress| ingress.stats().snapshot());
Json(StatsResponse {
http_active_requests: http.active_requests,
@@ -152,6 +156,44 @@ pub(crate) async fn stats(State(state): State<ApiState>) -> Json<StatsResponse>
})
}
pub(crate) async fn crawler_status(State(state): State<ApiState>) -> Json<CrawlerStatusResponse> {
let status = state.crawler.status();
Json(CrawlerStatusResponse {
enabled: status.enabled,
transitioning: status.transitioning,
})
}
pub(crate) async fn crawler_update(
State(state): State<ApiState>,
Json(request): Json<CrawlerUpdateRequest>,
) -> Result<Json<CrawlerStatusResponse>, ApiError> {
let previous = state.crawler.status().enabled;
if previous != request.enabled {
state
.crawler
.set_enabled(request.enabled)
.await
.map_err(|error| ApiError::internal(format!("切换 DHT 采集状态失败: {error}")))?;
}
let config = state.config.clone();
let enabled = request.enabled;
let persisted = tokio::task::spawn_blocking(move || config.set_dht_enabled(enabled))
.await
.map_err(|error| ApiError::internal(format!("采集状态保存任务失败: {error}")))?;
if let Err(error) = persisted {
if previous != request.enabled {
let _ = state.crawler.set_enabled(previous).await;
}
return Err(ApiError::internal(error.to_string()));
}
let status = state.crawler.status();
Ok(Json(CrawlerStatusResponse {
enabled: status.enabled,
transitioning: status.transitioning,
}))
}
pub(crate) async fn diagnostics_current(
State(state): State<ApiState>,
) -> Json<CurrentDiagnosticsResponse> {
@@ -208,7 +250,7 @@ pub(crate) async fn search(
State(state): State<ApiState>,
Query(request): Query<SearchRequest>,
) -> Result<Json<SearchPage>, ApiError> {
let verification = state.verification.clone();
let verification = state.crawler.verification();
if request.q.len() > 512 {
return Err(ApiError::bad_request("查询文本不能超过 512 字节"));
}
@@ -351,7 +393,7 @@ pub(crate) async fn torrent(
if request.file_offset > 1_000_000 {
return Err(ApiError::bad_request("file_offset 不能超过 1000000"));
}
let verification = state.verification.clone();
let verification = state.crawler.verification();
let info_hash =
InfoHash::from_str(&info_hash).map_err(|error| ApiError::bad_request(error.to_string()))?;
let record = tokio::task::spawn_blocking(move || state.repository.get_visible(info_hash))
+33 -8
View File
@@ -14,26 +14,23 @@ use axum::{
response::Response,
routing::get,
};
use dht_crawler::DhtRuntimeStats;
use tokio_util::sync::CancellationToken;
use tower_http::services::{ServeDir, ServeFile};
use crate::{
backup::BackupStats,
config::ConfigService,
crawler::pipeline::PersistenceIngress,
crawler::{pipeline::PersistenceIngress, runtime::CrawlerRuntime},
diagnostics::{DiagnosticsHandle, HttpStats},
disk_guard::DiskGuard,
verification::VerificationIngress,
};
#[derive(Clone)]
pub(crate) struct ApiState {
pub(crate) repository: Arc<dyn TorrentRepository>,
pub(crate) search: SearchEngine,
pub(crate) dht_stats: DhtRuntimeStats,
pub(crate) crawler: CrawlerRuntime,
pub(crate) persistence: PersistenceIngress,
pub(crate) verification: Option<VerificationIngress>,
pub(crate) disk_guard: DiskGuard,
pub(crate) backup_stats: BackupStats,
pub(crate) diagnostics: DiagnosticsHandle,
@@ -62,6 +59,10 @@ fn router(state: ApiState, web_dir: PathBuf) -> Router {
.route("/health", get(handlers::health))
.route("/ready", get(handlers::ready))
.route("/stats", get(handlers::stats))
.route(
"/crawler",
get(handlers::crawler_status).put(handlers::crawler_update),
)
.route("/diagnostics/current", get(handlers::diagnostics_current))
.route("/diagnostics/history", get(handlers::diagnostics_history))
.route(
@@ -96,7 +97,7 @@ mod tests {
body::{Body, to_bytes},
http::{Request, StatusCode},
};
use dht_crawler::DhtRuntimeStats;
use dht_crawler::DHTOptions;
use tempfile::TempDir;
use tower::ServiceExt;
@@ -154,6 +155,14 @@ mod tests {
disk_guard.clone(),
);
let verification = VerificationIngress::for_test(repository.clone(), 10);
let crawler = CrawlerRuntime::new(
DHTOptions::default(),
repository.clone(),
persistence.ingress.clone(),
disk_guard.clone(),
crate::config::VerificationConfig::default(),
);
crawler.set_verification_for_test(verification);
let web_dir = directory.path().join("web");
std::fs::create_dir_all(&web_dir).unwrap();
std::fs::write(web_dir.join("index.html"), "<main>DHT Search</main>").unwrap();
@@ -161,9 +170,8 @@ mod tests {
ApiState {
repository: repository_trait,
search,
dht_stats: DhtRuntimeStats::default(),
crawler,
persistence: persistence.ingress.clone(),
verification: Some(verification),
disk_guard,
backup_stats: BackupStats::default(),
diagnostics: DiagnosticsRuntime::disabled().handle(),
@@ -225,6 +233,23 @@ mod tests {
.is_some_and(|value| value >= 3)
);
let response = app
.clone()
.oneshot(
Request::builder()
.uri("/crawler")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::OK);
let json: serde_json::Value =
serde_json::from_slice(&to_bytes(response.into_body(), usize::MAX).await.unwrap())
.unwrap();
assert_eq!(json["enabled"], false);
assert_eq!(json["transitioning"], false);
let response = app
.clone()
.oneshot(
+5
View File
@@ -6,6 +6,11 @@ use crate::{
};
use serde::Deserialize;
#[derive(Debug, Deserialize)]
pub(crate) struct CrawlerUpdateRequest {
pub(crate) enabled: bool,
}
fn default_limit() -> usize {
20
}
+6
View File
@@ -14,6 +14,12 @@ pub(crate) struct ErrorResponse {
pub(crate) error: String,
}
#[derive(Debug, Serialize)]
pub(crate) struct CrawlerStatusResponse {
pub(crate) enabled: bool,
pub(crate) transitioning: bool,
}
#[derive(Debug, Serialize)]
pub(crate) struct StatsResponse {
pub(crate) http_active_requests: u64,
+19 -155
View File
@@ -1,28 +1,22 @@
// 负责连接采集存储索引和接口层并定义应用级启动顺序
use std::{
str::FromStr,
sync::Arc,
time::{Duration, SystemTime, UNIX_EPOCH},
};
use std::{sync::Arc, time::Duration};
use crate::{
domain::InfoHash,
search::SearchEngine,
storage::{RocksTorrentRepository, TorrentRepository},
};
use dht_crawler::DHTServer;
use tokio_util::sync::CancellationToken;
use crate::{
api::{self, ApiState},
backup::{self, BackupStats},
config::{AppConfig, ConfigService},
crawler::pipeline::PersistencePipeline,
crawler::{pipeline::PersistencePipeline, runtime::CrawlerRuntime},
diagnostics::{DiagnosticSources, DiagnosticsRuntime, HttpStats},
disk_guard::{self, DiskGuard},
error::AppError,
index_worker, monitor, shutdown, verification,
index_worker, monitor, shutdown,
};
pub(crate) async fn run(config: AppConfig, config_service: ConfigService) -> Result<(), AppError> {
@@ -73,121 +67,20 @@ pub(crate) async fn run(config: AppConfig, config_service: ConfigService) -> Res
(BackupStats::default(), None)
};
let options = config.dht_options();
let server = DHTServer::new(options.clone()).await?;
server.on_error(|error| tracing::error!(%error, "DHT 运行时错误"));
let sampled_repository = repository.clone();
let sampled_disk_guard = disk_guard.clone();
server.on_sampled_hashes(move |hashes| {
let repository = sampled_repository.clone();
let disk_guard = sampled_disk_guard.clone();
async move {
let Some(permit) = disk_guard.begin_admission() else {
return Vec::new();
};
let fallback = hashes.clone();
let info_hashes: Vec<_> = hashes.into_iter().map(InfoHash::from_bytes).collect();
let result = tokio::task::spawn_blocking(move || {
let _permit = permit;
repository.filter_unknown_and_observe(&info_hashes, unix_timestamp())
})
.await;
match result {
Ok(Ok(unknown)) => unknown
.into_iter()
.map(|info_hash| *info_hash.as_bytes())
.collect(),
Ok(Err(error)) => {
tracing::error!(%error, "采样 infohash 批量持久化去重失败");
fallback
}
Err(error) => {
tracing::error!(%error, "采样 infohash 批量去重任务异常");
fallback
}
}
}
});
let gate_repository = repository.clone();
let gate_disk_guard = disk_guard.clone();
server.on_metadata_fetch(move |hash| {
let repository = gate_repository.clone();
let disk_guard = gate_disk_guard.clone();
async move {
let Some(permit) = disk_guard.begin_admission() else {
return false;
};
let Ok(info_hash) = InfoHash::from_str(&hash) else {
tracing::warn!(%hash, "DHT 提供了无效 infohash");
return false;
};
let result = tokio::task::spawn_blocking(move || {
let _permit = permit;
repository.observe_existing(info_hash, unix_timestamp())
})
.await;
match result {
Ok(Ok(already_exists)) => !already_exists,
Ok(Err(error)) => {
tracing::error!(%error, %hash, "持久化去重查询失败");
true
}
Err(error) => {
tracing::error!(%error, %hash, "持久化去重任务异常");
true
}
}
}
});
let callback_ingress = ingress.clone();
server.on_torrent_with_ack(move |torrent| callback_ingress.try_enqueue(torrent));
server.on_metadata_fetch_complete(|completion| {
tracing::debug!(
info_hash = %completion.info_hash,
status = ?completion.status,
attempts = completion.attempts,
"Metadata 任务完成"
)
});
let verification_cancel = CancellationToken::new();
let (verification_fatal_tx, mut verification_fatal) = tokio::sync::oneshot::channel();
let mut verification_fatal_guard = None;
let (verification_ingress, verification_task) = if config.verification.enabled {
let (ingress, worker) = verification::start(
repository.clone(),
server.clone(),
config.verification.clone(),
disk_guard.clone(),
verification_cancel.clone(),
);
let cancel = verification_cancel.clone();
let task = tokio::spawn(async move {
let result = match worker.await {
Ok(result) => result,
Err(error) => Err(error.to_string()),
};
if !cancel.is_cancelled() {
let message = result
.as_ref()
.err()
.cloned()
.unwrap_or_else(|| "可用性验证 worker 意外停止".to_owned());
let _ = verification_fatal_tx.send(message);
}
result
});
(Some(ingress), Some(task))
} else {
verification_fatal_guard = Some(verification_fatal_tx);
(None, None)
};
let crawler = CrawlerRuntime::new(
config.dht_options(),
repository.clone(),
ingress.clone(),
disk_guard.clone(),
config.verification.clone(),
);
if config.dht.enabled {
crawler.set_enabled(true).await?;
}
tracing::info!(
dht_port = options.port,
dht_enabled = config.dht.enabled,
dht_port = config.dht.port,
data_dir = %config.data_dir.display(),
persistence_queue_capacity = config.persistence_queue_capacity,
"dht-search 启动"
@@ -195,7 +88,7 @@ pub(crate) async fn run(config: AppConfig, config_service: ConfigService) -> Res
let monitor_cancel = CancellationToken::new();
let monitor = tokio::spawn(monitor::run(
server.clone(),
crawler.stats_handle(),
ingress,
disk_guard.clone(),
config.stats_interval_secs,
@@ -222,7 +115,7 @@ pub(crate) async fn run(config: AppConfig, config_service: ConfigService) -> Res
DiagnosticSources {
repository: repository.clone(),
search: search.clone(),
dht: server.runtime_stats(),
dht: crawler.stats_handle(),
persistence: persistence.ingress.clone(),
disk_guard: disk_guard.clone(),
http: http_stats.clone(),
@@ -244,9 +137,8 @@ pub(crate) async fn run(config: AppConfig, config_service: ConfigService) -> Res
ApiState {
repository: repository.clone(),
search,
dht_stats: server.runtime_stats(),
crawler: crawler.clone(),
persistence: persistence.ingress.clone(),
verification: verification_ingress,
disk_guard: disk_guard.clone(),
backup_stats,
diagnostics: diagnostics.handle(),
@@ -265,7 +157,6 @@ pub(crate) async fn run(config: AppConfig, config_service: ConfigService) -> Res
tokio::pin!(run_duration);
let run_result = tokio::select! {
result = server.start() => result.map_err(AppError::from),
_ = shutdown::signal() => {
tracing::info!("收到退出信号");
Ok(())
@@ -282,10 +173,6 @@ pub(crate) async fn run(config: AppConfig, config_service: ConfigService) -> Res
let message = fatal.unwrap_or_else(|_| "索引 worker 意外停止".to_owned());
Err(AppError::IndexWorker(message))
}
fatal = &mut verification_fatal => {
let message = fatal.unwrap_or_else(|_| "可用性验证 worker 意外停止".to_owned());
Err(AppError::VerificationWorker(message))
}
result = &mut api_task => {
match result {
Ok(Ok(())) => Err(AppError::Config("HTTP 服务意外停止".to_owned())),
@@ -297,27 +184,11 @@ pub(crate) async fn run(config: AppConfig, config_service: ConfigService) -> Res
let mut shutdown_error = None;
diagnostics.request_shutdown();
verification_cancel.cancel();
backup_cancel.cancel();
server.shutdown();
crawler.shutdown().await;
if let Some(task) = backup_task {
let _ = task.await;
}
drop(verification_fatal_guard);
if let Some(task) = verification_task {
match task.await {
Ok(Ok(())) => {}
Ok(Err(error)) => {
remember_shutdown_error(&mut shutdown_error, AppError::VerificationWorker(error));
}
Err(error) => {
remember_shutdown_error(
&mut shutdown_error,
AppError::VerificationWorker(error.to_string()),
);
}
}
}
disk_cancel.cancel();
let _ = disk_task.await;
monitor_cancel.cancel();
@@ -376,10 +247,3 @@ fn remember_shutdown_error(slot: &mut Option<AppError>, error: AppError) {
tracing::error!(%error, "关闭阶段发生附加错误");
}
}
fn unix_timestamp() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs()
}
+2
View File
@@ -123,6 +123,7 @@ mod tests {
#[test]
fn defaults_produce_valid_dht_options() {
let config = resolve(AppConfigDto::default()).unwrap();
assert!(config.dht.enabled);
let options = config.dht_options();
assert_eq!(options.port, 12_313);
assert_eq!(options.metadata.max_worker_count, 400);
@@ -150,6 +151,7 @@ mod tests {
let decoded: AppConfigDto = toml::from_str(&encoded).unwrap();
assert_eq!(decoded.data_dir, dto.data_dir);
assert_eq!(decoded.dht.port, dto.dht.port);
assert_eq!(decoded.dht.enabled, dto.dht.enabled);
assert_eq!(decoded.http.listen, dto.http.listen);
assert_eq!(decoded.logging.file_prefix, dto.logging.file_prefix);
assert_eq!(decoded.diagnostics.database, dto.diagnostics.database);
+2
View File
@@ -31,6 +31,7 @@ pub(crate) struct AppConfigDto {
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(default, deny_unknown_fields)]
pub(crate) struct DhtConfig {
pub(crate) enabled: bool,
pub(crate) port: u16,
pub(crate) netmode: NetworkMode,
pub(crate) hash_queue_capacity: usize,
@@ -215,6 +216,7 @@ impl Default for MetadataLimitsConfig {
impl Default for DhtConfig {
fn default() -> Self {
Self {
enabled: true,
port: 12_313,
netmode: NetworkMode::Ipv4Only,
hash_queue_capacity: 20_000,
+48 -1
View File
@@ -123,10 +123,37 @@ impl ConfigService {
Ok(self.snapshot_from(&state))
}
pub(crate) fn set_dht_enabled(
&self,
enabled: bool,
) -> Result<ConfigSnapshot, ConfigServiceError> {
let mut state = self
.inner
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if state.config.dht.enabled == enabled {
return Ok(self.snapshot_from(&state));
}
let mut candidate = state.config.clone();
candidate.dht.enabled = enabled;
AppConfig::resolve(candidate.clone(), &self.inner.base)
.map_err(|error| ConfigServiceError::Validation(error.to_string()))?;
let candidate_revision = revision(&candidate)
.map_err(|error| ConfigServiceError::Persistence(error.to_string()))?;
self.inner
.store
.save(&candidate)
.map_err(|error| ConfigServiceError::Persistence(error.to_string()))?;
state.config = candidate;
state.revision = candidate_revision;
Ok(self.snapshot_from(&state))
}
fn snapshot_from(&self, state: &ConfigState) -> ConfigSnapshot {
ConfigSnapshot {
revision: state.revision.clone(),
restart_required: state.config != self.inner.startup,
restart_required: restart_required(&state.config, &self.inner.startup),
source: self.inner.store.path().to_string_lossy().into_owned(),
command_line_overrides: self.inner.command_line_overrides.clone(),
config: state.config.clone(),
@@ -134,6 +161,12 @@ impl ConfigService {
}
}
fn restart_required(current: &AppConfigDto, startup: &AppConfigDto) -> bool {
let mut comparable = current.clone();
comparable.dht.enabled = startup.dht.enabled;
comparable != *startup
}
fn revision(config: &AppConfigDto) -> Result<String, AppError> {
let canonical = toml::to_string(config)?;
Ok(blake3::hash(canonical.as_bytes()).to_hex().to_string())
@@ -176,6 +209,20 @@ mod tests {
assert_eq!(stored, Some(config));
}
#[test]
fn live_crawler_toggle_is_persisted_without_requiring_restart() {
let directory = TempDir::new().unwrap();
let service = service(&directory);
let updated = service.set_dht_enabled(false).unwrap();
assert!(!updated.config.dht.enabled);
assert!(!updated.restart_required);
let stored = TomlConfigStore::new(directory.path().join("service.toml"))
.load()
.unwrap()
.unwrap();
assert!(!stored.dht.enabled);
}
#[test]
fn stale_revision_cannot_overwrite_a_newer_update() {
let directory = TempDir::new().unwrap();
+1
View File
@@ -2,3 +2,4 @@
pub(crate) mod mapper;
pub(crate) mod pipeline;
pub(crate) mod runtime;
+327
View File
@@ -0,0 +1,327 @@
// 负责动态启停 DHT 采集运行时并暴露稳定的状态与统计句柄
use std::{
str::FromStr,
sync::{
Arc, RwLock,
atomic::{AtomicBool, Ordering},
},
time::{SystemTime, UNIX_EPOCH},
};
use dht_crawler::{DHTOptions, DHTServer, DhtRuntimeStats};
use tokio::task::JoinHandle;
use tokio_util::sync::CancellationToken;
use crate::{
config::VerificationConfig,
domain::InfoHash,
error::AppError,
storage::{RocksTorrentRepository, TorrentRepository},
verification::{self, VerificationIngress},
};
use super::pipeline::PersistenceIngress;
use crate::disk_guard::DiskGuard;
#[derive(Debug, Clone, Copy)]
pub(crate) struct CrawlerRuntimeStatus {
pub(crate) enabled: bool,
pub(crate) transitioning: bool,
}
#[derive(Clone, Default)]
pub(crate) struct CrawlerStats {
inner: Arc<RwLock<DhtRuntimeStats>>,
}
impl CrawlerStats {
pub(crate) fn current(&self) -> DhtRuntimeStats {
self.inner
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
fn replace(&self, stats: DhtRuntimeStats) {
*self
.inner
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = stats;
}
}
#[derive(Clone)]
pub(crate) struct CrawlerRuntime {
inner: Arc<CrawlerRuntimeInner>,
}
struct CrawlerRuntimeInner {
options: DHTOptions,
repository: Arc<RocksTorrentRepository>,
persistence: PersistenceIngress,
disk_guard: DiskGuard,
verification_config: VerificationConfig,
stats: CrawlerStats,
verification: RwLock<Option<VerificationIngress>>,
enabled: AtomicBool,
transitioning: AtomicBool,
state: tokio::sync::Mutex<RunningState>,
}
#[derive(Default)]
struct RunningState {
server: Option<DHTServer>,
server_task: Option<JoinHandle<()>>,
verification_cancel: Option<CancellationToken>,
verification_task: Option<JoinHandle<Result<(), String>>>,
}
impl CrawlerRuntime {
pub(crate) fn new(
options: DHTOptions,
repository: Arc<RocksTorrentRepository>,
persistence: PersistenceIngress,
disk_guard: DiskGuard,
verification_config: VerificationConfig,
) -> Self {
Self {
inner: Arc::new(CrawlerRuntimeInner {
options,
repository,
persistence,
disk_guard,
verification_config,
stats: CrawlerStats::default(),
verification: RwLock::new(None),
enabled: AtomicBool::new(false),
transitioning: AtomicBool::new(false),
state: tokio::sync::Mutex::new(RunningState::default()),
}),
}
}
pub(crate) fn status(&self) -> CrawlerRuntimeStatus {
CrawlerRuntimeStatus {
enabled: self.inner.enabled.load(Ordering::Acquire),
transitioning: self.inner.transitioning.load(Ordering::Acquire),
}
}
pub(crate) fn stats(&self) -> DhtRuntimeStats {
self.inner.stats.current()
}
pub(crate) fn stats_handle(&self) -> CrawlerStats {
self.inner.stats.clone()
}
pub(crate) fn verification(&self) -> Option<VerificationIngress> {
self.inner
.verification
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
pub(crate) async fn set_enabled(&self, enabled: bool) -> Result<(), AppError> {
self.inner.transitioning.store(true, Ordering::Release);
let result = if enabled {
self.start().await
} else {
self.stop().await
};
self.inner.transitioning.store(false, Ordering::Release);
result
}
pub(crate) async fn shutdown(&self) {
if let Err(error) = self.set_enabled(false).await {
tracing::error!(%error, "DHT 采集运行时关闭失败");
}
}
async fn start(&self) -> Result<(), AppError> {
let mut state = self.inner.state.lock().await;
if state.server.is_some() {
self.inner.enabled.store(true, Ordering::Release);
return Ok(());
}
let server = DHTServer::new(self.inner.options.clone()).await?;
configure_callbacks(
&server,
self.inner.repository.clone(),
self.inner.persistence.clone(),
self.inner.disk_guard.clone(),
);
let stats = server.runtime_stats();
if self.inner.verification_config.enabled {
let cancel = CancellationToken::new();
let (ingress, task) = verification::start(
self.inner.repository.clone(),
server.clone(),
self.inner.verification_config.clone(),
self.inner.disk_guard.clone(),
cancel.clone(),
);
*self
.inner
.verification
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(ingress);
state.verification_cancel = Some(cancel);
state.verification_task = Some(task);
}
let runner = server.clone();
state.server_task = Some(tokio::spawn(async move {
if let Err(error) = runner.start().await {
tracing::error!(%error, "DHT 采集运行时异常停止");
}
}));
state.server = Some(server);
self.inner.stats.replace(stats);
self.inner.enabled.store(true, Ordering::Release);
tracing::info!(dht_port = self.inner.options.port, "DHT 持续采集已启动");
Ok(())
}
async fn stop(&self) -> Result<(), AppError> {
let mut state = self.inner.state.lock().await;
if state.server.is_none() {
self.inner.enabled.store(false, Ordering::Release);
self.inner.stats.replace(DhtRuntimeStats::default());
return Ok(());
}
*self
.inner
.verification
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = None;
if let Some(cancel) = state.verification_cancel.take() {
cancel.cancel();
}
if let Some(server) = state.server.take() {
server.shutdown();
}
if let Some(task) = state.verification_task.take() {
match task.await {
Ok(Ok(())) => {}
Ok(Err(error)) => tracing::warn!(%error, "有效性验证运行时停止时返回错误"),
Err(error) => tracing::warn!(%error, "有效性验证运行时任务异常"),
}
}
if let Some(task) = state.server_task.take() {
if let Err(error) = task.await {
tracing::warn!(%error, "DHT 采集运行时任务异常");
}
}
self.inner.stats.replace(DhtRuntimeStats::default());
self.inner.enabled.store(false, Ordering::Release);
tracing::info!("DHT 持续采集已停止 当前仅提供本地索引和搜索");
Ok(())
}
#[cfg(test)]
pub(crate) fn set_verification_for_test(&self, ingress: VerificationIngress) {
*self
.inner
.verification
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(ingress);
}
}
fn configure_callbacks(
server: &DHTServer,
repository: Arc<RocksTorrentRepository>,
persistence: PersistenceIngress,
disk_guard: DiskGuard,
) {
server.on_error(|error| tracing::error!(%error, "DHT 运行时错误"));
let sampled_repository = repository.clone();
let sampled_disk_guard = disk_guard.clone();
server.on_sampled_hashes(move |hashes| {
let repository = sampled_repository.clone();
let disk_guard = sampled_disk_guard.clone();
async move {
let Some(permit) = disk_guard.begin_admission() else {
return Vec::new();
};
let fallback = hashes.clone();
let info_hashes: Vec<_> = hashes.into_iter().map(InfoHash::from_bytes).collect();
let result = tokio::task::spawn_blocking(move || {
let _permit = permit;
repository.filter_unknown_and_observe(&info_hashes, unix_timestamp())
})
.await;
match result {
Ok(Ok(unknown)) => unknown
.into_iter()
.map(|info_hash| *info_hash.as_bytes())
.collect(),
Ok(Err(error)) => {
tracing::error!(%error, "采样 infohash 批量持久化去重失败");
fallback
}
Err(error) => {
tracing::error!(%error, "采样 infohash 批量去重任务异常");
fallback
}
}
}
});
let gate_repository = repository;
let gate_disk_guard = disk_guard;
server.on_metadata_fetch(move |hash| {
let repository = gate_repository.clone();
let disk_guard = gate_disk_guard.clone();
async move {
let Some(permit) = disk_guard.begin_admission() else {
return false;
};
let Ok(info_hash) = InfoHash::from_str(&hash) else {
tracing::warn!(%hash, "DHT 提供了无效 infohash");
return false;
};
let result = tokio::task::spawn_blocking(move || {
let _permit = permit;
repository.observe_existing(info_hash, unix_timestamp())
})
.await;
match result {
Ok(Ok(already_exists)) => !already_exists,
Ok(Err(error)) => {
tracing::error!(%error, %hash, "持久化去重查询失败");
true
}
Err(error) => {
tracing::error!(%error, %hash, "持久化去重任务异常");
true
}
}
}
});
server.on_torrent_with_ack(move |torrent| persistence.try_enqueue(torrent));
server.on_metadata_fetch_complete(|completion| {
tracing::debug!(
info_hash = %completion.info_hash,
status = ?completion.status,
attempts = completion.attempts,
"Metadata 任务完成"
)
});
}
fn unix_timestamp() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs()
}
+8 -6
View File
@@ -20,11 +20,12 @@ use crate::{
search::SearchEngine,
storage::{RocksTorrentRepository, StorageDiagnostics as RocksDiagnostics},
};
use dht_crawler::DhtRuntimeStats;
use tokio_util::sync::CancellationToken;
use crate::{
config::DiagnosticsConfig, crawler::pipeline::PersistenceIngress, disk_guard::DiskGuard,
config::DiagnosticsConfig,
crawler::{pipeline::PersistenceIngress, runtime::CrawlerStats},
disk_guard::DiskGuard,
};
pub(crate) use http::HttpStats;
@@ -38,7 +39,7 @@ use store::{DiagnosticStore, DiagnosticStoreError};
pub(crate) struct DiagnosticSources {
pub(crate) repository: Arc<RocksTorrentRepository>,
pub(crate) search: SearchEngine,
pub(crate) dht: DhtRuntimeStats,
pub(crate) dht: CrawlerStats,
pub(crate) persistence: PersistenceIngress,
pub(crate) disk_guard: DiskGuard,
pub(crate) http: HttpStats,
@@ -259,8 +260,9 @@ async fn collect_loop(
}
fn collect_sample(sources: &DiagnosticSources, session_started_at: u64) -> DiagnosticSample {
let dht = sources.dht.snapshot();
let observability = sources.dht.observability_snapshot();
let runtime_stats = sources.dht.current();
let dht = runtime_stats.snapshot();
let observability = runtime_stats.observability_snapshot();
let persistence = sources.persistence.snapshot();
let disk = sources.disk_guard.snapshot();
let storage = sources.repository.diagnostics().unwrap_or_else(|error| {
@@ -362,7 +364,7 @@ mod tests {
DiagnosticSources {
repository,
search: SearchEngine::open(directory.path().join("tantivy")).unwrap(),
dht: DhtRuntimeStats::default(),
dht: CrawlerStats::default(),
persistence: persistence.ingress.clone(),
disk_guard,
http: HttpStats::default(),
-2
View File
@@ -26,8 +26,6 @@ pub(crate) enum AppError {
PersistenceWorker(String),
#[error("索引 worker 失败: {0}")]
IndexWorker(String),
#[error("可用性验证 worker 失败: {0}")]
VerificationWorker(String),
#[error("运行诊断失败: {0}")]
Diagnostics(String),
}
+8 -5
View File
@@ -3,13 +3,15 @@
use std::time::Duration;
use crate::domain::MetadataRejectionReason;
use dht_crawler::DHTServer;
use tokio_util::sync::CancellationToken;
use crate::{crawler::pipeline::PersistenceIngress, disk_guard::DiskGuard};
use crate::{
crawler::{pipeline::PersistenceIngress, runtime::CrawlerStats},
disk_guard::DiskGuard,
};
pub(crate) async fn run(
server: DHTServer,
dht_stats: CrawlerStats,
ingress: PersistenceIngress,
disk_guard: DiskGuard,
interval_secs: u64,
@@ -23,8 +25,9 @@ pub(crate) async fn run(
tokio::select! {
_ = cancel.cancelled() => break,
_ = interval.tick() => {
let dht = server.runtime_stats().snapshot();
let observability = server.runtime_stats().observability_snapshot();
let runtime_stats = dht_stats.current();
let dht = runtime_stats.snapshot();
let observability = runtime_stats.observability_snapshot();
let storage = ingress.snapshot();
let disk = disk_guard.snapshot();
let udp_tx_per_second = observability
+1
View File
@@ -59,6 +59,7 @@ Axum 根据配置中的 `http.web_dir` 提供静态资源和单页回退 不需
- 相同内容的不同 infohash 变体展示
- 磁力链接打开和复制
- 在系统诊断页统一展示 DHT 采集 索引 持久化和验证运行状态
- 使用右上角图标即时启停 DHT 持续采集并持久化离线模式
- 仅在系统诊断页打开时每秒刷新服务状态
- 独立展示进程 RocksDB Tantivy DHT 和队列诊断数据
- 使用懒加载 ECharts 展示内存存储压力 Metadata 吞吐以及 HTTP 请求错误趋势
+46 -2
View File
@@ -1,11 +1,16 @@
<script setup lang="ts">
import { onMounted, ref } from 'vue'
import { ChartLine, Database, Moon, Search, Settings, Sun } from '@lucide/vue'
import { computed, onMounted, ref } from 'vue'
import { ChartLine, Database, LoaderCircle, Moon, Radio, Search, Settings, Sun, WifiOff } from '@lucide/vue'
import { RouterLink, RouterView } from 'vue-router'
import { Button } from '@/components/ui/button'
import { getCrawlerStatus, updateCrawlerStatus } from '@/lib/api'
const darkMode = ref(false)
const crawlerEnabled = ref(false)
const crawlerLoaded = ref(false)
const crawlerChanging = ref(false)
const crawlerError = ref('')
const navigation = [
{ to: '/', label: '搜索', icon: Search },
@@ -17,6 +22,39 @@ function resetHome() {
window.dispatchEvent(new Event('dht-search-reset'))
}
const crawlerButtonLabel = computed(() => {
if (crawlerError.value) return crawlerError.value
if (!crawlerLoaded.value) return '正在读取 DHT 采集状态'
if (crawlerChanging.value) return crawlerEnabled.value ? '正在停止 DHT 采集' : '正在启动 DHT 采集'
return crawlerEnabled.value ? '停止 DHT 后台采集并进入离线模式' : '启动 DHT 后台采集'
})
async function loadCrawlerStatus() {
crawlerError.value = ''
try {
const status = await getCrawlerStatus()
crawlerEnabled.value = status.enabled
} catch (cause) {
crawlerError.value = cause instanceof Error ? cause.message : '无法读取 DHT 采集状态'
} finally {
crawlerLoaded.value = true
}
}
async function toggleCrawler() {
if (!crawlerLoaded.value || crawlerChanging.value) return
crawlerChanging.value = true
crawlerError.value = ''
try {
const status = await updateCrawlerStatus(!crawlerEnabled.value)
crawlerEnabled.value = status.enabled
} catch (cause) {
crawlerError.value = cause instanceof Error ? cause.message : '无法切换 DHT 采集状态'
} finally {
crawlerChanging.value = false
}
}
function toggleTheme() {
darkMode.value = !darkMode.value
document.documentElement.classList.toggle('dark', darkMode.value)
@@ -25,6 +63,7 @@ function toggleTheme() {
onMounted(() => {
darkMode.value = document.documentElement.classList.contains('dark')
void loadCrawlerStatus()
})
</script>
@@ -47,6 +86,11 @@ onMounted(() => {
</nav>
<div class="ml-auto flex items-center gap-1">
<Button size="icon" variant="ghost" :disabled="!crawlerLoaded || crawlerChanging" :aria-label="crawlerButtonLabel" :title="crawlerButtonLabel" :class="crawlerError ? 'text-destructive' : crawlerEnabled ? 'text-emerald-600 dark:text-emerald-400' : 'text-muted-foreground'" @click="toggleCrawler">
<LoaderCircle v-if="!crawlerLoaded || crawlerChanging" class="animate-spin" />
<Radio v-else-if="crawlerEnabled" />
<WifiOff v-else />
</Button>
<Button size="icon" variant="ghost" :aria-label="darkMode ? '切换到浅色主题' : '切换到深色主题'" :title="darkMode ? '浅色主题' : '深色主题'" @click="toggleTheme"><Sun v-if="darkMode" /><Moon v-else /></Button>
</div>
</div>
+13
View File
@@ -2,6 +2,7 @@ import type {
ContentVariants,
ConfigSnapshot,
ConfigUpdateRequest,
CrawlerStatus,
CurrentDiagnostics,
DiagnosticHistory,
SearchPage,
@@ -81,3 +82,15 @@ export function updateConfig(input: ConfigUpdateRequest, signal?: AbortSignal):
body: JSON.stringify(input),
})
}
export function getCrawlerStatus(signal?: AbortSignal): Promise<CrawlerStatus> {
return request<CrawlerStatus>('/crawler', signal)
}
export function updateCrawlerStatus(enabled: boolean, signal?: AbortSignal): Promise<CrawlerStatus> {
return request<CrawlerStatus>('/crawler', signal, {
method: 'PUT',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ enabled }),
})
}
+6
View File
@@ -256,6 +256,7 @@ export interface AppConfigDto {
max_path_depth: number
}
dht: {
enabled: boolean
port: number
netmode: NetworkMode
hash_queue_capacity: number
@@ -331,3 +332,8 @@ export interface ConfigUpdateRequest {
revision: string
config: AppConfigDto
}
export interface CrawlerStatus {
enabled: boolean
transitioning: boolean
}
+1
View File
@@ -18,6 +18,7 @@ export default defineConfig({
'/health': apiTarget,
'/ready': apiTarget,
'/stats': apiTarget,
'/crawler': apiTarget,
'/diagnostics': apiTarget,
'/config': apiTarget,
'/search': apiTarget,