From 3b997c8f40ea2f1deb5e2f646d88b8da20674735 Mon Sep 17 00:00:00 2001 From: chuan Date: Mon, 10 Aug 2026 21:45:32 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E7=B2=BE=E7=AE=80=E9=85=8D=E7=BD=AE?= =?UTF-8?q?=E7=AE=A1=E7=90=86=E4=B8=8E=E8=BF=90=E8=A1=8C=E9=BB=98=E8=AE=A4?= =?UTF-8?q?=E5=80=BC?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- Dockerfile | 2 +- TODOS.md | 11 +- config.toml | 43 +---- src/search/README.md | 34 ++-- src/search/src/api/mod.rs | 11 +- src/search/src/app.rs | 15 +- src/search/src/config.rs | 60 +++---- src/search/src/config/model.rs | 132 ++++++++-------- src/search/src/config/runtime.rs | 22 +-- src/search/src/crawler/pipeline.rs | 14 +- src/search/src/diagnostics/mod.rs | 9 +- src/search/src/disk_guard.rs | 103 +++++++----- src/search/src/entry.rs | 1 - src/search/src/telemetry.rs | 10 +- src/search/src/verification.rs | 7 +- src/web/README.md | 2 +- src/web/src/pages/ConfigPage.vue | 244 +++++++++++------------------ src/web/src/style.css | 3 + src/web/src/types/api.ts | 26 +-- 19 files changed, 294 insertions(+), 455 deletions(-) diff --git a/Dockerfile b/Dockerfile index 95e9856..ab554f7 100644 --- a/Dockerfile +++ b/Dockerfile @@ -67,4 +67,4 @@ EXPOSE 12313/udp STOPSIGNAL SIGTERM ENTRYPOINT ["/dht-search/dht-search"] -CMD ["--config", "/dht-search/config.toml", "--http-listen", "0.0.0.0:8080", "--web-dir", "/dht-search/web", "--console-logging", "--no-file-logging"] +CMD ["--config", "/dht-search/config.toml", "--http-listen", "0.0.0.0:8080", "--web-dir", "/dht-search/web"] diff --git a/TODOS.md b/TODOS.md index cc76698..9bc153e 100644 --- a/TODOS.md +++ b/TODOS.md @@ -183,6 +183,8 @@ - [x] 将采集索引持久化和验证运行状态集中到系统诊断页并仅在页面打开时每秒刷新 - [x] 将诊断时间范围放入历史趋势区域并精简图表卡片的单行摘要信息 - [x] 在 Web 顶栏提供持久化的 DHT 即时启停开关并保留离线索引搜索能力 +- [x] 将配置页精简为设置与过滤两个页签并合并全部常用参数 +- [x] 使用数值与单位选择编辑 Metadata 大小并自动换算内部字节数 - [x] 实现浏览器持久化明暗主题 - [x] 为加载空结果接口错误和失败重试提供明确界面状态 - [x] 将生产静态资源交给 Axum 提供并支持单页回退 @@ -221,18 +223,16 @@ - [ ] 统计真实数据的 infohash 重复率和内容重复率 - [x] 在统一 `config.toml` 中定义文件名和文件路径隐藏规则 -- [x] 支持精确前缀后缀包含通配符和正则匹配并限制规则复杂度 +- [x] 使用文件名和文件路径双文本框按行管理不区分大小写的通配符规则 - [x] 保留 RocksDB 原始文件列表并为详情统计搜索和内容聚合生成有效内容视图 - [x] 使用规则指纹在配置变化时重算内容组并从 RocksDB 重建 Tantivy - [x] 默认隐藏 BitComet padding 文件以及 `.pad` 和 `.____padding_file` 填充目录 - [x] 全部文件被隐藏的 Metadata 只保留原始记录且不进入公开索引 - [ ] 根据真实垃圾数据决定是否增加种子名称扩展名和大小准入规则 -- [ ] 增加按规则 ID 分类的隐藏文件命中指标 - [x] 定义可配置的 Metadata 最大大小文件数名称路径长度和目录层级限制 - [x] 识别空名称控制字符异常路径大小溢出总大小不一致和文件数量攻击 - [ ] 设计可解释的名称标准化规则 - [ ] 为模糊相似结果生成聚合候选但不自动删除 -- [ ] 支持黑名单规则版本和命中原因 - [x] 使用带规则指纹的 RocksDB 轻量拒绝记录阻止相同异常 infohash 重复下载 - [x] 保留按原因分类的过滤指标但避免保存名称和大文件列表 - [ ] 增加误判测试和边界数据集 @@ -290,8 +290,11 @@ - [x] 记录 Tantivy IndexWriter 内存和 commit 延迟 - [ ] 根据实测调整批量大小队列容量和并发 - [x] 增加带排空阶段恢复滞回和探测失败保护的磁盘只读降级策略 +- [x] 固定每分钟检查磁盘并按文件系统容量自动计算保护和恢复阈值 - [x] 增加在线 RocksDB 检查点保留上限只读校验和带旧库保留的离线恢复 -- [x] 增加可配置的终端日志滚动文件日志和保留文件上限 +- [x] 默认关闭自动检查点并从 Web 隐藏全部备份调度参数 +- [x] 始终保留终端日志并仅在 Web 暴露滚动文件日志开关 +- [x] 移除 Docker 对文件日志的命令行覆盖并隐藏无关的基础设施覆盖提示 - [x] 验证间歇运行和正常退出恢复 - [x] 完成本机约七小时真实持续运行并确认采集索引和搜索服务可用 - [ ] 验证二十四小时和七天连续运行 diff --git a/config.toml b/config.toml index fa220eb..e7b19c7 100644 --- a/config.toml +++ b/config.toml @@ -3,18 +3,11 @@ data_dir = "data" persistence_queue_capacity = 8192 stats_interval_secs = 10 -# run_duration_secs = 3600 index_batch_size = 1024 index_interval_millis = 5000 -[disk_guard] -enabled = true -check_interval_secs = 10 -minimum_free_bytes = 5368709120 -resume_free_bytes = 6442450944 - [backup] -enabled = true +enabled = false directory = "data/backups" interval_secs = 21600 retain_checkpoints = 3 @@ -31,7 +24,6 @@ queue_capacity = 128 [logging] directory = "data/logs" file_enabled = true -console_enabled = false rotation = "daily" retain_files = 7 file_prefix = "dht-search" @@ -44,37 +36,8 @@ max_path_bytes = 4096 max_path_depth = 64 [content_filter] -version = 1 - -[[content_filter.file_rules]] -id = "bitcomet-padding-file" -enabled = true -field = "file-name" -match = "prefix" -value = "_____padding_file_" -case_sensitive = false -action = "hide" -reason = "BitComet 分片边界填充文件" - -[[content_filter.file_rules]] -id = "generic-pad-directory" -enabled = true -field = "file-path" -match = "regex" -value = '(^|/)\.pad/' -case_sensitive = false -action = "hide" -reason = "客户端分片边界填充目录" - -[[content_filter.file_rules]] -id = "libtorrent-padding-directory" -enabled = true -field = "file-path" -match = "regex" -value = '(^|/)\.____padding_file/' -case_sensitive = false -action = "hide" -reason = "libtorrent 分片边界填充目录" +file_name_patterns = ["*_____padding_file_*"] +file_path_patterns = ["*.pad/*", "*.____padding_file/*"] [dht] enabled = false diff --git a/src/search/README.md b/src/search/README.md index 8875e7d..d685558 100644 --- a/src/search/README.md +++ b/src/search/README.md @@ -26,13 +26,13 @@ cargo run -p dht-search -- --config config.toml 相对目录以 `config.toml` 所在目录为基准解析 -也可以通过命令行覆盖数据目录和本次运行时长 +也可以通过命令行覆盖数据目录 ```powershell -cargo run -p dht-search -- --data-dir D:\data\dht-search --run-duration-secs 3600 +cargo run -p dht-search -- --data-dir D:\data\dht-search ``` -不设置 `run-duration-secs` 时服务持续运行直到收到 Ctrl+C SIGINT 或 SIGTERM +服务持续运行直到收到 Ctrl+C SIGINT 或 SIGTERM 生产 Web 页面需要先在 `src/web` 目录执行 `bun run build` Axum 会从 `http.web_dir` 提供构建结果 @@ -86,35 +86,29 @@ Windows 下索引每五秒批量提交 临时文件占用会自动指数退避 ### 磁盘空间保护 -应用默认每十秒检查 `data_dir` 所在磁盘的剩余空间 低于保护阈值时先停止接收新 Metadata DHT 状态更新索引任务和可用性验证任务 已经进入有界持久化队列的记录会继续排空 随后进入只读保护 +应用固定每分钟检查 `data_dir` 所在文件系统的总容量和剩余空间 低于保护阈值时先停止接收新 Metadata DHT 状态更新索引任务和可用性验证任务 已经进入有界持久化队列的记录会继续排空 随后进入只读保护 -只读保护期间 RocksDB 和 Tantivy 不再产生业务写入 现有搜索详情健康检查和 Web 页面仍然可用 剩余空间达到独立恢复阈值后自动恢复采集 使用两个阈值可以避免临界空间附近反复暂停和恢复 +只读保护始终启用且不需要用户配置 保护阈值取文件系统总容量的 5% 并限制在 512 MiB 到 20 GiB 之间 恢复缓冲取总容量的 1% 并限制在 256 MiB 到 5 GiB 之间 文件系统扩容后会自动采用新阈值 -| 配置项 | 默认值 | 作用 | -|---|---:|---| -| `disk_guard.enabled` | `true` | 是否启用磁盘空间保护 | -| `disk_guard.check_interval_secs` | `10` | 剩余空间检查间隔 | -| `disk_guard.minimum_free_bytes` | `5368709120` | 低于 5 GiB 时停止接收新任务 | -| `disk_guard.resume_free_bytes` | `6442450944` | 恢复到 6 GiB 时重新接受写入 | +只读保护期间 RocksDB 和 Tantivy 不再产生业务写入 现有搜索详情健康检查和 Web 页面仍然可用 剩余空间达到自动恢复阈值后重新接受写入 使用两个阈值可以避免临界空间附近反复暂停和恢复 磁盘空间探测失败时采用保守策略进入保护状态 `/stats` 返回 `disk_state` `disk_available_bytes` 阈值 活跃写入数 探测失败数 状态转换数和拒绝任务数 Web 运行状态使用绿色或黄色状态点展示正常与保护状态 ### 日志轮转和保留 -应用会在读取配置后初始化日志 默认只写入 `data/logs` 的滚动文件而不重复输出到终端 因此用脚本或后台进程启动时不需要再把标准错误重定向到长期增长的日志文件 +应用会在读取配置后初始化日志并始终输出到终端 因此 Docker 可以直接通过 `docker logs` 读取运行日志 文件日志默认同时写入 `data/logs` 并允许从 Web 高级设置关闭 | 配置项 | 默认值 | 作用 | |---|---:|---| | `logging.directory` | `data/logs` | 日志文件目录 相对主配置文件解析 | | `logging.file_enabled` | `true` | 启用滚动文件日志 | -| `logging.console_enabled` | `false` | 同时输出到当前终端 | | `logging.rotation` | `daily` | 轮转周期 支持 `minutely` `hourly` `daily` 和 `never` | | `logging.retain_files` | `7` | 最多保留的匹配日志文件数量 | | `logging.file_prefix` | `dht-search` | 日志文件名前缀 | 默认按天轮转时保留 7 个文件约等于保留最近 7 天 日志组件只清理同目录中同时匹配前缀和 `.log` 后缀的普通文件 不删除目录和符号链接 清理失败会输出错误但不会让服务退出 -开发时需要直接观察终端日志可以设置 `console_enabled = true` 文件日志和终端日志不能同时关闭 +关闭文件日志不会影响终端输出 ### RocksDB 检查点备份和恢复 @@ -122,13 +116,13 @@ RocksDB 是唯一权威数据源 应用使用 RocksDB 原生 Checkpoint API 在 | 配置项 | 默认值 | 作用 | |---|---:|---| -| `backup.enabled` | `true` | 是否启用自动检查点 | +| `backup.enabled` | `false` | 是否启用自动检查点 默认关闭且不在 Web 展示 | | `backup.directory` | `data/backups` | 检查点目录 相对主配置文件解析 | | `backup.interval_secs` | `21600` | 每 6 小时创建一次检查点 | | `backup.retain_checkpoints` | `3` | 保留最近 3 个自动检查点 | | `backup.create_on_start` | `true` | 每次启动后立即创建一次检查点 | -备份目录与数据目录位于同一磁盘时 RocksDB 会尽量通过硬链接减少复制开销 放到其他磁盘时可能复制全部数据库文件 创建前会检查备份磁盘剩余空间 磁盘保护期间自动跳过而不会阻塞服务 +需要自动检查点时可以直接修改 `config.toml` 备份目录与数据目录位于同一磁盘时 RocksDB 会尽量通过硬链接减少复制开销 放到其他磁盘时可能复制全部数据库文件 创建前会检查备份磁盘剩余空间 磁盘保护期间自动跳过而不会阻塞服务 `/stats` 返回检查点成功失败跳过清理数量 最近成功时间耗时和记录数量 自动清理只处理名称严格匹配 `checkpoint-` 加二十位时间戳的直接子目录 不会删除手工目录文件或符号链接 @@ -191,10 +185,10 @@ Web 右上角采集开关通过 `/crawler` 即时停止或重新创建 DHT 运 `/stats` 返回 `metadata_filtered` 总数以及 `metadata_filtered_*` 分类计数 Web 运行状态展示本次运行的过滤总数 -可以通过命令行覆盖独立数据目录并运行五分钟测试 不会污染正式数据目录 +可以通过命令行覆盖独立数据目录进行测试 不会污染正式数据目录 ```powershell -cargo run -p dht-search -- --config config.toml --data-dir data-filter-test --run-duration-secs 300 +cargo run -p dht-search -- --config config.toml --data-dir data-filter-test Invoke-RestMethod http://127.0.0.1:8080/stats | ConvertTo-Json -Depth 5 ``` @@ -204,9 +198,9 @@ Invoke-RestMethod http://127.0.0.1:8080/stats | ConvertTo-Json -Depth 5 RocksDB 始终保存完整原始 Metadata 隐藏规则不会删除文件或种子 修改或回滚规则后应用会根据规则指纹重新计算内容组并从 RocksDB 重建 Tantivy -每条规则包含稳定 `id` 开关 匹配字段 匹配方式 值 大小写选项和可读原因 当前字段支持 `file-name` 与 `file-path` 匹配方式支持 `exact` `prefix` `suffix` `contains` `wildcard` 和 `regex` 动作只允许安全的 `hide` +配置只包含 `file_name_patterns` 和 `file_path_patterns` 两组不区分大小写的通配符 Web 使用左右两个多行文本框编辑并按行切分 空行和重复规则自动忽略 -通配符中 `*` 表示任意长度字符 `?` 表示一个字符并匹配完整字段 正则表达式使用 Rust `regex` 语法 文件路径在匹配前统一使用 `/` 分隔符 +通配符中 `*` 表示任意长度字符 `?` 表示一个字符并匹配完整字段 文件路径在匹配前统一使用 `/` 分隔符 如果一个 Metadata 的全部文件都被隐藏 原始记录仍保留在 RocksDB 但不会进入搜索索引或公开详情 diff --git a/src/search/src/api/mod.rs b/src/search/src/api/mod.rs index bf03983..346c3fd 100644 --- a/src/search/src/api/mod.rs +++ b/src/search/src/api/mod.rs @@ -102,7 +102,7 @@ mod tests { use tower::ServiceExt; use crate::{ - config::{AppConfigDto, ConfigService, DiskGuardConfig, TomlConfigStore}, + config::{AppConfigDto, ConfigService, TomlConfigStore}, crawler::pipeline::PersistencePipeline, diagnostics::DiagnosticsRuntime, disk_guard::DiskGuard, @@ -144,10 +144,7 @@ mod tests { let search = SearchEngine::open(directory.path().join("tantivy")).unwrap(); search.index_pending(repository.as_ref(), 10, 20).unwrap(); let repository_trait: Arc = repository.clone(); - let disk_guard = DiskGuard::new(&DiskGuardConfig { - enabled: false, - ..DiskGuardConfig::default() - }); + let disk_guard = DiskGuard::new(); let persistence = PersistencePipeline::start( repository_trait.clone(), 4, @@ -295,8 +292,8 @@ mod tests { .unwrap(); let original_revision = config_snapshot["revision"].as_str().unwrap().to_owned(); let mut invalid_filter = config_snapshot["config"].clone(); - invalid_filter["content_filter"]["file_rules"][0]["match"] = serde_json::json!("regex"); - invalid_filter["content_filter"]["file_rules"][0]["value"] = serde_json::json!("("); + invalid_filter["content_filter"]["file_name_patterns"][0] = + serde_json::json!("x".repeat(1_025)); let response = app .clone() .oneshot( diff --git a/src/search/src/app.rs b/src/search/src/app.rs index 317d4f0..8267beb 100644 --- a/src/search/src/app.rs +++ b/src/search/src/app.rs @@ -22,7 +22,7 @@ use crate::{ pub(crate) async fn run(config: AppConfig, config_service: ConfigService) -> Result<(), AppError> { std::fs::create_dir_all(&config.data_dir)?; let _data_lock = backup::acquire_data_lock(&config.data_dir)?; - let disk_guard = DiskGuard::new(&config.disk_guard); + let disk_guard = DiskGuard::new(); let database_path = config.data_dir.join("rocksdb"); let metadata_limits = config.metadata_limits(); let content_filter = Arc::new(config.content_filter()?); @@ -50,7 +50,6 @@ pub(crate) async fn run(config: AppConfig, config_service: ConfigService) -> Res let disk_task = tokio::spawn(disk_guard::run( disk_guard.clone(), config.data_dir.clone(), - config.disk_guard.clone(), ingress.clone(), disk_cancel.clone(), )); @@ -148,23 +147,11 @@ pub(crate) async fn run(config: AppConfig, config_service: ConfigService) -> Res api_cancel.clone(), )); - let run_duration = async { - match config.run_duration_secs { - Some(seconds) => tokio::time::sleep(Duration::from_secs(seconds)).await, - None => std::future::pending().await, - } - }; - tokio::pin!(run_duration); - let run_result = tokio::select! { _ = shutdown::signal() => { tracing::info!("收到退出信号"); Ok(()) } - _ = &mut run_duration => { - tracing::info!("达到配置的运行时长"); - Ok(()) - } fatal = &mut persistence.fatal => { let message = fatal.unwrap_or_else(|_| "持久化 worker 意外停止".to_owned()); Err(AppError::PersistenceWorker(message)) diff --git a/src/search/src/config.rs b/src/search/src/config.rs index 0f1a70e..449cde4 100644 --- a/src/search/src/config.rs +++ b/src/search/src/config.rs @@ -16,8 +16,7 @@ use clap::Parser; use crate::error::AppError; pub(crate) use model::{ - AppConfigDto, BackupConfig, DiagnosticsConfig, DiskGuardConfig, LogRotation, LoggingConfig, - VerificationConfig, + AppConfigDto, BackupConfig, DiagnosticsConfig, LogRotation, LoggingConfig, VerificationConfig, }; pub(crate) use runtime::AppConfig; pub(crate) use service::{ConfigService, ConfigServiceError, ConfigSnapshot, ConfigUpdateRequest}; @@ -35,12 +34,8 @@ pub(crate) struct Cli { #[arg(long)] web_dir: Option, #[arg(long)] - console_logging: bool, - #[arg(long)] no_file_logging: bool, #[arg(long)] - run_duration_secs: Option, - #[arg(long)] restore_checkpoint: Option, } @@ -71,18 +66,10 @@ impl Cli { dto.http.web_dir = web_dir; command_line_overrides.push("http.web_dir".to_owned()); } - if self.console_logging { - dto.logging.console_enabled = true; - command_line_overrides.push("logging.console_enabled".to_owned()); - } if self.no_file_logging { dto.logging.file_enabled = false; command_line_overrides.push("logging.file_enabled".to_owned()); } - if self.run_duration_secs.is_some() { - dto.run_duration_secs = self.run_duration_secs; - command_line_overrides.push("run_duration_secs".to_owned()); - } let base = config_path .parent() .ok_or_else(|| AppError::Config("配置文件没有父目录".to_owned()))?; @@ -168,9 +155,7 @@ mod tests { data_dir: None, http_listen: None, web_dir: None, - console_logging: false, no_file_logging: false, - run_duration_secs: None, restore_checkpoint: None, } .load() @@ -199,9 +184,7 @@ mod tests { data_dir: None, http_listen: None, web_dir: None, - console_logging: false, no_file_logging: false, - run_duration_secs: None, restore_checkpoint: None, } .load() @@ -224,9 +207,7 @@ mod tests { data_dir: None, http_listen: Some(listen), web_dir: Some(PathBuf::from("/dht-search/web")), - console_logging: true, no_file_logging: true, - run_duration_secs: None, restore_checkpoint: None, } .load() @@ -236,14 +217,8 @@ mod tests { assert!(startup.app.http.web_dir.ends_with("dht-search/web")); assert_eq!( startup.config_service.snapshot().command_line_overrides, - [ - "http.listen", - "http.web_dir", - "logging.console_enabled", - "logging.file_enabled" - ] + ["http.listen", "http.web_dir", "logging.file_enabled"] ); - assert!(startup.app.logging.console_enabled); assert!(!startup.app.logging.file_enabled); } @@ -264,24 +239,29 @@ mod tests { #[test] fn invalid_embedded_content_filter_is_rejected() { let mut dto = AppConfigDto::default(); - dto.content_filter.file_rules[0].match_kind = crate::domain::FileMatchKind::Regex; - dto.content_filter.file_rules[0].value = "(".to_owned(); + dto.content_filter.file_name_patterns[0] = "x".repeat(1_025); assert!(matches!(resolve(dto), Err(AppError::ContentFilter(_)))); } #[test] - fn disk_resume_threshold_must_exceed_minimum() { + fn content_filter_patterns_ignore_blank_duplicate_and_case() { let mut dto = AppConfigDto::default(); - dto.disk_guard.resume_free_bytes = dto.disk_guard.minimum_free_bytes; - assert!(matches!(resolve(dto), Err(AppError::Config(_)))); - } - - #[test] - fn logging_requires_at_least_one_output() { - let mut dto = AppConfigDto::default(); - dto.logging.file_enabled = false; - dto.logging.console_enabled = false; - assert!(matches!(resolve(dto), Err(AppError::Config(_)))); + dto.content_filter.file_name_patterns = vec![ + String::new(), + "*PADDING_FILE*".to_owned(), + "*padding_file*".to_owned(), + ]; + let config = resolve(dto).unwrap(); + assert_eq!(config.content_filter.to_domain().file_rules.len(), 3); + let filter = config.content_filter().unwrap(); + assert!(filter.is_hidden(&crate::domain::TorrentFile { + path: "release/Padding_File_1".to_owned(), + size: 1, + })); + assert!(filter.is_hidden(&crate::domain::TorrentFile { + path: "release/.pad/1".to_owned(), + size: 1, + })); } #[test] diff --git a/src/search/src/config/model.rs b/src/search/src/config/model.rs index d7a5c31..6c83611 100644 --- a/src/search/src/config/model.rs +++ b/src/search/src/config/model.rs @@ -12,15 +12,13 @@ use crate::domain::{ #[serde(default, deny_unknown_fields)] pub(crate) struct AppConfigDto { pub(crate) data_dir: PathBuf, - pub(crate) content_filter: ContentFilterConfig, + pub(crate) content_filter: ContentFilterConfigDto, pub(crate) persistence_queue_capacity: usize, pub(crate) stats_interval_secs: u64, - pub(crate) run_duration_secs: Option, pub(crate) index_batch_size: usize, pub(crate) index_interval_millis: u64, pub(crate) metadata_limits: MetadataLimitsConfig, pub(crate) dht: DhtConfig, - pub(crate) disk_guard: DiskGuardConfig, pub(crate) backup: BackupConfig, pub(crate) diagnostics: DiagnosticsConfig, pub(crate) logging: LoggingConfig, @@ -28,6 +26,32 @@ pub(crate) struct AppConfigDto { pub(crate) verification: VerificationConfig, } +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub(crate) struct ContentFilterConfigDto { + pub(crate) file_name_patterns: Vec, + pub(crate) file_path_patterns: Vec, +} + +impl ContentFilterConfigDto { + pub(crate) fn to_domain(&self) -> ContentFilterConfig { + let mut file_rules = pattern_rules( + "file-name", + FileMatchField::FileName, + &self.file_name_patterns, + ); + file_rules.extend(pattern_rules( + "file-path", + FileMatchField::FilePath, + &self.file_path_patterns, + )); + ContentFilterConfig { + version: 1, + file_rules, + } + } +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(default, deny_unknown_fields)] pub(crate) struct DhtConfig { @@ -60,21 +84,11 @@ pub(crate) struct HttpConfig { pub(crate) web_dir: PathBuf, } -#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] -#[serde(default, deny_unknown_fields)] -pub(crate) struct DiskGuardConfig { - pub(crate) enabled: bool, - pub(crate) check_interval_secs: u64, - pub(crate) minimum_free_bytes: u64, - pub(crate) resume_free_bytes: u64, -} - #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(default, deny_unknown_fields)] pub(crate) struct LoggingConfig { pub(crate) directory: PathBuf, pub(crate) file_enabled: bool, - pub(crate) console_enabled: bool, pub(crate) rotation: LogRotation, pub(crate) retain_files: usize, pub(crate) file_prefix: String, @@ -148,12 +162,10 @@ impl Default for AppConfigDto { content_filter: default_content_filter(), persistence_queue_capacity: 8_192, stats_interval_secs: 10, - run_duration_secs: None, index_batch_size: 1_024, index_interval_millis: 5_000, metadata_limits: MetadataLimitsConfig::default(), dht: DhtConfig::default(), - disk_guard: DiskGuardConfig::default(), backup: BackupConfig::default(), diagnostics: DiagnosticsConfig::default(), logging: LoggingConfig::default(), @@ -163,44 +175,50 @@ impl Default for AppConfigDto { } } -fn default_content_filter() -> ContentFilterConfig { - ContentFilterConfig { - version: 1, - file_rules: vec![ - FileFilterRule { - id: "bitcomet-padding-file".to_owned(), - enabled: true, - field: FileMatchField::FileName, - match_kind: FileMatchKind::Prefix, - value: "_____padding_file_".to_owned(), - case_sensitive: false, - action: FileRuleAction::Hide, - reason: "BitComet 分片边界填充文件".to_owned(), - }, - FileFilterRule { - id: "generic-pad-directory".to_owned(), - enabled: true, - field: FileMatchField::FilePath, - match_kind: FileMatchKind::Regex, - value: r"(^|/)\.pad/".to_owned(), - case_sensitive: false, - action: FileRuleAction::Hide, - reason: "客户端分片边界填充目录".to_owned(), - }, - FileFilterRule { - id: "libtorrent-padding-directory".to_owned(), - enabled: true, - field: FileMatchField::FilePath, - match_kind: FileMatchKind::Regex, - value: r"(^|/)\.____padding_file/".to_owned(), - case_sensitive: false, - action: FileRuleAction::Hide, - reason: "libtorrent 分片边界填充目录".to_owned(), - }, - ], +fn default_content_filter() -> ContentFilterConfigDto { + ContentFilterConfigDto { + file_name_patterns: vec!["*_____padding_file_*".to_owned()], + file_path_patterns: vec!["*.pad/*".to_owned(), "*.____padding_file/*".to_owned()], } } +fn pattern_rules( + id_prefix: &str, + field: FileMatchField, + patterns: &[String], +) -> Vec { + let mut patterns: Vec<_> = patterns + .iter() + .map(|pattern| pattern.trim()) + .filter(|pattern| !pattern.is_empty()) + .map(|pattern| { + if field == FileMatchField::FilePath { + pattern.replace('\\', "/").to_lowercase() + } else { + pattern.to_lowercase() + } + }) + .collect(); + patterns.sort_unstable(); + patterns.dedup(); + patterns + .into_iter() + .map(|value| { + let fingerprint = blake3::hash(format!("{id_prefix}\0{value}").as_bytes()); + FileFilterRule { + id: format!("{id_prefix}-{fingerprint}"), + enabled: true, + field, + match_kind: FileMatchKind::Wildcard, + value, + case_sensitive: false, + action: FileRuleAction::Hide, + reason: String::new(), + } + }) + .collect() +} + impl Default for MetadataLimitsConfig { fn default() -> Self { Self { @@ -249,23 +267,11 @@ impl Default for HttpConfig { } } -impl Default for DiskGuardConfig { - fn default() -> Self { - Self { - enabled: true, - check_interval_secs: 10, - minimum_free_bytes: 5 * 1024 * 1024 * 1024, - resume_free_bytes: 6 * 1024 * 1024 * 1024, - } - } -} - impl Default for LoggingConfig { fn default() -> Self { Self { directory: PathBuf::from("data/logs"), file_enabled: true, - console_enabled: false, rotation: LogRotation::Daily, retain_files: 7, file_prefix: "dht-search".to_owned(), @@ -276,7 +282,7 @@ impl Default for LoggingConfig { impl Default for BackupConfig { fn default() -> Self { Self { - enabled: true, + enabled: false, directory: PathBuf::from("data/backups"), interval_secs: 6 * 60 * 60, retain_checkpoints: 3, diff --git a/src/search/src/config/runtime.rs b/src/search/src/config/runtime.rs index c5cbc8b..52e5366 100644 --- a/src/search/src/config/runtime.rs +++ b/src/search/src/config/runtime.rs @@ -28,7 +28,7 @@ impl AppConfig { } pub(crate) fn content_filter(&self) -> Result { - ContentFilter::compile(self.content_filter.clone()).map_err(AppError::from) + ContentFilter::compile(self.content_filter.to_domain()).map_err(AppError::from) } pub(crate) fn dht_options(&self) -> DHTOptions { @@ -95,7 +95,7 @@ impl AppConfig { } fn validate(&self) -> Result<(), AppError> { - ContentFilter::compile(self.content_filter.clone())?; + ContentFilter::compile(self.content_filter.to_domain())?; if self.persistence_queue_capacity == 0 { return Err(AppError::Config( "persistence_queue_capacity 必须大于零".to_owned(), @@ -106,11 +106,6 @@ impl AppConfig { "stats_interval_secs 必须大于零".to_owned(), )); } - if self.run_duration_secs == Some(0) { - return Err(AppError::Config( - "run_duration_secs 必须大于零或不设置".to_owned(), - )); - } if self.index_batch_size == 0 || self.index_interval_millis == 0 { return Err(AppError::Config( "索引批量大小和执行间隔必须大于零".to_owned(), @@ -141,14 +136,6 @@ impl AppConfig { "验证队列容量并发尝试数租约和轮询间隔必须大于零".to_owned(), )); } - if self.disk_guard.check_interval_secs == 0 - || self.disk_guard.minimum_free_bytes == 0 - || self.disk_guard.resume_free_bytes <= self.disk_guard.minimum_free_bytes - { - return Err(AppError::Config( - "磁盘检查间隔必须大于零且恢复阈值必须大于保护阈值".to_owned(), - )); - } if self.backup.interval_secs == 0 || self.backup.retain_checkpoints == 0 { return Err(AppError::Config( "备份间隔和检查点保留数量必须大于零".to_owned(), @@ -171,11 +158,6 @@ impl AppConfig { "检查点目录不能位于 RocksDB 数据库目录内部".to_owned(), )); } - if !self.logging.file_enabled && !self.logging.console_enabled { - return Err(AppError::Config( - "文件日志和终端日志不能同时关闭".to_owned(), - )); - } if self.logging.file_enabled && self.logging.retain_files == 0 { return Err(AppError::Config("日志保留文件数量必须大于零".to_owned())); } diff --git a/src/search/src/crawler/pipeline.rs b/src/search/src/crawler/pipeline.rs index 122f8b5..8a1af79 100644 --- a/src/search/src/crawler/pipeline.rs +++ b/src/search/src/crawler/pipeline.rs @@ -246,13 +246,8 @@ mod tests { use dht_crawler::FileInfo; use super::*; - use crate::config::DiskGuardConfig; - fn disk_guard() -> DiskGuard { - DiskGuard::new(&DiskGuardConfig { - enabled: false, - ..DiskGuardConfig::default() - }) + DiskGuard::new() } #[derive(Default)] @@ -420,12 +415,7 @@ mod tests { #[tokio::test] async fn read_only_guard_rejects_new_metadata_without_blocking_shutdown() { let repository = Arc::new(MemoryRepository::default()); - let guard = DiskGuard::new(&DiskGuardConfig { - enabled: true, - check_interval_secs: 1, - minimum_free_bytes: 100, - resume_free_bytes: 200, - }); + let guard = DiskGuard::with_thresholds_for_test(100, 200); guard.observe_available(99, 0); let pipeline = PersistencePipeline::start(repository.clone(), 1, MetadataLimits::default(), guard); diff --git a/src/search/src/diagnostics/mod.rs b/src/search/src/diagnostics/mod.rs index 6cd910b..d1374c1 100644 --- a/src/search/src/diagnostics/mod.rs +++ b/src/search/src/diagnostics/mod.rs @@ -334,9 +334,7 @@ mod tests { use tempfile::TempDir; use super::*; - use crate::{ - config::DiskGuardConfig, crawler::pipeline::PersistencePipeline, disk_guard::DiskGuard, - }; + use crate::{crawler::pipeline::PersistencePipeline, disk_guard::DiskGuard}; #[tokio::test] async fn runtime_persists_a_queryable_snapshot_and_stops_cleanly() { @@ -344,10 +342,7 @@ mod tests { let repository = Arc::new(RocksTorrentRepository::open(directory.path().join("rocksdb")).unwrap()); let repository_trait: Arc = repository.clone(); - let disk_guard = DiskGuard::new(&DiskGuardConfig { - enabled: false, - ..DiskGuardConfig::default() - }); + let disk_guard = DiskGuard::new(); let persistence = PersistencePipeline::start( repository_trait, 4, diff --git a/src/search/src/disk_guard.rs b/src/search/src/disk_guard.rs index 966452f..fb0182d 100644 --- a/src/search/src/disk_guard.rs +++ b/src/search/src/disk_guard.rs @@ -12,9 +12,16 @@ use std::{ use tokio_util::sync::CancellationToken; -use crate::{config::DiskGuardConfig, crawler::pipeline::PersistenceIngress}; +use crate::crawler::pipeline::PersistenceIngress; const UNKNOWN_AVAILABLE_BYTES: u64 = u64::MAX; +const CHECK_INTERVAL_SECS: u64 = 60; +const MIB: u64 = 1024 * 1024; +const GIB: u64 = 1024 * MIB; +const MINIMUM_FLOOR: u64 = 512 * MIB; +const MINIMUM_CEILING: u64 = 20 * GIB; +const RESUME_BUFFER_FLOOR: u64 = 256 * MIB; +const RESUME_BUFFER_CEILING: u64 = 5 * GIB; #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub(crate) enum DiskMode { @@ -52,9 +59,8 @@ pub(crate) struct DiskGuard { } struct DiskGuardInner { - enabled: bool, - minimum_free_bytes: u64, - resume_free_bytes: u64, + minimum_free_bytes: AtomicU64, + resume_free_bytes: AtomicU64, state: Mutex, available_bytes: AtomicU64, probe_failed: AtomicBool, @@ -75,11 +81,10 @@ pub(crate) struct DiskWritePermit { pub(crate) async fn run( guard: DiskGuard, data_dir: PathBuf, - config: DiskGuardConfig, persistence: PersistenceIngress, cancel: CancellationToken, ) { - let mut ticker = tokio::time::interval(Duration::from_secs(config.check_interval_secs)); + let mut ticker = tokio::time::interval(Duration::from_secs(CHECK_INTERVAL_SECS)); ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); loop { tokio::select! { @@ -92,12 +97,12 @@ pub(crate) async fn run( } impl DiskGuard { - pub(crate) fn new(config: &DiskGuardConfig) -> Self { + pub(crate) fn new() -> Self { + let (minimum_free_bytes, resume_free_bytes) = automatic_thresholds(100 * GIB); Self { inner: Arc::new(DiskGuardInner { - enabled: config.enabled, - minimum_free_bytes: config.minimum_free_bytes, - resume_free_bytes: config.resume_free_bytes, + minimum_free_bytes: AtomicU64::new(minimum_free_bytes), + resume_free_bytes: AtomicU64::new(resume_free_bytes), state: Mutex::new(GateState { mode: DiskMode::Normal, active_writes: 0, @@ -112,12 +117,37 @@ impl DiskGuard { } pub(crate) fn probe(&self, path: &Path, persistence_queue: usize) { - match fs2::available_space(path) { - Ok(available) => self.observe_available(available, persistence_queue), - Err(error) => self.observe_probe_error(&error, persistence_queue), + match (fs2::total_space(path), fs2::available_space(path)) { + (Ok(total), Ok(available)) => { + let (minimum, resume) = automatic_thresholds(total); + self.inner + .minimum_free_bytes + .store(minimum, Ordering::Relaxed); + self.inner + .resume_free_bytes + .store(resume, Ordering::Relaxed); + self.observe_available(available, persistence_queue); + } + (Err(error), _) | (_, Err(error)) => { + self.observe_probe_error(&error, persistence_queue) + } } } + #[cfg(test)] + pub(crate) fn with_thresholds_for_test(minimum: u64, resume: u64) -> Self { + let guard = Self::new(); + guard + .inner + .minimum_free_bytes + .store(minimum, Ordering::Relaxed); + guard + .inner + .resume_free_bytes + .store(resume, Ordering::Relaxed); + guard + } + pub(crate) fn begin_admission(&self) -> Option { let permit = self.begin_write(false); if permit.is_none() { @@ -153,8 +183,8 @@ impl DiskGuard { DiskGuardSnapshot { mode: state.mode, available_bytes: (available != UNKNOWN_AVAILABLE_BYTES).then_some(available), - minimum_free_bytes: self.inner.minimum_free_bytes, - resume_free_bytes: self.inner.resume_free_bytes, + minimum_free_bytes: self.inner.minimum_free_bytes.load(Ordering::Relaxed), + resume_free_bytes: self.inner.resume_free_bytes.load(Ordering::Relaxed), active_writes: state.active_writes, probe_failed: self.inner.probe_failed.load(Ordering::Relaxed), probe_failures: self.inner.probe_failures.load(Ordering::Relaxed), @@ -164,17 +194,6 @@ impl DiskGuard { } fn begin_write(&self, allow_draining: bool) -> Option { - if !self.inner.enabled { - let mut state = self - .inner - .state - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); - state.active_writes = state.active_writes.saturating_add(1); - return Some(DiskWritePermit { - inner: self.inner.clone(), - }); - } let mut state = self .inner .state @@ -196,13 +215,12 @@ impl DiskGuard { .available_bytes .store(available, Ordering::Relaxed); self.inner.probe_failed.store(false, Ordering::Relaxed); - if !self.inner.enabled { - return; - } - if available < self.inner.minimum_free_bytes { + let minimum_free_bytes = self.inner.minimum_free_bytes.load(Ordering::Relaxed); + let resume_free_bytes = self.inner.resume_free_bytes.load(Ordering::Relaxed); + if available < minimum_free_bytes { self.transition_to(DiskMode::Draining, Some(available), None); self.finish_draining(persistence_queue); - } else if available >= self.inner.resume_free_bytes { + } else if available >= resume_free_bytes { self.transition_to(DiskMode::Normal, Some(available), None); } else { self.finish_draining(persistence_queue); @@ -215,9 +233,6 @@ impl DiskGuard { .store(UNKNOWN_AVAILABLE_BYTES, Ordering::Relaxed); self.inner.probe_failed.store(true, Ordering::Relaxed); self.inner.probe_failures.fetch_add(1, Ordering::Relaxed); - if !self.inner.enabled { - return; - } self.transition_to(DiskMode::Draining, None, Some(error)); self.finish_draining(persistence_queue); } @@ -283,6 +298,12 @@ impl DiskGuard { } } +fn automatic_thresholds(total_bytes: u64) -> (u64, u64) { + let minimum = (total_bytes / 20).clamp(MINIMUM_FLOOR, MINIMUM_CEILING); + let resume_buffer = (total_bytes / 100).clamp(RESUME_BUFFER_FLOOR, RESUME_BUFFER_CEILING); + (minimum, minimum.saturating_add(resume_buffer)) +} + impl Drop for DiskWritePermit { fn drop(&mut self) { let mut state = self @@ -299,12 +320,7 @@ mod tests { use super::*; fn guard() -> DiskGuard { - DiskGuard::new(&DiskGuardConfig { - enabled: true, - check_interval_secs: 1, - minimum_free_bytes: 100, - resume_free_bytes: 200, - }) + DiskGuard::with_thresholds_for_test(100, 200) } #[test] @@ -343,4 +359,11 @@ mod tests { assert!(guard.begin_new_write().is_none()); assert!(guard.begin_drain_write().is_some()); } + + #[test] + fn thresholds_scale_with_capacity_and_remain_bounded() { + assert_eq!(automatic_thresholds(20 * GIB), (GIB, GIB + 256 * MIB)); + assert_eq!(automatic_thresholds(100 * GIB), (5 * GIB, 6 * GIB)); + assert_eq!(automatic_thresholds(1024 * GIB), (20 * GIB, 25 * GIB)); + } } diff --git a/src/search/src/entry.rs b/src/search/src/entry.rs index 392614e..8ace92d 100644 --- a/src/search/src/entry.rs +++ b/src/search/src/entry.rs @@ -16,7 +16,6 @@ pub async fn run_cli() -> std::process::ExitCode { let mut logging = startup.app.logging.clone(); if restore_requested { logging.file_enabled = false; - logging.console_enabled = true; } let _telemetry = match telemetry::init(&logging) { Ok(guard) => guard, diff --git a/src/search/src/telemetry.rs b/src/search/src/telemetry.rs index 8970885..e000e94 100644 --- a/src/search/src/telemetry.rs +++ b/src/search/src/telemetry.rs @@ -15,12 +15,10 @@ pub(crate) struct TelemetryGuard { pub(crate) fn init(config: &LoggingConfig) -> Result { let filter = EnvFilter::try_from_default_env() .unwrap_or_else(|_| EnvFilter::new("warn,dht_search=info,dht_crawler=info")); - let console_layer = config.console_enabled.then(|| { - tracing_subscriber::fmt::layer() - .with_target(true) - .with_ansi(std::io::IsTerminal::is_terminal(&std::io::stderr())) - .with_writer(std::io::stderr) - }); + let console_layer = tracing_subscriber::fmt::layer() + .with_target(true) + .with_ansi(std::io::IsTerminal::is_terminal(&std::io::stderr())) + .with_writer(std::io::stderr); let (file_layer, file_guard) = if config.file_enabled { let appender = build_file_appender(config)?; let (writer, guard) = tracing_appender::non_blocking(appender); diff --git a/src/search/src/verification.rs b/src/search/src/verification.rs index bab07a3..91568fc 100644 --- a/src/search/src/verification.rs +++ b/src/search/src/verification.rs @@ -15,8 +15,6 @@ use dht_crawler::DHTServer; use tokio::task::{JoinHandle, JoinSet}; use tokio_util::sync::CancellationToken; -#[cfg(test)] -use crate::config::DiskGuardConfig; use crate::{ config::VerificationConfig, disk_guard::{DiskGuard, DiskWritePermit}, @@ -104,10 +102,7 @@ impl VerificationIngress { repository, capacity, VerificationStats::default(), - DiskGuard::new(&DiskGuardConfig { - enabled: false, - ..DiskGuardConfig::default() - }), + DiskGuard::new(), ) } diff --git a/src/web/README.md b/src/web/README.md index 2409d3e..4accc07 100644 --- a/src/web/README.md +++ b/src/web/README.md @@ -63,7 +63,7 @@ Axum 根据配置中的 `http.web_dir` 提供静态资源和单页回退 不需 - 仅在系统诊断页打开时每秒刷新服务状态 - 独立展示进程 RocksDB Tantivy DHT 和队列诊断数据 - 使用懒加载 ECharts 展示内存存储压力 Metadata 吞吐以及 HTTP 请求错误趋势 -- 使用 reka-ui 页签按职责分组的完整配置查看编辑校验保存和重启提示 +- 使用设置与过滤两个精简页签管理常用配置 Metadata 单位换算和按行通配符 - 搜索诊断配置使用互不冲突的单页路由 - 加载 空结果 接口错误 重试和移动端适配 diff --git a/src/web/src/pages/ConfigPage.vue b/src/web/src/pages/ConfigPage.vue index bf91a4e..77a1d49 100644 --- a/src/web/src/pages/ConfigPage.vue +++ b/src/web/src/pages/ConfigPage.vue @@ -1,108 +1,55 @@