perf(search): coalesce repeated index refreshes

This commit is contained in:
chuan
2026-08-11 12:24:42 +08:00
parent e74ebf3a95
commit 753534f53c
17 changed files with 630 additions and 394 deletions
+4
View File
@@ -17,6 +17,8 @@ RocksDB 是唯一权威数据源 Tantivy 索引可以从 RocksDB 完整重建
Tantivy 使用代际影子索引处理 Schema 文档格式损坏和缺失等全量重建 用户过滤规则通过可恢复扫描只增量更新真正变化的搜索文档 两类进度都可以在搜索提示与系统诊断页查看
重复发现只会精确更新 RocksDB 同一个内容组每六小时最多触发一次 Tantivy 动态状态刷新 新内容过滤变化和可用性验证仍会立即更新索引 搜索结果返回前会从 RocksDB 补全当前精确状态 因此界面数字保持最新但最近发现次数和热度的全局排序允许最多六小时的索引快照延迟
应用全部配置和内容隐藏规则统一位于 [`config.toml`](config.toml)
Web 右上角的无线电图标可以即时停止或恢复 DHT 持续采集 状态会写回 `dht.enabled` 关闭后不会建立 DHT 和 Metadata 网络任务 但现有 RocksDB 数据仍会继续补建索引并提供本地搜索
@@ -151,6 +153,8 @@ bun run build
每条记录包含 2048 个文件的极端基准中 450 个内容文档的 Tantivy 索引为 19.84 MiB 平均每文档 46.2 KiB 峰值内存为 113.08 MiB
运行诊断中的 `index_refresh_scheduled``index_refresh_suppressed` 分别表示重复发现实际安排和被时间桶合并的索引刷新数 `index_documents_written``index_documents_skipped` 用于确认待索引任务最终是否产生 Tantivy 文档写入
## 相关文档
- 当前实施状态和后续计划见 [`TODOS.md`](TODOS.md)
+61 -344
View File
@@ -1,386 +1,103 @@
# DHT 元数据搜索服务计划
# DHT 元数据搜索服务待办
本文档记录项目当前规划实施顺序和完成状态
本文档记录尚未完成且确定需要继续推进的工作
它是随需求实现结果性能数据和部署条件持续调整的活文档
已经验收的功能不再保留在此文件中 实现现状和使用方式统一记录在 `README.md`
## 维护规则
- 已经通过验收的任务使用 `[x]` 标记
- 正在规划但尚未完成的任务使用 `[ ]` 标记
- 需求变化时允许新增删除拆分合并或调整阶段顺序
- 调整计划时同步修改任务说明依赖关系和验收标准
- 不因代码已经存在就标记完成必须满足对应验收标准
- 发现原方案不合适时记录新决策并更新后续阶段
- 每次完成一个可交付功能时同步更新本文档
- 只记录仍需实施或仍需验收的任务
- 完成并通过相应验证后直接删除任务
- 需求或技术决策变化时同步调整任务和验收标准
- 不保留已经放弃或当前没有实际需求的预留功能
- 涉及 Tantivy Schema 的改动尽量合并以减少全量重建次数
## 当前技术方向
- `src/crawler` 负责可复用的 DHT 协议节点发现 Peer 查找和 Metadata 下载
- `src/search` 负责持久化去重索引搜索接口配置和运行生命周期
- `src/web` 负责最终用户搜索诊断和配置管理界面
- RocksDB 保存权威数据去重信息和任务状态
- Tantivy 保存可以从 RocksDB 重建的搜索索引
- Axum 提供搜索详情统计和健康检查接口
- 用户配置使用强类型 DTO 表达 TOML 只是当前持久化适配器
- SQLite 保存有明确保留上限且可安全删除的运行诊断历史
- 所有长期任务通过有界队列和背压控制资源占用
如果实际运行证明 RocksDB 的构建部署或资源成本不合适可以重新评估 redb SQLite 或其他存储方案
## 阶段零 项目基础
## 阶段一 搜索索引空间优化
### 目标
建立清晰的 workspace 边界开发规则和可持续验证的基础库
降低真实数据持续增长时的 Tantivy 磁盘占用并避免磁盘保护过早停止采集
### 任务
- [x] 将 crawler search 和 web 源码统一收纳到根目录 `src`
- [x] 使用当前 Git 配置统一作者仓库许可证和 edition 元数据
- [x] 编写 `AGENTS.md` 记录架构边界和开发约定
- [x] 将最终应用与可复用 DHT 基础库分离
- [x] 删除被内容组索引替代的旧状态和无效兼容代码
- [x] 收紧仅供测试或存储内部使用的接口和依赖
- [x] 按运行统计调度队列 Peer 竞速和 RocksDB 生命周期职责拆分核心大模块
- [x] 实现 BEP-51 `sample_infohashes` 主动发现
- [x] 实现主动 Peer 查找和 Metadata 获取
- [x] 验证远程公网设备能够持续获取 Metadata
- [x] 确认 Xray 全局代理会影响 Metadata TCP 连接并完成旁路验证
- [x] 保持基础库测试通过
- [ ] 在远端真实数据上验证六小时动态状态分桶对 Tantivy 文档重写和删除比例的改善
- [ ] 对比优化前后的每小时索引增长提交文档数和段合并压力
- [ ] 测量标题别名文件名完整路径 N-Gram fast field 和 stored field 的实际空间占比
- [ ] 根据测量结果设计尽量不降低常用搜索效果的索引精简方案
- [ ] 评估大型种子文件采样数量和路径文本预算对搜索覆盖率与索引体积的影响
- [ ] 明确优化后的单文档平均占用重建峰值空间和预期最大可容纳种子数量
### 验收标准
- [x] `cargo check --workspace --all-targets` 通过
- [x] `dht-crawler` 单元测试通过
- [x] `opencodes` 参考项目不参与 workspace 构建
- [ ] 使用远端真实数据或等比例数据对比优化前后的索引体积
- [ ] 普通文本通配符正则扩展名精确哈希和详情文件匹配测试保持通过
- [ ] 影子索引重建期间旧索引继续提供搜索且完成后能够原子切换
- [ ] 重建峰值空间不会触发磁盘只读保护
## 阶段一 本地持久化基础
## 阶段二 搜索排序与热度完善
### 目标
建立跨重启保留的权威数据源并完成精确去重和内容聚合基础
让热度随时间和活跃状态变化并让完全同分的搜索结果保持稳定顺序
### 任务
- [x] 定义二十字节 `InfoHash` 类型和十六进制转换
- [x] 定义 `TorrentRecord` `TorrentFile` 和内容组索引状态
- [x] 校验名称文件列表文件总大小和 infohash
- [x] 使用 BLAKE3 计算规范化内容指纹
- [x] 规范化 Unicode 路径分隔符大小写和文件顺序
- [x] 保留真实子目录避免内容指纹碰撞
- [x] 定义 RocksDB 二进制键空间
- [x] 实现数据库格式检查
- [x] 实现 infohash 精确查询和存在性判断
- [x] 实现新记录 WriteBatch 原子写入
- [x] 实现重复 infohash 的 `last_seen` `seen_count` 和 Peer 更新
- [x] 实现相同内容不同 infohash 的聚合映射
- [x] 实现待索引记录查询和索引完成标记
- [x] 配置 Bloom Filter LZ4 压缩和有限 block cache
- [x] 准备 Windows 本地 RocksDB 构建所需的 libclang
- [x] 将本地构建工具目录排除出 Git
### 验收标准
- [x] 数据库关闭并重新打开后记录仍可读取
- [x] 重复写入不会创建第二条 torrent 记录
- [x] 重复写入会正确增加发现次数
- [x] 相同内容的不同 infohash 可以独立保存并聚合查询
- [x] Metadata 主体内容映射和待索引标记原子写入
- [x] RocksDB 功能测试通过
- [x] `dht-search` Clippy `-D warnings` 通过
## 阶段二 采集持久化闭环
### 目标
让 DHT 获取的真实 Metadata 自动进入有界持久化管线并支持安全停止和重新启动
### 任务
- [x] 定义应用配置结构和默认配置文件
- [x] 支持通过配置指定固定数据目录
- [x] 支持配置 DHT 端口并发队列容量和 Metadata 限制
- [x] 初始化 RocksDB repository 并处理启动错误
- [x]`TorrentInfo` callback 转换为 `TorrentRecord`
- [x] 建立有界持久化队列并实现背压
- [x] 使用专用阻塞任务执行 RocksDB 操作避免阻塞 Tokio worker
- [x] 在 Metadata 下载前查询持久化 infohash 状态减少重复下载
- [x] 在 BEP-51 Peer Lookup 前批量查询 RocksDB 并更新已有 infohash 发现状态
- [x] 将已存在记录更新为再次发现而不是重复创建
- [x] 增加接收写入重复拒绝失败和队列深度指标
- [x] 实现 `Ctrl+C` `SIGINT``SIGTERM` 优雅退出
- [x] 退出时停止接收新任务并排空或持久化剩余任务
- [x] 支持重新启动后继续使用原数据库
- [x] 将 example 运行方式替换为正式 `dht-search` 二进制
- [x] 支持通过运行时长参数进行间歇运行
### 验收标准
- [x] 本地运行可以持续向 RocksDB 写入真实 Metadata
- [x] 停止并重启后旧 infohash 不会作为新记录重复写入
- [x] 队列达到容量时内存不继续无界增长
- [x] 正常退出后已接受的任务不会静默丢失
- [ ] 远程设备运行一小时没有持续内存增长
- [x] 记录采集速度重复率数据库增长和写入延迟
## 阶段三 Tantivy 搜索索引
### 目标
让持久化 Metadata 支持快速全文搜索过滤排序和索引恢复
### 任务
- [x] 定义 Tantivy schema 和索引版本
- [x] 索引名称文件路径扩展名 infohash 和内容指纹
- [x] 将大小文件数时间和发现次数定义为 fast fields
- [x] 设计中英文数字和文件名子串 tokenizer
- [x] 实现待索引任务批量消费
- [x] 实现按数量和时间间隔批量 commit
- [x] commit 成功后原子更新 RocksDB 索引状态
- [x] 实现关键词短语和精确 infohash 查询
- [x] 实现大小时间扩展名文件数热度和可用性过滤
- [x] 实现大小范围和扩展名过滤
- [x] 实现相关性时间热度大小和发现次数排序
- [x] 建立带时间衰减的 DHT 活跃度分数和用户可读等级
- [x] 实现分页并限制最大翻页成本
- [x] 实现相同 `content_key` 结果精确折叠和变体分页
- [x] 实现从 RocksDB 全量重建 Tantivy 索引
- [x] 支持索引结构不兼容时直接重建
- [x] 使用影子索引保留旧搜索并在完整校验后原子切换
- [x] 持久化影子索引构建清单并支持跨重启继续重建
- [x] 暴露种子内容组待索引数量和全量重建状态进度
- [x] 按文件大小为大型种子选择最多 2048 个文件并限制路径文本预算
- [x] 优化完整名称别名文件名路径的相关性权重并使用热度时间稳定同分结果
- [x] 将标题文件名和路径拆分为有限 N-Gram 与路径分词策略并移除无用位置索引
- [x] 搜索写入器按需创建使影子重建期间旧活动索引保持纯查询占用
### 验收标准
- [x] 新写入记录在目标延迟内可搜索
- [x] 搜索索引删除后可以从 RocksDB 完整重建
- [x] 索引过程中异常退出不会永久丢失文档
- [x] 全量重建期间旧索引继续提供完整旧结果且新数据在原子切换后可见
- [x] 百万级测试数据常用查询延迟达到 `README.md` 记录的目标
- [ ] 使用远端真实数据验证紧凑索引的最终体积重建峰值内存和稳态内存
## 阶段四 HTTP 搜索服务
### 目标
提供稳定可验证并且资源受限的搜索和详情接口
### 任务
- [x] 使用 Axum 建立 HTTP 服务
- [x] 实现 `/health``/ready` 接口
- [x] 实现 `/stats` 运行状态接口
- [x] 实现 `/search` 搜索过滤和分页接口
- [x] 实现 `/torrents/{infohash}` 详情接口
- [x] 实现 `/contents/{content_key}` 内容变体接口
- [x] 使用 Bun Vue TypeScript Vite Tailwind CSS 和 shadcn-vue 建立 Web 基础环境
- [x] 实现简单现代并适配移动端的单页搜索界面
- [x] 接入搜索排序分页详情内容变体和磁力链接复制
- [x] 详情文件分页优先展示匹配搜索条件的文件并继承外层大小或名称排序
- [x] 实现名称别名和文件路径的有限状态自动机正则搜索
- [x] 根据输入语法自动识别普通文本通配符和正则表达式并移除独立模式开关
- [x] 品牌入口可清除搜索查询排序分页和详情状态并返回主页
- [x] 实现种子详情文件列表后端分页并限制浏览器单页节点数量
- [x] 保留种子详情顶部结构并使用 reka-ui 数字分页重构文件条目
- [x] 支持用户选择并持久化文件列表每页数量
- [x] 统一搜索结果数字分页并支持用户选择每页数量
- [x] 将采集索引持久化和验证运行状态集中到系统诊断页并仅在页面打开时每秒刷新
- [x] 将诊断时间范围放入历史趋势区域并精简图表卡片的单行摘要信息
- [x] 在 Web 顶栏提供持久化的 DHT 即时启停开关并保留离线索引搜索能力
- [x] 将配置页精简为设置与过滤两个页签并合并全部常用参数
- [x] 使用数值与单位选择编辑 Metadata 大小并自动换算内部字节数
- [x] 实现浏览器持久化明暗主题
- [x] 为加载空结果接口错误和失败重试提供明确界面状态
- [x] 将生产静态资源交给 Axum 提供并支持单页回退
- [x] 提供 Windows 一键启动后端和 Web 开发服务的脚本
- [x] 修复一键启动脚本只停止父进程导致 Vite 子进程残留的问题
- [x] 定义统一错误响应
- [x] 限制查询长度分页大小和最大 offset
- [x] 增加请求延迟错误率和并发指标
- [x] 增加搜索详情字段和按需验证入队 API 端到端测试
- [x] 增加搜索过滤折叠精确哈希和变体接口测试
### 验收标准
- [x] API 能搜索真实采集数据
- [x] 非法参数返回稳定的客户端错误
- [x] 搜索查询在独立阻塞任务执行不会阻塞异步 worker
- [x] 健康检查能区分进程存活和服务可用
## 阶段五 质量过滤和重复内容控制
### 目标
减少垃圾数据和重复展示同时避免不可恢复的误删
### 任务
- [x] 将种子可用性定义为最近通过 DHT 找到并完成 BitTorrent 握手
- [x] 区分 Metadata 结构有效和 swarm 当前可用性
- [x] 实现未验证活跃和可能失效三态模型
- [x] 实现详情高优先级和搜索普通优先级的仅按需验证
- [x] 使用持久化有界验证队列租约恢复去重和失败退避
- [x] Metadata 与验证握手共享 TCP 建连总预算
- [x] 开发阶段清理测试数据库并以当前数据结构重新采集
- [x] 搜索和详情接口返回热度与可用性数据
- [x] `/stats` 返回验证队列发现握手成功失败和拒绝指标
- [ ] 统计真实数据的 infohash 重复率和内容重复率
- [x] 在统一 `config.toml` 中定义种子标题和内部文件隐藏规则
- [x] 使用标题与内部文件双文本框按行管理不区分大小写的通配符和 `regex:` 规则
- [x] 保留 RocksDB 原始文件列表并将用户过滤从内容指纹和内容组身份中解耦
- [x] 使用内容组过滤投影哈希只增量更新真正变化的 Tantivy 文档
- [x] 持久化过滤扫描游标并支持连续修改采用最新规则和跨重启恢复
- [x] 在搜索页和系统诊断页展示过滤基线扫描更新提交和失败状态
- [x] 默认隐藏 BitComet padding 文件以及 `.pad``.____padding_file` 填充目录
- [x] 全部文件被隐藏的 Metadata 只保留原始记录且不进入公开索引
- [ ] 根据真实垃圾数据决定是否增加种子名称扩展名和大小准入规则
- [x] 定义可配置的 Metadata 最大大小文件数名称路径长度和目录层级限制
- [x] 识别空名称控制字符异常路径大小溢出总大小不一致和文件数量攻击
- [ ] 设计可解释的名称标准化规则
- [ ] 为模糊相似结果生成聚合候选但不自动删除
- [x] 使用带规则指纹的 RocksDB 轻量拒绝记录阻止相同异常 infohash 重复下载
- [x] 保留按原因分类的过滤指标但避免保存名称和大文件列表
- [ ] 增加误判测试和边界数据集
### 验收标准
- [x] 搜索和详情响应不等待 DHT 或 Peer 网络验证
- [x] 进程重启后已接受的验证任务能够通过租约恢复
- [x] 一次验证失败不会删除记录或标记为绝对失效
- [x] 数据清空后能够建立新的内容组搜索文档
- [x] 精确重复不会重复下载和重复展示
- [x] 内容重复可以折叠并保留全部 infohash
- [x] 过滤规则可以配置更新和回滚且不会删除原始 Metadata
- [ ] 模糊去重不会直接造成数据丢失
## 阶段六 性能资源和长期运行
### 目标
以真实数据验证持续运行时的吞吐延迟磁盘放大和资源上限
### 任务
- [x] 增加 `find_node` Peer Lookup 新目标和 Metadata 建连的显式配置
- [x] 为主动 `find_node` `get_peers``sample_infohashes` 增加共享 UDP 查询总预算
- [x] 为 Metadata TCP 建连增加独立每秒速率限制
- [x] 使用保守网络预算定位早期本机断网问题
- [x] 根据本机首次验证将主动 UDP 从 `40/s` 下调至 `10/s` 并将 Metadata 建连从 `5/s` 下调至 `2/s` 完成故障隔离
- [x] 对照 Bitmagnet 默认并发建立受全局预算和有界队列保护的激进配置
- [x] 根据两轮一分钟资源测试将 Bitmagnet 等效激进配置设为应用和运行模板默认值
- [x] 修复 Windows 临时索引文件占用导致整个服务退出的问题
- [x] 验证极保守配置运行三分钟不影响同机代理网络并安全退出
- [x] 将 BEP-51 采样准入压力反向传递到采样查询调度
- [x] 实现样本来源节点单点 `get_peers` 优先和失败后有限递归降级
- [x] 将 BEP-51 最大在途请求和采样失败后的迭代回退暴露为应用配置
- [x] 使用 Bitmagnet 等效并发完成一分钟资源测试并确认主网卡无丢包无错误且持久化无积压
- [x] 将新发现但尚未验证的节点按地址稳定分流到有界 BEP-51 通道并避免与 `find_node` 重复探测
- [x] 增加直接采样候选队列请求响应重复过滤丢弃和 Metadata 成功来源转化指标
- [x] 使用相同网络预算对旧快照采样和新节点直接采样完成十分钟对比
- [x] 确认直接采样在相近 UDP 流量下 Metadata 成功数提高约百分之六点六且单条成功 UDP 成本降低约百分之六点六
- [x] 完成首轮三分钟对比并验证 Peer Lookup UDP 从 `278` 降至 `254` 且网络稳定
- [ ] 通过多轮或更长时间运行评估随机 DHT 样本下的 Metadata 成功率
- [ ] 设计可在查询时计算热度的索引字段避免权重和等级调整触发再次重建
- [ ] 为完全同分的搜索结果增加稳定的最终排序键
- [ ] 将热度字段稳定排序键和索引空间优化合并为一次 Tantivy Schema 更新
- [ ] 统计按需验证的 Peer 发现率握手成功率和平均验证耗时
- [x] 完成首轮真实按需验证并确认旧记录两次握手均成功更新为活跃
- [x] 验证新 Metadata 记录直接继承成功来源 Peer 的活跃状态
- [x] 验证启用按需可用性功能后保守预算运行三分钟并安全停止
- [ ] 根据真实验证数据校准热度权重等级边界和失败退避时间
- [ ] 根据公网设备长期实测设计超时率自动降速
- [x] 建立可重复的采集存储索引和查询规模基准工具
- [x] 完成一万十万和一百万条 release 基线并记录查询 P50 P95 P99
- [x] 记录每条元数据和每个索引文档的平均磁盘占用
- [x] 记录百万级基准进程峰值内存和索引总吞吐
- [x] 记录 RocksDB block cache memtable 和 compaction 指标
- [x] 记录 Tantivy IndexWriter 内存和 commit 延迟
- [ ] 根据实测调整批量大小队列容量和并发
- [x] 增加带排空阶段恢复滞回和探测失败保护的磁盘只读降级策略
- [x] 固定每分钟检查磁盘并按文件系统容量自动计算保护和恢复阈值
- [x] 增加在线 RocksDB 检查点保留上限只读校验和带旧库保留的离线恢复
- [x] 默认关闭自动检查点并从 Web 隐藏全部备份调度参数
- [x] 始终保留终端日志并仅在 Web 暴露滚动文件日志开关
- [x] 移除 Docker 对文件日志的命令行覆盖并隐藏无关的基础设施覆盖提示
- [x] 验证间歇运行和正常退出恢复
- [x] 完成本机约七小时真实持续运行并确认采集索引和搜索服务可用
- [ ] 验证二十四小时和七天连续运行
- [ ] 根据规模决定是否继续使用 RocksDB
### 验收标准
- [ ] 内存使用在目标上限内稳定
- [ ] 队列和缓存不会随运行时间无限增长
- [x] 磁盘不足时能够拒绝新任务排空持久化队列并保留搜索能力
- [x] 备份可以在独立目录恢复并从 RocksDB 重建索引搜索
- [ ] 连续运行期间没有数据格式损坏和不可恢复任务
- [ ] 热度能够随新发现可用性验证和时间衰减产生可观察变化
- [ ] 相同查询和相同数据重复执行时结果顺序保持一致
- [ ] 后续调整热度权重或等级边界不需要全量重建索引
## 阶段七 部署和运维
## 阶段三 数据质量统计与回归保护
### 目标
让应用可以在公网 Linux 设备上重复构建部署监控停止和恢复
用真实数据量化重复和过滤效果并防止规则调整造成误判
### 任务
- [x] 固化 Linux 目标构建方式和 RocksDB 构建依赖
- [x] 生成 release 二进制并使用 SHA-256 校验部署
- [x] 定义 Docker 配置数据日志索引和静态资源目录布局
- [x] 编写 Debian 三阶段最小运行镜像并使用非 root 用户完成构建启动和重启恢复测试
- [x] 编写 Compose 配置管理端口数据卷重启策略和文件句柄限制
- [x] Docker 将根目录唯一 `config.toml` 映射到容器并使用单独数据卷保存全部运行数据
- [x] 使用专用低权限 UID 运行验证
- [x] 在公网 Debian 设备验证镜像传输内部 Docker 网络 Caddy 反向代理 Compose 重启数据复用和优雅停止
- [x] 固化 Xray 环境下的最小范围网络旁路
- [x] 为 Docker DHT 容器增加独立直连网络和持久化 Xray 来源网段旁路
- [x] 验证旁路后一分钟内 Metadata 成功入库且 Xray 文件描述符和主机 UDP socket 恢复稳定
- [x] 验证高并发运行需要 `LimitNOFILE=65536`
- [ ] 实现启动前数据目录权限检查
- [ ] 实现优雅升级和回滚流程
- [ ] 编写备份恢复和故障排查文档
- [ ] 统计真实数据的 infohash 重复率和不同 infohash 的内容重复率
- [ ] 建立包含正常 Metadata 垃圾标题填充文件异常路径和极端文件数量的边界数据集
- [ ] 增加标题规则内部文件规则通配符正则和全部文件隐藏场景的误判测试
- [ ] 根据真实垃圾数据决定是否需要增加新的种子准入规则
### 验收标准
- [ ] 新设备可以按 Docker 文档完成部署
- [x] 容器重建并复用数据卷不会丢失已提交数据和诊断历史
- [ ] Xray 旁路只影响爬虫进程
- [ ] 更新失败时可以恢复上一版本二进制和数据
- [ ] 重复率和过滤结果可以通过诊断数据或离线工具复现
- [ ] 过滤规则不会删除 RocksDB 中的原始 Metadata
- [ ] 正常种子不会因为垃圾过滤规则被错误丢弃或错误聚合
## 阶段四 长期运行与部署收尾
### 目标
在版本保持不变的情况下完成长期运行验收并改善部署失败时的错误提示
### 任务
- [ ] 在应用启动前检查数据目录配置文件和关键子目录的读写权限
- [ ] 使用同一版本完成二十四小时连续运行
- [ ] 使用同一版本完成七天连续运行
- [ ] 记录连续运行期间的私有内存队列深度 Metadata 成功率磁盘增长和索引提交状态
- [ ]`README.md` 补充简洁的备份恢复和常见故障排查步骤
### 验收标准
- [ ] 权限错误在打开 RocksDB 或 Tantivy 前返回明确的路径和操作错误
- [ ] 私有内存和有界队列进入稳定区间且不随时间无限增长
- [ ] 连续运行期间没有数据格式损坏索引提交失败或不可恢复任务
- [ ] 容器重建后能够继续使用原 RocksDB Tantivy 和诊断历史
## 当前下一步
- [x] 按职责拆分应用编排索引 worker 运行监控搜索查询构建和 Tantivy 文档映射
- [x] 按领域边界拆分 infohash Metadata 校验内容聚合和种子活跃状态
- [x] 将 DHT 响应限流器从服务器编排中提取为独立组合组件
- [x] 将规模基准拆分为参数数据集工作负载采样报告和编排模块
- [x] 明确单元组件集成端到端和性能测试层级并增加公开 API 集成测试
- [x] 实现可回滚的无效文件过滤并重建有效内容聚合和搜索索引
- [x] 将 crawler search 和 web 统一迁移到根目录 `src` 并修复构建脚本文档路径
- [x] 使用领域 Metadata DTO 切断领域层对 DHT 传输 DTO 的直接依赖
- [x] 将配置拆分为可序列化 DTO TOML 读取适配器和运行时解析结果
- [x] 使用独立 SQLite 建立有界运行诊断历史存储
- [x] 采集进程 RocksDB Tantivy DHT 队列和磁盘资源快照
- [x] 在系统诊断中持续展示已保存种子总量并以历史趋势替代低价值 HTTP 请求图表
- [x] 提供当前诊断快照和原始或分钟历史查询接口
- [x] 将 HTTP 请求并发客户端错误服务端错误和延迟分布写入诊断历史
- [x] 增加配置查询完整校验原子保存并发修订和统一重启提示
- [x] 增加 Web 诊断页和配置管理页
- [x] 精简诊断页即时指标并由趋势图承担重复的资源和吞吐数据
- [x] 在系统诊断中区分 Peer 下载结果和 Metadata 内容有效性并展示 Peer 失败原因构成
- [x] 使用页签拆分配置分类并隐藏底层配置存储位置
- [x] 将主配置和内容过滤规则合并为唯一 `config.toml`
- [x] 将 Docker 测试和性能基线核心内容合并到根目录 `README.md`
- [x] 将 Xray 透明代理异常环境的判断处理和回滚方案独立到 `docs`
- [x] 实现可恢复的影子索引原子切换和 Web 重建进度
- [x] 将大型种子文件路径覆盖扩大到按大小选择的 2048 个文件和 256 KiB
- [x] 使用文本优先热度时间同分的稳定相关性排序
先部署并观察六小时动态状态分桶 使用诊断中的安排刷新抑制刷新实际写入和跳过写入数量确认效果
完成二十四小时持续运行并继续观察私有内存 Metadata 成功率候选队列深度和每条成功 Metadata 的网络成本
完成真实运行对比后再测量 Tantivy 字段空间占用并确定索引精简方案
RocksDB Tantivy HTTP 和进程资源指标已经接入诊断历史 后续根据长期实测继续调优
下一步持续观察公网设备 Metadata 成功率 Xray 文件描述符主机 socket 和磁盘增长是否长期稳定
Schema 调整时同时加入动态热度所需字段和稳定排序键 最终只执行一次影子索引重建
+35 -3
View File
@@ -52,7 +52,10 @@ pub(crate) async fn stats(State(state): State<ApiState>) -> Json<StatsResponse>
.crawler
.verification()
.map(|ingress| ingress.stats().snapshot());
let index = state.search.status(state.repository.index_inventory());
let storage = state.repository.index_inventory();
let refresh_diagnostics = state.repository.index_refresh_diagnostics();
let search_diagnostics = state.search.diagnostics();
let index = state.search.status(storage);
Json(StatsResponse {
http_active_requests: http.active_requests,
http_requests: http.requests,
@@ -144,6 +147,10 @@ pub(crate) async fn stats(State(state): State<ApiState>) -> Json<StatsResponse>
persistence_rejected_full: persistence.rejected_full,
persistence_queue: persistence.queue_depth,
indexed_documents: state.search.num_docs(),
index_refresh_scheduled: refresh_diagnostics.scheduled,
index_refresh_suppressed: refresh_diagnostics.suppressed,
index_documents_written: search_diagnostics.documents_written,
index_documents_skipped: search_diagnostics.documents_skipped,
index,
filter: state.filter.status(),
verification_queue: verification.map_or(0, |stats| stats.queue_depth),
@@ -314,8 +321,9 @@ pub(crate) async fn search(
} else {
None
};
let page = tokio::task::spawn_blocking(move || {
state.search.search_with(SearchOptions {
let search = state.search.clone();
let mut page = tokio::task::spawn_blocking(move || {
search.search_with(SearchOptions {
query,
mode: request.mode,
offset: request.offset,
@@ -338,6 +346,30 @@ pub(crate) async fn search(
.await
.map_err(|error| ApiError::internal(error.to_string()))?
.map_err(|error| ApiError::bad_request(error.to_string()))?;
if !page.hits.is_empty() {
let repository = state.repository.clone();
page = tokio::task::spawn_blocking(move || {
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
for hit in &mut page.hits {
let Ok(decoded) = hex::decode(&hit.content_key) else {
continue;
};
let Ok(content_key) = <[u8; 32]>::try_from(decoded) else {
continue;
};
if let Some(group) = repository.content_group(&content_key, now)? {
hit.refresh_from_group(&group);
}
}
Ok::<_, crate::storage::StorageError>(page)
})
.await
.map_err(|error| ApiError::internal(format!("搜索结果补全任务失败: {error}")))?
.map_err(|error| ApiError::internal(error.to_string()))?;
}
if let Some(verification) = &verification {
let hashes = page
.hits
+5
View File
@@ -358,6 +358,9 @@ mod tests {
.unwrap();
assert_eq!(response.status(), StatusCode::CONFLICT);
repository.observe_existing(record.info_hash, 30).unwrap();
assert!(repository.pending_index(10).unwrap().is_empty());
let response = app
.clone()
.oneshot(
@@ -376,6 +379,8 @@ mod tests {
assert!(json["hits"][0]["heat"]["score"].is_number());
assert_eq!(json["hits"][0]["availability"]["status"], "unknown");
assert_eq!(json["hits"][0]["variant_count"], 2);
assert_eq!(json["hits"][0]["last_seen"], 30);
assert_eq!(json["hits"][0]["seen_count"], 3);
assert_eq!(repository.verification_queue_len().unwrap(), 1);
let response = app
+4
View File
@@ -109,6 +109,10 @@ pub(crate) struct StatsResponse {
pub(crate) persistence_rejected_full: u64,
pub(crate) persistence_queue: usize,
pub(crate) indexed_documents: u64,
pub(crate) index_refresh_scheduled: u64,
pub(crate) index_refresh_suppressed: u64,
pub(crate) index_documents_written: u64,
pub(crate) index_documents_skipped: u64,
pub(crate) index: IndexStatus,
pub(crate) filter: FilterStatus,
pub(crate) verification_queue: u64,
+2
View File
@@ -302,6 +302,8 @@ mod tests {
_: &[u8; 32],
_: u64,
_: [u8; 32],
_: [u8; 32],
_: u64,
_: bool,
) -> Result<bool, StorageError> {
Ok(true)
+4
View File
@@ -283,6 +283,8 @@ fn collect_sample(sources: &DiagnosticSources, session_started_at: u64) -> Diagn
live_sst_bytes: storage.live_sst_bytes,
running_compactions: storage.running_compactions,
estimated_keys: storage.estimated_keys,
index_refresh_scheduled: storage.index_refresh_scheduled,
index_refresh_suppressed: storage.index_refresh_suppressed,
},
search: SearchDiagnostics {
documents: search.documents,
@@ -292,6 +294,8 @@ fn collect_sample(sources: &DiagnosticSources, session_started_at: u64) -> Diagn
last_commit_at: search.last_commit_at,
last_commit_duration_millis: search.last_commit_duration_millis,
last_commit_documents: search.last_commit_documents,
documents_written: search.documents_written,
documents_skipped: search.documents_skipped,
},
http: sources.http.snapshot(),
runtime: RuntimeDiagnostics {
+8
View File
@@ -33,6 +33,10 @@ pub(crate) struct StorageDiagnostics {
pub(crate) live_sst_bytes: Option<u64>,
pub(crate) running_compactions: Option<u64>,
pub(crate) estimated_keys: Option<u64>,
#[serde(default)]
pub(crate) index_refresh_scheduled: u64,
#[serde(default)]
pub(crate) index_refresh_suppressed: u64,
}
#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, PartialEq, Eq)]
@@ -44,6 +48,10 @@ pub(crate) struct SearchDiagnostics {
pub(crate) last_commit_at: Option<u64>,
pub(crate) last_commit_duration_millis: u64,
pub(crate) last_commit_documents: u64,
#[serde(default)]
pub(crate) documents_written: u64,
#[serde(default)]
pub(crate) documents_skipped: u64,
}
#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, PartialEq, Eq)]
+37 -3
View File
@@ -48,6 +48,8 @@ pub struct SearchDiagnostics {
pub last_commit_at: Option<u64>,
pub last_commit_duration_millis: u64,
pub last_commit_documents: u64,
pub documents_written: u64,
pub documents_skipped: u64,
}
struct SearchInner {
@@ -61,6 +63,8 @@ struct SearchInner {
last_commit_at: AtomicU64,
last_commit_duration_millis: AtomicU64,
last_commit_documents: AtomicU64,
documents_written: AtomicU64,
documents_skipped: AtomicU64,
}
impl SearchEngine {
@@ -155,6 +159,8 @@ impl SearchEngine {
last_commit_at: AtomicU64::new(0),
last_commit_duration_millis: AtomicU64::new(0),
last_commit_documents: AtomicU64::new(0),
documents_written: AtomicU64::new(0),
documents_skipped: AtomicU64::new(0),
}),
})
}
@@ -163,6 +169,18 @@ impl SearchEngine {
if documents.is_empty() {
return Ok(());
}
let write_count = documents
.iter()
.filter(|document| document.requires_write)
.count();
let skipped_count = documents.len().saturating_sub(write_count);
self.inner.documents_skipped.fetch_add(
skipped_count.min(u64::MAX as usize) as u64,
AtomicOrdering::Relaxed,
);
if write_count == 0 {
return Ok(());
}
let started = Instant::now();
let fields = self.inner.fields;
let result = (|| {
@@ -179,7 +197,7 @@ impl SearchEngine {
);
}
let writer = writer.as_mut().expect("search writer was initialized");
for document in documents {
for document in documents.iter().filter(|document| document.requires_write) {
writer.delete_term(Term::from_field_text(
fields.content_key,
&hex::encode(document.content_key),
@@ -202,7 +220,11 @@ impl SearchEngine {
.last_commit_at
.store(unix_timestamp(), AtomicOrdering::Relaxed);
self.inner.last_commit_documents.store(
documents.len().min(u64::MAX as usize) as u64,
write_count.min(u64::MAX as usize) as u64,
AtomicOrdering::Relaxed,
);
self.inner.documents_written.fetch_add(
write_count.min(u64::MAX as usize) as u64,
AtomicOrdering::Relaxed,
);
} else {
@@ -221,6 +243,9 @@ impl SearchEngine {
.map(|group| IndexDocument {
content_key: group.content_key,
projection_hash: crate::storage::projection_hash(Some(&group)),
dynamic_projection_hash: crate::storage::dynamic_projection_hash(Some(&group)),
refresh_bucket: crate::storage::refresh_bucket(group.last_seen),
requires_write: true,
group: Some(group),
})
.collect();
@@ -258,6 +283,8 @@ impl SearchEngine {
.inner
.last_commit_documents
.load(AtomicOrdering::Relaxed),
documents_written: self.inner.documents_written.load(AtomicOrdering::Relaxed),
documents_skipped: self.inner.documents_skipped.load(AtomicOrdering::Relaxed),
}
}
@@ -270,7 +297,9 @@ impl SearchEngine {
let tasks = repository.pending_index(limit)?;
let mut documents = Vec::with_capacity(tasks.len());
for task in &tasks {
documents.push(repository.index_document(&task.content_key, now)?);
let mut document = repository.index_document(&task.content_key, now)?;
document.requires_write |= task.force_write;
documents.push(document);
}
self.index_documents(&documents)?;
for (task, document) in tasks.iter().zip(&documents) {
@@ -278,6 +307,8 @@ impl SearchEngine {
&task.content_key,
task.revision,
document.projection_hash,
document.dynamic_projection_hash,
document.refresh_bucket,
document.group.is_some(),
)?;
}
@@ -507,6 +538,9 @@ mod tests {
content_key: record.content_key,
group: None,
projection_hash: [9; 32],
dynamic_projection_hash: [8; 32],
refresh_bucket: 0,
requires_write: true,
}])
.unwrap();
+20 -1
View File
@@ -2,7 +2,7 @@
use serde::{Deserialize, Serialize};
use crate::domain::{AvailabilityStatus, Heat, HeatLevel};
use crate::domain::{AvailabilityStatus, ContentGroup, Heat, HeatLevel};
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Deserialize, Serialize)]
#[serde(rename_all = "snake_case")]
@@ -62,6 +62,25 @@ pub struct SearchHit {
pub availability: AvailabilitySummary,
}
impl SearchHit {
pub(crate) fn refresh_from_group(&mut self, group: &ContentGroup) {
self.info_hash = group.representative.info_hash.to_string();
self.name.clone_from(&group.representative.name);
self.total_size = group.representative.total_size;
self.file_count = group.representative.original_file_count();
self.first_seen = group.first_seen;
self.last_seen = group.last_seen;
self.seen_count = group.seen_count;
self.variant_count = group.variant_count;
self.heat = group.heat;
self.availability = AvailabilitySummary {
status: group.availability.status,
last_verified_at: group.availability.last_verified_at,
reachable_peers: group.availability.reachable_peers,
};
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct AvailabilitySummary {
pub status: AvailabilityStatus,
+4 -4
View File
@@ -6,12 +6,12 @@ mod repository;
#[cfg(feature = "rocksdb-storage")]
mod rocks;
#[cfg(test)]
pub(crate) use repository::projection_hash;
pub use repository::{
CheckpointSummary, ContentGroupTask, ContentVariants, IndexDocument, IndexInventory,
StorageDiagnostics, StorageError, TorrentRepository, UpsertOutcome, VerificationEnqueueOutcome,
VerificationPriority, VerificationRequest,
IndexRefreshDiagnostics, StorageDiagnostics, StorageError, TorrentRepository, UpsertOutcome,
VerificationEnqueueOutcome, VerificationPriority, VerificationRequest,
};
#[cfg(test)]
pub(crate) use repository::{dynamic_projection_hash, projection_hash, refresh_bucket};
#[cfg(feature = "rocksdb-storage")]
pub use rocks::RocksTorrentRepository;
+73
View File
@@ -53,10 +53,18 @@ pub trait TorrentRepository: Send + Sync {
) -> Result<IndexDocument, StorageError> {
let group = self.content_group(content_key, now)?;
let projection_hash = projection_hash(group.as_ref());
let dynamic_projection_hash = dynamic_projection_hash(group.as_ref());
let refresh_bucket = group.as_ref().map_or_else(
|| refresh_bucket(now),
|group| refresh_bucket(group.last_seen),
);
Ok(IndexDocument {
content_key: *content_key,
group,
projection_hash,
dynamic_projection_hash,
refresh_bucket,
requires_write: true,
})
}
@@ -65,6 +73,8 @@ pub trait TorrentRepository: Send + Sync {
content_key: &[u8; 32],
revision: u64,
projection_hash: [u8; 32],
dynamic_projection_hash: [u8; 32],
refresh_bucket: u64,
visible: bool,
) -> Result<bool, StorageError>;
@@ -102,6 +112,10 @@ pub trait TorrentRepository: Send + Sync {
fn index_inventory(&self) -> IndexInventory {
IndexInventory::default()
}
fn index_refresh_diagnostics(&self) -> IndexRefreshDiagnostics {
IndexRefreshDiagnostics::default()
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, serde::Serialize)]
@@ -111,10 +125,17 @@ pub struct IndexInventory {
pub pending_documents: u64,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, serde::Serialize)]
pub struct IndexRefreshDiagnostics {
pub scheduled: u64,
pub suppressed: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ContentGroupTask {
pub content_key: [u8; 32],
pub revision: u64,
pub force_write: bool,
}
#[derive(Debug, Clone, PartialEq)]
@@ -122,6 +143,9 @@ pub struct IndexDocument {
pub content_key: [u8; 32],
pub group: Option<ContentGroup>,
pub projection_hash: [u8; 32],
pub dynamic_projection_hash: [u8; 32],
pub refresh_bucket: u64,
pub requires_write: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
@@ -143,6 +167,8 @@ pub struct StorageDiagnostics {
pub live_sst_bytes: Option<u64>,
pub running_compactions: Option<u64>,
pub estimated_keys: Option<u64>,
pub index_refresh_scheduled: u64,
pub index_refresh_suppressed: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
@@ -196,6 +222,53 @@ pub(crate) fn projection_hash(group: Option<&ContentGroup>) -> [u8; 32] {
*hasher.finalize().as_bytes()
}
pub(crate) const INDEX_REFRESH_INTERVAL_SECS: u64 = 6 * 3_600;
pub(crate) fn refresh_bucket(timestamp: u64) -> u64 {
timestamp / INDEX_REFRESH_INTERVAL_SECS
}
pub(crate) fn dynamic_projection_hash(group: Option<&ContentGroup>) -> [u8; 32] {
let mut hasher = blake3::Hasher::new();
match group {
None => {
hasher.update(b"hidden");
}
Some(group) => {
hasher.update(b"visible");
hasher.update(&refresh_bucket(group.last_seen).to_be_bytes());
hasher.update(&discovery_bucket(group.seen_count).to_be_bytes());
hasher.update(&[group.heat.score / 5]);
hasher.update(&[availability_number(group.availability.status)]);
hasher.update(&[peer_bucket(group.availability.reachable_peers)]);
hasher.update(
&refresh_bucket(group.availability.last_verified_at.unwrap_or(0)).to_be_bytes(),
);
}
}
*hasher.finalize().as_bytes()
}
fn discovery_bucket(value: u64) -> u32 {
if value == 0 { 0 } else { value.ilog2() }
}
fn peer_bucket(value: u32) -> u8 {
match value {
0 => 0,
1 => 1,
value => value.ilog2().saturating_add(1).min(u32::from(u8::MAX)) as u8,
}
}
fn availability_number(status: crate::domain::AvailabilityStatus) -> u8 {
match status {
crate::domain::AvailabilityStatus::Unknown => 0,
crate::domain::AvailabilityStatus::Active => 1,
crate::domain::AvailabilityStatus::PossiblyStale => 2,
}
}
#[derive(Debug, thiserror::Error)]
pub enum StorageError {
#[cfg(feature = "rocksdb-storage")]
+205 -33
View File
@@ -20,9 +20,10 @@ use super::{
verification_locator_key, verification_task_key, verification_task_prefix,
},
repository::{
ContentGroupTask, ContentVariants, IndexDocument, IndexInventory, StorageError,
TorrentRepository, UpsertOutcome, VerificationEnqueueOutcome, VerificationPriority,
VerificationRequest, projection_hash,
ContentGroupTask, ContentVariants, IndexDocument, IndexInventory, IndexRefreshDiagnostics,
StorageError, TorrentRepository, UpsertOutcome, VerificationEnqueueOutcome,
VerificationPriority, VerificationRequest, dynamic_projection_hash, projection_hash,
refresh_bucket,
},
};
@@ -30,6 +31,8 @@ mod filter_migration;
mod lifecycle;
const DEFAULT_BLOCK_CACHE_BYTES: usize = 64 * 1024 * 1024;
const UNKNOWN_DYNAMIC_PROJECTION: [u8; 32] = [0; 32];
const BASELINED_DYNAMIC_PROJECTION: [u8; 32] = [u8::MAX; 32];
type VerificationTaskEntry = (Box<[u8]>, InfoHash);
#[derive(Debug, Clone, Copy, serde::Serialize, serde::Deserialize)]
@@ -41,6 +44,8 @@ struct ContentGroupState {
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct ProjectionState {
hash: [u8; 32],
dynamic_hash: [u8; 32],
refresh_bucket: u64,
visible: bool,
}
@@ -51,6 +56,8 @@ pub struct RocksTorrentRepository {
canonical_filter: Arc<ContentFilter>,
content_filter: RwLock<Arc<ContentFilter>>,
inventory: InventoryState,
index_refresh_scheduled: AtomicU64,
index_refresh_suppressed: AtomicU64,
}
#[derive(Default)]
@@ -67,6 +74,12 @@ struct InventoryDelta {
pending_documents: i64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum RefreshDecision {
Scheduled,
Suppressed,
}
impl RocksTorrentRepository {
fn current_filter(&self) -> Arc<ContentFilter> {
self.content_filter
@@ -129,7 +142,13 @@ impl RocksTorrentRepository {
inserted: bool,
) -> Result<InventoryDelta, StorageError> {
let current = self.group_state(content_key)?;
let pending = self.db.get(pending_index_key(content_key))?.is_some();
let pending_value = self.db.get(pending_index_key(content_key))?;
let pending = pending_value.is_some();
let force_write = pending_value
.as_deref()
.map(decode_pending_revision)
.transpose()?
.is_some_and(|(_, force_write)| force_write);
let mut state = current.unwrap_or(ContentGroupState {
revision: 0,
member_count: 0,
@@ -139,7 +158,14 @@ impl RocksTorrentRepository {
state.member_count = state.member_count.saturating_add(1);
}
batch.put(content_group_key(content_key), Self::encode_group(state)?);
batch.put(pending_index_key(content_key), state.revision.to_be_bytes());
if force_write {
batch.put(
pending_index_key(content_key),
encode_forced_pending_revision(state.revision),
);
} else {
batch.put(pending_index_key(content_key), state.revision.to_be_bytes());
}
Ok(InventoryDelta {
searchable_groups: u64::from(current.is_none()),
pending_documents: i64::from(!pending),
@@ -147,6 +173,42 @@ impl RocksTorrentRepository {
})
}
fn schedule_observation_refresh(
&self,
batch: &mut WriteBatch,
content_key: &[u8; 32],
observed_at: u64,
) -> Result<RefreshDecision, StorageError> {
let Some(mut state) = self.projection_state(content_key)? else {
return Ok(RefreshDecision::Scheduled);
};
let bucket = refresh_bucket(observed_at);
if state.dynamic_hash == UNKNOWN_DYNAMIC_PROJECTION {
state.dynamic_hash = BASELINED_DYNAMIC_PROJECTION;
state.refresh_bucket = bucket;
batch.put(filter_projection_key(content_key), encode_projection(state));
return Ok(RefreshDecision::Suppressed);
}
if bucket <= state.refresh_bucket {
return Ok(RefreshDecision::Suppressed);
}
state.refresh_bucket = bucket;
batch.put(filter_projection_key(content_key), encode_projection(state));
Ok(RefreshDecision::Scheduled)
}
fn record_refresh_decision(&self, decision: RefreshDecision) {
match decision {
RefreshDecision::Scheduled => {
self.index_refresh_scheduled.fetch_add(1, Ordering::Relaxed);
}
RefreshDecision::Suppressed => {
self.index_refresh_suppressed
.fetch_add(1, Ordering::Relaxed);
}
}
}
fn inventory_snapshot(&self) -> IndexInventory {
IndexInventory {
stored_torrents: self.inventory.stored_torrents.load(Ordering::Relaxed),
@@ -305,12 +367,25 @@ impl TorrentRepository for RocksTorrentRepository {
current.observe_again(observation.last_seen, &observation.source_peers);
let mut batch = WriteBatch::default();
batch.put(torrent_key(current.info_hash), Self::encode(&current)?);
let delta = if current.searchable {
self.dirty_group(&mut batch, &current.content_key, false)?
let (delta, refresh) = if current.searchable {
let refresh = self.schedule_observation_refresh(
&mut batch,
&current.content_key,
current.last_seen,
)?;
let delta = if refresh == RefreshDecision::Scheduled {
self.dirty_group(&mut batch, &current.content_key, false)?
} else {
InventoryDelta::default()
};
(delta, Some(refresh))
} else {
InventoryDelta::default()
(InventoryDelta::default(), None)
};
self.commit_inventory_batch(batch, delta)?;
if let Some(refresh) = refresh {
self.record_refresh_decision(refresh);
}
return Ok(UpsertOutcome::Updated {
seen_count: current.seen_count,
});
@@ -393,12 +468,25 @@ impl TorrentRepository for RocksTorrentRepository {
record.observe_again(observed_at, &[]);
let mut batch = WriteBatch::default();
batch.put(torrent_key(info_hash), Self::encode(&record)?);
let delta = if record.searchable {
self.dirty_group(&mut batch, &record.content_key, false)?
let (delta, refresh) = if record.searchable {
let refresh = self.schedule_observation_refresh(
&mut batch,
&record.content_key,
record.last_seen,
)?;
let delta = if refresh == RefreshDecision::Scheduled {
self.dirty_group(&mut batch, &record.content_key, false)?
} else {
InventoryDelta::default()
};
(delta, Some(refresh))
} else {
InventoryDelta::default()
(InventoryDelta::default(), None)
};
self.commit_inventory_batch(batch, delta)?;
if let Some(refresh) = refresh {
self.record_refresh_decision(refresh);
}
Ok(true)
}
@@ -425,7 +513,7 @@ impl TorrentRepository for RocksTorrentRepository {
let mut unknown = Vec::with_capacity(info_hashes.len());
let mut batch = WriteBatch::default();
let mut updated = 0_usize;
let mut dirty_groups = std::collections::BTreeSet::new();
let mut observed_groups = std::collections::BTreeMap::new();
for ((((info_hash, key), rejected_key), record), rejection) in info_hashes
.iter()
.zip(&keys)
@@ -438,7 +526,12 @@ impl TorrentRepository for RocksTorrentRepository {
record.observe_again(observed_at, &[]);
batch.put(key, Self::encode(&record)?);
if record.searchable {
dirty_groups.insert(record.content_key);
observed_groups
.entry(record.content_key)
.and_modify(|last_seen: &mut u64| {
*last_seen = (*last_seen).max(record.last_seen);
})
.or_insert(record.last_seen);
}
updated += 1;
continue;
@@ -456,16 +549,25 @@ impl TorrentRepository for RocksTorrentRepository {
}
if updated > 0 {
let mut delta = InventoryDelta::default();
for content_key in dirty_groups {
let group_delta = self.dirty_group(&mut batch, &content_key, false)?;
delta.searchable_groups = delta
.searchable_groups
.saturating_add(group_delta.searchable_groups);
delta.pending_documents = delta
.pending_documents
.saturating_add(group_delta.pending_documents);
let mut refreshes = Vec::with_capacity(observed_groups.len());
for (content_key, last_seen) in observed_groups {
let refresh =
self.schedule_observation_refresh(&mut batch, &content_key, last_seen)?;
if refresh == RefreshDecision::Scheduled {
let group_delta = self.dirty_group(&mut batch, &content_key, false)?;
delta.searchable_groups = delta
.searchable_groups
.saturating_add(group_delta.searchable_groups);
delta.pending_documents = delta
.pending_documents
.saturating_add(group_delta.pending_documents);
}
refreshes.push(refresh);
}
self.commit_inventory_batch(batch, delta)?;
for refresh in refreshes {
self.record_refresh_decision(refresh);
}
}
Ok(unknown)
}
@@ -484,14 +586,11 @@ impl TorrentRepository for RocksTorrentRepository {
let Some(content_key) = decode_pending_content_key(&key) else {
break;
};
let revision = value
.as_ref()
.try_into()
.map(u64::from_be_bytes)
.map_err(|_| StorageError::CorruptContentGroup)?;
let (revision, force_write) = decode_pending_revision(&value)?;
tasks.push(ContentGroupTask {
content_key,
revision,
force_write,
});
if tasks.len() == limit {
break;
@@ -515,9 +614,24 @@ impl TorrentRepository for RocksTorrentRepository {
now: u64,
) -> Result<IndexDocument, StorageError> {
let group = self.content_group(content_key, now)?;
let projection_hash = projection_hash(group.as_ref());
let dynamic_projection_hash = dynamic_projection_hash(group.as_ref());
let bucket = group.as_ref().map_or_else(
|| refresh_bucket(now),
|group| refresh_bucket(group.last_seen),
);
let previous = self.projection_state(content_key)?;
let requires_write = previous.is_none_or(|state| {
state.hash != projection_hash
|| state.dynamic_hash != dynamic_projection_hash
|| state.visible != group.is_some()
});
Ok(IndexDocument {
content_key: *content_key,
projection_hash: projection_hash(group.as_ref()),
projection_hash,
dynamic_projection_hash,
refresh_bucket: bucket,
requires_write,
group,
})
}
@@ -527,6 +641,8 @@ impl TorrentRepository for RocksTorrentRepository {
content_key: &[u8; 32],
revision: u64,
projection_hash: [u8; 32],
dynamic_projection_hash: [u8; 32],
refresh_bucket: u64,
visible: bool,
) -> Result<bool, StorageError> {
let _guard = self
@@ -547,6 +663,8 @@ impl TorrentRepository for RocksTorrentRepository {
filter_projection_key(content_key),
encode_projection(ProjectionState {
hash: projection_hash,
dynamic_hash: dynamic_projection_hash,
refresh_bucket,
visible,
}),
);
@@ -585,7 +703,7 @@ impl TorrentRepository for RocksTorrentRepository {
let state: ContentGroupState = rmp_serde::from_slice(&value)?;
batch.put(
pending_index_key(&content_key),
state.revision.to_be_bytes(),
encode_forced_pending_revision(state.revision),
);
batch_len += 1;
total += 1;
@@ -814,24 +932,49 @@ impl TorrentRepository for RocksTorrentRepository {
fn index_inventory(&self) -> IndexInventory {
self.inventory_snapshot()
}
fn index_refresh_diagnostics(&self) -> IndexRefreshDiagnostics {
IndexRefreshDiagnostics {
scheduled: self.index_refresh_scheduled.load(Ordering::Relaxed),
suppressed: self.index_refresh_suppressed.load(Ordering::Relaxed),
}
}
}
fn encode_projection(state: ProjectionState) -> [u8; 33] {
let mut bytes = [0_u8; 33];
fn encode_projection(state: ProjectionState) -> [u8; 73] {
let mut bytes = [0_u8; 73];
bytes[0] = u8::from(state.visible);
bytes[1..].copy_from_slice(&state.hash);
bytes[1..33].copy_from_slice(&state.hash);
bytes[33..65].copy_from_slice(&state.dynamic_hash);
bytes[65..].copy_from_slice(&state.refresh_bucket.to_be_bytes());
bytes
}
fn decode_projection(bytes: &[u8]) -> Result<ProjectionState, StorageError> {
if bytes.len() != 33 || bytes[0] > 1 {
if !matches!(bytes.len(), 33 | 73) || bytes[0] > 1 {
return Err(StorageError::CorruptContentGroup);
}
Ok(ProjectionState {
visible: bytes[0] == 1,
hash: bytes[1..]
hash: bytes[1..33]
.try_into()
.map_err(|_| StorageError::CorruptContentGroup)?,
dynamic_hash: if bytes.len() == 73 {
bytes[33..65]
.try_into()
.map_err(|_| StorageError::CorruptContentGroup)?
} else {
[0; 32]
},
refresh_bucket: if bytes.len() == 73 {
u64::from_be_bytes(
bytes[65..]
.try_into()
.map_err(|_| StorageError::CorruptContentGroup)?,
)
} else {
0
},
})
}
@@ -847,6 +990,35 @@ fn read_counter(db: &DB, key: &[u8]) -> Result<Option<u64>, StorageError> {
.transpose()
}
fn encode_forced_pending_revision(revision: u64) -> [u8; 9] {
let mut bytes = [0_u8; 9];
bytes[0] = 1;
bytes[1..].copy_from_slice(&revision.to_be_bytes());
bytes
}
fn decode_pending_revision(bytes: &[u8]) -> Result<(u64, bool), StorageError> {
match bytes {
bytes if bytes.len() == 8 => Ok((
u64::from_be_bytes(
bytes
.try_into()
.map_err(|_| StorageError::CorruptContentGroup)?,
),
false,
)),
[1, revision @ ..] if revision.len() == 8 => Ok((
u64::from_be_bytes(
revision
.try_into()
.map_err(|_| StorageError::CorruptContentGroup)?,
),
true,
)),
_ => Err(StorageError::CorruptContentGroup),
}
}
fn count_prefix<const N: usize>(db: &DB, prefix: [u8; N]) -> Result<u64, StorageError> {
let iterator = db.iterator(IteratorMode::From(&prefix, Direction::Forward));
let mut count = 0_u64;
@@ -7,7 +7,11 @@ use rocksdb::{Direction, IteratorMode, WriteBatch};
use super::{ProjectionState, RocksTorrentRepository, encode_projection};
use crate::{
domain::{ContentFilter, ContentGroup, ContentGroupBuilder},
storage::{StorageError, TorrentRepository, keys::*, repository::projection_hash},
storage::{
StorageError, TorrentRepository,
keys::*,
repository::{dynamic_projection_hash, projection_hash, refresh_bucket},
},
};
const COUNTER_BYTES: usize = std::mem::size_of::<u64>();
@@ -92,6 +96,11 @@ impl RocksTorrentRepository {
filter_projection_key(&content_key),
encode_projection(ProjectionState {
hash,
dynamic_hash: dynamic_projection_hash(group.as_ref()),
refresh_bucket: group.as_ref().map_or_else(
|| refresh_bucket(now),
|group| refresh_bucket(group.last_seen),
),
visible: group.is_some(),
}),
)?;
+5 -1
View File
@@ -12,7 +12,7 @@ use rocksdb::{
};
use std::{
path::Path,
sync::{Arc, Mutex},
sync::{Arc, Mutex, atomic::Ordering},
};
impl RocksTorrentRepository {
@@ -55,6 +55,8 @@ impl RocksTorrentRepository {
canonical_filter: Arc::new(ContentFilter::legacy_default()),
content_filter: std::sync::RwLock::new(content_filter),
inventory: Default::default(),
index_refresh_scheduled: Default::default(),
index_refresh_suppressed: Default::default(),
};
repository.initialize_format()?;
repository.initialize_inventory(false)?;
@@ -75,6 +77,8 @@ impl RocksTorrentRepository {
.db
.property_int_value("rocksdb.num-running-compactions")?,
estimated_keys: self.db.property_int_value("rocksdb.estimate-num-keys")?,
index_refresh_scheduled: self.index_refresh_scheduled.load(Ordering::Relaxed),
index_refresh_suppressed: self.index_refresh_suppressed.load(Ordering::Relaxed),
})
}
+145 -4
View File
@@ -56,7 +56,7 @@ fn index_inventory_is_exact_and_survives_reopen() {
let task = repository.pending_index(1).unwrap()[0];
assert!(
repository
.mark_indexed(&task.content_key, task.revision, [0; 32], true)
.mark_indexed(&task.content_key, task.revision, [0; 32], [0; 32], 0, true)
.unwrap()
);
assert_eq!(repository.index_inventory().pending_documents, 0);
@@ -132,6 +132,8 @@ fn changed_filter_enqueues_only_a_changed_group_without_rekeying_metadata() {
&task.content_key,
task.revision,
document.projection_hash,
document.dynamic_projection_hash,
document.refresh_bucket,
document.group.is_some(),
)
.unwrap();
@@ -177,6 +179,8 @@ fn filter_scan_with_unchanged_projection_does_not_reindex() {
&task.content_key,
task.revision,
document.projection_hash,
document.dynamic_projection_hash,
document.refresh_bucket,
true,
)
.unwrap();
@@ -321,6 +325,98 @@ fn existing_hash_can_be_observed_without_downloading_metadata_again() {
assert_eq!(stored.seen_count, 2);
}
#[test]
fn repeated_observations_schedule_at_most_one_refresh_per_bucket() {
let directory = TempDir::new().unwrap();
let repository = RocksTorrentRepository::open(directory.path()).unwrap();
let started_at = crate::storage::repository::INDEX_REFRESH_INTERVAL_SECS * 10 + 10;
let record = test_record(16, started_at);
repository.upsert(record.clone()).unwrap();
let task = repository.pending_index(1).unwrap()[0];
let document = repository
.index_document(&record.content_key, started_at)
.unwrap();
assert!(
repository
.mark_indexed(
&record.content_key,
task.revision,
document.projection_hash,
document.dynamic_projection_hash,
document.refresh_bucket,
true,
)
.unwrap()
);
repository
.observe_existing(record.info_hash, started_at + 60)
.unwrap();
assert!(repository.pending_index(10).unwrap().is_empty());
let next_bucket = started_at + crate::storage::repository::INDEX_REFRESH_INTERVAL_SECS;
repository
.observe_existing(record.info_hash, next_bucket)
.unwrap();
let scheduled = repository.pending_index(10).unwrap();
assert_eq!(scheduled.len(), 1);
let revision = scheduled[0].revision;
repository
.observe_existing(record.info_hash, next_bucket + 60)
.unwrap();
assert_eq!(repository.pending_index(10).unwrap()[0].revision, revision);
let diagnostics = repository.diagnostics().unwrap();
assert_eq!(diagnostics.index_refresh_scheduled, 1);
assert_eq!(diagnostics.index_refresh_suppressed, 2);
}
#[test]
fn legacy_projection_is_baselined_before_scheduling_a_later_bucket() {
let directory = TempDir::new().unwrap();
let repository = RocksTorrentRepository::open(directory.path()).unwrap();
let started_at = crate::storage::repository::INDEX_REFRESH_INTERVAL_SECS * 20 + 10;
let record = test_record(17, started_at);
repository.upsert(record.clone()).unwrap();
let task = repository.pending_index(1).unwrap()[0];
let document = repository
.index_document(&record.content_key, started_at)
.unwrap();
repository
.mark_indexed(
&record.content_key,
task.revision,
document.projection_hash,
document.dynamic_projection_hash,
document.refresh_bucket,
true,
)
.unwrap();
let mut legacy = [0_u8; 33];
legacy[0] = 1;
legacy[1..].copy_from_slice(&document.projection_hash);
repository
.db
.put(filter_projection_key(&record.content_key), legacy)
.unwrap();
repository
.observe_existing(record.info_hash, started_at + 60)
.unwrap();
repository
.observe_existing(record.info_hash, started_at + 120)
.unwrap();
assert!(repository.pending_index(10).unwrap().is_empty());
repository
.observe_existing(
record.info_hash,
started_at + crate::storage::repository::INDEX_REFRESH_INTERVAL_SECS,
)
.unwrap();
assert_eq!(repository.pending_index(10).unwrap().len(), 1);
}
#[test]
fn rejection_survives_restart_and_blocks_the_same_rule() {
let directory = TempDir::new().unwrap();
@@ -422,7 +518,14 @@ fn marking_indexed_is_atomic_with_removing_pending_marker() {
let task = repository.pending_index(10).unwrap()[0];
assert!(
repository
.mark_indexed(&record.content_key, task.revision, [0; 32], true)
.mark_indexed(
&record.content_key,
task.revision,
[0; 32],
[0; 32],
0,
true,
)
.unwrap()
);
@@ -441,7 +544,14 @@ fn stale_index_revision_cannot_clear_a_newer_update() {
assert!(
!repository
.mark_indexed(&record.content_key, stale.revision, [0; 32], true)
.mark_indexed(
&record.content_key,
stale.revision,
[0; 32],
[0; 32],
0,
true,
)
.unwrap()
);
let current = repository.pending_index(10).unwrap()[0];
@@ -477,7 +587,7 @@ fn full_reindex_restores_pending_markers_for_every_record() {
repository.upsert(second.clone()).unwrap();
for task in repository.pending_index(10).unwrap() {
repository
.mark_indexed(&task.content_key, task.revision, [0; 32], true)
.mark_indexed(&task.content_key, task.revision, [0; 32], [0; 32], 0, true)
.unwrap();
}
assert!(repository.pending_index(10).unwrap().is_empty());
@@ -486,6 +596,37 @@ fn full_reindex_restores_pending_markers_for_every_record() {
assert_eq!(repository.pending_index(10).unwrap().len(), 2);
}
#[test]
fn observation_preserves_forced_write_during_full_reindex() {
let directory = TempDir::new().unwrap();
let repository = RocksTorrentRepository::open(directory.path()).unwrap();
let record = test_record(18, 10);
repository.upsert(record.clone()).unwrap();
let task = repository.pending_index(1).unwrap()[0];
let document = repository.index_document(&record.content_key, 20).unwrap();
repository
.mark_indexed(
&record.content_key,
task.revision,
document.projection_hash,
document.dynamic_projection_hash,
document.refresh_bucket,
true,
)
.unwrap();
repository.prepare_full_reindex().unwrap();
repository
.observe_existing(
record.info_hash,
crate::storage::repository::INDEX_REFRESH_INTERVAL_SECS * 2,
)
.unwrap();
let pending = repository.pending_index(1).unwrap()[0];
assert!(pending.force_write);
}
#[test]
fn verification_queue_is_persistent_prioritized_and_updates_record_atomically() {
let directory = TempDir::new().unwrap();
+8
View File
@@ -133,6 +133,10 @@ export interface ServiceStats {
persistence_updated: number
persistence_queue: number
indexed_documents: number
index_refresh_scheduled: number
index_refresh_suppressed: number
index_documents_written: number
index_documents_skipped: number
index: IndexStatus
filter: FilterStatus
verification_queue: number
@@ -185,6 +189,8 @@ export interface StorageDiagnostics {
live_sst_bytes: number | null
running_compactions: number | null
estimated_keys: number | null
index_refresh_scheduled: number
index_refresh_suppressed: number
}
export interface SearchDiagnostics {
@@ -195,6 +201,8 @@ export interface SearchDiagnostics {
last_commit_at: number | null
last_commit_duration_millis: number
last_commit_documents: number
documents_written: number
documents_skipped: number
}
export interface RuntimeDiagnostics {