Compare commits
4
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
2df0af5b82 | ||
|
|
4f219ef92f | ||
|
|
c622bc6247 | ||
|
|
d91cdbc6c9 |
@@ -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 下实际产物。
|
||||||
|
|
||||||
## 对结构化摘要结果执行确定性过滤规则
|
## 对结构化摘要结果执行确定性过滤规则
|
||||||
|
|
||||||
|
|||||||
@@ -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. 记录区
|
||||||
|
|||||||
@@ -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,6 +314,7 @@ 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:
|
||||||
|
|
||||||
|
|||||||
@@ -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 任务和路径拼接仓库来驱动。**
|
||||||
@@ -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,在恢复点和前置产物都满足时,从最近可恢复点继续执行;否则明确拒绝恢复,不做隐式魔法补救。**
|
||||||
@@ -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]
|
||||||
@@ -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(
|
||||||
|
|||||||
@@ -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,
|
||||||
)
|
)
|
||||||
|
|||||||
Reference in New Issue
Block a user