Compare commits
4
Commits
3a85d47f00
...
4a02894c43
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
4a02894c43 | ||
|
|
d9173fbe0f | ||
|
|
7563aa8fea | ||
|
|
72a6853c03 |
@@ -1,83 +1,217 @@
|
||||
# TODO
|
||||
# TODO - Reader MCP 正式化
|
||||
|
||||
## 当前状态
|
||||
> 本文件用于架构与 Codex 协作同步。
|
||||
>
|
||||
> 规则:
|
||||
> - `TODO` = 未开始
|
||||
> - `DOING` = 正在进行
|
||||
> - `DONE` = 已完成
|
||||
> - 每次只允许一个最高优先级主任务处于 `DOING`
|
||||
|
||||
项目当前已经进入“可交付给 OpenClaw 调用”的阶段。
|
||||
## 0. 协作约束
|
||||
|
||||
当前主链路:
|
||||
开始编码前必须阅读:
|
||||
|
||||
`FreshRSS 未读 -> RSS 内容提取 -> LLM 总结 -> 规则过滤 -> OpenClaw delivery payload`
|
||||
|
||||
当前已经完成:
|
||||
|
||||
- [x] FreshRSS `greader` API 接入
|
||||
- [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`
|
||||
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`
|
||||
5. `docs/openclaw/openclaw-handoff.md`
|
||||
|
||||
---
|
||||
|
||||
## P0 - 交接前后最优先
|
||||
## 1. 当前主任务
|
||||
|
||||
- [x] 为 OpenClaw 补齐交接文档
|
||||
- [x] 将 MCP 工具作为统一生产入口
|
||||
- [x] 将默认输出收敛为最小必要文件
|
||||
- [ ] 设计 OpenClaw webhook / delivery payload 的主动推送方式
|
||||
- [ ] 明确 OpenClaw 侧如何注册和启动本 MCP 服务
|
||||
### [DONE][P0] 建立 run-state 运行态基础设施
|
||||
|
||||
目标:
|
||||
- 给 freshrss pipeline 引入正式 run state
|
||||
- 即使失败或中断,也能留下明确运行真相
|
||||
|
||||
要求:
|
||||
- 新增 `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. 后续任务队列
|
||||
|
||||
- [ ] 设计“人工确认后再沉淀知识库”的状态流转
|
||||
- [ ] 收敛 `paywall` 误判规则,降低中文文本误报
|
||||
- [ ] 细化过滤规则并引入更多个性化上下文
|
||||
- [ ] 将 `keyword-cleanup-review` skill 接入周期性执行流程,产出别名/停用词/兴趣词建议
|
||||
- [ ] 增加清洗前后效果对比报告,验证配置调整是否真的改善过滤质量
|
||||
- [ ] 增加批量 run 的保留策略与历史清理策略
|
||||
- [ ] 为 OpenClaw 补一份更正式的 MCP 调用示例和接线说明
|
||||
### [DONE][P1] 增加 MCP 状态查询接口 `get_run_status`
|
||||
|
||||
目标:
|
||||
- 可通过 MCP 查询 run 状态
|
||||
|
||||
要求:
|
||||
- 输入 `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` 的批处理能力
|
||||
- [ ] 让 OpenClaw 聚合候选内容并生成日级摘要
|
||||
- [ ] 将日级摘要写入知识库,并同步生成面向用户的日报消息
|
||||
- [ ] 支持更多 `content_kind`
|
||||
- [ ] 增加提取缓存、重试和更细粒度日志
|
||||
- [ ] 整理历史 rerun 目录与调试产物保留策略
|
||||
目标:
|
||||
- 查看近期 runs
|
||||
|
||||
要求:
|
||||
- 支持按 workflow / status / latest_n 过滤
|
||||
|
||||
完成情况:
|
||||
- 已通过 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. 设计“人工确认后再沉淀知识库”的状态流转
|
||||
3. 收敛规则误判,尤其是 `paywall` 相关启发式
|
||||
完成情况:
|
||||
- 已通过 MCP 暴露 `list_run_artifacts`
|
||||
- 对新 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`
|
||||
- `docs/openclaw/openclaw-candidate-input-field-spec.md`
|
||||
- `docs/openclaw/openclaw-delivery-payload-spec.md`
|
||||
- `docs/design/daily-keyword-index-design.md`
|
||||
- `skills/keyword-cleanup-review/SKILL.md`
|
||||
- `docs/current/context-reset-brief.md`
|
||||
目标:
|
||||
- 按 run_id 读取 delivery payload
|
||||
|
||||
完成情况:
|
||||
- 已通过 MCP 暴露 `get_delivery_payload`
|
||||
- 查询优先复用 `run-state.json` 已注册 artifacts,其次回退标准产物路径、`run-report.json` 引用和 run 目录扫描
|
||||
- 返回补充了 `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
|
||||
- 中间产物复用
|
||||
- 分阶段补跑
|
||||
@@ -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 为下游编排层的工作流服务。**
|
||||
@@ -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 工作流服务”,再做更多工具;不要反过来先堆接口名。
|
||||
@@ -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",
|
||||
]
|
||||
@@ -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"))
|
||||
@@ -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)
|
||||
@@ -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)
|
||||
@@ -15,6 +15,11 @@ from summary_mcp.models.filtering import FilterContext, FilterInput
|
||||
from summary_mcp.models.item import Item
|
||||
from summary_mcp.models.llm_result import LlmSummaryResult
|
||||
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.article_summary import ArticleSummaryConfig, summarize_selected_articles
|
||||
|
||||
@@ -119,6 +124,40 @@ def run_freshrss_openclaw_pipeline(
|
||||
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()
|
||||
def generate_article_summaries(
|
||||
*,
|
||||
|
||||
@@ -1,14 +1,8 @@
|
||||
from __future__ import annotations
|
||||
|
||||
# FreshRSS 全链路管道:拉取未读条目 -> 内容提取 -> LLM 摘要 -> 规则过滤 ->
|
||||
# 构建 OpenClaw delivery payload -> 写盘 -> 词元统计 -> 标记已读。
|
||||
# 生产入口:run_freshrss_pipeline(),由 MCP 工具 run_freshrss_openclaw_pipeline 调用。
|
||||
|
||||
import json
|
||||
import os
|
||||
from datetime import date, datetime, timezone
|
||||
|
||||
UTC = timezone.utc
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
@@ -32,8 +26,10 @@ from summary_mcp.models.openclaw_delivery import (
|
||||
build_openclaw_digest_brief,
|
||||
)
|
||||
from summary_mcp.models.summary_io import ExtractionInput
|
||||
from summary_mcp.runtime import RunStore
|
||||
|
||||
|
||||
UTC = timezone.utc
|
||||
REPO_ROOT = Path(__file__).resolve().parents[3]
|
||||
OUTPUT_ROOT = REPO_ROOT / "outputs"
|
||||
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_DAILY_DIR = DATA_ROOT / "daily"
|
||||
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:
|
||||
@@ -56,7 +60,6 @@ def _load_json(path: Path) -> dict[str, Any]:
|
||||
|
||||
|
||||
def _load_required_env(name: str, value: str | None) -> str:
|
||||
# 优先使用传入的 value,否则读取同名环境变量;两者均缺失时抛出 RuntimeError
|
||||
if value:
|
||||
return value
|
||||
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
|
||||
|
||||
|
||||
def _process_item(
|
||||
*,
|
||||
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,不写盘)
|
||||
def _build_item_context(*, index: int, item: Any, resolved_output_dir: Path, debug_artifacts: bool) -> dict[str, Any]:
|
||||
item_key = f"item-{index:02d}"
|
||||
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"
|
||||
@@ -102,10 +87,6 @@ def _process_item(
|
||||
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")
|
||||
|
||||
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_key": item_key,
|
||||
"item_id": item.item_id,
|
||||
@@ -117,115 +98,27 @@ def _process_item(
|
||||
if debug_artifacts:
|
||||
item_report["paths"] = {
|
||||
"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,
|
||||
"filter": str(filter_path) if filter_path else None,
|
||||
"article_candidate": str(record_path) if record_path else None,
|
||||
"openclaw_candidate": str(openclaw_path) if openclaw_path else None,
|
||||
}
|
||||
|
||||
# 步骤 3:内容提取(RSS 内联内容 or 回源抓取);extracted.json 始终写盘
|
||||
extraction = extract_content(ExtractionInput(item=item))
|
||||
extracted_payload = extraction.model_dump(mode="json")
|
||||
if extracted_path is not None:
|
||||
_save_json(extracted_path, extracted_payload)
|
||||
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
|
||||
return item_report
|
||||
|
||||
# 步骤 4:LLM 摘要循环,失败时最多重试 max_retries 次
|
||||
item_report["status"] = "extracted"
|
||||
summary_exit_code, summary_payload, summary_report = run_loop_payload(
|
||||
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
|
||||
return {
|
||||
"item": item,
|
||||
"item_key": item_key,
|
||||
"item_path": item_path,
|
||||
"extracted_path": extracted_path,
|
||||
"summary_output": summary_output,
|
||||
"filter_path": filter_path,
|
||||
"record_path": record_path,
|
||||
"openclaw_path": openclaw_path,
|
||||
"item_report": item_report,
|
||||
"extraction": None,
|
||||
"extracted_payload": None,
|
||||
"summary_payload": None,
|
||||
}
|
||||
|
||||
|
||||
def _build_run_report(
|
||||
@@ -233,8 +126,8 @@ def _build_run_report(
|
||||
resolved_run_id: str,
|
||||
started_at: datetime,
|
||||
limit: int,
|
||||
items: list,
|
||||
delivered_candidates: list,
|
||||
items: list[Any],
|
||||
delivered_candidates: list[OpenClawCandidateInput],
|
||||
marked_count: int,
|
||||
mark_read: bool,
|
||||
debug_artifacts: bool,
|
||||
@@ -244,7 +137,6 @@ def _build_run_report(
|
||||
keyword_index_result: dict[str, Any],
|
||||
item_reports: list[dict[str, Any]],
|
||||
) -> dict[str, Any]:
|
||||
# 统计各状态计数,组装 run report dict
|
||||
status_counts: dict[str, int] = {}
|
||||
for item_report in item_reports:
|
||||
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(
|
||||
*,
|
||||
api_base_url: str | None = None,
|
||||
@@ -293,139 +191,413 @@ def run_freshrss_pipeline(
|
||||
delivery_date: date | None = None,
|
||||
output_dir: Path | None = None,
|
||||
) -> dict[str, Any]:
|
||||
# --- 阶段 1:初始化 run_id、输出路径、凭证 ---
|
||||
started_at = datetime.now(tz=UTC)
|
||||
resolved_output_dir = output_dir or default_output_dir()
|
||||
run_stamp = started_at.strftime("%Y%m%d-%H%M%S")
|
||||
resolved_run_id = run_id or f"freshrss-pipeline-{run_stamp}"
|
||||
resolved_delivery_date = delivery_date or datetime.now(tz=UTC).date()
|
||||
|
||||
resolved_api_base_url = _load_required_env("FRESHRSS_API_BASE_URL", api_base_url)
|
||||
resolved_username = _load_required_env("FRESHRSS_USERNAME", username)
|
||||
resolved_api_password = _load_required_env("FRESHRSS_API_PASSWORD", api_password)
|
||||
resolved_llm_api_key, resolved_llm_model, resolved_llm_api_url = resolve_llm_settings(
|
||||
api_key=llm_api_key,
|
||||
model=llm_model,
|
||||
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"
|
||||
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")
|
||||
|
||||
# --- 阶段 2:登录 FreshRSS,拉取未读条目,写原始 payload ---
|
||||
client = FreshRSSClient(
|
||||
api_base_url=resolved_api_base_url,
|
||||
username=resolved_username,
|
||||
api_password=resolved_api_password,
|
||||
timeout_seconds=timeout_seconds,
|
||||
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,
|
||||
)
|
||||
auth_token = client.client_login()
|
||||
payload = client.fetch_stream_contents(
|
||||
auth_token=auth_token,
|
||||
stream_id=stream_id,
|
||||
limit=limit,
|
||||
continuation=continuation,
|
||||
exclude_targets=[] if include_read else [READ_TAG],
|
||||
)
|
||||
entries = payload.get("items")
|
||||
if not isinstance(entries, list):
|
||||
raise RuntimeError("FreshRSS stream response does not contain an items array.")
|
||||
run_store.save()
|
||||
|
||||
_save_json(raw_output, payload)
|
||||
|
||||
items = [map_entry_to_item(entry) for entry in entries]
|
||||
if items_list_output is not None:
|
||||
_save_json(items_list_output, [item.model_dump(mode="json") for item in items])
|
||||
|
||||
loaded_rules = load_filter_rules(resolved_rules_path)
|
||||
# context 优先使用直接传入的 dict,其次读取 context_path 文件,两者均缺失则使用空 context
|
||||
if context is not None:
|
||||
filter_context = FilterContext.model_validate(context)
|
||||
elif context_path is not None:
|
||||
filter_context = FilterContext.model_validate(_load_json(context_path))
|
||||
else:
|
||||
filter_context = FilterContext()
|
||||
|
||||
# --- 阶段 3:逐条处理(提取 -> LLM 摘要 -> 规则过滤 -> 候选构建) ---
|
||||
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] = []
|
||||
item_reports: list[dict[str, Any]] = []
|
||||
|
||||
for index, item in enumerate(items, start=1):
|
||||
result = _process_item(
|
||||
index=index,
|
||||
item=item,
|
||||
resolved_output_dir=resolved_output_dir,
|
||||
resolved_prompt_path=resolved_prompt_path,
|
||||
resolved_run_id=resolved_run_id,
|
||||
debug_artifacts=debug_artifacts,
|
||||
loaded_rules=loaded_rules,
|
||||
filter_context=filter_context,
|
||||
max_retries=max_retries,
|
||||
timeout_seconds=timeout_seconds,
|
||||
resolved_llm_api_key=resolved_llm_api_key,
|
||||
resolved_llm_model=resolved_llm_model,
|
||||
resolved_llm_api_url=resolved_llm_api_url,
|
||||
)
|
||||
if result.get("status") == "delivered":
|
||||
delivered_candidates.append(result.pop("_candidate"))
|
||||
external_id = result.pop("_external_id", None)
|
||||
if external_id:
|
||||
delivered_item_ids.append(external_id)
|
||||
item_reports.append(result)
|
||||
|
||||
# --- 阶段 4+5:构建 delivery payload 并持久化词元索引 ---
|
||||
delivery_payload, digest_brief_output, keyword_index_result = _build_and_persist_delivery(
|
||||
delivered_candidates=delivered_candidates,
|
||||
resolved_run_id=resolved_run_id,
|
||||
resolved_delivery_date=resolved_delivery_date,
|
||||
delivery_output=delivery_output,
|
||||
)
|
||||
|
||||
# --- 阶段 6:标记已读,汇总报告,返回结果 ---
|
||||
keyword_index_result: dict[str, Any] = {}
|
||||
marked_count = 0
|
||||
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)
|
||||
marked_count = len({item_id for item_id in delivered_item_ids if item_id})
|
||||
|
||||
report = _build_run_report(
|
||||
resolved_run_id=resolved_run_id,
|
||||
started_at=started_at,
|
||||
limit=limit,
|
||||
items=items,
|
||||
delivered_candidates=delivered_candidates,
|
||||
marked_count=marked_count,
|
||||
mark_read=mark_read,
|
||||
debug_artifacts=debug_artifacts,
|
||||
raw_output=raw_output,
|
||||
delivery_output=delivery_output,
|
||||
digest_brief_output=digest_brief_output,
|
||||
keyword_index_result=keyword_index_result,
|
||||
item_reports=item_reports,
|
||||
)
|
||||
_save_json(report_output, report)
|
||||
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_username = _load_required_env("FRESHRSS_USERNAME", username)
|
||||
resolved_api_password = _load_required_env("FRESHRSS_API_PASSWORD", api_password)
|
||||
resolved_llm_api_key, resolved_llm_model, resolved_llm_api_url = resolve_llm_settings(
|
||||
api_key=llm_api_key,
|
||||
model=llm_model,
|
||||
api_url=llm_api_url,
|
||||
)
|
||||
|
||||
return {
|
||||
"run_id": resolved_run_id,
|
||||
"output_dir": str(resolved_output_dir),
|
||||
"raw_output": str(raw_output),
|
||||
"delivery_output": str(delivery_output),
|
||||
"digest_brief_output": str(digest_brief_output),
|
||||
"report_output": str(report_output),
|
||||
"keyword_index": keyword_index_result,
|
||||
"pulled_count": len(items),
|
||||
"delivered_count": len(delivered_candidates),
|
||||
"marked_read_count": marked_count,
|
||||
"status_counts": report["status_counts"],
|
||||
"debug_artifacts": debug_artifacts,
|
||||
"delivery_payload": delivery_payload.model_dump(mode="json"),
|
||||
"items": item_reports,
|
||||
}
|
||||
client = FreshRSSClient(
|
||||
api_base_url=resolved_api_base_url,
|
||||
username=resolved_username,
|
||||
api_password=resolved_api_password,
|
||||
timeout_seconds=timeout_seconds,
|
||||
)
|
||||
auth_token = client.client_login()
|
||||
payload = client.fetch_stream_contents(
|
||||
auth_token=auth_token,
|
||||
stream_id=stream_id,
|
||||
limit=limit,
|
||||
continuation=continuation,
|
||||
exclude_targets=[] if include_read else [READ_TAG],
|
||||
)
|
||||
entries = payload.get("items")
|
||||
if not isinstance(entries, list):
|
||||
raise RuntimeError("FreshRSS stream response does not contain an items array.")
|
||||
|
||||
_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]
|
||||
if items_list_output is not None:
|
||||
_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)
|
||||
if context is not None:
|
||||
filter_context = FilterContext.model_validate(context)
|
||||
elif context_path is not None:
|
||||
filter_context = FilterContext.model_validate(_load_json(context_path))
|
||||
else:
|
||||
filter_context = FilterContext()
|
||||
|
||||
run_store.finish_stage(
|
||||
FETCH_STAGE,
|
||||
outputs={
|
||||
"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):
|
||||
item_context = _build_item_context(
|
||||
index=index,
|
||||
item=item,
|
||||
resolved_output_dir=resolved_output_dir,
|
||||
debug_artifacts=debug_artifacts,
|
||||
)
|
||||
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,
|
||||
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
|
||||
summary_failed_count += 1
|
||||
else:
|
||||
item_context["summary_payload"] = summary_payload
|
||||
item_report["status"] = "summarized"
|
||||
summary_success_count += 1
|
||||
|
||||
run_store.update_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,
|
||||
},
|
||||
)
|
||||
|
||||
if debug_artifacts and (resolved_output_dir / "summary").exists():
|
||||
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:
|
||||
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})
|
||||
|
||||
report = _build_run_report(
|
||||
resolved_run_id=resolved_run_id,
|
||||
started_at=started_at,
|
||||
limit=limit,
|
||||
items=items,
|
||||
delivered_candidates=delivered_candidates,
|
||||
marked_count=marked_count,
|
||||
mark_read=mark_read,
|
||||
debug_artifacts=debug_artifacts,
|
||||
raw_output=raw_output,
|
||||
delivery_output=delivery_output,
|
||||
digest_brief_output=digest_brief_output,
|
||||
keyword_index_result=keyword_index_result,
|
||||
item_reports=item_reports,
|
||||
)
|
||||
_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 {
|
||||
"run_id": resolved_run_id,
|
||||
"output_dir": str(resolved_output_dir),
|
||||
"raw_output": str(raw_output),
|
||||
"delivery_output": str(delivery_output),
|
||||
"digest_brief_output": str(digest_brief_output),
|
||||
"report_output": str(report_output),
|
||||
"keyword_index": keyword_index_result,
|
||||
"pulled_count": len(items),
|
||||
"delivered_count": len(delivered_candidates),
|
||||
"marked_read_count": marked_count,
|
||||
"status_counts": report["status_counts"],
|
||||
"debug_artifacts": debug_artifacts,
|
||||
"delivery_payload": delivery_payload.model_dump(mode="json"),
|
||||
"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:
|
||||
|
||||
Reference in New Issue
Block a user