Compare commits

..
4 Commits
9 changed files with 1563 additions and 27 deletions
+50 -13
View File
@@ -1,6 +1,6 @@
# 内容提取 MCP # Reader MCP Workflow Service
这是一个基于 Python 的 MCP 项目骨架,用于完成文章内容提取、结构化摘要校验、确定性过滤,以及 Markdown 输出落盘。 reader 当前已经收口为面向 OpenClaw 的 MCP workflow service。正式能力边界以 FreshRSS 日报工作流为准:启动 run、写入 `run-state.json`、查询运行状态、读取结构化结果,以及最小可用的 `resume_run`。CLI 仍保留,但定位为 debug / fallback,而不是正式集成入口。
## 运行 ## 运行
@@ -9,14 +9,50 @@ pip install -e .
summary-mcp summary-mcp
``` ```
服务当前暴露 5 个工具: 服务当前暴露 11 个工具。
正式 workflow service 相关工具:
- `run_freshrss_openclaw_pipeline`
- `get_run_status`
- `list_runs`
- `list_run_artifacts`
- `get_delivery_payload`
- `get_run_report`
- `resume_run`
单步处理 / 调试相关工具:
- `extract_url_content` - `extract_url_content`
- `extract_item_content` - `extract_item_content`
- `filter_summary_result` - `filter_summary_result`
- `run_freshrss_openclaw_pipeline`
- `generate_article_summaries` - `generate_article_summaries`
## 正式能力边界
- 当前正式 workflow 只有 `freshrss_daily_digest`
- 每次 FreshRSS 主流水线 run 都会在 `outputs/freshrss/rerun/<run_dir>/run-state.json` 落地运行真相
- OpenClaw 正式读取结果应优先使用 `get_delivery_payload` 与 `get_run_report`,而不是自己拼输出目录路径
- `digest-brief.json` 当前会随主流水线产出,但还没有独立的 MCP 读取工具;如需定位它,应通过 `list_run_artifacts` 或 `get_run_report` 返回的信息发现
- `run_freshrss_openclaw_pipeline` 与 `resume_run` 当前都是同步 MCP 调用;仓库里还没有后台队列 / worker / 异步任务管理
## OpenClaw 推荐调用路径
1. 调用 `run_freshrss_openclaw_pipeline` 启动正式日报 run,并保存返回的 `run_id`
2. 后续所有状态判断都基于 `get_run_status(run_id)` 或 `list_runs(...)`
3. 需要看产物列表时用 `list_run_artifacts(run_id)`,不要在 OpenClaw 里硬编码 `outputs/freshrss/rerun/...`
4. 需要消费正式结果时优先用 `get_delivery_payload(run_id)` 与 `get_run_report(run_id)`
5. 仅当 `resume_run` 的最小恢复范围满足时,才对失败 run 调用 `resume_run(run_id)`;否则应重启一个新 run
## `resume_run` 当前最小范围
- 只支持带有效 `run-state.json` 的 run
- 只支持 workflow `freshrss_daily_digest`
- 恢复时继续沿用原 `run_id`,不会新建 retry run
- 当前支持的恢复起点只有:`generate_summaries`、`apply_filters`、`build_delivery_payload`、`write_run_report`
- 当前明确不支持从 `fetch_feed`、`extract_articles` 恢复;这类失败应新开 run
- 恢复前会校验关键中间产物是否齐备,缺失时直接返回不可恢复,而不会自动回退到更早 stage
## 单篇文章总结后处理(可选使用独立 LLM) ## 单篇文章总结后处理(可选使用独立 LLM)
### 生产环境推荐输入 ### 生产环境推荐输入
@@ -34,7 +70,7 @@ summary-mcp
相关能力: 相关能力:
- `article-summary` MCP 工具(基于已有 extracted payload 做单篇总结) - `generate_article_summaries` MCP 工具(基于已有 extracted payload 做单篇总结)
- `scripts/run_article_summaries.py` CLI 辅助脚本 - `scripts/run_article_summaries.py` CLI 辅助脚本
## 校验 LLM 摘要结果 ## 校验 LLM 摘要结果
@@ -111,13 +147,14 @@ python scripts/run_freshrss_pipeline.py ^
--mark-read --mark-read
``` ```
这是当前推荐的**正式生产入口**。默认只写出这些产物: 这条 CLI 与 MCP `run_freshrss_openclaw_pipeline` 共用同一条主流水线逻辑,但正式生产集成应优先走 MCP;CLI 仅用于本地 debug / fallback。默认会写出这些产物:
- `outputs/freshrss/rerun/<timestamp>/raw/freshrss.raw.json` - `outputs/freshrss/rerun/<run_dir>/run-state.json`
- `outputs/freshrss/rerun/<timestamp>/candidates/openclaw-delivery-payload.json` - `outputs/freshrss/rerun/<run_dir>/raw/freshrss.raw.json`
- `outputs/freshrss/rerun/<timestamp>/candidates/digest-brief.json`(给 OpenClaw 生成 public digest 用的轻量输入,仅包含 `keep` 候选) - `outputs/freshrss/rerun/<run_dir>/candidates/openclaw-delivery-payload.json`
- `outputs/freshrss/rerun/<timestamp>/run-report.json` - `outputs/freshrss/rerun/<run_dir>/candidates/digest-brief.json`(给 OpenClaw 生成 public digest 用的轻量输入,仅包含 `keep` 候选)
- `outputs/freshrss/rerun/<timestamp>/extracted/item-XX.extracted.json`(每篇一份) - `outputs/freshrss/rerun/<run_dir>/run-report.json`
- `outputs/freshrss/rerun/<run_dir>/extracted/item-XX.extracted.json`(每篇一份)
同时还会更新每日关键词索引运行数据: 同时还会更新每日关键词索引运行数据:
@@ -127,8 +164,8 @@ python scripts/run_freshrss_pipeline.py ^
主流水线默认**不会**产出批量级的 `freshrss.extracted.json`。 主流水线默认**不会**产出批量级的 `freshrss.extracted.json`。
如果你需要更多逐条中间产物,例如标准化 items、摘要结果、过滤决策、candidate record、candidate input,可以加 `--debug-artifacts`。 如果你需要更多逐条中间产物,例如标准化 items、摘要结果、过滤决策、candidate record、candidate input,可以加 `--debug-artifacts`。
当 OpenClaw 接入这个 MCP 服务后,应直接调用 `run_freshrss_openclaw_pipeline` 来获得同样行为。 当 OpenClaw 接入这个 MCP 服务后,应以 `run_id` 作为稳定句柄:先调用 `run_freshrss_openclaw_pipeline`,再通过 `get_run_status` / `get_delivery_payload` / `get_run_report` 读取状态与结果,而不是直接拼接目录路径。
在排查复杂问题时,也可以把 `debug_artifacts=true` 打开。 在排查复杂问题时,也可以把 `debug_artifacts=true` 打开,并结合 `list_run_artifacts` 查看该 run 下实际产物。
## 对结构化摘要结果执行确定性过滤规则 ## 对结构化摘要结果执行确定性过滤规则
+21 -2
View File
@@ -168,7 +168,7 @@
--- ---
### [TODO][P2] 设计并实现 `resume_run` ### [DONE][P2] 设计并实现 `resume_run`
目标: 目标:
- 基于 `run-state.json` 和现有中间产物继续执行 - 基于 `run-state.json` 和现有中间产物继续执行
@@ -176,6 +176,14 @@
说明: 说明:
- 先做最小可用恢复 - 先做最小可用恢复
- 暂不追求任意 stage 任意重入 - 暂不追求任意 stage 任意重入
- 设计约束已补充到 `plans/resume-run-minimal-design.md`
- 第一版只支持 freshrss workflow 且仅支持有 `run-state.json` 的 run
- 第一版仅考虑从最近可恢复点继续;`fetch_feed` / `extract_articles` 暂不支持恢复
进展备注:
- 2026-04-07:已完成 `resume_run` minimal design 与现有 runtime/workflow/server 代码对齐分析,开始实现最小恢复链路。
- 2026-04-07:已完成 `resume_run` 最小实现编码,新增 runtime 恢复服务并接入 MCP server;当前进入设计对齐与本地自检。
- 2026-04-07:已完成 `resume_run` 架构对齐与本地自检;已验证 `write_run_report` 可恢复,且 `extract_articles` 会被明确拒绝恢复。
--- ---
@@ -194,12 +202,23 @@
--- ---
### [TODO][P3] 更新 README / handoff / docs,明确 MCP 为正式入口 ### [DONE][P3] 更新 README / handoff / docs,明确 MCP 为正式入口
目标: 目标:
- 把生产建议从 CLI 迁移到 MCP - 把生产建议从 CLI 迁移到 MCP
- CLI 明确降级为 debug / fallback - CLI 明确降级为 debug / fallback
完成情况:
- 已更新 `README.md`,补齐 reader 作为正式 MCP workflow service 的当前能力边界、推荐调用路径、已支持 tools 与最小 `resume_run` 范围
- 已更新 `docs/openclaw/openclaw-handoff.md`,明确 OpenClaw 应优先通过 MCP 读取 run 状态与结果,不再自己拼接 reader 输出路径
改动文件:
- `README.md`
- `docs/openclaw/openclaw-handoff.md`
遗留风险:
- 当前仍无独立 `get_digest_brief` tool;若下游确实需要该产物,仍应先通过 `list_run_artifacts` / `get_run_report` 发现,而不是写死路径
--- ---
## 3. 记录区 ## 3. 记录区
+96 -11
View File
@@ -16,6 +16,7 @@ This repository is responsible only for upstream reading-pipeline work:
- rule-based filtering - rule-based filtering
- OpenClaw delivery payload generation - OpenClaw delivery payload generation
- selected-article summary capability based on existing extracted text - selected-article summary capability based on existing extracted text
- run-state persistence, status lookup, result lookup, and minimal resume for the FreshRSS workflow
This repository should **not** take over downstream orchestration responsibilities that belong to OpenClaw / skills, such as: This repository should **not** take over downstream orchestration responsibilities that belong to OpenClaw / skills, such as:
@@ -26,11 +27,34 @@ This repository should **not** take over downstream orchestration responsibiliti
## Production Entrypoint ## Production Entrypoint
OpenClaw should call the MCP tool: reader 当前正式工作流服务入口是 MCP tool:
- `run_freshrss_openclaw_pipeline` - `run_freshrss_openclaw_pipeline`
This is the canonical entrypoint for production use. It starts the only formally supported workflow today: `freshrss_daily_digest`.
OpenClaw should treat the returned `run_id` as the only stable handle for follow-up reads. Do not hand-build `outputs/freshrss/rerun/...` paths in OpenClaw.
## Supported MCP Tools
Current MCP tools: 11 total.
Workflow service tools:
- `run_freshrss_openclaw_pipeline`
- `get_run_status`
- `list_runs`
- `list_run_artifacts`
- `get_delivery_payload`
- `get_run_report`
- `resume_run`
Single-step / debug tools:
- `extract_url_content`
- `extract_item_content`
- `filter_summary_result`
- `generate_article_summaries`
## Required Environment Variables ## Required Environment Variables
@@ -68,9 +92,28 @@ Start the MCP server:
summary-mcp summary-mcp
``` ```
## Recommended MCP Workflow
Recommended production path:
1. Call `run_freshrss_openclaw_pipeline` and persist the returned `run_id`
2. Use `get_run_status(run_id)` as the authoritative run-state read for status, stage, artifacts, and recovery
3. Use `list_runs(...)` when OpenClaw needs recent-run discovery or high-level inspection
4. Use `list_run_artifacts(run_id)` when OpenClaw needs to inspect what this run actually produced
5. Use `get_delivery_payload(run_id)` and `get_run_report(run_id)` as the formal result-reading APIs
6. Use `resume_run(run_id)` only when the run falls inside the minimal supported resume scope
OpenClaw should not directly derive or hardcode:
- `outputs/freshrss/rerun/<run_dir>/run-state.json`
- `outputs/freshrss/rerun/<run_dir>/candidates/openclaw-delivery-payload.json`
- `outputs/freshrss/rerun/<run_dir>/run-report.json`
If filesystem access is needed for debugging, consume only paths returned by MCP such as `output_dir`, `artifact.path`, `delivery_output`, or `report_output`.
## Recommended MCP Call ## Recommended MCP Call
Recommended production call: Recommended production start call:
```json ```json
{ {
@@ -89,17 +132,36 @@ Recommended semantics:
- Use `mark_read=false` only for debug, test, or validation runs. - Use `mark_read=false` only for debug, test, or validation runs.
- Keep `debug_artifacts=false` for routine production runs. - Keep `debug_artifacts=false` for routine production runs.
- Set `debug_artifacts=true` only when troubleshooting a bad batch. - Set `debug_artifacts=true` only when troubleshooting a bad batch.
- Treat the returned `run_id` as the stable identifier for all follow-up MCP reads.
- If no real `openclaw-delivery-payload.json` was produced, OpenClaw should stop instead of generating a digest from placeholders or examples. - If no real `openclaw-delivery-payload.json` was produced, OpenClaw should stop instead of generating a digest from placeholders or examples.
## Formal Capability Boundary
reader 当前正式 MCP workflow service 的边界如下:
- formal workflow: only `freshrss_daily_digest`
- run truth: every FreshRSS run writes `run-state.json`
- state query tools: `get_run_status`, `list_runs`, `list_run_artifacts`
- result read tools: `get_delivery_payload`, `get_run_report`
- `digest-brief.json` is generated and registered as an artifact, but there is no standalone `get_digest_brief` tool yet
- `run_freshrss_openclaw_pipeline` and `resume_run` are synchronous MCP calls today; there is no background queue / worker model yet
- `generate_article_summaries` is supported, but it is outside the formal `resume_run` scope and not part of the FreshRSS workflow-state model
Historical compatibility note:
- `get_run_status` / `list_runs` / `list_run_artifacts` can still infer basic state for older run directories without `run-state.json`
- `resume_run` does **not** support those inferred historical runs; it requires a valid `run-state.json`
## What The Tool Returns ## What The Tool Returns
Primary return fields: Primary return fields from `run_freshrss_openclaw_pipeline`:
- `run_id` - `run_id`
- `output_dir` - `output_dir`
- `raw_output` - `raw_output`
- `delivery_output` - `delivery_output`
- `report_output` - `report_output`
- `digest_brief_output`
- `pulled_count` - `pulled_count`
- `delivered_count` - `delivered_count`
- `marked_read_count` - `marked_read_count`
@@ -110,16 +172,20 @@ Primary return fields:
Optional: Optional:
- `items` - `items`
- Returned only when `include_item_reports=true` - returned only when `include_item_reports=true`
Follow-up structured reads should use MCP tools rather than re-reading these files directly.
## Minimal Output Files ## Minimal Output Files
By default the pipeline writes only: By default the pipeline writes these core artifacts:
- `outputs/freshrss/rerun/<run_id>/raw/freshrss.raw.json` - `outputs/freshrss/rerun/<run_dir>/run-state.json`
- `outputs/freshrss/rerun/<run_id>/candidates/openclaw-delivery-payload.json` - `outputs/freshrss/rerun/<run_dir>/raw/freshrss.raw.json`
- `outputs/freshrss/rerun/<run_id>/run-report.json` - `outputs/freshrss/rerun/<run_dir>/candidates/openclaw-delivery-payload.json`
- `outputs/freshrss/rerun/<run_id>/extracted/item-XX.extracted.json` (one per item) - `outputs/freshrss/rerun/<run_dir>/candidates/digest-brief.json`
- `outputs/freshrss/rerun/<run_dir>/run-report.json`
- `outputs/freshrss/rerun/<run_dir>/extracted/item-XX.extracted.json` (one per item)
It also updates local runtime keyword data: It also updates local runtime keyword data:
@@ -130,6 +196,24 @@ Per-item extracted files live under `extracted/` and are always written.
If `debug_artifacts=true`, the pipeline additionally writes normalized items, summaries, filter decisions, candidate records, and candidate inputs. If `debug_artifacts=true`, the pipeline additionally writes normalized items, summaries, filter decisions, candidate records, and candidate inputs.
The main pipeline does not emit a batch-level `freshrss.extracted.json` file by default. The main pipeline does not emit a batch-level `freshrss.extracted.json` file by default.
## `resume_run` Minimal Scope
`resume_run` currently supports only the minimum resume contract:
- only runs with a valid `run-state.json`
- only workflow `freshrss_daily_digest`
- resume in place on the original `run_id`
- supported resume points: `generate_summaries`, `apply_filters`, `build_delivery_payload`, `write_run_report`
- unsupported resume points: `fetch_feed`, `extract_articles`
- if required artifacts are missing, the tool returns a non-resumable response instead of silently falling back to an earlier stage
Artifact expectations by resume point:
- `generate_summaries`: requires `raw/freshrss.raw.json` and `extracted/`
- `apply_filters`: requires the above plus per-item summary outputs
- `build_delivery_payload`: requires per-item candidate inputs consistent with filter-stage output
- `write_run_report`: requires `candidates/openclaw-delivery-payload.json`; if `mark_read=true`, raw input must still be present
## Payload Specs ## Payload Specs
OpenClaw payload field specs live here: OpenClaw payload field specs live here:
@@ -230,9 +314,10 @@ The minimum you need to give OpenClaw is:
- the MCP server startup command - the MCP server startup command
- the required environment variables in the target environment - the required environment variables in the target environment
- the instruction to call `run_freshrss_openclaw_pipeline` - the instruction to call `run_freshrss_openclaw_pipeline`
- the rule that follow-up state/result reads must go through MCP tools first, not handwritten filesystem paths
If OpenClaw will also participate in keyword-governance review, additionally point it to: If OpenClaw will also participate in keyword-governance review, additionally point it to:
- `docs/design/daily-keyword-index-design.md` - `docs/design/daily-keyword-index-design.md`
- `skills/keyword-cleanup-review/SKILL.md` - `skills/keyword-cleanup-review/SKILL.md`
- `scripts/apply_term_suggestions.py` - `scripts/apply_term_suggestions.py`
@@ -0,0 +1,307 @@
# OpenClaw → reader MCP 标准编排流程
## 1. 文档目的
本文档定义 OpenClaw 在正式环境中如何调用 reader 作为上游 MCP workflow service。
目标不是描述 reader 内部实现,而是明确 OpenClaw 的编排动作:
- 什么时候启动新 run
- 什么时候查询状态
- 什么时候读取结果
- 什么时候尝试恢复
- 什么时候直接新开 run
- 什么时候需要人工介入
本文档基于 reader 当前**已真实落地**的能力编写,不描述尚未实现的未来接口。
---
## 2. 当前 reader 已正式支持的 MCP 能力
当前可用能力:
- `run_freshrss_openclaw_pipeline`
- `get_run_status`
- `list_runs`
- `list_run_artifacts`
- `get_delivery_payload`
- `get_run_report`
- `resume_run`
其中:
- `run_freshrss_openclaw_pipeline` 是当前正式启动入口
- `get_run_status` / `list_runs` / `list_run_artifacts` 用于观测
- `get_delivery_payload` / `get_run_report` 用于读取正式结果
- `resume_run` 用于最小恢复能力
---
## 3. 编排基本原则
### 3.1 OpenClaw 不再手拼路径
OpenClaw 不应再自己拼 reader 输出路径来判断运行状态或读取核心结果。
优先使用 MCP:
- 查状态 → `get_run_status`
- 读 payload → `get_delivery_payload`
- 读 report → `get_run_report`
- 做恢复 → `resume_run`
只有在排障/人工核查时,才回退到直接看 reader run 目录。
### 3.2 reader 是上游 workflow engine
reader 负责:
- FreshRSS 拉取
- 内容提取
- 摘要
- 过滤
- payload 生成
- run 状态记录
- 最小恢复
OpenClaw 负责:
- 触发执行
- 轮询状态
- 读取结果
- 生成 digest markdown
- Hugo 发布
- 聊天汇报
- 用户确认精选
- IMA 编排
### 3.3 默认生产语义
正式生产运行默认:
- `mark_read=true`
- `debug_artifacts=false`
- 只在 debug/test/validation 时显式放宽
---
## 4. 标准 Happy Path
### Step 1: 启动新 run
调用:
- `run_freshrss_openclaw_pipeline`
推荐参数示例:
```json
{
"limit": 5,
"mark_read": true,
"include_read": false,
"debug_artifacts": false,
"timeout_seconds": 60,
"max_retries": 2
}
```
期望:
- 获得 `run_id`
- 获得 `output_dir`
- 获得初始结果摘要
如果启动阶段直接抛错:
- 直接判为启动失败
- 不进入后续查询
### Step 2: 查询运行状态
调用:
- `get_run_status(run_id=...)`
根据返回:
- `status=running` → 继续轮询
- `status=success` → 进入结果读取
- `status=failed` → 进入失败处理
- `status=partial` → 视为未完成,优先看 `recovery` 和当前阶段
### Step 3: 读取正式结果
成功后读取:
- `get_delivery_payload(run_id=...)`
- `get_run_report(run_id=...)`
后续 OpenClaw 编排应以这两个接口为正式结果源,而不是自己拼路径读取 JSON。
### Step 4: 进入下游编排
OpenClaw 在拿到正式 payload / report 后,继续执行:
- public/internal digest 生成
- Hugo 发布
- 聊天汇报
- 用户确认精选
- IMA 沉淀
---
## 5. 状态 → 动作映射
| reader 状态 | OpenClaw 动作 |
|---|---|
| `running` | 继续轮询 `get_run_status` |
| `success` | 读取 `get_delivery_payload` 和 `get_run_report` |
| `failed` 且 `recovery.resumable=true` | 评估是否调用 `resume_run` |
| `failed` 且 `recovery.resumable=false` | 直接判失败,通常新开 run 或人工介入 |
| `partial` | 先读状态详情和 recovery,再决定继续等 / 恢复 / 人工介入 |
---
## 6. 失败处理与恢复决策
### 6.1 什么时候优先尝试 `resume_run`
满足以下条件时,优先考虑恢复而不是新开 run:
- `get_run_status` 返回 `failed`
- `recovery.resumable=true`
- 当前 run 对应的是 freshrss workflow
- 当前失败点在 reader 第一版支持的恢复范围内
### 6.2 `resume_run` 当前支持范围
当前最小实现仅支持:
- 仅对带 `run-state.json` 的 freshrss run
- 仅从最近可恢复点继续
- 支持的恢复点:
- `generate_summaries`
- `apply_filters`
- `build_delivery_payload`
- `write_run_report`
明确不支持:
- `fetch_feed`
- `extract_articles`
### 6.3 什么时候不要恢复,直接新开 run
以下情况不建议 `resume_run`:
- `recovery.resumable=false`
- run 没有 `run-state.json`
- 失败点是 `fetch_feed` 或 `extract_articles`
- 恢复所需关键产物缺失
- 恢复点语义不明确或结果存在明显漂移风险
这时更合理的动作通常是:
- 直接新开 run
- 或人工介入排查
### 6.4 什么时候需要人工介入
出现以下任一情况时,建议人工介入:
- 连续恢复失败
- `get_run_status` 与实际产物明显不一致
- payload/report 结构不符合预期
- 恢复依赖的关键文件缺失且原因不明
- FreshRSS / LLM / 外部环境异常
---
## 7. 读取结果的标准动作
### 7.1 `get_delivery_payload`
用途:
- 获取正式交付给 OpenClaw 的 payload
- 后续 digest 生成应以该返回为准
OpenClaw 应做:
- 读取后直接进入 digest 生成
- 不再自己拼 `candidates/openclaw-delivery-payload.json`
### 7.2 `get_run_report`
用途:
- 获取 run 的结果摘要与关键元信息
- 用于状态判断、排障、补充上下文
OpenClaw 应做:
- 作为诊断与编排辅助信息读取
- 不把 raw file path 解析逻辑继续散落到 skill 里
---
## 8. 最小编排动作表
### 8.1 标准生产执行
1. 调 `run_freshrss_openclaw_pipeline`
2. 拿 `run_id`
3. 轮询 `get_run_status`
4. 若 `success`:
- 调 `get_delivery_payload`
- 调 `get_run_report`
5. 进入 digest / Hugo / chat / IMA 下游编排
### 8.2 失败恢复执行
1. 调 `get_run_status`
2. 若 `failed && recovery.resumable=true`:
- 调 `resume_run`
3. 恢复后再次:
- 调 `get_run_status`
- 若成功,再读 payload / report
4. 若恢复失败或明确不可恢复:
- 新开 run 或人工介入
---
## 9. 不推荐做法
以下做法不应再作为正式主路径:
- 让 OpenClaw 直接长时间 `exec` reader CLI 作为主要生产入口
- 让 OpenClaw 自己拼 reader 输出路径来判断成功/失败
- 让 OpenClaw 自己读取 `outputs/.../*.json` 作为正式结果源
- 在未确认 `resume_run` 支持范围外的失败点上强行恢复
CLI 现在的定位是:
- debug
- fallback
- 人工排障
而不是正式生产主入口。
---
## 10. 当前已知局限
- `resume_run` 仍是最小实现,不支持任意 stage 任意重入
- 历史无 `run-state.json` 的 run 不支持正式恢复
- 极旧 run 的结果读取仍可能依赖保守目录扫描
- `write_run_report` 若涉及重新 `mark_read`,仍依赖 FreshRSS 环境和可用凭据
---
## 11. 一句话结论
OpenClaw 当前应把 reader 当作正式 MCP workflow service 使用:
**启动用 `run_freshrss_openclaw_pipeline`,观测用 `get_run_status`,结果读取用 `get_delivery_payload` / `get_run_report`,恢复仅在 `resume_run` 最小支持范围内启用;不要再把 reader 当成长 CLI 任务和路径拼接仓库来驱动。**
+276
View File
@@ -0,0 +1,276 @@
# resume_run 最小恢复语义设计
## 1. 背景
前面已经完成:
- run-state 基础设施
- 状态查询接口
- 结果读取接口
reader 现在已经具备“运行真相 + 查询 + 结果读取”能力。下一步是补 `resume_run`,让失败后的 run 可以从已有中间产物继续,而不是完全重跑。
但 `resume_run` 是第一项真正涉及“重新进入执行流程”的能力,复杂度高于前几轮。因此本轮设计必须进一步收窄范围,避免一次性把任意 stage 重入、复杂后台管理、任务调度全拉进来。
---
## 2. 本轮目标
只做 **最小可用的 `resume_run`**:
> 基于已有 `run-state.json`,让 freshrss pipeline 能从“最近可恢复点”继续执行。
本轮不追求:
- 任意 stage 任意重入
- 历史 run 的通用恢复
- 无 `run-state.json` run 的恢复
- 多任务后台管理
- 队列 / 数据库 / worker
---
## 3. 支持范围
### 3.1 仅支持 workflow
只支持:
- `freshrss_daily_digest`
- 或当前 `run-state.workflow` 对应的 freshrss pipeline workflow
不支持:
- article summary 独立恢复
- 其他未来 workflow
### 3.2 仅支持有 run-state 的 run
调用 `resume_run(run_id=...)` 时,必须满足:
- run 对应目录存在
- `run-state.json` 存在
- `run-state.json` 可解析
否则直接返回不可恢复错误,而不是尝试猜目录结构。
---
## 4. 最近可恢复点定义
本轮采用保守定义:
### 4.1 恢复依据
优先使用 `run-state.recovery.resume_from_stage`。
如果没有该字段,使用:
- `state.status == failed` 时:失败 stage
- 否则:拒绝恢复
### 4.2 只支持以下恢复点
第一版仅支持从下面几类 stage 恢复:
1. `generate_summaries`
2. `apply_filters`
3. `build_delivery_payload`
4. `write_run_report`
### 4.3 明确不支持的恢复点
第一版暂不支持从以下位置恢复:
1. `fetch_feed`
2. `extract_articles`
原因:
- 这两个阶段更依赖外部抓取与逐条内容处理过程
- 恢复语义更复杂
- 容易和 FreshRSS 读状态、副作用、原始输入不一致问题缠在一起
如果 run 停在这两个阶段:
- `resume_run` 返回 `resumable=false`
- 并明确建议重新触发新 run,而不是恢复
---
## 5. 恢复前置条件
### 5.1 必须存在的中间产物
按恢复点要求最小前置产物:
#### 从 `generate_summaries` 恢复
必须至少有:
- raw output
- extracted outputs(或足以驱动 summary 的当前输入)
#### 从 `apply_filters` 恢复
必须至少有:
- extracted outputs
- summary outputs
#### 从 `build_delivery_payload` 恢复
必须至少有:
- filter / candidate 所需输入已齐备
#### 从 `write_run_report` 恢复
必须至少有:
- delivery payload 已存在
### 5.2 缺失产物处理
如果恢复点所需的关键产物缺失:
- 直接返回不可恢复
- 返回字段中注明缺失产物名
- 不自动降级到更早 stage
原因:
- 第一版先避免隐式魔法恢复
- 让行为更可预测
---
## 6. 执行语义
### 6.1 resume_run 的行为
调用 `resume_run(run_id)` 后:
1. 读取 `run-state.json`
2. 校验 workflow、status、recovery 信息
3. 判断恢复点是否在本轮支持范围内
4. 校验该恢复点所需产物是否齐备
5. 复用现有 freshrss pipeline / runtime,从该恢复点之后继续执行
6. 更新原 run 的 `run-state.json`
### 6.2 不新建 run_id
第一版恢复时:
- **继续沿用原 run_id**
- 不新建子 run / shadow run / retry run
原因:
- 保持恢复行为简单直观
- 避免多 run 关系管理复杂化
### 6.3 状态更新
恢复开始时:
- `status` 置回 `running`
- `current_stage` 置为恢复点
- `error` 清理为 null
- 在必要时更新 `recovery` 信息
恢复成功时:
- 正常走到 `success`
恢复失败时:
- 正常写入新的失败信息
- 保留新的 `run-state.json`
---
## 7. MCP 接口建议
### 7.1 新增 tool
建议新增:
- `resume_run`
### 7.2 输入
```json
{
"run_id": "freshrss-pipeline-20260407-010236"
}
```
第一版不加 `from_stage`,避免人为覆盖恢复点逻辑。
### 7.3 输出
建议输出:
```json
{
"run_id": "freshrss-pipeline-20260407-010236",
"workflow": "freshrss_daily_digest",
"resumed": true,
"resume_from_stage": "build_delivery_payload",
"status": "success",
"output_dir": "outputs/freshrss/rerun/freshrss-pipeline-20260407-010236",
"delivery_payload": {...},
"run_report": {...},
"message": "Run resumed from build_delivery_payload and completed successfully."
}
```
### 7.4 不可恢复时输出
```json
{
"run_id": "freshrss-pipeline-20260407-010236",
"resumed": false,
"status": "failed",
"resume_from_stage": "extract_articles",
"message": "This run cannot be resumed from extract_articles in the current minimal implementation.",
"missing_artifacts": []
}
```
---
## 8. 代码组织建议
优先在现有 runtime 体系内最小扩展:
- `run_store.py`
- 补载入/更新辅助能力(如还缺)
- `query_service.py`
- 可复用读取 run-state / artifact / report / payload 的查询能力
- 新增或补充 runtime service
- 例如 `resume_service.py` 或在现有 runtime 层增加 resume 逻辑
- `server.py`
- 新增 MCP tool `resume_run`
关键原则:
- 不要重写整条 freshrss pipeline
- 应尽量让 pipeline 能接受“从某 stage 之后继续”的最小参数
---
## 9. 明确不做的事
本轮明确不做:
1. `rerun_stage`
2. 指定任意 `from_stage`
3. 多 workflow 通用恢复框架
4. 恢复时自动修补缺失产物
5. 历史无 run-state run 的恢复
6. 后台异步恢复任务管理
---
## 10. 一句话结论
`resume_run` 第一版只做一件事:
**对已有 `run-state.json` 的 freshrss run,在恢复点和前置产物都满足时,从最近可恢复点继续执行;否则明确拒绝恢复,不做隐式魔法补救。**
+798
View File
@@ -0,0 +1,798 @@
from __future__ import annotations
from datetime import date, datetime, timezone
from pathlib import Path
from typing import Any
from summary_mcp.core.keyword_index import persist_keyword_indexes
from summary_mcp.core.summary_loop import resolve_llm_settings, run_loop_payload
from summary_mcp.filters.engine import evaluate_filter_rules, load_filter_rules
from summary_mcp.integrations.freshrss import FreshRSSClient, map_entry_to_item
from summary_mcp.models.article_candidate import (
CandidateMetadata,
CandidateSourceRefs,
OpenClawCandidateInput,
build_article_candidate_record,
build_openclaw_candidate_input,
candidate_id_for,
)
from summary_mcp.models.filtering import FilterContext, FilterInput
from summary_mcp.models.llm_result import LlmSummaryResult
from summary_mcp.models.openclaw_delivery import (
OpenClawDeliveryPayload,
build_openclaw_delivery_payload,
build_openclaw_digest_brief,
)
from summary_mcp.models.summary_io import ExtractionOutput
from summary_mcp.workflows.freshrss_pipeline import (
DEFAULT_PROMPT_PATH,
DEFAULT_RULES_PATH,
DEFAULT_TERM_ALIASES_PATH,
DEFAULT_TERM_DAILY_DIR,
DEFAULT_TERM_STATS_PATH,
DEFAULT_TERM_STOPWORDS_PATH,
DELIVERY_STAGE,
EXTRACT_STAGE,
FETCH_STAGE,
FILTER_STAGE,
REPORT_STAGE,
REPO_ROOT,
SUMMARY_STAGE,
WORKFLOW_NAME,
_build_item_context,
_final_run_status,
_load_json,
_load_required_env,
_save_json,
)
from .query_service import _normalize_repo_path, _resolve_repo_path, _resolve_run_record
from .run_store import RunStore
from .state_models import RunState, StageState
SUPPORTED_RESUME_STAGES = {
SUMMARY_STAGE,
FILTER_STAGE,
DELIVERY_STAGE,
REPORT_STAGE,
}
UNSUPPORTED_RESUME_STAGES = {
FETCH_STAGE,
EXTRACT_STAGE,
}
UTC = timezone.utc
DEFAULT_STREAM_ID = "user/-/state/com.google/reading-list"
def resume_run(*, run_id: str) -> dict[str, Any]:
record = _resolve_run_record(run_id)
base_response = {
"run_id": record["run_id"],
"workflow": record["workflow"],
"status": record["status"],
"output_dir": _normalize_repo_path(record["run_dir"]),
}
if record["state_source"] != "run_state" or not isinstance(record.get("state"), RunState):
return {
**base_response,
"resumed": False,
"resume_from_stage": None,
"message": "This run cannot be resumed because run-state.json is missing or could not be loaded.",
"missing_artifacts": [],
}
run_store = RunStore.load(path=record["run_dir"] / "run-state.json", repo_root=REPO_ROOT)
state = run_store.state
resume_from_stage = _resolve_resume_from_stage(state)
if state.workflow != WORKFLOW_NAME:
return {
**base_response,
"resumed": False,
"resume_from_stage": resume_from_stage,
"message": f"This run cannot be resumed because workflow '{state.workflow}' is not supported by the minimal resume_run implementation.",
"missing_artifacts": [],
}
if resume_from_stage is None:
return {
**base_response,
"resumed": False,
"resume_from_stage": None,
"message": "This run does not expose a recoverable stage in run-state.json.",
"missing_artifacts": [],
}
if resume_from_stage in UNSUPPORTED_RESUME_STAGES:
return {
**base_response,
"resumed": False,
"resume_from_stage": resume_from_stage,
"message": f"This run cannot be resumed from {resume_from_stage} in the current minimal implementation.",
"missing_artifacts": [],
}
if resume_from_stage not in SUPPORTED_RESUME_STAGES:
return {
**base_response,
"resumed": False,
"resume_from_stage": resume_from_stage,
"message": f"This run cannot be resumed because stage '{resume_from_stage}' is not supported.",
"missing_artifacts": [],
}
missing_artifacts = _validate_resume_artifacts(record=record, state=state, resume_from_stage=resume_from_stage)
if missing_artifacts:
return {
**base_response,
"resumed": False,
"resume_from_stage": resume_from_stage,
"message": f"This run cannot be resumed from {resume_from_stage} because required artifacts are missing.",
"missing_artifacts": missing_artifacts,
}
try:
result = _resume_freshrss_run(record=record, run_store=run_store, resume_from_stage=resume_from_stage)
except Exception as error:
failed_stage = run_store.state.current_stage or resume_from_stage
run_store.fail_stage(failed_stage, error=error)
return {
**base_response,
"resumed": True,
"resume_from_stage": resume_from_stage,
"status": run_store.state.status,
"message": f"Run resumed from {resume_from_stage} but failed again at {failed_stage}: {error}",
"missing_artifacts": [],
"error_summary": run_store.state.error.model_dump(mode="json") if run_store.state.error else None,
}
return {
"run_id": run_store.state.run_id,
"workflow": run_store.state.workflow,
"resumed": True,
"resume_from_stage": resume_from_stage,
"status": result["status"],
"output_dir": _normalize_repo_path(record["run_dir"]),
"message": f"Run resumed from {resume_from_stage} and completed with status {result['status']}.",
"delivery_payload": result.get("delivery_payload"),
"run_report": result.get("run_report"),
"delivery_output": result.get("delivery_output"),
"report_output": result.get("report_output"),
"digest_brief_output": result.get("digest_brief_output"),
"keyword_index": result.get("keyword_index"),
"pulled_count": result.get("pulled_count"),
"delivered_count": result.get("delivered_count"),
"marked_read_count": result.get("marked_read_count"),
"status_counts": result.get("status_counts"),
"missing_artifacts": [],
}
def _resume_freshrss_run(*, record: dict[str, Any], run_store: RunStore, resume_from_stage: str) -> dict[str, Any]:
state = run_store.state
run_dir = record["run_dir"]
config = _build_resume_config(state=state, run_dir=run_dir)
items = _load_items(config["raw_output"]) if config["raw_output"] is not None else []
item_contexts = _load_item_contexts(
items=items,
run_dir=run_dir,
debug_artifacts=config["debug_artifacts"],
include_summaries=resume_from_stage in {FILTER_STAGE, DELIVERY_STAGE},
include_candidates=resume_from_stage == DELIVERY_STAGE,
)
item_reports = [context["item_report"] for context in item_contexts]
delivered_candidates = [context["candidate"] for context in item_contexts if context.get("candidate") is not None]
if resume_from_stage == SUMMARY_STAGE:
_run_summary_stage(run_store=run_store, item_contexts=item_contexts, config=config)
_run_filter_stage(run_store=run_store, item_contexts=item_contexts, config=config)
delivered_candidates = [context["candidate"] for context in item_contexts if context.get("candidate") is not None]
delivery_payload, keyword_index_result = _run_delivery_stage(
run_store=run_store,
delivered_candidates=delivered_candidates,
config=config,
)
elif resume_from_stage == FILTER_STAGE:
_run_filter_stage(run_store=run_store, item_contexts=item_contexts, config=config)
delivered_candidates = [context["candidate"] for context in item_contexts if context.get("candidate") is not None]
delivery_payload, keyword_index_result = _run_delivery_stage(
run_store=run_store,
delivered_candidates=delivered_candidates,
config=config,
)
elif resume_from_stage == DELIVERY_STAGE:
delivery_payload, keyword_index_result = _run_delivery_stage(
run_store=run_store,
delivered_candidates=delivered_candidates,
config=config,
)
else:
delivery_payload = OpenClawDeliveryPayload.model_validate(_load_json(config["delivery_output"]))
candidate_by_id = {candidate.candidate_id: candidate for candidate in delivery_payload.candidates}
item_contexts = _load_item_contexts(
items=items,
run_dir=run_dir,
debug_artifacts=config["debug_artifacts"],
delivery_candidates_by_id=candidate_by_id,
)
item_reports = [context["item_report"] for context in item_contexts]
delivered_candidates = list(delivery_payload.candidates)
keyword_index_result = _build_keyword_index_result(
delivery_payload=delivery_payload,
keyword_daily_output=_stage_output(state, DELIVERY_STAGE, "keyword_daily_output"),
keyword_stats_output=_stage_output(state, DELIVERY_STAGE, "keyword_stats_output"),
)
item_reports = [context["item_report"] for context in item_contexts]
delivered_item_ids = [
context["item"].external_id
for context in item_contexts
if context.get("candidate") is not None and context["item"].external_id
]
report, marked_count = _run_report_stage(
run_store=run_store,
items=items,
item_reports=item_reports,
delivered_candidates=delivered_candidates,
delivered_item_ids=delivered_item_ids,
delivery_payload=delivery_payload,
keyword_index_result=keyword_index_result,
config=config,
)
final_status = _final_run_status(item_reports)
run_store.finish_run(status=final_status)
return {
"status": final_status,
"delivery_payload": delivery_payload.model_dump(mode="json"),
"run_report": report,
"delivery_output": _normalize_repo_path(config["delivery_output"]),
"report_output": _normalize_repo_path(config["report_output"]),
"digest_brief_output": _normalize_repo_path(config["digest_brief_output"]),
"keyword_index": keyword_index_result,
"pulled_count": report["pulled_count"],
"delivered_count": report["delivered_count"],
"marked_read_count": marked_count,
"status_counts": report["status_counts"],
}
def _build_resume_config(*, state: RunState, run_dir: Path) -> dict[str, Any]:
input_payload = state.input if isinstance(state.input, dict) else {}
prompt_path = Path(input_payload.get("prompt")) if input_payload.get("prompt") else DEFAULT_PROMPT_PATH
rules_path = Path(input_payload.get("rules")) if input_payload.get("rules") else DEFAULT_RULES_PATH
delivery_date_value = input_payload.get("delivery_date")
resolved_delivery_date = date.fromisoformat(delivery_date_value) if isinstance(delivery_date_value, str) else datetime.now(tz=UTC).date()
context = input_payload.get("context") if isinstance(input_payload.get("context"), dict) else None
context_path_value = input_payload.get("context_path")
context_path = Path(context_path_value) if isinstance(context_path_value, str) and context_path_value.strip() else None
raw_output = _find_artifact_path(run_dir=run_dir, state=state, artifact_name="raw_output", relative_path=Path("raw/freshrss.raw.json"))
delivery_output = _find_artifact_path(
run_dir=run_dir,
state=state,
artifact_name="delivery_payload",
relative_path=Path("candidates/openclaw-delivery-payload.json"),
) or (run_dir / "candidates" / "openclaw-delivery-payload.json")
digest_brief_output = _find_artifact_path(
run_dir=run_dir,
state=state,
artifact_name="digest_brief",
relative_path=Path("candidates/digest-brief.json"),
) or (run_dir / "candidates" / "digest-brief.json")
return {
"run_dir": run_dir,
"raw_output": raw_output,
"delivery_output": delivery_output,
"digest_brief_output": digest_brief_output,
"report_output": run_dir / "run-report.json",
"prompt_path": prompt_path,
"rules_path": rules_path,
"context": context,
"context_path": context_path,
"limit": int(input_payload.get("limit") or _stage_output(state, FETCH_STAGE, "pulled_count") or 0),
"mark_read": bool(input_payload.get("mark_read", False)),
"include_read": bool(input_payload.get("include_read", False)),
"debug_artifacts": bool(input_payload.get("debug_artifacts", False)),
"continuation": input_payload.get("continuation"),
"stream_id": str(input_payload.get("stream_id") or DEFAULT_STREAM_ID),
"timeout_seconds": float(input_payload.get("timeout_seconds") or 60.0),
"max_retries": int(input_payload.get("max_retries") or 2),
"delivery_date": resolved_delivery_date,
}
def _load_item_contexts(
*,
items: list[Any],
run_dir: Path,
debug_artifacts: bool,
include_summaries: bool = False,
include_candidates: bool = False,
delivery_candidates_by_id: dict[str, OpenClawCandidateInput] | None = None,
) -> list[dict[str, Any]]:
contexts: list[dict[str, Any]] = []
for index, item in enumerate(items, start=1):
context = _build_item_context(
index=index,
item=item,
resolved_output_dir=run_dir,
debug_artifacts=debug_artifacts,
)
context["candidate"] = None
extraction_path = context["extracted_path"]
if extraction_path.exists():
extracted_payload = _load_json(extraction_path)
extraction = ExtractionOutput.model_validate(extracted_payload)
context["extraction"] = extraction
context["extracted_payload"] = extracted_payload
if not extraction.success or extraction.article is None:
context["item_report"]["status"] = "extract_failed"
context["item_report"]["error"] = extraction.error.model_dump(mode="json") if extraction.error else None
else:
context["item_report"]["status"] = "extracted"
if include_summaries and context["summary_output"] is not None and context["summary_output"].exists():
summary_payload = _load_json(context["summary_output"])
context["summary_payload"] = summary_payload
context["item_report"]["status"] = "summarized"
elif include_summaries and context.get("extraction") is not None and context["extraction"].success:
context["item_report"]["status"] = "summary_failed"
candidate: OpenClawCandidateInput | None = None
if include_candidates and context["openclaw_path"] is not None and context["openclaw_path"].exists():
candidate = OpenClawCandidateInput.model_validate(_load_json(context["openclaw_path"]))
elif delivery_candidates_by_id is not None and context.get("extraction") is not None and context["extraction"].success:
candidate_id = candidate_id_for(item, context["extraction"].article)
candidate = delivery_candidates_by_id.get(candidate_id)
if candidate is None and context["item_report"]["status"] == "extracted":
context["item_report"]["status"] = "summary_failed"
if candidate is not None:
context["candidate"] = candidate
context["item_report"]["status"] = "delivered"
context["item_report"]["selection_decision"] = candidate.selection_decision
context["item_report"]["candidate_id"] = candidate.candidate_id
contexts.append(context)
return contexts
def _run_summary_stage(*, run_store: RunStore, item_contexts: list[dict[str, Any]], config: dict[str, Any]) -> None:
extraction_success_count = sum(1 for context in item_contexts if context.get("extraction") is not None and context["extraction"].success)
run_store.start_stage(
SUMMARY_STAGE,
outputs={
"expected_items": extraction_success_count,
"completed_items": 0,
"success_count": 0,
"failed_count": 0,
},
)
summary_success_count = 0
summary_failed_count = 0
resolved_llm_api_key, resolved_llm_model, resolved_llm_api_url = resolve_llm_settings()
for item_context in [context for context in item_contexts if context.get("extraction") is not None and context["extraction"].success]:
summary_exit_code, summary_payload, summary_report = run_loop_payload(
extracted_payload=item_context["extracted_payload"],
prompt_path=config["prompt_path"],
output_path=item_context["summary_output"],
max_retries=config["max_retries"],
timeout_seconds=config["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_context["item_report"]["status"] = "summary_failed"
if summary_report is not None:
item_context["item_report"]["summary_errors"] = summary_report.errors
summary_failed_count += 1
else:
item_context["summary_payload"] = summary_payload
item_context["item_report"]["status"] = "summarized"
summary_success_count += 1
run_store.update_stage(
SUMMARY_STAGE,
outputs={
"expected_items": extraction_success_count,
"completed_items": summary_success_count + summary_failed_count,
"success_count": summary_success_count,
"failed_count": summary_failed_count,
},
)
if config["debug_artifacts"] and (config["run_dir"] / "summary").exists():
run_store.register_artifact(name="summary_dir", path=config["run_dir"] / "summary", kind="directory", stage=SUMMARY_STAGE)
run_store.finish_stage(
SUMMARY_STAGE,
outputs={
"expected_items": extraction_success_count,
"completed_items": summary_success_count + summary_failed_count,
"success_count": summary_success_count,
"failed_count": summary_failed_count,
},
)
def _run_filter_stage(*, run_store: RunStore, item_contexts: list[dict[str, Any]], config: dict[str, Any]) -> None:
loaded_rules = load_filter_rules(config["rules_path"])
filter_context = _load_filter_context(config)
summary_success_count = sum(1 for context in item_contexts if context.get("summary_payload") is not None)
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 [context for context in item_contexts if context.get("summary_payload") is not None]:
item = item_context["item"]
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,
)
if item_context["filter_path"] is not None:
_save_json(item_context["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(item_context["filter_path"]) if item_context["filter_path"] else None,
),
metadata=CandidateMetadata(
generated_at=datetime.now(tz=UTC),
producer="resume_run",
run_id=run_store.state.run_id,
),
)
candidate = build_openclaw_candidate_input(record)
item_context["candidate"] = candidate
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"], candidate.model_dump(mode="json"))
item_context["item_report"]["status"] = "delivered"
item_context["item_report"]["selection_decision"] = decision.decision
item_context["item_report"]["candidate_id"] = candidate.candidate_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
candidate_count = sum(1 for context in item_contexts if context.get("candidate") is not None)
run_store.update_stage(
FILTER_STAGE,
outputs={
"expected_items": summary_success_count,
"completed_items": filter_completed_count,
"candidate_count": candidate_count,
"keep_count": keep_count,
"review_count": review_count,
"drop_count": drop_count,
},
)
if config["debug_artifacts"] and (config["run_dir"] / "candidates").exists():
run_store.register_artifact(name="candidate_dir", path=config["run_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": sum(1 for context in item_contexts if context.get("candidate") is not None),
"keep_count": keep_count,
"review_count": review_count,
"drop_count": drop_count,
},
)
def _run_delivery_stage(
*,
run_store: RunStore,
delivered_candidates: list[OpenClawCandidateInput],
config: dict[str, Any],
) -> tuple[OpenClawDeliveryPayload, dict[str, Any]]:
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=run_store.state.run_id,
for_date=config["delivery_date"],
)
_save_json(config["delivery_output"], delivery_payload.model_dump(mode="json"))
run_store.register_artifact(name="delivery_payload", path=config["delivery_output"], kind="json", stage=DELIVERY_STAGE)
digest_brief = build_openclaw_digest_brief(delivery_payload)
_save_json(config["digest_brief_output"], digest_brief.model_dump(mode="json"))
run_store.register_artifact(name="digest_brief", path=config["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(config["delivery_output"]),
"digest_brief_output": str(config["digest_brief_output"]),
"keyword_daily_output": str(keyword_index_result["daily_output"]),
"keyword_stats_output": str(keyword_index_result["stats_output"]),
},
)
return delivery_payload, keyword_index_result
def _run_report_stage(
*,
run_store: RunStore,
items: list[Any],
item_reports: list[dict[str, Any]],
delivered_candidates: list[OpenClawCandidateInput],
delivered_item_ids: list[str],
delivery_payload: OpenClawDeliveryPayload,
keyword_index_result: dict[str, Any],
config: dict[str, Any],
) -> tuple[dict[str, Any], int]:
run_store.start_stage(REPORT_STAGE, outputs={"mark_read_requested": config["mark_read"]})
marked_count = 0
if config["mark_read"] and delivered_item_ids:
resolved_api_base_url = _load_required_env("FRESHRSS_API_BASE_URL", None)
resolved_username = _load_required_env("FRESHRSS_USERNAME", None)
resolved_api_password = _load_required_env("FRESHRSS_API_PASSWORD", None)
client = FreshRSSClient(
api_base_url=resolved_api_base_url,
username=resolved_username,
api_password=resolved_api_password,
timeout_seconds=config["timeout_seconds"],
)
auth_token = client.client_login()
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_resume_run_report(
state=run_store.state,
items=items,
item_reports=item_reports,
delivered_candidates=delivered_candidates,
marked_count=marked_count,
delivery_payload=delivery_payload,
keyword_index_result=keyword_index_result,
config=config,
)
_save_json(config["report_output"], report)
run_store.register_artifact(name="run_report", path=config["report_output"], kind="json", stage=REPORT_STAGE)
run_store.finish_stage(
REPORT_STAGE,
outputs={
"marked_read_count": marked_count,
"report_output": str(config["report_output"]),
},
)
return report, marked_count
def _build_resume_run_report(
*,
state: RunState,
items: list[Any],
item_reports: list[dict[str, Any]],
delivered_candidates: list[OpenClawCandidateInput],
marked_count: int,
delivery_payload: OpenClawDeliveryPayload,
keyword_index_result: dict[str, Any],
config: dict[str, Any],
) -> dict[str, Any]:
status_counts: dict[str, int] = {}
for item_report in item_reports:
status = str(item_report["status"])
status_counts[status] = status_counts.get(status, 0) + 1
pulled_count = len(items)
if not items:
pulled_count = int(_stage_output(state, FETCH_STAGE, "pulled_count") or len(item_reports))
if not status_counts:
failed_count = int(_stage_output(state, SUMMARY_STAGE, "failed_count") or 0)
if failed_count:
status_counts["summary_failed"] = failed_count
if delivered_candidates:
status_counts["delivered"] = len(delivered_candidates)
report: dict[str, Any] = {
"run_id": state.run_id,
"started_at": state.started_at.isoformat(),
"completed_at": datetime.now(tz=UTC).isoformat(),
"requested_limit": config["limit"],
"pulled_count": pulled_count,
"delivered_count": len(delivered_candidates),
"marked_read_count": marked_count,
"mark_read_requested": config["mark_read"],
"debug_artifacts": config["debug_artifacts"],
"raw_output": str(config["raw_output"]) if config["raw_output"] is not None else None,
"delivery_output": str(config["delivery_output"]),
"digest_brief_output": str(config["digest_brief_output"]),
"keyword_index": keyword_index_result,
"status_counts": status_counts,
"items": item_reports,
}
return report
def _build_keyword_index_result(
*,
delivery_payload: OpenClawDeliveryPayload,
keyword_daily_output: str | None,
keyword_stats_output: str | None,
) -> dict[str, Any]:
result: dict[str, Any] = {
"source": "openclaw_delivery_payload",
"candidate_count": len(delivery_payload.candidates),
}
if isinstance(keyword_daily_output, str) and keyword_daily_output.strip():
result["daily_output"] = keyword_daily_output
if isinstance(keyword_stats_output, str) and keyword_stats_output.strip():
result["stats_output"] = keyword_stats_output
return result
def _load_filter_context(config: dict[str, Any]) -> FilterContext:
if config["context"] is not None:
return FilterContext.model_validate(config["context"])
if config["context_path"] is not None and config["context_path"].exists():
return FilterContext.model_validate(_load_json(config["context_path"]))
return FilterContext()
def _validate_resume_artifacts(*, record: dict[str, Any], state: RunState, resume_from_stage: str) -> list[str]:
run_dir = record["run_dir"]
missing_artifacts: list[str] = []
raw_output = _find_artifact_path(run_dir=run_dir, state=state, artifact_name="raw_output", relative_path=Path("raw/freshrss.raw.json"))
extracted_dir = _find_artifact_path(run_dir=run_dir, state=state, artifact_name="extracted_dir", relative_path=Path("extracted"))
if resume_from_stage in {SUMMARY_STAGE, FILTER_STAGE, DELIVERY_STAGE} and raw_output is None:
missing_artifacts.append("raw/freshrss.raw.json")
if resume_from_stage in {SUMMARY_STAGE, FILTER_STAGE} and extracted_dir is None:
missing_artifacts.append("extracted/")
if missing_artifacts:
return missing_artifacts
items = _load_items(raw_output) if raw_output is not None else []
item_contexts = _load_item_contexts(
items=items,
run_dir=run_dir,
debug_artifacts=bool(state.input.get("debug_artifacts", False)),
include_summaries=resume_from_stage in {FILTER_STAGE, DELIVERY_STAGE},
include_candidates=resume_from_stage == DELIVERY_STAGE,
)
if resume_from_stage == SUMMARY_STAGE:
for context in item_contexts:
if not context["extracted_path"].exists():
missing_artifacts.append(_normalize_repo_path(context["extracted_path"]))
if resume_from_stage == FILTER_STAGE:
for context in item_contexts:
if context.get("extraction") is not None and context["extraction"].success:
if context["summary_output"] is None or not context["summary_output"].exists():
missing_artifacts.append(
_normalize_repo_path(context["summary_output"] or (run_dir / "summary" / context["item_key"] / "result.loop.json"))
)
if resume_from_stage == DELIVERY_STAGE:
expected_candidate_count = int(_stage_output(state, FILTER_STAGE, "candidate_count") or 0)
actual_candidate_count = sum(1 for context in item_contexts if context.get("candidate") is not None)
if expected_candidate_count != actual_candidate_count:
for context in item_contexts:
if context.get("summary_payload") is not None and (context["openclaw_path"] is None or not context["openclaw_path"].exists()):
missing_artifacts.append(
_normalize_repo_path(context["openclaw_path"] or (run_dir / "candidates" / f"{context['item_key']}.openclaw-candidate-input.json"))
)
if resume_from_stage == REPORT_STAGE:
delivery_output = _find_artifact_path(
run_dir=run_dir,
state=state,
artifact_name="delivery_payload",
relative_path=Path("candidates/openclaw-delivery-payload.json"),
)
if delivery_output is None:
missing_artifacts.append("candidates/openclaw-delivery-payload.json")
elif bool(state.input.get("mark_read", False)) and raw_output is None:
missing_artifacts.append("raw/freshrss.raw.json")
deduped_missing_artifacts: list[str] = []
for artifact in missing_artifacts:
if artifact not in deduped_missing_artifacts:
deduped_missing_artifacts.append(artifact)
return deduped_missing_artifacts
def _resolve_resume_from_stage(state: RunState) -> str | None:
if state.recovery.resume_from_stage:
return state.recovery.resume_from_stage
if state.status == "failed":
failed_stage = next((stage.name for stage in reversed(state.stages) if stage.status == "failed"), None)
return failed_stage
return None
def _stage_output(state: RunState, stage_name: str, key: str) -> Any:
stage_state = _find_stage_state(state, stage_name)
if stage_state is None:
return None
return stage_state.outputs.get(key)
def _find_stage_state(state: RunState, stage_name: str) -> StageState | None:
for stage in state.stages:
if stage.name == stage_name:
return stage
return None
def _find_artifact_path(*, run_dir: Path, state: RunState, artifact_name: str, relative_path: Path) -> Path | None:
for artifact in state.artifacts:
if artifact.name != artifact_name:
continue
resolved = _resolve_repo_path(artifact.path)
if resolved.exists():
return resolved
path = run_dir / relative_path
if path.exists():
return path
return None
def _load_items(raw_output: Path) -> list[Any]:
payload = _load_json(raw_output)
entries = payload.get("items")
if not isinstance(entries, list):
raise RuntimeError("FreshRSS raw output does not contain an items array.")
return [map_entry_to_item(entry) for entry in entries]
+1 -1
View File
@@ -168,7 +168,7 @@ class RunStore:
if self.state.status == "running": if self.state.status == "running":
resume_from_stage = self.state.current_stage resume_from_stage = self.state.current_stage
elif self.state.status == "failed": elif self.state.status == "failed":
failed_stage = next((stage.name for stage in self.state.stages if stage.status == "failed"), None) failed_stage = next((stage.name for stage in reversed(self.state.stages) if stage.status == "failed"), None)
resume_from_stage = failed_stage or self.state.current_stage resume_from_stage = failed_stage or self.state.current_stage
self.state.recovery = RecoveryState( self.state.recovery = RecoveryState(
+7
View File
@@ -20,6 +20,7 @@ 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 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_run_artifacts as load_run_artifacts
from summary_mcp.runtime import list_runs as load_runs from summary_mcp.runtime import list_runs as load_runs
from summary_mcp.runtime.resume_service import resume_run as resume_existing_run
from summary_mcp.workflows import run_freshrss_pipeline from summary_mcp.workflows import run_freshrss_pipeline
from summary_mcp.workflows.article_summary import ArticleSummaryConfig, summarize_selected_articles from summary_mcp.workflows.article_summary import ArticleSummaryConfig, summarize_selected_articles
@@ -158,6 +159,12 @@ def get_run_report(run_id: str) -> dict:
return load_run_report(run_id=run_id) return load_run_report(run_id=run_id)
@mcp.tool()
def resume_run(run_id: str) -> dict:
"""Resume a failed or interrupted FreshRSS workflow run from its latest supported recovery point."""
return resume_existing_run(run_id=run_id)
@mcp.tool() @mcp.tool()
def generate_article_summaries( def generate_article_summaries(
*, *,
@@ -219,6 +219,13 @@ def run_freshrss_pipeline(
"debug_artifacts": debug_artifacts, "debug_artifacts": debug_artifacts,
"continuation": continuation, "continuation": continuation,
"stream_id": stream_id, "stream_id": stream_id,
"timeout_seconds": timeout_seconds,
"max_retries": max_retries,
"delivery_date": resolved_delivery_date.isoformat(),
"prompt": str(resolved_prompt_path),
"rules": str(resolved_rules_path),
"context": context,
"context_path": str(context_path) if context_path else None,
}, },
repo_root=REPO_ROOT, repo_root=REPO_ROOT,
) )