Compare commits

..
4 Commits
10 changed files with 2599 additions and 307 deletions
+192 -58
View File
@@ -1,83 +1,217 @@
# TODO # TODO - Reader MCP 正式化
## 当前状态 > 本文件用于架构与 Codex 协作同步。
>
> 规则:
> - `TODO` = 未开始
> - `DOING` = 正在进行
> - `DONE` = 已完成
> - 每次只允许一个最高优先级主任务处于 `DOING`
项目当前已经进入“可交付给 OpenClaw 调用”的阶段。 ## 0. 协作约束
当前主链路: 开始编码前必须阅读:
`FreshRSS 未读 -> RSS 内容提取 -> LLM 总结 -> 规则过滤 -> OpenClaw delivery payload` 1. `plans/reader-mcp-architecture-design.md`
2. `plans/reader-mcp-implementation-plan.md`
当前已经完成: 3. 本文件
4. `plans/issues/2026-04-06-reader-digest-sigterm.md`
- [x] FreshRSS `greader` API 接入 5. `docs/openclaw/openclaw-handoff.md`
- [x] `entry -> item` 标准化映射
- [x] RSS-first 内容提取策略
- [x] LLM 摘要与校验闭环
- [x] 第一版规则引擎
- [x] `ArticleCandidateRecord` / `OpenClawCandidateInput` 分层
- [x] `OpenClawDeliveryPayload` 批量投递结构
- [x] 最终 payload 成功后才标记 FreshRSS 已读
- [x] MCP 工具 `run_freshrss_openclaw_pipeline`
- [x] 默认精简输出模式
- [x] 日报级 `keywords` 词元库与周期性词元清洗 skill 设计完成
- [x] 日报级 `keywords` 词元库与全局词频统计实现完成
- [x] `keyword-cleanup-review` skill 骨架与 review bundle 脚本实现完成
- [x] 词元清洗低复杂治理层落地:`term_cleanup_policy` / `term_watchlist` / `term_change_log`
- [x] 已支持人工确认采纳建议并写入 `term_watchlist` / `term_change_log`
--- ---
## P0 - 交接前后最优先 ## 1. 当前主任务
- [x] 为 OpenClaw 补齐交接文档 ### [DONE][P0] 建立 run-state 运行态基础设施
- [x] 将 MCP 工具作为统一生产入口
- [x] 将默认输出收敛为最小必要文件 目标:
- [ ] 设计 OpenClaw webhook / delivery payload 的主动推送方式 - 给 freshrss pipeline 引入正式 run state
- [ ] 明确 OpenClaw 侧如何注册和启动本 MCP 服务 - 即使失败或中断,也能留下明确运行真相
要求:
- 新增 `RunState / StageState / ArtifactRecord` 模型
- 在 `outputs/freshrss/rerun/<run_id>/run-state.json` 持久化
- 至少覆盖以下 stages:
- `fetch_feed`
- `extract_articles`
- `generate_summaries`
- `apply_filters`
- `build_delivery_payload`
- `write_run_report`
- 失败时写入失败阶段与错误摘要
- 不破坏现有输出目录兼容性
建议文件:
- `src/summary_mcp/runtime/state_models.py`
- `src/summary_mcp/runtime/run_store.py`
- `src/summary_mcp/workflows/...`
完成标准:
- 跑一次 pipeline 后,无论成功失败,都存在 `run-state.json`
- 文件中可看出当前/最后阶段、整体状态、关键 artifacts
进展备注:
- 2026-04-07:架构设计文档已建立;开始进入实现阶段。
- 2026-04-07:已新增 `runtime` 包骨架,落地 `RunState / StageState / ArtifactRecord` 与文件存储接口。
- 2026-04-07:已将 `run-state.json` 接入 `freshrss` 主流程,按阶段持续写入状态与关键 artifacts。
- 2026-04-07:已完成成功/失败路径自检,确认 `run-state.json` 在两类路径下都保留且不改变既有对外返回字段。
--- ---
## P1 - 下一阶段推进 ## 2. 后续任务队列
- [ ] 设计“人工确认后再沉淀知识库”的状态流转 ### [DONE][P1] 增加 MCP 状态查询接口 `get_run_status`
- [ ] 收敛 `paywall` 误判规则,降低中文文本误报
- [ ] 细化过滤规则并引入更多个性化上下文 目标:
- [ ] 将 `keyword-cleanup-review` skill 接入周期性执行流程,产出别名/停用词/兴趣词建议 - 可通过 MCP 查询 run 状态
- [ ] 增加清洗前后效果对比报告,验证配置调整是否真的改善过滤质量
- [ ] 增加批量 run 的保留策略与历史清理策略 要求:
- [ ] 为 OpenClaw 补一份更正式的 MCP 调用示例和接线说明 - 输入 `run_id`
- 返回 status / current_stage / completed_stages / failed_stage / artifacts / recovery
完成情况:
- 已通过 MCP 暴露 `get_run_status`
- 优先读取 `run-state.json`;对无 `run-state.json` 的历史 run 兼容基于现有 run 目录与 `run-report.json` 推断状态
- 返回补充了 `progress` / `output_dir` / `state_source`,便于 OpenClaw 稳定消费且不必手拼路径
改动文件:
- `src/summary_mcp/runtime/query_service.py`
- `src/summary_mcp/runtime/run_store.py`
- `src/summary_mcp/runtime/__init__.py`
- `src/summary_mcp/server.py`
遗留风险:
- 历史 run 若缺少 `run-state.json`,其阶段状态只能基于现有目录与 `run-report.json` 做保守推断
--- ---
## P2 - 后续增强 ### [DONE][P1] 增加 MCP 查询接口 `list_runs`
- [ ] 将 `Markdown sink` 进一步降级为 debug / fallback 能力 目标:
- [ ] 增加按天聚合 `ArticleCandidateRecord` 的批处理能力 - 查看近期 runs
- [ ] 让 OpenClaw 聚合候选内容并生成日级摘要
- [ ] 将日级摘要写入知识库,并同步生成面向用户的日报消息 要求:
- [ ] 支持更多 `content_kind` - 支持按 workflow / status / latest_n 过滤
- [ ] 增加提取缓存、重试和更细粒度日志
- [ ] 整理历史 rerun 目录与调试产物保留策略 完成情况:
- 已通过 MCP 暴露 `list_runs`
- 支持按 `workflow` / `status` / `latest_n` 过滤近期 runs
- 返回 `run_id`、`status`、`progress`、`recovery`、`output_dir` 与 `state_source`
改动文件:
- `src/summary_mcp/runtime/query_service.py`
- `src/summary_mcp/runtime/__init__.py`
- `src/summary_mcp/server.py`
遗留风险:
- 当前按文件系统扫描 `outputs/freshrss/rerun/` 聚合,规模继续增大时可能需要再评估缓存或索引,但本阶段先保持文件系统真相
--- ---
## 当前建议的下一步 ### [DONE][P1] 增加 MCP 查询接口 `list_run_artifacts`
优先做这三件事: 目标:
- 统一列出 run 下 artifact
1. 将 `keyword-cleanup-review` skill 接入周期性执行流程 完成情况:
2. 设计“人工确认后再沉淀知识库”的状态流转 - 已通过 MCP 暴露 `list_run_artifacts`
3. 收敛规则误判,尤其是 `paywall` 相关启发式 - 对新 run 优先返回 `run-state.json` 已注册 artifacts,并补充 run 目录扫描发现的标准产物
- 对历史 run 直接基于现有 run 目录发现标准产物,保持兼容
改动文件:
- `src/summary_mcp/runtime/query_service.py`
- `src/summary_mcp/runtime/__init__.py`
- `src/summary_mcp/server.py`
遗留风险:
- 当前仅补充扫描固定的一组标准产物;未注册且不在标准集合内的调试文件不会进入稳定 artifact 列表
--- ---
## 交接时优先阅读 ### [DONE][P1] 增加 MCP 结果读取接口 `get_delivery_payload`
- `README.md` 目标:
- `docs/openclaw/openclaw-handoff.md` - 按 run_id 读取 delivery payload
- `docs/openclaw/openclaw-candidate-input-field-spec.md`
- `docs/openclaw/openclaw-delivery-payload-spec.md` 完成情况:
- `docs/design/daily-keyword-index-design.md` - 已通过 MCP 暴露 `get_delivery_payload`
- `skills/keyword-cleanup-review/SKILL.md` - 查询优先复用 `run-state.json` 已注册 artifacts,其次回退标准产物路径、`run-report.json` 引用和 run 目录扫描
- `docs/current/context-reset-brief.md` - 返回补充了 `artifact`、`payload_schema_version`、`generated_at`、`delivery_date`、`candidate_count`、`stats`,并保留完整 `payload`
改动文件:
- `src/summary_mcp/runtime/query_service.py`
- `src/summary_mcp/server.py`
- `src/summary_mcp/runtime/__init__.py`
遗留风险:
- 历史 run 的 payload 若既未注册也不在标准路径下,只能依赖 `run-report.json` 引用或目录扫描做兼容发现
---
### [DONE][P1] 增加 MCP 结果读取接口 `get_run_report`
目标:
- 按 run_id 读取 run report
完成情况:
- 已通过 MCP 暴露 `get_run_report`
- 查询优先复用 `run-state.json` 已注册 artifacts,其次回退标准产物路径、历史 `run-report.json` 固定位置和 run 目录扫描
- 返回补充了 `artifact`、核心计数摘要、规范化后的关键产物路径与 `keyword_index`,并保留完整 `report`
改动文件:
- `src/summary_mcp/runtime/query_service.py`
- `src/summary_mcp/runtime/__init__.py`
- `src/summary_mcp/server.py`
遗留风险:
- 历史 run 的 `run-report.json` 内嵌路径仍保留原始绝对路径于 `report` 字段,当前仅在顶层摘要字段做规范化,避免改变历史文件真相
---
### [TODO][P2] 设计并实现 `resume_run`
目标:
- 基于 `run-state.json` 和现有中间产物继续执行
说明:
- 先做最小可用恢复
- 暂不追求任意 stage 任意重入
---
### [TODO][P2] 评估 `rerun_stage` 是否值得进入第一阶段
目标:
- 在 `resume_run` 之后评估是否继续增加更细粒度补跑能力
---
### [TODO][P2] 调整 `run_freshrss_openclaw_pipeline` 内部实现以复用 runtime
目标:
- 保持外部兼容
- 内部不再是黑箱长函数
---
### [TODO][P3] 更新 README / handoff / docs,明确 MCP 为正式入口
目标:
- 把生产建议从 CLI 迁移到 MCP
- CLI 明确降级为 debug / fallback
---
## 3. 记录区
### 已完成记录
- 2026-04-07:新增架构设计文档 `plans/reader-mcp-architecture-design.md`
- 2026-04-07:新增实施计划文档 `plans/reader-mcp-implementation-plan.md`
- 2026-04-07:完成 `freshrss` pipeline 的 run-state 基础设施,新增 `runtime` 包并覆盖关键 stages 状态持久化。
### 风险提醒
- 不要在第一阶段引入复杂任务队列
- 不要让 CLI 和 MCP 背后变成两套独立逻辑
- 若实现偏离架构,先更新 `plans/` 再改代码
@@ -0,0 +1,112 @@
# reader 日报链路在 OpenClaw/Feishu 外层执行中被 SIGTERM 截断
## 背景
2026-04-06 在 OpenClaw 中执行 reader 日报流程时,出现多次“前半段有产物、最终产物缺失”的现象。
典型表现:
- 能生成 `raw/freshrss.raw.json`
- 能部分生成 `extracted/item-xx.extracted.json`
- 但经常拿不到:
- `run-report.json`
- `candidates/openclaw-delivery-payload.json`
- `candidates/digest-brief.json`
- 外层日志多次出现 `Exec failed (..., signal SIGTERM)`
## 已确认结论
### 1. RSS 抓取与排序正常
已确认:
- FreshRSS API 正常
- 原始 raw 数据能拉到
- 默认排序正常(默认相当于 `r=d`,新到旧)
因此问题不在:
- RSS 接口
- 鉴权
- 排序规则
### 2. reader 前半段模块正常
已确认:
- 单篇 summary 能成功
- 单篇 filter + candidate 构建能成功
因此问题不在:
- summary 模块整体损坏
- candidate 构建整体损坏
### 3. 真正问题在长链路执行方式
更合理的判断是:
> 整条 reader 日报 pipeline 是长串行任务,在 OpenClaw / Feishu 当前这条外层执行链路里,容易被外层执行环境提前 `SIGTERM`。
也就是说:
- 不是 reader 总是自己抛 Python 异常退出
- 更多是脚本尚未跑完,外层执行会话先被终止
## 关键认知
当天很多排障与补跑步骤,实际不是通过稳定常驻的 MCP 服务在跑,而是直接执行 reader 仓库里的 Python 脚本:
- `python scripts/run_freshrss_pipeline.py`
- `python scripts/run_article_summaries.py`
因此更准确地说:
> 当前 reader 正式日报运行入口偏 CLI/脚本模式,而不是稳定 MCP 服务调用模式。
这也是为什么长任务更容易受外层 exec 生命周期影响。
## 为什么前几次没问题
可能原因:
1. 之前任务更短、内容更轻,刚好能在外层执行环境截断前跑完
2. 之前不是链路天然稳,而是还没撞上边界条件
3. 当前正式链路缺少稳健的断点恢复能力,因此一旦遇到较重任务,就暴露出问题
## 当天动作记录
### 做过的排查
- 确认 FreshRSS 默认排序
- 确认 raw 文件能生成
- 确认 extracted 文件能部分生成
- 单独验证单篇 summary 成功
- 单独验证单篇 filter / candidate 成功
### 做过的临时修复
当天曾尝试加入:
- `--resume`
- 中间产物复用
- 断点恢复思路
目的是降低 SIGTERM 后的损失。
### 当前状态
这些临时代码修改已全部回滚,reader 工作区已恢复干净。
## 当天结果
虽然正式链路异常,但通过手工恢复推进,最终仍补齐了:
- `candidates/openclaw-delivery-payload.json`
- `candidates/digest-brief.json`
- `run-report.json`
当天最终保留文章:
- 1
- 3
## 后续建议
### 短期
- CLI 仅保留为 debug / fallback
- 不再把正式生产日报流主要建立在长脚本入口上
### 中期
- 将 reader 作为正式 MCP 服务
- OpenClaw 正式通过 MCP tool 调用:
- `run_freshrss_openclaw_pipeline`
- `generate_article_summaries`
### 长期
让 reader 的正式能力具备:
- run_id
- status 查询
- resume
- 中间产物复用
- 分阶段补跑
+766
View File
@@ -0,0 +1,766 @@
# Reader 正式 MCP 服务架构设计
## 1. 背景
当前 reader 仓库已经具备 MCP 服务入口(`src/summary_mcp/server.py`),并暴露了:
- `extract_url_content`
- `extract_item_content`
- `filter_summary_result`
- `run_freshrss_openclaw_pipeline`
- `generate_article_summaries`
但从实际生产使用方式看,reader 的正式日报链路仍然**偏向 CLI/脚本模式**,尤其是:
- `python scripts/run_freshrss_pipeline.py`
- `python scripts/run_article_summaries.py`
这导致在 OpenClaw / Feishu 外层执行环境下,长链路任务容易因为 exec 生命周期、超时或外层中断而失败。
2026-04-06 的事故已经说明:
- `raw` 能产出
- `extracted` 能部分产出
- 但 `run-report.json` / `openclaw-delivery-payload.json` / `digest-brief.json` 经常缺失
- 外层日志出现 `signal SIGTERM`
因此,reader 需要从“可被 exec 调起的仓库”升级为“**正式的 MCP 工作流服务**”。
---
## 2. 设计目标
本次架构设计的目标不是重做 reader 的业务逻辑,而是把现有能力收口为稳定的服务接口。
目标如下:
1. **MCP 成为正式入口**
- OpenClaw 与上层编排默认通过 MCP tool 调用 reader
- CLI 降级为 debug / fallback 入口
2. **引入稳定的运行态抽象**
- 统一 `run_id`
- 统一 `status`
- 统一 `stage`
- 统一 `artifacts`
3. **支持长链路可观测与恢复**
- 可查询当前运行状态
- 可查询失败阶段
- 可基于已有中间产物 resume / rerun
4. **让 OpenClaw 消费结构化结果,而不是硬编码目录细节**
- OpenClaw 不再依赖 reader 的脚本 stdout 作为唯一信号
- OpenClaw 尽量不直接拼接 reader 的输出目录路径
5. **保持最小重构成本**
- 先用文件系统持久化 run state
- 先不引入复杂任务队列 / 数据库 / 多 worker 平台
- 先让单机、单实例、顺序执行场景稳定起来
---
## 3. 核心结论
一句话总结:
> Reader 应该被设计为“有状态的工作流 MCP 服务”,而不是“包了一层 MCP 壳的长 CLI 命令”。
换句话说:
- **错误方向**:`reader_run(command="python scripts/run_xxx.py ...")`
- **正确方向**:`start_run` + `get_run_status` + `get_artifact` + `resume_run`
协议层只是入口,真正的关键是把 reader 内部抽象成:
- run
- stage
- status
- artifact
- recovery
---
## 4. 服务定位
### 4.1 Reader 负责什么
reader 作为 MCP 服务,负责上游阅读工作流本身:
- FreshRSS 拉取
- 内容提取
- LLM 摘要
- 规则过滤
- candidate / delivery payload 生成
- 单篇总结生成
- 中间产物保存
- 运行状态记录
- 恢复与补跑
### 4.2 Reader 不负责什么
以下继续由 OpenClaw / skill 层负责:
- Hugo 发布
- 对话汇报
- 用户确认精选
- IMA 上传编排
- 最终对人类的日常交互
### 4.3 边界结论
- **reader 是上游引擎 / workflow service**
- **OpenClaw 是下游编排层 / orchestration layer**
这个边界与当前 `docs/openclaw/openclaw-handoff.md` 的原则保持一致,但会进一步强化“reader 通过正式 MCP 接口暴露工作流状态”,而不是只暴露“跑完后的结果”。
---
## 5. 当前问题分析
### 5.1 现状问题
当前 reader 的 MCP 工具里虽然已经有 `run_freshrss_openclaw_pipeline`,但其调用语义仍然更像:
- 一次性同步执行整个长链路
- 直接返回最终结果
- 对外隐藏中间状态
这会带来几个问题:
1. 长任务运行时,外层必须一直等
2. 一旦外层中断,状态观测困难
3. 难以精确判断失败发生在哪个阶段
4. resume / rerun 只能靠脚本层补丁式处理
5. OpenClaw 不容易构建“先触发,后查询,再消费”的稳定编排流
### 5.2 根本问题
根本问题不是“有没有 MCP server”,而是:
> **reader 还没有被真正建模成一个有状态的 workflow service。**
当前更接近:
- MCP 暴露了几个函数
- 但长链路执行模型仍是 CLI thinking
因此改造重点不应放在“多加几个 tool 名字”,而应该放在:
- 状态机
- 运行记录
- 恢复机制
- 标准产物注册
---
## 6. 目标架构
建议的正式架构如下:
```text
OpenClaw / Skill Layer
-> MCP Client Calls
-> Reader MCP Server
-> Workflow Service Layer
-> Workflow Runtime / Stage Engine
-> Existing Reader Core Modules
- FreshRSS pull
- extraction
- summary
- filter
- candidate builder
- article summary
-> Run Store (filesystem-backed)
-> Artifact Store (existing outputs directory)
```
### 6.1 分层说明
#### A. MCP Server Layer
职责:
- tool 注册
- 输入校验
- 输出包装
- 对外暴露稳定 API
不负责:
- 大段业务逻辑
- 状态推进细节
- 复杂流程编排
#### B. Workflow Service Layer
职责:
- 接收“启动日报任务”“恢复任务”“查询状态”等请求
- 管理 run 生命周期
- 将调用分发给 runtime/stage engine
#### C. Workflow Runtime / Stage Engine
职责:
- 执行各阶段
- 记录阶段状态
- 收集中间产物
- 更新运行态文件
- 支持 resume / rerun
#### D. Reader Core Modules
继续复用现有业务能力:
- `summary_mcp.core.*`
- `summary_mcp.workflows.*`
- 现有脚本中成熟的逻辑
这层应该尽量保持“业务逻辑纯净”,不要耦合 MCP 协议。
---
## 7. 核心抽象
### 7.1 Run
每一次完整的 reader 工作流执行,都对应一个 `run_id`。
建议语义:
- `run_id` 是 reader 工作流的一等公民
- 所有状态、产物、恢复都围绕 `run_id` 建立
- 所有下游消费都优先通过 `run_id` 取结果,而不是手拼路径
建议输出目录继续沿用现有结构:
```text
outputs/freshrss/rerun/<run_id>/
```
### 7.2 Stage
建议将 reader 的正式工作流拆为明确阶段:
1. `fetch_feed`
2. `extract_articles`
3. `generate_summaries`
4. `apply_filters`
5. `build_delivery_payload`
6. `write_run_report`
7. `completed`
对于单篇总结链路,可单独建另一类 workflow,或者作为独立 run type。
### 7.3 Status
建议统一使用:
- `queued`
- `running`
- `partial`
- `success`
- `failed`
- `cancelled`
其中:
- `partial` 表示已有阶段成功,但整体尚未完成或部分失败
- `failed` 表示本次 run 已终止且未达到最终成功
### 7.4 Artifact
所有关键输出都应被注册为 artifact,而不只是“写到了某个目录”。
artifact 至少应包含:
- `name`
- `path`
- `kind`
- `stage`
- `exists`
- `created_at`
- `metadata`
关键 artifact 示例:
- `raw_output`
- `extracted_dir`
- `delivery_payload`
- `digest_brief`
- `run_report`
- `single_summary_dir`
---
## 8. 运行态持久化设计
### 8.1 为什么要单独持久化 run state
仅靠目录中是否存在某些 json 文件,不足以稳定表达:
- 当前是否还在跑
- 跑到了哪个阶段
- 哪个阶段失败
- 是否可以 resume
- 哪些 artifact 已确认可用
因此必须增加显式的运行态文件。
### 8.2 建议文件
建议在每个 run 目录下引入:
```text
outputs/freshrss/rerun/<run_id>/run-state.json
```
### 8.3 run-state.json 建议结构
```json
{
"run_id": "20260407-093000",
"workflow": "freshrss_daily_digest",
"run_type": "daily_digest",
"status": "running",
"current_stage": "extract_articles",
"started_at": "2026-04-07T09:30:00+08:00",
"updated_at": "2026-04-07T09:31:10+08:00",
"finished_at": null,
"input": {
"limit": 5,
"mark_read": true,
"include_read": false,
"debug_artifacts": false
},
"stages": [
{
"name": "fetch_feed",
"status": "success",
"started_at": "2026-04-07T09:30:00+08:00",
"finished_at": "2026-04-07T09:30:05+08:00",
"outputs": {
"pulled_count": 5,
"raw_output": "outputs/freshrss/rerun/20260407-093000/raw/freshrss.raw.json"
},
"error": null
},
{
"name": "extract_articles",
"status": "running",
"started_at": "2026-04-07T09:30:05+08:00",
"finished_at": null,
"outputs": {
"completed_items": 3,
"expected_items": 5
},
"error": null
}
],
"artifacts": [
{
"name": "raw_output",
"kind": "json",
"stage": "fetch_feed",
"path": "outputs/freshrss/rerun/20260407-093000/raw/freshrss.raw.json",
"exists": true
}
],
"error": null,
"recovery": {
"resumable": true,
"resume_from_stage": "extract_articles",
"last_success_stage": "fetch_feed"
}
}
```
### 8.4 设计原则
- 运行中每完成一个 stage,就更新一次 `run-state.json`
- 不要求数据库,先以文件系统为准
- 所有对外 status 查询优先读 `run-state.json`
- 其他 output 文件仍然可以保留现有格式和目录结构
---
## 9. MCP Tool 设计
不建议继续把正式能力收口成一个“巨型同步工具”。
建议拆成以下 MCP tools。
### 9.1 任务启动类
#### `start_freshrss_digest_run`
作用:
- 启动一条新的日报工作流
- 默认返回 `run_id`,而不是要求调用方一直同步等待到最后
输入示例:
```json
{
"limit": 5,
"mark_read": true,
"include_read": false,
"debug_artifacts": false,
"timeout_seconds": 60,
"max_retries": 2
}
```
输出示例:
```json
{
"run_id": "20260407-093000",
"status": "queued",
"workflow": "freshrss_daily_digest",
"output_dir": "outputs/freshrss/rerun/20260407-093000"
}
```
#### 兼容策略
第一阶段可保留现有 `run_freshrss_openclaw_pipeline`,但其语义逐步调整为:
- 内部复用新的 workflow runtime
- 可选 `wait=true/false`
- 默认建议 `wait=false`
---
### 9.2 状态查询类
#### `get_run_status`
作用:
- 查询 run 当前状态
- 返回 current stage、completed stages、关键 artifact、失败信息
输入:
```json
{
"run_id": "20260407-093000"
}
```
输出字段建议:
- `run_id`
- `workflow`
- `status`
- `current_stage`
- `started_at`
- `updated_at`
- `finished_at`
- `progress`
- `completed_stages`
- `failed_stage`
- `error_summary`
- `artifacts`
- `recovery`
#### `list_runs`
作用:
- 支持按日期、状态、workflow 类型查看最近 runs
输入示例:
```json
{
"workflow": "freshrss_daily_digest",
"status": "failed",
"latest_n": 10
}
```
---
### 9.3 结果读取类
#### `get_delivery_payload`
作用:
- 按 `run_id` 获取 delivery payload
- 返回结构化对象,不要求调用方自己读文件
#### `get_digest_brief`
作用:
- 按 `run_id` 获取 digest brief
- 可支持 `mode=public|internal`
#### `get_run_report`
作用:
- 返回 run-report 内容
- 供 OpenClaw / 调试 / 运维查看
#### `get_extracted_article`
作用:
- 获取单篇 extracted 内容
- 输入:`run_id + item_id` 或 `run_id + item_index`
#### `list_run_artifacts`
作用:
- 统一列出 run 下当前注册的 artifact
---
### 9.4 恢复与补跑类
#### `resume_run`
作用:
- 基于已有 `run-state.json` 和中间产物,从可恢复点继续
输入示例:
```json
{
"run_id": "20260407-093000"
}
```
#### `rerun_stage`
作用:
- 指定从某个阶段开始重跑
输入示例:
```json
{
"run_id": "20260407-093000",
"stage": "generate_summaries",
"force": true
}
```
#### `explain_run_failure`
作用:
- 给出适合 OpenClaw / 人类阅读的失败解释
- 不只是 Python stacktrace
---
### 9.5 单篇总结相关工具
现有 `generate_article_summaries` 可继续保留,但建议长期也纳入 run 模型。
后续可扩展为:
- `start_article_summary_run`
- `get_article_summary_status`
- `get_article_summary_outputs`
短期内,如果单篇总结执行耗时可控,也可以先保持同步接口。
---
## 10. 向后兼容策略
为了降低迁移成本,不建议一次性砍掉现有接口。
### 10.1 保留现有 MCP 工具
短期保留:
- `run_freshrss_openclaw_pipeline`
- `generate_article_summaries`
但内部逐步改为调用新的 workflow runtime。
### 10.2 调整 `run_freshrss_openclaw_pipeline` 语义
建议演进成:
- `wait=true` 时:保留当前类似同步行为
- `wait=false` 时:只返回 `run_id`
- 若未显式指定,生产建议默认 `wait=false`
### 10.3 CLI 的新定位
CLI 继续保留,但只作为:
- debug
- fallback
- 本地排障
- 开发验证
CLI 最好只是新 runtime 的薄包装,而不是另一套独立实现。
---
## 11. 与 OpenClaw 的集成方式
### 11.1 旧模式
OpenClaw:
- 直接 exec reader 脚本
- 等待整条 CLI 跑完
- 根据 stdout 或目录文件判断是否成功
### 11.2 新模式
OpenClaw:
1. 调用 `start_freshrss_digest_run`
2. 获得 `run_id`
3. 周期性调用 `get_run_status`
4. status = `success` 后调用 `get_delivery_payload`
5. 再执行下游 Hugo / chat / IMA 编排
这样会有几个明显好处:
- 上游 reader 运行态可观察
- OpenClaw 不必绑死在长 CLI 会话上
- 中途失败可以明确知道失败位置
- 下游编排只依赖结构化结果
---
## 12. 最小实现方案(MVP)
为了尽快落地,不建议一次做到最重。
### Phase 1:先做内部状态化
目标:
- 在现有 pipeline 内引入 `run_id`
- 新增 `run-state.json`
- 固化 stage 切分
- 所有关键输出注册成 artifact
- 将 resume 所需信息写入 recovery 字段
这一步完成后,即使对外接口还没完全变化,内部也已经不再是“黑箱长函数”。
### Phase 2:新增 MCP 状态查询接口
目标:
- 新增 `get_run_status`
- 新增 `get_delivery_payload`
- 新增 `list_runs`
- 新增 `resume_run`
这一步完成后,OpenClaw 就可以逐步改走正式服务调用。
### Phase 3:调整现有生产接入
目标:
- OpenClaw 默认不再 exec `scripts/run_freshrss_pipeline.py`
- OpenClaw 默认走 MCP run + status + payload 模式
- 将 CLI 降级为 debug/fallback
---
## 13. 目录与模块建议
建议新增一层 workflow runtime 模块,例如:
```text
src/summary_mcp/
server.py
runtime/
run_store.py
artifact_store.py
state_models.py
workflow_service.py
stage_runner.py
workflows/
freshrss_pipeline.py
article_summary.py
core/
...
```
### 建议职责
- `runtime/state_models.py`
- RunState / StageState / Artifact models
- `runtime/run_store.py`
- 读写 `run-state.json`
- `runtime/artifact_store.py`
- artifact 注册与查询
- `runtime/workflow_service.py`
- start / status / resume / rerun 核心服务
- `runtime/stage_runner.py`
- 阶段推进与失败捕获
这样可以保持:
- 协议层清晰
- 运行态管理独立
- 业务逻辑复用现有 workflows
---
## 14. 风险与注意事项
### 14.1 不要把“异步”理解成“一定要上复杂队列”
当前阶段,异步的核心不是上 Redis/Celery,而是:
- 有 `run_id`
- 有状态文件
- 可以先启动、后查询
单机单进程也完全可以做到。
### 14.2 不要让 OpenClaw 继续依赖 reader 内部目录细节
OpenClaw 可以知道输出目录存在,但不应继续以“自己拼 `outputs/.../*.json`”作为主交互方式。
正式模式下,优先通过 MCP 取:
- status
- payload
- report
- artifact list
### 14.3 不要保留两套行为漂移的实现
如果 CLI 和 MCP 背后各跑各的逻辑,后续一定会漂移。
正确做法:
- 先统一 runtime
- 再让 CLI / MCP 都调用同一套 runtime
---
## 15. 最终建议
### 架构判断
Reader 现在已经不再是一个简单脚本仓库,而是:
- 有明确上游输入
- 有稳定工作流
- 有中间产物
- 有下游消费者
- 有恢复与补跑需求
因此它应该正式升级为:
> **Reader Workflow MCP Service**
### 最终建议清单
1. 把 `run_id / stage / status / artifact / recovery` 作为正式核心抽象
2. 在 `outputs/freshrss/rerun/<run_id>/` 下新增 `run-state.json`
3. 新增 `get_run_status / list_runs / get_delivery_payload / resume_run`
4. 让现有 `run_freshrss_openclaw_pipeline` 内部复用新 runtime
5. 让 OpenClaw 逐步从 exec 切到 MCP 调用
6. CLI 保留,但降级为 debug/fallback
---
## 16. 一句话结论
Reader 的正式生产能力不应再主要依赖长 CLI exec,而应演进为:
**以 MCP 为正式入口、以 run-state 为运行真相、以 artifact 为交付契约、以 OpenClaw 为下游编排层的工作流服务。**
+209
View File
@@ -0,0 +1,209 @@
# Reader MCP 实施计划
## 目标
把 reader 从“生产上主要依赖长 CLI/exec”推进到“以 MCP 为正式入口的工作流服务”。
## 协作分工
- **架构与方向**:由我负责
- **编码实现**:由 Codex 负责
- **同步机制**:通过 `plans/` 下规划文档 + `TODO.md` 保持进度与方向一致
## 当前权威文档
实现前,必须先读:
1. `plans/reader-mcp-architecture-design.md`
2. `plans/reader-mcp-implementation-plan.md`
3. `TODO.md`
4. `plans/issues/2026-04-06-reader-digest-sigterm.md`
5. `docs/openclaw/openclaw-handoff.md`
如实现细节与旧文档冲突,以:
1. 最新架构设计文档
2. 最新实施计划
3. TODO 当前项
为准。
---
## 里程碑
### M1:运行态落地(最优先)
目标:让现有 freshrss pipeline 拥有明确 run state。
交付:
- 新增 `run-state.json` 持久化能力
- 定义 `RunState / StageState / ArtifactRecord` 模型
- 现有 pipeline 按 stage 更新状态
- 保持现有产物目录兼容
完成标准:
- 任意一次 run 都能产出 `outputs/freshrss/rerun/<run_id>/run-state.json`
- 中途失败时也能看到失败阶段和已有 artifacts
---
### M2:状态查询接口
目标:通过 MCP 查询 run 状态,不再只靠 CLI 或手看目录。
交付:
- 新增 `get_run_status`
- 新增 `list_runs`
- 新增 `list_run_artifacts`
完成标准:
- OpenClaw 可通过 MCP 查询 run 当前状态
- 不需要直接读磁盘路径判断是否成功
---
### M3:结果读取接口
目标:通过 MCP 获取结果内容,而不是自己拼文件路径。
交付:
- 新增 `get_delivery_payload`
- 新增 `get_run_report`
- 视情况新增 `get_digest_brief`
- 视情况新增 `get_extracted_article`
完成标准:
- OpenClaw 只要知道 `run_id`,就能获取关键结果
---
### M4:恢复与补跑
目标:reader 具备正式的恢复机制。
交付:
- 新增 `resume_run`
- 视情况新增 `rerun_stage`
- 在 `run-state.json` 中记录 recovery 信息
完成标准:
- 至少支持从最近成功 stage 之后继续执行
- 能给出可恢复/不可恢复的明确判断
---
### M5:生产入口切换
目标:OpenClaw 正式从 exec 模式切到 MCP 模式。
交付:
- 保留 CLI 作为 debug/fallback
- 生产推荐入口改为 MCP run + status + payload
- 更新 handoff / README / docs
完成标准:
- 正式流程默认不再依赖长 CLI exec
---
## 编码原则
1. **优先复用现有业务逻辑**
- 不要重写成熟的 FreshRSS / extract / summary / filter 逻辑
- 优先抽 runtime 层把现有逻辑包起来
2. **先状态化,再协议扩展**
- 先把 run-state 跑通
- 再补 MCP tools
3. **保持向后兼容**
- `run_freshrss_openclaw_pipeline` 先保留
- 内部逐步改为调用新的 runtime
4. **CLI 降级,不删除**
- 仍保留 debug/fallback 价值
- 但不要让 CLI 和 MCP 背后出现两套逻辑
5. **小步提交**
- 每完成一个可验证的小目标就提交
- 不要攒一个超大变更
---
## 建议目录改造
建议新增:
```text
src/summary_mcp/runtime/
state_models.py
run_store.py
artifact_store.py
workflow_service.py
stage_runner.py
```
说明:
- 目录名可微调
- 但必须把“运行态管理”从现有 workflows 中抽出来,避免继续黑箱化
---
## Codex 工作方式要求
Codex 每次开始前:
1. 先读 `plans/reader-mcp-architecture-design.md`
2. 再读本文件
3. 再读 `TODO.md`
4. 只处理 TODO 中 `TODO` 状态的当前优先项
Codex 每完成一项后:
1. 更新 `TODO.md`
2. 在 TODO 对应项下补:
- 完成情况
- 改动文件
- 遗留风险
3. 如实现偏离原架构,必须先更新 `plans/` 文档,再继续代码
---
## 决策规则
如果出现以下情况:
- 需要新增 MCP tool,但架构设计中未定义
- 需要改动现有 pipeline 关键语义
- 需要引入数据库 / 队列 / 线程池 / 后台 worker
- 需要改变 OpenClaw 与 reader 的边界
则不允许 Codex自行拍板,必须先回写到:
- `plans/reader-mcp-architecture-design.md`
- 或新增 `plans/issues/*.md`
由架构层确认后再继续。
---
## 当前实现顺序(强约束)
按以下顺序推进:
1. `run-state.json` 模型与持久化
2. pipeline 中 stage 状态更新
3. `get_run_status`
4. `list_runs` / `list_run_artifacts`
5. `get_delivery_payload` / `get_run_report`
6. `resume_run`
7. 再考虑 `rerun_stage`
不要一上来就做复杂异步后台队列。
---
## 一句话执行口径
先把 reader 做成“有运行真相的 MCP 工作流服务”,再做更多工具;不要反过来先堆接口名。
+17
View File
@@ -0,0 +1,17 @@
from .query_service import get_delivery_payload, get_run_report, get_run_status, list_run_artifacts, list_runs
from .run_store import RunStore
from .state_models import ArtifactRecord, RecoveryState, RunError, RunState, StageState
__all__ = [
"ArtifactRecord",
"get_delivery_payload",
"get_run_report",
"RecoveryState",
"RunError",
"RunState",
"RunStore",
"StageState",
"get_run_status",
"list_run_artifacts",
"list_runs",
]
+595
View File
@@ -0,0 +1,595 @@
from __future__ import annotations
# Compatibility note:
# historical runs may have a timestamp directory name that differs from the recorded run_id,
# so status queries resolve both identifiers without changing the existing output layout.
import json
from datetime import datetime
from pathlib import Path
from typing import Any
from .run_store import RunStore
from .state_models import ArtifactRecord, RunError, RunState, StageState
REPO_ROOT = Path(__file__).resolve().parents[3]
RUNS_ROOT = REPO_ROOT / "outputs" / "freshrss" / "rerun"
DEFAULT_WORKFLOW = "freshrss_daily_digest"
DEFAULT_RUN_TYPE = "daily_digest"
DEFAULT_STAGES = [
"fetch_feed",
"extract_articles",
"generate_summaries",
"apply_filters",
"build_delivery_payload",
"write_run_report",
]
DISCOVERED_ARTIFACTS = [
{"name": "run_state", "kind": "json", "stage": "runtime_state", "relative_path": Path("run-state.json")},
{"name": "run_report", "kind": "json", "stage": "write_run_report", "relative_path": Path("run-report.json")},
{"name": "raw_output", "kind": "json", "stage": "fetch_feed", "relative_path": Path("raw/freshrss.raw.json")},
{"name": "items_output", "kind": "json", "stage": "fetch_feed", "relative_path": Path("items/freshrss.items.json")},
{"name": "extracted_dir", "kind": "directory", "stage": "extract_articles", "relative_path": Path("extracted")},
{"name": "summary_dir", "kind": "directory", "stage": "generate_summaries", "relative_path": Path("summary")},
{"name": "candidate_dir", "kind": "directory", "stage": "apply_filters", "relative_path": Path("candidates")},
{
"name": "delivery_payload",
"kind": "json",
"stage": "build_delivery_payload",
"relative_path": Path("candidates/openclaw-delivery-payload.json"),
},
{
"name": "digest_brief",
"kind": "json",
"stage": "build_delivery_payload",
"relative_path": Path("candidates/digest-brief.json"),
},
]
def get_run_status(*, run_id: str) -> dict[str, Any]:
record = _resolve_run_record(run_id)
return _build_status_response(record)
def list_runs(
*,
workflow: str | None = None,
status: str | None = None,
latest_n: int = 20,
) -> dict[str, Any]:
if latest_n <= 0:
raise ValueError("latest_n must be greater than 0.")
records = []
for run_dir in _iter_run_dirs():
record = _build_run_record(run_dir)
if workflow is not None and record["workflow"] != workflow:
continue
if status is not None and record["status"] != status:
continue
records.append(record)
records.sort(key=_record_sort_key, reverse=True)
selected_records = records[:latest_n]
return {
"runs": [_build_list_response(record) for record in selected_records],
"count": len(selected_records),
"filters": {
"workflow": workflow,
"status": status,
"latest_n": latest_n,
},
}
def list_run_artifacts(*, run_id: str) -> dict[str, Any]:
record = _resolve_run_record(run_id)
artifacts = _collect_artifacts(record)
return {
"run_id": record["run_id"],
"run_dir": record["run_dir"].name,
"workflow": record["workflow"],
"status": record["status"],
"output_dir": _normalize_repo_path(record["run_dir"]),
"state_source": record["state_source"],
"artifact_count": len(artifacts),
"artifacts": artifacts,
}
def get_delivery_payload(*, run_id: str) -> dict[str, Any]:
record = _resolve_run_record(run_id)
artifact = _resolve_artifact(
record,
artifact_name="delivery_payload",
relative_path=Path("candidates/openclaw-delivery-payload.json"),
fallback_stage="build_delivery_payload",
fallback_kind="json",
report_field="delivery_output",
)
payload = _load_json(_resolve_repo_path(artifact["path"]))
candidates = payload.get("candidates")
candidate_count = len(candidates) if isinstance(candidates, list) else 0
stats = payload.get("stats") if isinstance(payload.get("stats"), dict) else None
response = _build_run_lookup_response(record)
response.update(
{
"artifact": artifact,
"payload_schema_version": payload.get("schema_version"),
"generated_at": payload.get("generated_at"),
"delivery_date": payload.get("date"),
"candidate_count": candidate_count,
"stats": stats,
"payload": payload,
}
)
return response
def get_run_report(*, run_id: str) -> dict[str, Any]:
record = _resolve_run_record(run_id)
artifact = _resolve_artifact(
record,
artifact_name="run_report",
relative_path=Path("run-report.json"),
fallback_stage="write_run_report",
fallback_kind="json",
report_field="report_output",
)
report = _load_json(_resolve_repo_path(artifact["path"]))
items = report.get("items")
item_count = len(items) if isinstance(items, list) else 0
response = _build_run_lookup_response(record)
response.update(
{
"artifact": artifact,
"started_at": report.get("started_at"),
"completed_at": report.get("completed_at"),
"requested_limit": report.get("requested_limit"),
"pulled_count": report.get("pulled_count"),
"delivered_count": report.get("delivered_count"),
"marked_read_count": report.get("marked_read_count"),
"mark_read_requested": report.get("mark_read_requested"),
"debug_artifacts": report.get("debug_artifacts"),
"raw_output": _normalize_optional_repo_path(report.get("raw_output")),
"delivery_output": _normalize_optional_repo_path(report.get("delivery_output")),
"digest_brief_output": _normalize_optional_repo_path(report.get("digest_brief_output")),
"status_counts": report.get("status_counts"),
"keyword_index": _normalize_keyword_index(report.get("keyword_index")),
"item_count": item_count,
"report": report,
}
)
return response
def _resolve_run_record(run_id: str) -> dict[str, Any]:
for run_dir in _iter_run_dirs():
record = _build_run_record(run_dir)
if run_id in record["aliases"]:
return record
raise FileNotFoundError(f"Run not found for run_id: {run_id}")
def _iter_run_dirs() -> list[Path]:
if not RUNS_ROOT.exists():
return []
return sorted((path for path in RUNS_ROOT.iterdir() if path.is_dir()), key=lambda path: path.name, reverse=True)
def _build_run_record(run_dir: Path) -> dict[str, Any]:
state_path = run_dir / "run-state.json"
report_path = run_dir / "run-report.json"
if state_path.exists():
run_store = RunStore.load(path=state_path, repo_root=REPO_ROOT)
state = run_store.state
run_id = state.run_id
return {
"run_id": run_id,
"workflow": state.workflow,
"run_type": state.run_type,
"status": state.status,
"current_stage": state.current_stage,
"started_at": state.started_at.isoformat(),
"updated_at": state.updated_at.isoformat(),
"finished_at": state.finished_at.isoformat() if state.finished_at else None,
"stages": [_stage_to_dict(stage) for stage in state.stages],
"error": _error_to_dict(state.error),
"recovery": state.recovery.model_dump(mode="json"),
"run_dir": run_dir,
"state_source": "run_state",
"aliases": {run_id, run_dir.name},
"state": state,
"report": _load_json(report_path) if report_path.exists() else None,
}
report = _load_json(report_path) if report_path.exists() else None
run_id = str(report.get("run_id")) if isinstance(report, dict) and report.get("run_id") else run_dir.name
inferred_record = _infer_run_record_from_directory(run_dir=run_dir, report=report, run_id=run_id)
inferred_record["aliases"] = {run_id, run_dir.name}
return inferred_record
def _infer_run_record_from_directory(*, run_dir: Path, report: dict[str, Any] | None, run_id: str) -> dict[str, Any]:
discovered_artifacts = _discover_artifacts(run_dir)
completed_stages = [artifact["stage"] for artifact in discovered_artifacts if artifact["stage"] in DEFAULT_STAGES]
deduped_completed_stages = []
for stage_name in DEFAULT_STAGES:
if stage_name in completed_stages:
deduped_completed_stages.append(stage_name)
if report is not None:
status = _infer_status_from_report(report)
started_at = _maybe_iso(report.get("started_at")) or _parse_run_dir_timestamp(run_dir.name)
updated_at = _maybe_iso(report.get("completed_at")) or started_at
finished_at = _maybe_iso(report.get("completed_at")) or updated_at
stages = [
_stage_dict(name=stage_name, status="success")
for stage_name in DEFAULT_STAGES
if stage_name in {"fetch_feed", "extract_articles", "generate_summaries", "apply_filters", "build_delivery_payload", "write_run_report"}
]
error = None
else:
status = "failed"
started_at = _parse_run_dir_timestamp(run_dir.name)
updated_at = started_at
finished_at = updated_at if discovered_artifacts else None
stages = [_stage_dict(name=stage_name, status="success") for stage_name in deduped_completed_stages]
next_stage = _infer_failed_stage(deduped_completed_stages)
if next_stage is not None:
stages.append(_stage_dict(name=next_stage, status="failed"))
error = {
"type": "InferredRunState",
"message": "run-state.json is missing; status inferred from existing run directory.",
"stage": next_stage,
"details": {},
}
else:
error = {
"type": "InferredRunState",
"message": "run-state.json is missing and no completed stage could be confirmed.",
"stage": None,
"details": {},
}
return {
"run_id": run_id,
"workflow": DEFAULT_WORKFLOW,
"run_type": DEFAULT_RUN_TYPE,
"status": status,
"current_stage": None,
"started_at": started_at,
"updated_at": updated_at,
"finished_at": finished_at,
"stages": stages,
"error": error,
"recovery": {
"resumable": False,
"resume_from_stage": None,
"last_success_stage": deduped_completed_stages[-1] if deduped_completed_stages else None,
},
"run_dir": run_dir,
"state_source": "directory_inference",
"state": None,
"report": report,
}
def _build_status_response(record: dict[str, Any]) -> dict[str, Any]:
completed_stages = [stage["name"] for stage in record["stages"] if stage["status"] == "success"]
failed_stage = next((stage["name"] for stage in record["stages"] if stage["status"] == "failed"), None)
return {
"run_id": record["run_id"],
"run_dir": record["run_dir"].name,
"workflow": record["workflow"],
"run_type": record["run_type"],
"status": record["status"],
"current_stage": record["current_stage"],
"started_at": record["started_at"],
"updated_at": record["updated_at"],
"finished_at": record["finished_at"],
"output_dir": _normalize_repo_path(record["run_dir"]),
"progress": _build_progress(record["stages"]),
"completed_stages": completed_stages,
"failed_stage": failed_stage,
"error_summary": record["error"],
"artifacts": _collect_artifacts(record),
"recovery": record["recovery"],
"state_source": record["state_source"],
}
def _build_run_lookup_response(record: dict[str, Any]) -> dict[str, Any]:
return {
"run_id": record["run_id"],
"run_dir": record["run_dir"].name,
"workflow": record["workflow"],
"run_type": record["run_type"],
"status": record["status"],
"output_dir": _normalize_repo_path(record["run_dir"]),
"state_source": record["state_source"],
}
def _build_list_response(record: dict[str, Any]) -> dict[str, Any]:
progress = _build_progress(record["stages"])
return {
"run_id": record["run_id"],
"run_dir": record["run_dir"].name,
"workflow": record["workflow"],
"run_type": record["run_type"],
"status": record["status"],
"current_stage": record["current_stage"],
"started_at": record["started_at"],
"updated_at": record["updated_at"],
"finished_at": record["finished_at"],
"output_dir": _normalize_repo_path(record["run_dir"]),
"progress": progress,
"recovery": record["recovery"],
"artifact_count": len(_collect_artifacts(record)),
"state_source": record["state_source"],
}
def _collect_artifacts(record: dict[str, Any]) -> list[dict[str, Any]]:
artifacts: list[dict[str, Any]] = []
seen_names: set[str] = set()
state = record.get("state")
if isinstance(state, RunState):
for artifact in state.artifacts:
artifact_dict = _artifact_from_state_record(artifact)
artifacts.append(artifact_dict)
seen_names.add(artifact_dict["name"])
for artifact in _discover_artifacts(record["run_dir"]):
if artifact["name"] in seen_names:
continue
artifacts.append(artifact)
seen_names.add(artifact["name"])
artifacts.sort(key=lambda artifact: (artifact["stage"], artifact["name"]))
return artifacts
def _resolve_artifact(
record: dict[str, Any],
*,
artifact_name: str,
relative_path: Path,
fallback_stage: str,
fallback_kind: str,
report_field: str | None = None,
) -> dict[str, Any]:
state = record.get("state")
if isinstance(state, RunState):
for artifact in state.artifacts:
if artifact.name != artifact_name:
continue
if _artifact_exists(artifact.path):
return _artifact_from_state_record(artifact)
standard_path = record["run_dir"] / relative_path
if standard_path.exists():
return _build_resolved_artifact(
name=artifact_name,
path=standard_path,
kind=fallback_kind,
stage=fallback_stage,
source="standard_path",
)
report = record.get("report")
if report_field and isinstance(report, dict):
report_path_value = report.get(report_field)
if isinstance(report_path_value, str) and report_path_value.strip():
report_path = _resolve_repo_path(report_path_value)
if report_path.exists():
return _build_resolved_artifact(
name=artifact_name,
path=report_path,
kind=fallback_kind,
stage=fallback_stage,
source="run_report_reference",
)
for path in sorted(record["run_dir"].rglob(relative_path.name)):
if path.is_file():
return _build_resolved_artifact(
name=artifact_name,
path=path,
kind=fallback_kind,
stage=fallback_stage,
source="directory_scan",
)
available_artifacts = [artifact["name"] for artifact in _collect_artifacts(record)]
available_summary = ", ".join(available_artifacts) if available_artifacts else "none"
raise FileNotFoundError(
f"Artifact '{artifact_name}' not found for run_id: {record['run_id']} "
f"(status={record['status']}, run_dir={record['run_dir'].name}, available_artifacts={available_summary})"
)
def _build_resolved_artifact(
*,
name: str,
path: Path,
kind: str,
stage: str,
source: str,
) -> dict[str, Any]:
return {
"name": name,
"kind": kind,
"stage": stage,
"path": _normalize_repo_path(path),
"exists": True,
"created_at": _iso_from_stat(path),
"metadata": {},
"source": source,
}
def _discover_artifacts(run_dir: Path) -> list[dict[str, Any]]:
artifacts = []
for spec in DISCOVERED_ARTIFACTS:
path = run_dir / spec["relative_path"]
if not path.exists():
continue
artifacts.append(
{
"name": spec["name"],
"kind": spec["kind"],
"stage": spec["stage"],
"path": _normalize_repo_path(path),
"exists": True,
"created_at": _iso_from_stat(path),
"metadata": {},
"source": "directory_scan",
}
)
return artifacts
def _artifact_from_state_record(artifact: ArtifactRecord) -> dict[str, Any]:
return {
"name": artifact.name,
"kind": artifact.kind,
"stage": artifact.stage,
"path": _normalize_repo_path(artifact.path),
"exists": _artifact_exists(artifact.path),
"created_at": artifact.created_at.isoformat(),
"metadata": artifact.metadata,
"source": "run_state",
}
def _artifact_exists(path_value: str) -> bool:
path = _resolve_repo_path(path_value)
return path.exists()
def _build_progress(stages: list[dict[str, Any]]) -> dict[str, Any]:
expected_stage_names = list(DEFAULT_STAGES)
stage_status_by_name = {stage["name"]: stage["status"] for stage in stages}
for stage_name in stage_status_by_name:
if stage_name not in expected_stage_names:
expected_stage_names.append(stage_name)
completed_count = sum(1 for stage_name in expected_stage_names if stage_status_by_name.get(stage_name) == "success")
running_count = sum(1 for stage_name in expected_stage_names if stage_status_by_name.get(stage_name) == "running")
failed_count = sum(1 for stage_name in expected_stage_names if stage_status_by_name.get(stage_name) == "failed")
pending_count = sum(1 for stage_name in expected_stage_names if stage_status_by_name.get(stage_name, "pending") == "pending")
return {
"completed_stage_count": completed_count,
"running_stage_count": running_count,
"failed_stage_count": failed_count,
"pending_stage_count": pending_count,
"total_stage_count": len(expected_stage_names),
}
def _infer_status_from_report(report: dict[str, Any]) -> str:
status_counts = report.get("status_counts")
if isinstance(status_counts, dict) and any(key in status_counts for key in {"extract_failed", "summary_failed"}):
return "partial"
return "success"
def _infer_failed_stage(completed_stages: list[str]) -> str | None:
for stage_name in DEFAULT_STAGES:
if stage_name not in completed_stages:
return stage_name
return None
def _stage_to_dict(stage: StageState) -> dict[str, Any]:
return {
"name": stage.name,
"status": stage.status,
"started_at": stage.started_at.isoformat() if stage.started_at else None,
"finished_at": stage.finished_at.isoformat() if stage.finished_at else None,
"outputs": stage.outputs,
"error": _error_to_dict(stage.error),
}
def _stage_dict(*, name: str, status: str) -> dict[str, Any]:
return {
"name": name,
"status": status,
"started_at": None,
"finished_at": None,
"outputs": {},
"error": None,
}
def _error_to_dict(error: RunError | dict[str, Any] | None) -> dict[str, Any] | None:
if error is None:
return None
if isinstance(error, RunError):
return error.model_dump(mode="json")
return error
def _normalize_repo_path(path_value: str | Path) -> str:
path = _resolve_repo_path(path_value)
try:
return str(path.relative_to(REPO_ROOT))
except ValueError:
return str(path)
def _resolve_repo_path(path_value: str | Path) -> Path:
path = Path(path_value)
if path.is_absolute():
return path
return REPO_ROOT / path
def _normalize_optional_repo_path(path_value: Any) -> str | None:
if not isinstance(path_value, str) or not path_value.strip():
return None
return _normalize_repo_path(path_value)
def _normalize_keyword_index(keyword_index: Any) -> dict[str, Any] | None:
if not isinstance(keyword_index, dict):
return None
normalized = dict(keyword_index)
normalized["daily_output"] = _normalize_optional_repo_path(keyword_index.get("daily_output"))
normalized["stats_output"] = _normalize_optional_repo_path(keyword_index.get("stats_output"))
return normalized
def _iso_from_stat(path: Path) -> str:
return datetime.fromtimestamp(path.stat().st_mtime).astimezone().isoformat()
def _parse_run_dir_timestamp(run_dir_name: str) -> str | None:
try:
return datetime.strptime(run_dir_name, "%Y%m%d-%H%M%S").astimezone().isoformat()
except ValueError:
return None
def _maybe_iso(value: Any) -> str | None:
if isinstance(value, str):
try:
return datetime.fromisoformat(value).isoformat()
except ValueError:
return value
return None
def _record_sort_key(record: dict[str, Any]) -> tuple[str, str]:
return (record.get("updated_at") or "", record["run_dir"].name)
def _load_json(path: Path) -> dict[str, Any]:
return json.loads(path.read_text(encoding="utf-8-sig"))
+190
View File
@@ -0,0 +1,190 @@
from __future__ import annotations
import json
from datetime import datetime
from pathlib import Path
from typing import Any
from .state_models import ArtifactRecord, RecoveryState, RunError, RunState, StageState
class RunStore:
def __init__(self, *, path: Path, state: RunState, repo_root: Path | None = None) -> None:
self.path = path
self.state = state
self.repo_root = repo_root
@classmethod
def create(
cls,
*,
path: Path,
run_id: str,
workflow: str,
run_type: str,
started_at: datetime,
input_payload: dict[str, Any] | None = None,
repo_root: Path | None = None,
) -> "RunStore":
state = RunState(
run_id=run_id,
workflow=workflow,
run_type=run_type,
status="running",
current_stage=None,
started_at=started_at,
updated_at=started_at,
input=input_payload or {},
recovery=RecoveryState(),
)
return cls(path=path, state=state, repo_root=repo_root)
@classmethod
def load(cls, *, path: Path, repo_root: Path | None = None) -> "RunStore":
state = RunState.model_validate(json.loads(path.read_text(encoding="utf-8-sig")))
return cls(path=path, state=state, repo_root=repo_root)
def save(self) -> None:
self.state.updated_at = datetime.now(tz=self.state.started_at.tzinfo)
self.path.parent.mkdir(parents=True, exist_ok=True)
self.path.write_text(
json.dumps(self.state.model_dump(mode="json"), ensure_ascii=False, indent=2),
encoding="utf-8",
)
def start_stage(self, name: str, *, outputs: dict[str, Any] | None = None) -> None:
stage = self._get_or_create_stage(name)
now = datetime.now(tz=self.state.started_at.tzinfo)
stage.status = "running"
stage.started_at = stage.started_at or now
stage.finished_at = None
if outputs:
stage.outputs.update(outputs)
stage.error = None
self.state.current_stage = name
self.state.status = "running"
self.state.error = None
self._refresh_recovery()
self.save()
def update_stage(self, name: str, *, outputs: dict[str, Any] | None = None) -> None:
stage = self._get_or_create_stage(name)
if outputs:
stage.outputs.update(outputs)
self.state.current_stage = name
self._refresh_recovery()
self.save()
def finish_stage(self, name: str, *, outputs: dict[str, Any] | None = None) -> None:
stage = self._get_or_create_stage(name)
now = datetime.now(tz=self.state.started_at.tzinfo)
stage.status = "success"
stage.started_at = stage.started_at or now
stage.finished_at = now
if outputs:
stage.outputs.update(outputs)
stage.error = None
if self.state.current_stage == name:
self.state.current_stage = None
self._refresh_recovery()
self.save()
def fail_stage(
self,
name: str,
*,
error: BaseException | RunError,
outputs: dict[str, Any] | None = None,
) -> None:
stage = self._get_or_create_stage(name)
now = datetime.now(tz=self.state.started_at.tzinfo)
stage.status = "failed"
stage.started_at = stage.started_at or now
stage.finished_at = now
if outputs:
stage.outputs.update(outputs)
stage_error = error if isinstance(error, RunError) else self._build_error(error, stage=name)
stage.error = stage_error
self.state.current_stage = name
self.state.status = "failed"
self.state.error = stage_error
self.state.finished_at = now
self._refresh_recovery()
self.save()
def register_artifact(
self,
*,
name: str,
path: Path,
kind: str,
stage: str,
metadata: dict[str, Any] | None = None,
) -> None:
artifact = ArtifactRecord(
name=name,
path=self._normalize_path(path),
kind=kind,
stage=stage,
exists=path.exists(),
created_at=datetime.now(tz=self.state.started_at.tzinfo),
metadata=metadata or {},
)
existing = next((item for item in self.state.artifacts if item.name == name), None)
if existing is None:
self.state.artifacts.append(artifact)
else:
existing.path = artifact.path
existing.kind = artifact.kind
existing.stage = artifact.stage
existing.exists = artifact.exists
existing.created_at = artifact.created_at
existing.metadata = artifact.metadata
self.save()
def finish_run(self, *, status: str) -> None:
now = datetime.now(tz=self.state.started_at.tzinfo)
self.state.status = status
self.state.current_stage = None
self.state.finished_at = now
self.state.error = None
self._refresh_recovery()
self.save()
def _get_or_create_stage(self, name: str) -> StageState:
for stage in self.state.stages:
if stage.name == name:
return stage
stage = StageState(name=name)
self.state.stages.append(stage)
return stage
def _refresh_recovery(self) -> None:
successful_stages = [stage.name for stage in self.state.stages if stage.status == "success"]
last_success_stage = successful_stages[-1] if successful_stages else None
resume_from_stage = None
if self.state.status == "running":
resume_from_stage = self.state.current_stage
elif self.state.status == "failed":
failed_stage = next((stage.name for stage in self.state.stages if stage.status == "failed"), None)
resume_from_stage = failed_stage or self.state.current_stage
self.state.recovery = RecoveryState(
resumable=self.state.status in {"running", "failed"} and resume_from_stage is not None,
resume_from_stage=resume_from_stage,
last_success_stage=last_success_stage,
)
def _build_error(self, error: BaseException, *, stage: str | None = None) -> RunError:
return RunError(type=type(error).__name__, message=str(error), stage=stage)
def _normalize_path(self, path: Path) -> str:
resolved_path = path.resolve()
if self.repo_root is not None:
try:
return str(resolved_path.relative_to(self.repo_root.resolve()))
except ValueError:
return str(path)
return str(path)
+58
View File
@@ -0,0 +1,58 @@
from __future__ import annotations
from datetime import datetime
from typing import Any, Literal
from pydantic import BaseModel, Field
RunStatus = Literal["running", "success", "partial", "failed"]
StageStatus = Literal["pending", "running", "success", "failed"]
class RunError(BaseModel):
type: str
message: str
stage: str | None = None
details: dict[str, Any] = Field(default_factory=dict)
class ArtifactRecord(BaseModel):
name: str
path: str
kind: str
stage: str
exists: bool = True
created_at: datetime
metadata: dict[str, Any] = Field(default_factory=dict)
class StageState(BaseModel):
name: str
status: StageStatus = "pending"
started_at: datetime | None = None
finished_at: datetime | None = None
outputs: dict[str, Any] = Field(default_factory=dict)
error: RunError | None = None
class RecoveryState(BaseModel):
resumable: bool = False
resume_from_stage: str | None = None
last_success_stage: str | None = None
class RunState(BaseModel):
run_id: str
workflow: str
run_type: str
status: RunStatus = "running"
current_stage: str | None = None
started_at: datetime
updated_at: datetime
finished_at: datetime | None = None
input: dict[str, Any] = Field(default_factory=dict)
stages: list[StageState] = Field(default_factory=list)
artifacts: list[ArtifactRecord] = Field(default_factory=list)
error: RunError | None = None
recovery: RecoveryState = Field(default_factory=RecoveryState)
+39
View File
@@ -15,6 +15,11 @@ from summary_mcp.models.filtering import FilterContext, FilterInput
from summary_mcp.models.item import Item from summary_mcp.models.item import Item
from summary_mcp.models.llm_result import LlmSummaryResult from summary_mcp.models.llm_result import LlmSummaryResult
from summary_mcp.models.summary_io import ExtractionInput from summary_mcp.models.summary_io import ExtractionInput
from summary_mcp.runtime import get_delivery_payload as load_delivery_payload
from summary_mcp.runtime import get_run_status as load_run_status
from summary_mcp.runtime import get_run_report as load_run_report
from summary_mcp.runtime import list_run_artifacts as load_run_artifacts
from summary_mcp.runtime import list_runs as load_runs
from summary_mcp.workflows import run_freshrss_pipeline from summary_mcp.workflows import run_freshrss_pipeline
from summary_mcp.workflows.article_summary import ArticleSummaryConfig, summarize_selected_articles from summary_mcp.workflows.article_summary import ArticleSummaryConfig, summarize_selected_articles
@@ -119,6 +124,40 @@ def run_freshrss_openclaw_pipeline(
return result return result
@mcp.tool()
def get_run_status(run_id: str) -> dict:
"""Get the current status of a workflow run by run_id."""
return load_run_status(run_id=run_id)
@mcp.tool()
def list_runs(
workflow: str | None = None,
status: str | None = None,
latest_n: int = 20,
) -> dict:
"""List recent workflow runs with optional workflow/status filters."""
return load_runs(workflow=workflow, status=status, latest_n=latest_n)
@mcp.tool()
def list_run_artifacts(run_id: str) -> dict:
"""List registered and discovered artifacts for a workflow run."""
return load_run_artifacts(run_id=run_id)
@mcp.tool()
def get_delivery_payload(run_id: str) -> dict:
"""Get the structured OpenClaw delivery payload for a workflow run."""
return load_delivery_payload(run_id=run_id)
@mcp.tool()
def get_run_report(run_id: str) -> dict:
"""Get the structured run report for a workflow run."""
return load_run_report(run_id=run_id)
@mcp.tool() @mcp.tool()
def generate_article_summaries( def generate_article_summaries(
*, *,
+345 -173
View File
@@ -1,14 +1,8 @@
from __future__ import annotations from __future__ import annotations
# FreshRSS 全链路管道:拉取未读条目 -> 内容提取 -> LLM 摘要 -> 规则过滤 ->
# 构建 OpenClaw delivery payload -> 写盘 -> 词元统计 -> 标记已读。
# 生产入口:run_freshrss_pipeline(),由 MCP 工具 run_freshrss_openclaw_pipeline 调用。
import json import json
import os import os
from datetime import date, datetime, timezone from datetime import date, datetime, timezone
UTC = timezone.utc
from pathlib import Path from pathlib import Path
from typing import Any from typing import Any
@@ -32,8 +26,10 @@ from summary_mcp.models.openclaw_delivery import (
build_openclaw_digest_brief, build_openclaw_digest_brief,
) )
from summary_mcp.models.summary_io import ExtractionInput from summary_mcp.models.summary_io import ExtractionInput
from summary_mcp.runtime import RunStore
UTC = timezone.utc
REPO_ROOT = Path(__file__).resolve().parents[3] REPO_ROOT = Path(__file__).resolve().parents[3]
OUTPUT_ROOT = REPO_ROOT / "outputs" OUTPUT_ROOT = REPO_ROOT / "outputs"
FRESHRSS_OUTPUT_ROOT = OUTPUT_ROOT / "freshrss" FRESHRSS_OUTPUT_ROOT = OUTPUT_ROOT / "freshrss"
@@ -44,6 +40,14 @@ DEFAULT_TERM_ALIASES_PATH = REPO_ROOT / "configs" / "term_aliases.json"
DEFAULT_TERM_STOPWORDS_PATH = REPO_ROOT / "configs" / "term_stopwords.json" DEFAULT_TERM_STOPWORDS_PATH = REPO_ROOT / "configs" / "term_stopwords.json"
DEFAULT_TERM_DAILY_DIR = DATA_ROOT / "daily" DEFAULT_TERM_DAILY_DIR = DATA_ROOT / "daily"
DEFAULT_TERM_STATS_PATH = DATA_ROOT / "term_stats.json" DEFAULT_TERM_STATS_PATH = DATA_ROOT / "term_stats.json"
WORKFLOW_NAME = "freshrss_daily_digest"
RUN_TYPE = "daily_digest"
FETCH_STAGE = "fetch_feed"
EXTRACT_STAGE = "extract_articles"
SUMMARY_STAGE = "generate_summaries"
FILTER_STAGE = "apply_filters"
DELIVERY_STAGE = "build_delivery_payload"
REPORT_STAGE = "write_run_report"
def _save_json(path: Path, payload: dict[str, Any] | list[Any]) -> None: def _save_json(path: Path, payload: dict[str, Any] | list[Any]) -> None:
@@ -56,7 +60,6 @@ def _load_json(path: Path) -> dict[str, Any]:
def _load_required_env(name: str, value: str | None) -> str: def _load_required_env(name: str, value: str | None) -> str:
# 优先使用传入的 value,否则读取同名环境变量;两者均缺失时抛出 RuntimeError
if value: if value:
return value return value
env_value = os.environ.get(name) env_value = os.environ.get(name)
@@ -75,25 +78,7 @@ def _maybe_path(enabled: bool, path: Path) -> Path | None:
return path if enabled else None return path if enabled else None
def _process_item( def _build_item_context(*, index: int, item: Any, resolved_output_dir: Path, debug_artifacts: bool) -> dict[str, Any]:
*,
index: int,
item: Any,
resolved_output_dir: Path,
resolved_prompt_path: Path,
resolved_run_id: str,
debug_artifacts: bool,
loaded_rules: list,
filter_context: Any,
max_retries: int,
timeout_seconds: float,
resolved_llm_api_key: str,
resolved_llm_model: str,
resolved_llm_api_url: str,
) -> dict[str, Any]:
# 处理单条 item:提取 -> LLM 摘要 -> 规则过滤 -> 候选构建
# 返回 item_report dict;delivered 时额外携带 _candidate/_external_id 供调用方解包
# 步骤 1:确定各中间文件路径(debug_artifacts=False 时大部分路径为 None,不写盘)
item_key = f"item-{index:02d}" item_key = f"item-{index:02d}"
item_path = _maybe_path(debug_artifacts, resolved_output_dir / "items" / f"{item_key}.item.json") item_path = _maybe_path(debug_artifacts, resolved_output_dir / "items" / f"{item_key}.item.json")
extracted_path = resolved_output_dir / "extracted" / f"{item_key}.extracted.json" extracted_path = resolved_output_dir / "extracted" / f"{item_key}.extracted.json"
@@ -102,10 +87,6 @@ def _process_item(
record_path = _maybe_path(debug_artifacts, resolved_output_dir / "candidates" / f"{item_key}.article-candidate-record.json") record_path = _maybe_path(debug_artifacts, resolved_output_dir / "candidates" / f"{item_key}.article-candidate-record.json")
openclaw_path = _maybe_path(debug_artifacts, resolved_output_dir / "candidates" / f"{item_key}.openclaw-candidate-input.json") openclaw_path = _maybe_path(debug_artifacts, resolved_output_dir / "candidates" / f"{item_key}.openclaw-candidate-input.json")
if item_path is not None:
_save_json(item_path, item.model_dump(mode="json"))
# 步骤 2:初始化 item_report,记录基础元信息;debug 模式下附加各中间文件路径
item_report: dict[str, Any] = { item_report: dict[str, Any] = {
"item_key": item_key, "item_key": item_key,
"item_id": item.item_id, "item_id": item.item_id,
@@ -117,115 +98,27 @@ def _process_item(
if debug_artifacts: if debug_artifacts:
item_report["paths"] = { item_report["paths"] = {
"item": str(item_path) if item_path else None, "item": str(item_path) if item_path else None,
"extracted": str(extracted_path) if extracted_path else None, "extracted": str(extracted_path),
"summary": str(summary_output) if summary_output else None, "summary": str(summary_output) if summary_output else None,
"filter": str(filter_path) if filter_path else None, "filter": str(filter_path) if filter_path else None,
"article_candidate": str(record_path) if record_path else None, "article_candidate": str(record_path) if record_path else None,
"openclaw_candidate": str(openclaw_path) if openclaw_path else None, "openclaw_candidate": str(openclaw_path) if openclaw_path else None,
} }
# 步骤 3:内容提取(RSS 内联内容 or 回源抓取);extracted.json 始终写盘 return {
extraction = extract_content(ExtractionInput(item=item)) "item": item,
extracted_payload = extraction.model_dump(mode="json") "item_key": item_key,
if extracted_path is not None: "item_path": item_path,
_save_json(extracted_path, extracted_payload) "extracted_path": extracted_path,
if not extraction.success or extraction.article is None: "summary_output": summary_output,
item_report["status"] = "extract_failed" "filter_path": filter_path,
item_report["error"] = extraction.error.model_dump(mode="json") if extraction.error else None "record_path": record_path,
return item_report "openclaw_path": openclaw_path,
"item_report": item_report,
# 步骤 4:LLM 摘要循环,失败时最多重试 max_retries 次 "extraction": None,
item_report["status"] = "extracted" "extracted_payload": None,
summary_exit_code, summary_payload, summary_report = run_loop_payload( "summary_payload": None,
extracted_payload=extracted_payload, }
prompt_path=resolved_prompt_path,
output_path=summary_output,
max_retries=max_retries,
timeout_seconds=timeout_seconds,
api_key=resolved_llm_api_key,
model=resolved_llm_model,
api_url=resolved_llm_api_url,
)
if summary_exit_code != 0 or summary_payload is None:
item_report["status"] = "summary_failed"
if summary_report is not None:
item_report["summary_errors"] = summary_report.errors
return item_report
# 步骤 5:规则引擎过滤,产出 keep/review/drop 决策及 digest_rank
summary = LlmSummaryResult.model_validate(summary_payload)
decision = evaluate_filter_rules(
FilterInput(item=item, article=extraction.article, summary=summary, context=filter_context),
loaded_rules,
)
if filter_path is not None:
_save_json(filter_path, decision.model_dump(mode="json"))
# 步骤 6:构建 ArticleCandidateRecord 和 OpenClawCandidateInput
record = build_article_candidate_record(
summary=summary,
article=extraction.article,
filter_result=decision,
item=item,
source_refs=CandidateSourceRefs(
item_path=str(item_path) if item_path else None,
extracted_path=str(extracted_path) if extracted_path else None,
summary_path=str(summary_output) if summary_output else None,
filter_path=str(filter_path) if filter_path else None,
),
metadata=CandidateMetadata(
generated_at=datetime.now(tz=UTC),
producer="run_freshrss_pipeline",
run_id=resolved_run_id,
),
)
openclaw_input = build_openclaw_candidate_input(record)
if record_path is not None:
_save_json(record_path, record.model_dump(mode="json"))
if openclaw_path is not None:
_save_json(openclaw_path, openclaw_input.model_dump(mode="json"))
# 步骤 7:标记 delivered,附加临时键 _candidate/_external_id 供主函数解包后加入投递列表
item_report["status"] = "delivered"
item_report["selection_decision"] = decision.decision
item_report["candidate_id"] = openclaw_input.candidate_id
item_report["_candidate"] = openclaw_input
item_report["_external_id"] = item.external_id
return item_report
def _build_and_persist_delivery(
*,
delivered_candidates: list[OpenClawCandidateInput],
resolved_run_id: str,
resolved_delivery_date: Any,
delivery_output: Path,
) -> tuple[OpenClawDeliveryPayload, Path, dict[str, Any]]:
# 按 digest_rank 排序,构建 delivery payload;同时派生一个更轻量的 digest-brief.json 供日报生成使用。
delivered_candidates.sort(key=lambda candidate: candidate.digest_rank, reverse=True)
delivery_payload = build_openclaw_delivery_payload(
delivered_candidates,
run_id=resolved_run_id,
for_date=resolved_delivery_date,
)
_save_json(delivery_output, delivery_payload.model_dump(mode="json"))
digest_brief = build_openclaw_digest_brief(delivery_payload)
digest_brief_output = delivery_output.with_name("digest-brief.json")
_save_json(digest_brief_output, digest_brief.model_dump(mode="json"))
keyword_index_result = persist_keyword_indexes(
delivery_payload.candidates,
for_date=delivery_payload.date,
digest_id=delivery_payload.run_id,
source="openclaw_delivery_payload",
daily_dir=DEFAULT_TERM_DAILY_DIR,
stats_path=DEFAULT_TERM_STATS_PATH,
aliases_path=DEFAULT_TERM_ALIASES_PATH,
stopwords_path=DEFAULT_TERM_STOPWORDS_PATH,
)
return delivery_payload, digest_brief_output, keyword_index_result
def _build_run_report( def _build_run_report(
@@ -233,8 +126,8 @@ def _build_run_report(
resolved_run_id: str, resolved_run_id: str,
started_at: datetime, started_at: datetime,
limit: int, limit: int,
items: list, items: list[Any],
delivered_candidates: list, delivered_candidates: list[OpenClawCandidateInput],
marked_count: int, marked_count: int,
mark_read: bool, mark_read: bool,
debug_artifacts: bool, debug_artifacts: bool,
@@ -244,7 +137,6 @@ def _build_run_report(
keyword_index_result: dict[str, Any], keyword_index_result: dict[str, Any],
item_reports: list[dict[str, Any]], item_reports: list[dict[str, Any]],
) -> dict[str, Any]: ) -> dict[str, Any]:
# 统计各状态计数,组装 run report dict
status_counts: dict[str, int] = {} status_counts: dict[str, int] = {}
for item_report in item_reports: for item_report in item_reports:
status = str(item_report["status"]) status = str(item_report["status"])
@@ -269,6 +161,12 @@ def _build_run_report(
} }
def _final_run_status(item_reports: list[dict[str, Any]]) -> str:
if any(item_report.get("status") in {"extract_failed", "summary_failed"} for item_report in item_reports):
return "partial"
return "success"
def run_freshrss_pipeline( def run_freshrss_pipeline(
*, *,
api_base_url: str | None = None, api_base_url: str | None = None,
@@ -293,13 +191,51 @@ def run_freshrss_pipeline(
delivery_date: date | None = None, delivery_date: date | None = None,
output_dir: Path | None = None, output_dir: Path | None = None,
) -> dict[str, Any]: ) -> dict[str, Any]:
# --- 阶段 1:初始化 run_id、输出路径、凭证 ---
started_at = datetime.now(tz=UTC) started_at = datetime.now(tz=UTC)
resolved_output_dir = output_dir or default_output_dir() resolved_output_dir = output_dir or default_output_dir()
run_stamp = started_at.strftime("%Y%m%d-%H%M%S") run_stamp = started_at.strftime("%Y%m%d-%H%M%S")
resolved_run_id = run_id or f"freshrss-pipeline-{run_stamp}" resolved_run_id = run_id or f"freshrss-pipeline-{run_stamp}"
resolved_delivery_date = delivery_date or datetime.now(tz=UTC).date() resolved_delivery_date = delivery_date or datetime.now(tz=UTC).date()
resolved_prompt_path = prompt or DEFAULT_PROMPT_PATH
resolved_rules_path = rules or DEFAULT_RULES_PATH
raw_output = resolved_output_dir / "raw" / "freshrss.raw.json"
delivery_output = resolved_output_dir / "candidates" / "openclaw-delivery-payload.json"
digest_brief_output = delivery_output.with_name("digest-brief.json")
report_output = resolved_output_dir / "run-report.json"
run_state_output = resolved_output_dir / "run-state.json"
items_list_output = _maybe_path(debug_artifacts, resolved_output_dir / "items" / "freshrss.items.json")
run_store = RunStore.create(
path=run_state_output,
run_id=resolved_run_id,
workflow=WORKFLOW_NAME,
run_type=RUN_TYPE,
started_at=started_at,
input_payload={
"limit": limit,
"mark_read": mark_read,
"include_read": include_read,
"debug_artifacts": debug_artifacts,
"continuation": continuation,
"stream_id": stream_id,
},
repo_root=REPO_ROOT,
)
run_store.save()
client: FreshRSSClient | None = None
auth_token: str | None = None
items: list[Any] = []
item_contexts: list[dict[str, Any]] = []
item_reports: list[dict[str, Any]] = []
delivered_candidates: list[OpenClawCandidateInput] = []
delivered_item_ids: list[str] = []
keyword_index_result: dict[str, Any] = {}
marked_count = 0
try:
run_store.start_stage(FETCH_STAGE, outputs={"output_dir": str(resolved_output_dir)})
resolved_api_base_url = _load_required_env("FRESHRSS_API_BASE_URL", api_base_url) resolved_api_base_url = _load_required_env("FRESHRSS_API_BASE_URL", api_base_url)
resolved_username = _load_required_env("FRESHRSS_USERNAME", username) resolved_username = _load_required_env("FRESHRSS_USERNAME", username)
resolved_api_password = _load_required_env("FRESHRSS_API_PASSWORD", api_password) resolved_api_password = _load_required_env("FRESHRSS_API_PASSWORD", api_password)
@@ -309,14 +245,6 @@ def run_freshrss_pipeline(
api_url=llm_api_url, api_url=llm_api_url,
) )
resolved_prompt_path = prompt or DEFAULT_PROMPT_PATH
resolved_rules_path = rules or DEFAULT_RULES_PATH
raw_output = resolved_output_dir / "raw" / "freshrss.raw.json"
delivery_output = resolved_output_dir / "candidates" / "openclaw-delivery-payload.json"
report_output = resolved_output_dir / "run-report.json"
items_list_output = _maybe_path(debug_artifacts, resolved_output_dir / "items" / "freshrss.items.json")
# --- 阶段 2:登录 FreshRSS,拉取未读条目,写原始 payload ---
client = FreshRSSClient( client = FreshRSSClient(
api_base_url=resolved_api_base_url, api_base_url=resolved_api_base_url,
username=resolved_username, username=resolved_username,
@@ -336,13 +264,14 @@ def run_freshrss_pipeline(
raise RuntimeError("FreshRSS stream response does not contain an items array.") raise RuntimeError("FreshRSS stream response does not contain an items array.")
_save_json(raw_output, payload) _save_json(raw_output, payload)
run_store.register_artifact(name="raw_output", path=raw_output, kind="json", stage=FETCH_STAGE)
items = [map_entry_to_item(entry) for entry in entries] items = [map_entry_to_item(entry) for entry in entries]
if items_list_output is not None: if items_list_output is not None:
_save_json(items_list_output, [item.model_dump(mode="json") for item in items]) _save_json(items_list_output, [item.model_dump(mode="json") for item in items])
run_store.register_artifact(name="items_output", path=items_list_output, kind="json", stage=FETCH_STAGE)
loaded_rules = load_filter_rules(resolved_rules_path) loaded_rules = load_filter_rules(resolved_rules_path)
# context 优先使用直接传入的 dict,其次读取 context_path 文件,两者均缺失则使用空 context
if context is not None: if context is not None:
filter_context = FilterContext.model_validate(context) filter_context = FilterContext.model_validate(context)
elif context_path is not None: elif context_path is not None:
@@ -350,46 +279,276 @@ def run_freshrss_pipeline(
else: else:
filter_context = FilterContext() filter_context = FilterContext()
# --- 阶段 3:逐条处理(提取 -> LLM 摘要 -> 规则过滤 -> 候选构建) --- run_store.finish_stage(
delivered_candidates: list[OpenClawCandidateInput] = [] FETCH_STAGE,
delivered_item_ids: list[str] = [] outputs={
item_reports: list[dict[str, Any]] = [] "pulled_count": len(items),
"raw_output": str(raw_output),
"items_output": str(items_list_output) if items_list_output else None,
},
)
run_store.start_stage(
EXTRACT_STAGE,
outputs={
"expected_items": len(items),
"completed_items": 0,
"success_count": 0,
"failed_count": 0,
},
)
extracted_success_count = 0
extracted_failed_count = 0
for index, item in enumerate(items, start=1): for index, item in enumerate(items, start=1):
result = _process_item( item_context = _build_item_context(
index=index, index=index,
item=item, item=item,
resolved_output_dir=resolved_output_dir, resolved_output_dir=resolved_output_dir,
resolved_prompt_path=resolved_prompt_path,
resolved_run_id=resolved_run_id,
debug_artifacts=debug_artifacts, debug_artifacts=debug_artifacts,
loaded_rules=loaded_rules, )
filter_context=filter_context, item_contexts.append(item_context)
item_report = item_context["item_report"]
item_reports.append(item_report)
item_path = item_context["item_path"]
if item_path is not None:
_save_json(item_path, item.model_dump(mode="json"))
extraction = extract_content(ExtractionInput(item=item))
extracted_payload = extraction.model_dump(mode="json")
item_context["extraction"] = extraction
item_context["extracted_payload"] = extracted_payload
_save_json(item_context["extracted_path"], extracted_payload)
run_store.register_artifact(name="extracted_dir", path=resolved_output_dir / "extracted", kind="directory", stage=EXTRACT_STAGE)
if not extraction.success or extraction.article is None:
item_report["status"] = "extract_failed"
item_report["error"] = extraction.error.model_dump(mode="json") if extraction.error else None
extracted_failed_count += 1
else:
item_report["status"] = "extracted"
extracted_success_count += 1
run_store.update_stage(
EXTRACT_STAGE,
outputs={
"expected_items": len(items),
"completed_items": extracted_success_count + extracted_failed_count,
"success_count": extracted_success_count,
"failed_count": extracted_failed_count,
},
)
run_store.finish_stage(
EXTRACT_STAGE,
outputs={
"expected_items": len(items),
"completed_items": extracted_success_count + extracted_failed_count,
"success_count": extracted_success_count,
"failed_count": extracted_failed_count,
"extracted_dir": str(resolved_output_dir / "extracted"),
},
)
run_store.start_stage(
SUMMARY_STAGE,
outputs={
"expected_items": extracted_success_count,
"completed_items": 0,
"success_count": 0,
"failed_count": 0,
},
)
summary_success_count = 0
summary_failed_count = 0
summary_candidates = [ctx for ctx in item_contexts if ctx["extraction"] is not None and ctx["extraction"].success]
for item_context in summary_candidates:
item_report = item_context["item_report"]
summary_exit_code, summary_payload, summary_report = run_loop_payload(
extracted_payload=item_context["extracted_payload"],
prompt_path=resolved_prompt_path,
output_path=item_context["summary_output"],
max_retries=max_retries, max_retries=max_retries,
timeout_seconds=timeout_seconds, timeout_seconds=timeout_seconds,
resolved_llm_api_key=resolved_llm_api_key, api_key=resolved_llm_api_key,
resolved_llm_model=resolved_llm_model, model=resolved_llm_model,
resolved_llm_api_url=resolved_llm_api_url, api_url=resolved_llm_api_url,
) )
if result.get("status") == "delivered": if summary_exit_code != 0 or summary_payload is None:
delivered_candidates.append(result.pop("_candidate")) item_report["status"] = "summary_failed"
external_id = result.pop("_external_id", None) if summary_report is not None:
if external_id: item_report["summary_errors"] = summary_report.errors
delivered_item_ids.append(external_id) summary_failed_count += 1
item_reports.append(result) else:
item_context["summary_payload"] = summary_payload
item_report["status"] = "summarized"
summary_success_count += 1
# --- 阶段 4+5:构建 delivery payload 并持久化词元索引 --- run_store.update_stage(
delivery_payload, digest_brief_output, keyword_index_result = _build_and_persist_delivery( SUMMARY_STAGE,
delivered_candidates=delivered_candidates, outputs={
resolved_run_id=resolved_run_id, "expected_items": extracted_success_count,
resolved_delivery_date=resolved_delivery_date, "completed_items": summary_success_count + summary_failed_count,
delivery_output=delivery_output, "success_count": summary_success_count,
"failed_count": summary_failed_count,
},
) )
# --- 阶段 6:标记已读,汇总报告,返回结果 --- if debug_artifacts and (resolved_output_dir / "summary").exists():
marked_count = 0 run_store.register_artifact(name="summary_dir", path=resolved_output_dir / "summary", kind="directory", stage=SUMMARY_STAGE)
run_store.finish_stage(
SUMMARY_STAGE,
outputs={
"expected_items": extracted_success_count,
"completed_items": summary_success_count + summary_failed_count,
"success_count": summary_success_count,
"failed_count": summary_failed_count,
},
)
run_store.start_stage(
FILTER_STAGE,
outputs={
"expected_items": summary_success_count,
"completed_items": 0,
"candidate_count": 0,
"keep_count": 0,
"review_count": 0,
"drop_count": 0,
},
)
filter_completed_count = 0
keep_count = 0
review_count = 0
drop_count = 0
for item_context in [ctx for ctx in item_contexts if ctx["summary_payload"] is not None]:
item = item_context["item"]
item_report = item_context["item_report"]
extraction = item_context["extraction"]
summary = LlmSummaryResult.model_validate(item_context["summary_payload"])
decision = evaluate_filter_rules(
FilterInput(item=item, article=extraction.article, summary=summary, context=filter_context),
loaded_rules,
)
filter_path = item_context["filter_path"]
if filter_path is not None:
_save_json(filter_path, decision.model_dump(mode="json"))
record = build_article_candidate_record(
summary=summary,
article=extraction.article,
filter_result=decision,
item=item,
source_refs=CandidateSourceRefs(
item_path=str(item_context["item_path"]) if item_context["item_path"] else None,
extracted_path=str(item_context["extracted_path"]),
summary_path=str(item_context["summary_output"]) if item_context["summary_output"] else None,
filter_path=str(filter_path) if filter_path else None,
),
metadata=CandidateMetadata(
generated_at=datetime.now(tz=UTC),
producer="run_freshrss_pipeline",
run_id=resolved_run_id,
),
)
openclaw_input = build_openclaw_candidate_input(record)
item_context["candidate"] = openclaw_input
if item_context["record_path"] is not None:
_save_json(item_context["record_path"], record.model_dump(mode="json"))
if item_context["openclaw_path"] is not None:
_save_json(item_context["openclaw_path"], openclaw_input.model_dump(mode="json"))
item_report["status"] = "delivered"
item_report["selection_decision"] = decision.decision
item_report["candidate_id"] = openclaw_input.candidate_id
delivered_candidates.append(openclaw_input)
if item.external_id:
delivered_item_ids.append(item.external_id)
if decision.decision == "keep":
keep_count += 1
elif decision.decision == "review":
review_count += 1
elif decision.decision == "drop":
drop_count += 1
filter_completed_count += 1
run_store.update_stage(
FILTER_STAGE,
outputs={
"expected_items": summary_success_count,
"completed_items": filter_completed_count,
"candidate_count": len(delivered_candidates),
"keep_count": keep_count,
"review_count": review_count,
"drop_count": drop_count,
},
)
if debug_artifacts and (resolved_output_dir / "candidates").exists():
run_store.register_artifact(name="candidate_dir", path=resolved_output_dir / "candidates", kind="directory", stage=FILTER_STAGE)
run_store.finish_stage(
FILTER_STAGE,
outputs={
"expected_items": summary_success_count,
"completed_items": filter_completed_count,
"candidate_count": len(delivered_candidates),
"keep_count": keep_count,
"review_count": review_count,
"drop_count": drop_count,
},
)
run_store.start_stage(DELIVERY_STAGE, outputs={"candidate_count": len(delivered_candidates)})
delivered_candidates.sort(key=lambda candidate: candidate.digest_rank, reverse=True)
delivery_payload = build_openclaw_delivery_payload(
delivered_candidates,
run_id=resolved_run_id,
for_date=resolved_delivery_date,
)
_save_json(delivery_output, delivery_payload.model_dump(mode="json"))
run_store.register_artifact(name="delivery_payload", path=delivery_output, kind="json", stage=DELIVERY_STAGE)
digest_brief = build_openclaw_digest_brief(delivery_payload)
_save_json(digest_brief_output, digest_brief.model_dump(mode="json"))
run_store.register_artifact(name="digest_brief", path=digest_brief_output, kind="json", stage=DELIVERY_STAGE)
keyword_index_result = persist_keyword_indexes(
delivery_payload.candidates,
for_date=delivery_payload.date,
digest_id=delivery_payload.run_id,
source="openclaw_delivery_payload",
daily_dir=DEFAULT_TERM_DAILY_DIR,
stats_path=DEFAULT_TERM_STATS_PATH,
aliases_path=DEFAULT_TERM_ALIASES_PATH,
stopwords_path=DEFAULT_TERM_STOPWORDS_PATH,
)
run_store.register_artifact(
name="keyword_daily_index",
path=Path(str(keyword_index_result["daily_output"])),
kind="json",
stage=DELIVERY_STAGE,
)
run_store.register_artifact(
name="keyword_stats_index",
path=Path(str(keyword_index_result["stats_output"])),
kind="json",
stage=DELIVERY_STAGE,
)
run_store.finish_stage(
DELIVERY_STAGE,
outputs={
"candidate_count": len(delivered_candidates),
"delivery_output": str(delivery_output),
"digest_brief_output": str(digest_brief_output),
"keyword_daily_output": str(keyword_index_result["daily_output"]),
"keyword_stats_output": str(keyword_index_result["stats_output"]),
},
)
run_store.start_stage(REPORT_STAGE, outputs={"mark_read_requested": mark_read})
if mark_read and delivered_item_ids: if mark_read and delivered_item_ids:
# 仅标记成功投递(delivered)的 item;drop/review 的 item 保持未读状态
client.mark_items_as_read(auth_token=auth_token, item_ids=delivered_item_ids) client.mark_items_as_read(auth_token=auth_token, item_ids=delivered_item_ids)
marked_count = len({item_id for item_id in delivered_item_ids if item_id}) marked_count = len({item_id for item_id in delivered_item_ids if item_id})
@@ -409,7 +568,16 @@ def run_freshrss_pipeline(
item_reports=item_reports, item_reports=item_reports,
) )
_save_json(report_output, report) _save_json(report_output, report)
run_store.register_artifact(name="run_report", path=report_output, kind="json", stage=REPORT_STAGE)
run_store.finish_stage(
REPORT_STAGE,
outputs={
"marked_read_count": marked_count,
"report_output": str(report_output),
},
)
run_store.finish_run(status=_final_run_status(item_reports))
return { return {
"run_id": resolved_run_id, "run_id": resolved_run_id,
"output_dir": str(resolved_output_dir), "output_dir": str(resolved_output_dir),
@@ -426,6 +594,10 @@ def run_freshrss_pipeline(
"delivery_payload": delivery_payload.model_dump(mode="json"), "delivery_payload": delivery_payload.model_dump(mode="json"),
"items": item_reports, "items": item_reports,
} }
except Exception as error:
failed_stage = run_store.state.current_stage or FETCH_STAGE
run_store.fail_stage(failed_stage, error=error)
raise
def read_delivery_payload(path: Path) -> OpenClawDeliveryPayload: def read_delivery_payload(path: Path) -> OpenClawDeliveryPayload: