diff --git a/README.md b/README.md index 0fd99ca..d7c0e1f 100644 --- a/README.md +++ b/README.md @@ -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 可以在仓库根目录执行 diff --git a/TODOS.md b/TODOS.md index beda17a..cc76698 100644 --- a/TODOS.md +++ b/TODOS.md @@ -182,6 +182,7 @@ - [x] 统一搜索结果数字分页并支持用户选择每页数量 - [x] 将采集索引持久化和验证运行状态集中到系统诊断页并仅在页面打开时每秒刷新 - [x] 将诊断时间范围放入历史趋势区域并精简图表卡片的单行摘要信息 +- [x] 在 Web 顶栏提供持久化的 DHT 即时启停开关并保留离线索引搜索能力 - [x] 实现浏览器持久化明暗主题 - [x] 为加载空结果接口错误和失败重试提供明确界面状态 - [x] 将生产静态资源交给 Axum 提供并支持单页回退 diff --git a/config.toml b/config.toml index aa0b4df..fa220eb 100644 --- a/config.toml +++ b/config.toml @@ -77,6 +77,7 @@ action = "hide" reason = "libtorrent 分片边界填充目录" [dht] +enabled = false port = 12313 netmode = "ipv4-only" hash_queue_capacity = 20000 diff --git a/src/search/README.md b/src/search/README.md index 300c30d..8875e7d 100644 --- a/src/search/README.md +++ b/src/search/README.md @@ -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` 暴露到不受信任的公网入口 diff --git a/src/search/src/api/handlers.rs b/src/search/src/api/handlers.rs index 6a8e904..b462a85 100644 --- a/src/search/src/api/handlers.rs +++ b/src/search/src/api/handlers.rs @@ -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 { } pub(crate) async fn stats(State(state): State) -> Json { - 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) -> Json }) } +pub(crate) async fn crawler_status(State(state): State) -> Json { + let status = state.crawler.status(); + Json(CrawlerStatusResponse { + enabled: status.enabled, + transitioning: status.transitioning, + }) +} + +pub(crate) async fn crawler_update( + State(state): State, + Json(request): Json, +) -> Result, 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, ) -> Json { @@ -208,7 +250,7 @@ pub(crate) async fn search( State(state): State, Query(request): Query, ) -> Result, 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)) diff --git a/src/search/src/api/mod.rs b/src/search/src/api/mod.rs index 48cb1ac..bf03983 100644 --- a/src/search/src/api/mod.rs +++ b/src/search/src/api/mod.rs @@ -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, pub(crate) search: SearchEngine, - pub(crate) dht_stats: DhtRuntimeStats, + pub(crate) crawler: CrawlerRuntime, pub(crate) persistence: PersistenceIngress, - pub(crate) verification: Option, 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"), "
DHT Search
").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( diff --git a/src/search/src/api/request.rs b/src/search/src/api/request.rs index 93ee677..7c57388 100644 --- a/src/search/src/api/request.rs +++ b/src/search/src/api/request.rs @@ -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 } diff --git a/src/search/src/api/response.rs b/src/search/src/api/response.rs index 1f98fda..0e2a20a 100644 --- a/src/search/src/api/response.rs +++ b/src/search/src/api/response.rs @@ -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, diff --git a/src/search/src/app.rs b/src/search/src/app.rs index 286a2be..317d4f0 100644 --- a/src/search/src/app.rs +++ b/src/search/src/app.rs @@ -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, error: AppError) { tracing::error!(%error, "关闭阶段发生附加错误"); } } - -fn unix_timestamp() -> u64 { - SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_secs() -} diff --git a/src/search/src/config.rs b/src/search/src/config.rs index 9f9ac37..0f1a70e 100644 --- a/src/search/src/config.rs +++ b/src/search/src/config.rs @@ -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); diff --git a/src/search/src/config/model.rs b/src/search/src/config/model.rs index 173e065..d7a5c31 100644 --- a/src/search/src/config/model.rs +++ b/src/search/src/config/model.rs @@ -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, diff --git a/src/search/src/config/service.rs b/src/search/src/config/service.rs index bab1c03..1e5033c 100644 --- a/src/search/src/config/service.rs +++ b/src/search/src/config/service.rs @@ -123,10 +123,37 @@ impl ConfigService { Ok(self.snapshot_from(&state)) } + pub(crate) fn set_dht_enabled( + &self, + enabled: bool, + ) -> Result { + 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 { 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(); diff --git a/src/search/src/crawler/mod.rs b/src/search/src/crawler/mod.rs index cf3b9ec..3345fac 100644 --- a/src/search/src/crawler/mod.rs +++ b/src/search/src/crawler/mod.rs @@ -2,3 +2,4 @@ pub(crate) mod mapper; pub(crate) mod pipeline; +pub(crate) mod runtime; diff --git a/src/search/src/crawler/runtime.rs b/src/search/src/crawler/runtime.rs new file mode 100644 index 0000000..0a155c2 --- /dev/null +++ b/src/search/src/crawler/runtime.rs @@ -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>, +} + +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, +} + +struct CrawlerRuntimeInner { + options: DHTOptions, + repository: Arc, + persistence: PersistenceIngress, + disk_guard: DiskGuard, + verification_config: VerificationConfig, + stats: CrawlerStats, + verification: RwLock>, + enabled: AtomicBool, + transitioning: AtomicBool, + state: tokio::sync::Mutex, +} + +#[derive(Default)] +struct RunningState { + server: Option, + server_task: Option>, + verification_cancel: Option, + verification_task: Option>>, +} + +impl CrawlerRuntime { + pub(crate) fn new( + options: DHTOptions, + repository: Arc, + 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 { + 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, + 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() +} diff --git a/src/search/src/diagnostics/mod.rs b/src/search/src/diagnostics/mod.rs index d17f10c..6cd910b 100644 --- a/src/search/src/diagnostics/mod.rs +++ b/src/search/src/diagnostics/mod.rs @@ -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, 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(), diff --git a/src/search/src/error.rs b/src/search/src/error.rs index 06057be..6453ff5 100644 --- a/src/search/src/error.rs +++ b/src/search/src/error.rs @@ -26,8 +26,6 @@ pub(crate) enum AppError { PersistenceWorker(String), #[error("索引 worker 失败: {0}")] IndexWorker(String), - #[error("可用性验证 worker 失败: {0}")] - VerificationWorker(String), #[error("运行诊断失败: {0}")] Diagnostics(String), } diff --git a/src/search/src/monitor.rs b/src/search/src/monitor.rs index 12a01fb..1ca9304 100644 --- a/src/search/src/monitor.rs +++ b/src/search/src/monitor.rs @@ -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 diff --git a/src/web/README.md b/src/web/README.md index d254581..2409d3e 100644 --- a/src/web/README.md +++ b/src/web/README.md @@ -59,6 +59,7 @@ Axum 根据配置中的 `http.web_dir` 提供静态资源和单页回退 不需 - 相同内容的不同 infohash 变体展示 - 磁力链接打开和复制 - 在系统诊断页统一展示 DHT 采集 索引 持久化和验证运行状态 +- 使用右上角图标即时启停 DHT 持续采集并持久化离线模式 - 仅在系统诊断页打开时每秒刷新服务状态 - 独立展示进程 RocksDB Tantivy DHT 和队列诊断数据 - 使用懒加载 ECharts 展示内存存储压力 Metadata 吞吐以及 HTTP 请求错误趋势 diff --git a/src/web/src/App.vue b/src/web/src/App.vue index 78fde3d..19202b4 100644 --- a/src/web/src/App.vue +++ b/src/web/src/App.vue @@ -1,11 +1,16 @@ @@ -47,6 +86,11 @@ onMounted(() => {
+
diff --git a/src/web/src/lib/api.ts b/src/web/src/lib/api.ts index 372d50a..8d38b8a 100644 --- a/src/web/src/lib/api.ts +++ b/src/web/src/lib/api.ts @@ -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 { + return request('/crawler', signal) +} + +export function updateCrawlerStatus(enabled: boolean, signal?: AbortSignal): Promise { + return request('/crawler', signal, { + method: 'PUT', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ enabled }), + }) +} diff --git a/src/web/src/types/api.ts b/src/web/src/types/api.ts index 11afad2..2dfc24b 100644 --- a/src/web/src/types/api.ts +++ b/src/web/src/types/api.ts @@ -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 +} diff --git a/src/web/vite.config.ts b/src/web/vite.config.ts index 71d7dc0..3d21b70 100644 --- a/src/web/vite.config.ts +++ b/src/web/vite.config.ts @@ -18,6 +18,7 @@ export default defineConfig({ '/health': apiTarget, '/ready': apiTarget, '/stats': apiTarget, + '/crawler': apiTarget, '/diagnostics': apiTarget, '/config': apiTarget, '/search': apiTarget,