diff --git a/README.md b/README.md index 69fe154..6f70f95 100644 --- a/README.md +++ b/README.md @@ -9,7 +9,7 @@ pip install -e . summary-mcp ``` -服务当前暴露 11 个工具。 +服务当前暴露 14 个工具。 正式 workflow service 相关工具: @@ -26,7 +26,10 @@ summary-mcp - `extract_url_content` - `extract_item_content` - `filter_summary_result` -- `generate_article_summaries` +- `generate_article_summaries`(同步模式) +- `start_article_summary_job` +- `get_article_summary_job_status` +- `get_article_summary_job_result` ## 正式能力边界 @@ -34,7 +37,9 @@ summary-mcp - 每次 FreshRSS 主流水线 run 都会在 `outputs/freshrss/rerun//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 / 异步任务管理 +- `run_freshrss_openclaw_pipeline` 与 `resume_run` 当前都是同步 MCP 调用;仓库里还没有针对 FreshRSS 主流程的后台队列 / worker / 异步任务管理 +- 单篇总结已补上最小异步 job 形态:`start_article_summary_job` / `get_article_summary_job_status` / `get_article_summary_job_result` +- `generate_article_summaries` 仍保留,但定位是同步 debug 路径,而不是 OpenClaw 的正式生产集成入口 ## OpenClaw 推荐调用路径 @@ -70,8 +75,10 @@ summary-mcp 相关能力: -- `generate_article_summaries` MCP 工具(基于已有 extracted payload 做单篇总结) +- `start_article_summary_job` / `get_article_summary_job_status` / `get_article_summary_job_result`(正式推荐的最小异步 job 路径) +- `generate_article_summaries` MCP 工具(同步 debug 路径) - `scripts/run_article_summaries.py` CLI 辅助脚本 +- `scripts/run_article_summary_job.py` 后台 runner 入口 ## 校验 LLM 摘要结果 @@ -257,7 +264,7 @@ python scripts/build_keyword_index.py ^ - `skills/keyword-cleanup-review/` -构建给 LLM skill 使用的评审数据包(review bundle): +构建给关键词治理流程使用的评审数据包(review bundle,临时工作文件): ```bash python skills/keyword-cleanup-review/scripts/build_review_bundle.py ^ @@ -266,17 +273,44 @@ python skills/keyword-cleanup-review/scripts/build_review_bundle.py ^ --output outputs/term_index/review/keyword-cleanup-bundle.json ``` -现在这个评审数据包(review bundle)还会额外携带治理上下文: +这个评审数据包(review bundle)会额外携带治理上下文: - 来自 `configs/term_cleanup_policy.json` 的清理阈值 - 当前 watch list(`configs/term_watchlist.json`) - 最近已应用的变更(`configs/term_change_log.json`) -这个 skill 只负责生成 review 输入与建议,不会自动修改: +接下来可以把 bundle 渲染成正式建议产物(默认走确定性规则,不把 LLM 作为默认路径): -- `term_aliases` -- `term_stopwords` -- `filter_context.personal.json` +```bash +python scripts/generate_term_cleanup_suggestions.py ^ + --bundle outputs/term_index/review/keyword-cleanup-bundle.json +``` + +默认只生成: + +- `outputs/term_index/review/term-cleanup-suggestions-YYYY-MM-DD.json`(唯一正式建议产物,建议短期保留) + +如需人工审阅展示稿,再显式加: + +```bash +python scripts/generate_term_cleanup_suggestions.py ^ + --bundle outputs/term_index/review/keyword-cleanup-bundle.json ^ + --emit-markdown +``` + +这时才会额外生成: + +- `outputs/term_index/review/term-cleanup-suggestions-YYYY-MM-DD.md`(临时展示稿,可按需生成,不必作为长期资产保留) + +其中本轮最小版本优先覆盖 `interest_keyword_suggestions` 和 `watch_terms` 主链路;`alias_suggestions` / `stopword_suggestions` 先保持保守。 + +产物保留策略建议: + +- `data/term_index/daily/YYYY-MM-DD.json`、`data/term_index/term_stats.json` 作为事实层长期保留 +- `configs/filter_context.personal.json`、`configs/term_watchlist.json`、`configs/term_aliases.json`、`configs/term_stopwords.json`、`configs/term_change_log.json` 作为状态层长期保留 +- `term-cleanup-suggestions-YYYY-MM-DD.json` 作为正式建议产物短期保留(最近少量几份或仅保留已应用过的) +- `keyword-cleanup-bundle.json` 仅作为临时工作文件,默认只保留当前最新一份 +- `term-cleanup-suggestions-YYYY-MM-DD.md` 仅作为临时展示稿,优先按需生成,不建议默认长期归档 如果你想先预览已接受建议,再决定是否写配置文件: @@ -324,12 +358,26 @@ python scripts/run_article_summaries.py ^ 也可以通过 `summary_mcp.server` 暴露的 MCP 工具 `generate_article_summaries` 调用: - `extracted_path`(string):单篇 extracted JSON 路径(例如 `outputs/freshrss/rerun//extracted/item-01.extracted.json`),或者包含 `results` 数组的 batch extracted JSON -- `selected_ids`(array of strings):要总结的一个或多个 `item_id`。如果传空数组,则对文件中的全部条目做总结 +- `selected_ids`(array of strings):必填,至少传一个 `item_id`;如果传入的 ID 在 extracted payload 中一个都匹配不到,会直接报错 - `output_dir`(optional string):Markdown 输出目录;若不传,则默认写到 extracted 文件旁边的 `single_summaries/` 目录 - `llm_api_key` / `llm_model` / `llm_api_url`(optional strings):单篇总结 LLM 的覆盖配置;不传时会按前文规则回退到 `ARTICLE_SUMMARY_*` 或主 `LLM_*` 该工具返回一个 JSON 数组,内容为生成好的 Markdown 文件路径。 +OpenClaw / 正式集成建议优先走异步 job: + +- 调 `start_article_summary_job` 启动任务,立即拿到 `job_id` +- 轮询 `get_article_summary_job_status(job_id)`,直到 `status` 进入 `success` 或 `failed` +- 成功后调用 `get_article_summary_job_result(job_id)` 读取 `written_paths` 与结构化结果 +- 失败时优先查看 `error_summary` 与 job 目录中的 `job-report.json` + +异步 job 状态目录固定落在 `outputs/freshrss/article_summary_jobs//`,最小会包含: + +- `run-state.json` +- `input.json` +- `result.json`(成功时) +- `job-report.json` + 单篇总结使用独立 prompt:`outputs/prompts/article-summary-prompt.txt`。 它与日报 prompt 完全独立,输出的是中文结构化知识笔记,包含这些部分: diff --git a/TODO.md b/TODO.md index bab5c9d..0ae5613 100644 --- a/TODO.md +++ b/TODO.md @@ -187,6 +187,47 @@ --- +### [DOING][P1] 单篇总结改为最小真异步 job + +目标: +- 解决 `generate_article_summaries` 在 OpenClaw → MCP 同步链路里易 timeout 的问题 +- 将单篇总结正式升级为可启动、可轮询、可读取结果的异步 job + +要求: +- 新增最小异步接口: + - `start_article_summary_job` + - `get_article_summary_job_status` + - `get_article_summary_job_result` +- 状态目录固定落到: + - `outputs/freshrss/article_summary_jobs//` +- 至少包含: + - `run-state.json` + - `input.json` + - `result.json`(成功时) + - `job-report.json` +- 执行模型优先使用后台子进程,不使用线程 +- 业务逻辑继续复用 `summarize_selected_articles(...)`,不要重写正文总结核心逻辑 +- 对 OpenClaw / reader-digest-flow 而言,异步 job 成功后应可继续接 IMA 沉淀闭环 + +当前进展: +- 2026-04-10:已完成方案文档 `plans/article-summary-async-job-plan.md` +- 2026-04-10:已落地最小代码骨架: + - `src/summary_mcp/runtime/article_summary_jobs.py` + - `scripts/run_article_summary_job.py` + - `src/summary_mcp/server.py` 已新增 3 个 async job tools +- 2026-04-10:已用真实 extracted 文件验证最小异步链路可跑通,job 能成功进入 `running → success`,并可读回结果 +- 2026-04-10:已补 README / OpenClaw handoff 文档,并对外统一为 `start_article_summary_job` / `get_article_summary_job_status` / `get_article_summary_job_result` +- 2026-04-10:已完成聚焦自检: + - MCP tool 注册名校验通过 + - stubbed async job 成功路径通过 + - stubbed async job 失败路径通过,`error_summary` 与 `job-report.json` 可回读 + +下一步: +- 将 `reader-digest-flow` 正式默认路径切到 async job +- 用真实 LLM 配置再做一次非 stub 的服务端冒烟验证 + +--- + ### [TODO][P2] 评估 `rerun_stage` 是否值得进入第一阶段 目标: diff --git a/docs/openclaw/openclaw-handoff.md b/docs/openclaw/openclaw-handoff.md index ebe4a3c..acbf0de 100644 --- a/docs/openclaw/openclaw-handoff.md +++ b/docs/openclaw/openclaw-handoff.md @@ -37,7 +37,7 @@ OpenClaw should treat the returned `run_id` as the only stable handle for follow ## Supported MCP Tools -Current MCP tools: 11 total. +Current MCP tools: 14 total, including the FreshRSS workflow set plus article-summary async job tools. Workflow service tools: @@ -49,12 +49,34 @@ Workflow service tools: - `get_run_report` - `resume_run` +Article-summary tools: + +- `start_article_summary_job` +- `get_article_summary_job_status` +- `get_article_summary_job_result` + Single-step / debug tools: - `extract_url_content` - `extract_item_content` - `filter_summary_result` -- `generate_article_summaries` +- `generate_article_summaries`(同步模式,仅适合轻量调试) + +## Recommended Selected-Article Flow + +For OpenClaw selected-article follow-up, prefer the async job path: + +1. Call `start_article_summary_job` with a real extracted file path plus a non-empty `selected_ids` list. +2. Poll `get_article_summary_job_status` until `status` becomes `success` or `failed`. +3. On success, call `get_article_summary_job_result` and continue downstream processing from `written_paths`. +4. Use `generate_article_summaries` only as a synchronous debug fallback, not as the default production path. + +Job state is written under `outputs/freshrss/article_summary_jobs//` and will minimally contain: + +- `run-state.json` +- `input.json` +- `result.json` on success +- `job-report.json` ## Required Environment Variables @@ -144,8 +166,9 @@ reader 当前正式 MCP workflow service 的边界如下: - 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 +- `run_freshrss_openclaw_pipeline` and `resume_run` are synchronous MCP calls today; there is no background queue / worker model for the FreshRSS daily workflow yet +- article summary now has a minimal asynchronous job model with `start_article_summary_job` / `get_article_summary_job_status` / `get_article_summary_job_result` +- `generate_article_summaries` is still supported, but it is a synchronous debug path and outside the formal `resume_run` scope Historical compatibility note: diff --git a/plans/article-summary-async-job-plan.md b/plans/article-summary-async-job-plan.md new file mode 100644 index 0000000..90752a8 --- /dev/null +++ b/plans/article-summary-async-job-plan.md @@ -0,0 +1,757 @@ +# 单篇总结异步 job 最小版落地方案 + +## 1. 背景与问题定义 + +当前 reader 已经把日更 FreshRSS 主流程做成了带 `run-state.json` 的 runtime 模型: + +- `src/summary_mcp/runtime/state_models.py` +- `src/summary_mcp/runtime/run_store.py` +- `src/summary_mcp/runtime/query_service.py` +- `src/summary_mcp/workflows/freshrss_pipeline.py` + +这条主链路已经具备: + +- run / stage / artifact 的结构化状态 +- MCP 查询接口:`get_run_status` / `list_runs` / `list_run_artifacts` +- 结果读取接口:`get_delivery_payload` / `get_run_report` + +但“单篇总结”这条线目前还是同步调用: + +- 核心逻辑:`src/summary_mcp/workflows/article_summary.py` +- MCP 暴露:`src/summary_mcp/server.py` 中的 `generate_article_summaries` +- CLI 辅助:`scripts/run_article_summaries.py` + +现状问题已经很明确: + +- 在 OpenClaw → MCP tool 这条链路里,`generate_article_summaries` 可能因为 tool 调用时长而 timeout +- 但 reader 项目本体在 `.venv` 下直接跑 article summary,大约 37.5 秒即可成功 +- 这说明问题不一定在 summary 业务本身,而更可能在“同步工具调用 + 上层等待模型”这个包装层 + +所以目标不是先继续调 timeout,而是把单篇总结也纳入 **真正异步、可轮询、可落盘、可恢复基本状态** 的最小 job 模型里。 + +--- + +## 2. 为什么同步 MCP 不适合这一步 + +`generate_article_summaries` 当前在 `server.py` 里直接同步执行 `summarize_selected_articles(...)`,调用方必须一直阻塞等待,直到: + +1. 读取 extracted payload +2. 调用 LLM 生成总结 +3. 可选 repair retry +4. 渲染 Markdown +5. 写文件完成 +6. MCP tool 返回生成路径 + +这个模式对“几十秒级、依赖外部 LLM、可能重试”的任务不稳,核心问题有三层: + +### 2.1 tool 调用时长不可控 + +`run_loop_payload()` 内部会发起外部 HTTP 请求,还可能做 repair retry。即便单次平均 37.5 秒,也已经接近很多上层编排系统的心理和技术超时边界。 + +### 2.2 调用方看不到中间状态 + +现在如果卡住,调用方只能等: + +- 不知道是在读输入 +- 不知道是在调 LLM +- 不知道是在重试 +- 不知道是否已经写出部分结果 + +这也是同步接口最烦的点:失败时只能看到“tool timeout / tool failed”,而不是“业务跑到哪一步了”。 + +### 2.3 与 reader 已有 runtime 风格不一致 + +FreshRSS 主流程已经是“run truth + status query + artifact read”的思路,而单篇总结仍然是黑箱同步函数。继续维持两套风格,只会让 SOP 更复杂: + +- 日报主链路用 `run_id` +- 单篇总结却要么同步等,要么退回 CLI fallback + +这不利于后续把 `reader-digest-flow` 稳定成正式 SOP。 + +--- + +## 3. 本轮目标:最小真异步,不做大而全 + +这次只做 **单篇总结异步 job 最小版**,目标是: + +> 让 OpenClaw 或其他调用方能先“启动单篇总结 job”,立即拿到 `job_id`,再通过状态接口轮询,最后读取输出文件/结果。 + +### 3.1 本轮必须做到的范围 + +1. 新增单篇总结 job 的 start/status/result 最小接口 +2. job 真正在后台执行,而不是 MCP 请求线程里阻塞到完成 +3. 状态落盘到文件,遵循 reader 当前 runtime 风格 +4. 复用现有 `summarize_selected_articles` 逻辑,不重写业务 +5. 输出仍然是现有 Markdown 文件,不改知识内容 schema + +### 3.2 本轮明确不做 + +1. **不做通用队列系统** +2. **不做数据库** +3. **不做多 worker / 分布式调度** +4. **不做取消 job / kill job** +5. **不做并发配额控制** +6. **不做 resume/retry from stage** +7. **不把 article summary 一次性并入 freshrss `resume_run` 体系** +8. **不改 summary prompt / validator / 输出格式** +9. **不处理批量高吞吐场景优化** + +一句话:这轮只解“同步 tool 容易 timeout,但业务本身能跑完”这个问题,不顺手扩成任务调度平台。 + +--- + +## 4. 接入当前 reader 结构的建议 + +### 4.1 复用现有 runtime 设计,但单独建 article summary job 命名空间 + +不建议把 article summary job 粗暴塞进现有 `freshrss_daily_digest` run 查询里混用一个 schema;更合适的是: + +- 复用 `RunState / StageState / ArtifactRecord / RunStore` 这套思维 +- 但给单篇总结定义独立 workflow 名称与存储目录 + +建议: + +- workflow: `article_summary_job` +- run_type: `article_summary` +- output root: `outputs/freshrss/article_summary_jobs//` + +这样有几个好处: + +- 不污染 `outputs/freshrss/rerun/` +- 语义清楚:这是独立 job,不是假装自己是日报 rerun +- 查询和排查时更直观 + +### 4.2 job 与结果 Markdown 解耦 + +job 目录只负责: + +- 状态文件 +- 输入快照 +- artifact 索引 +- 执行报告 + +真正生成的总结 Markdown,仍然可以写到用户指定的 `output_dir`(或默认 `single_summaries/`)。 + +这样不破坏当前下游 SOP: + +- `reader-digest-flow` 依然从原来的单篇总结输出目录拿 `.md` +- job 目录只提供状态与索引,不强迫下游改结果路径约定 + +--- + +## 5. 新增工具 / API 设计 + +本轮建议新增 3 个 MCP tool,名字尽量和当前 runtime 风格一致。 + +## 5.1 `start_article_summary_job` + +### 作用 + +启动一个后台 job,立即返回 `job_id`,不等待总结完成。 + +### 输入建议 + +```json +{ + "extracted_path": "outputs/freshrss/rerun//extracted/item-01.extracted.json", + "selected_ids": ["12345"], + "output_dir": "outputs/freshrss/single_summaries/2026-04-10", + "max_retries": 2, + "timeout_seconds": 120, + "llm_api_key": null, + "llm_model": null, + "llm_api_url": null +} +``` + +### 返回建议 + +```json +{ + "job_id": "article-summary-20260410-144500-ab12cd34", + "workflow": "article_summary_job", + "run_type": "article_summary", + "status": "running", + "output_dir": "outputs/freshrss/article_summary_jobs/article-summary-20260410-144500-ab12cd34", + "message": "Article summary job started successfully. Use get_article_summary_job_status to poll progress." +} +``` + +### 约束建议 + +- `selected_ids` 第一版允许多个,但建议由 OpenClaw 每次只传一篇或少量篇,避免一个 job 干太多事 +- `extracted_path` 必须存在,否则直接拒绝启动 +- `output_dir` 不传则按当前默认逻辑推导 + +--- + +## 5.2 `get_article_summary_job_status` + +### 作用 + +查询 job 当前状态、阶段、进度、错误摘要、已注册 artifact。 + +### 输入 + +```json +{ + "job_id": "article-summary-20260410-144500-ab12cd34" +} +``` + +### 返回建议 + +```json +{ + "job_id": "article-summary-20260410-144500-ab12cd34", + "workflow": "article_summary_job", + "run_type": "article_summary", + "status": "running", + "current_stage": "generate_markdown", + "started_at": "2026-04-10T14:45:00+08:00", + "updated_at": "2026-04-10T14:45:23+08:00", + "finished_at": null, + "output_dir": "outputs/freshrss/article_summary_jobs/article-summary-20260410-144500-ab12cd34", + "progress": { + "completed_stage_count": 2, + "running_stage_count": 1, + "failed_stage_count": 0, + "pending_stage_count": 1, + "total_stage_count": 4 + }, + "artifacts": [...], + "error_summary": null +} +``` + +--- + +## 5.3 `get_article_summary_job_result` + +### 作用 + +当 job 成功后,返回结构化结果,供 OpenClaw 继续下游 IMA 沉淀。 + +### 输入 + +```json +{ + "job_id": "article-summary-20260410-144500-ab12cd34" +} +``` + +### 返回建议 + +```json +{ + "job_id": "article-summary-20260410-144500-ab12cd34", + "status": "success", + "written_paths": [ + "outputs/freshrss/single_summaries/2026-04-10/某篇文章标题.md" + ], + "artifact": { + "name": "job_result", + "path": "outputs/freshrss/article_summary_jobs/article-summary-20260410-144500-ab12cd34/result.json", + "kind": "json", + "stage": "write_result" + }, + "result": { + "selected_ids": ["12345"], + "written_paths": [ + "outputs/freshrss/single_summaries/2026-04-10/某篇文章标题.md" + ] + } +} +``` + +### 行为建议 + +- 若 job 还没完成,返回当前状态 + 提示“not ready” +- 若 job 失败,返回失败摘要,不硬抛文件不存在异常 + +--- + +## 6. job 状态文件设计 + +建议直接复用现有 `RunState` 模型,不另造一套 schema。 + +job 目录示例: + +```text +outputs/freshrss/article_summary_jobs/ + article-summary-20260410-144500-ab12cd34/ + run-state.json + input.json + result.json + job-report.json +``` + +## 6.1 `run-state.json` + +建议直接沿用当前字段: + +```json +{ + "run_id": "article-summary-20260410-144500-ab12cd34", + "workflow": "article_summary_job", + "run_type": "article_summary", + "status": "running", + "current_stage": "generate_markdown", + "started_at": "2026-04-10T14:45:00+08:00", + "updated_at": "2026-04-10T14:45:23+08:00", + "finished_at": null, + "input": { + "extracted_path": "outputs/freshrss/rerun//extracted/item-01.extracted.json", + "selected_ids": ["12345"], + "output_dir": "outputs/freshrss/single_summaries/2026-04-10", + "max_retries": 2, + "timeout_seconds": 120 + }, + "stages": [...], + "artifacts": [...], + "error": null, + "recovery": { + "resumable": false, + "resume_from_stage": null, + "last_success_stage": "load_input" + } +} +``` + +### 第一版 stage 建议 + +建议只切 4 个 stage,够看即可: + +1. `prepare_job` + - 校验输入 + - 解析路径 + - 写 `input.json` + +2. `load_input` + - 读取 extracted payload + - 确认 `selected_ids` 可匹配条目 + +3. `generate_markdown` + - 调 `summarize_selected_articles(...)` + - 这是主要耗时阶段 + +4. `write_result` + - 写 `result.json` + - 注册输出 artifact + +这里不要把 LLM 调用再拆更多细 stage,否则最小版反而过度设计。 + +--- + +## 6.2 `input.json` + +作用:保留启动请求快照,便于排查。 + +建议内容与 `run-state.input` 基本一致。 + +--- + +## 6.3 `result.json` + +成功时写: + +```json +{ + "job_id": "article-summary-20260410-144500-ab12cd34", + "selected_ids": ["12345"], + "written_paths": [ + "outputs/freshrss/single_summaries/2026-04-10/某篇文章标题.md" + ], + "completed_at": "2026-04-10T14:45:41+08:00" +} +``` + +失败时可以不写,或只写失败快照都行。最小版建议: + +- 成功写 `result.json` +- 失败只依赖 `run-state.json` + +避免双份失败状态不一致。 + +--- + +## 7. 执行模型建议:优先子进程,不建议线程 + +### 7.1 推荐:子进程后台执行 + +最小真异步推荐模型: + +- `start_article_summary_job` 负责: + - 创建 job 目录 + - 初始化 `run-state.json` + - 通过 `subprocess.Popen(...)` 启动一个独立 Python 进程执行 job runner + - 立即返回 `job_id` + +后台 runner 再去: + +- 读取 `input.json` +- 用 `RunStore` 更新状态 +- 调用 `summarize_selected_articles(...)` +- 写 `result.json` +- finish/fail run + +### 7.2 为什么不推荐线程 + +虽然线程实现看起来更省事,但不适合作为 reader 的正式最小异步落地: + +1. **MCP server 进程重启后线程直接丢失** +2. 线程状态不天然可恢复,容易出现“状态文件还在 running,但线程没了” +3. 未来要做健康检查/孤儿 job 检测时,线程模型更难收口 + +### 7.3 为什么子进程更贴当前项目风格 + +reader 现在本来就偏“文件产物 + runtime 状态真相”风格。子进程模式有天然优势: + +- 和 CLI/fallback 思维一致 +- 进程边界清楚 +- `run-state.json` 由实际执行者写,职责清晰 +- 将来如果要做 orphan detection / stale running job 修复,也容易补 + +### 7.4 本轮不做进程管理增强 + +最小版里,不要求: + +- 记录 PID 后做 kill/cancel +- 自动清理僵尸进程 +- 守护进程/worker 池 + +但建议在 `input.json` 或 `run-state.input` 里附带: + +- `launcher_pid` +- `runner_command` + +方便排障。 + +--- + +## 8. 与现有 `summarize_selected_articles` 的复用关系 + +核心原则:**不重写总结业务,只包一层 job runner。** + +### 8.1 直接复用的部分 + +`src/summary_mcp/workflows/article_summary.py` 已经做了: + +- 读取 extracted payload +- 根据 `selected_ids` 找条目 +- 调用 `run_loop_payload(...)` +- 渲染 Markdown +- 写到 `output_dir` +- 返回 `list[Path]` + +这些都继续用。 + +### 8.2 最小新增建议 + +建议只新增一层 runtime/service,例如: + +- `src/summary_mcp/runtime/article_summary_jobs.py` + +职责: + +- 生成 `job_id` +- 创建 job 目录 +- 初始化 `RunStore` +- 启动 runner 子进程 +- 提供 status/result 查询 +- 在 runner 里调用 `summarize_selected_articles` + +### 8.3 是否需要改 `summarize_selected_articles` + +最小版尽量少改,只建议加两类低风险增强: + +1. **可选输入校验增强** + - 如果 `selected_ids` 一个都匹配不到,显式报错 + - 避免“成功返回空列表”却让 job 看起来像成功 + +2. **可选 hook / telemetry(非必须)** + - 如果后面需要更细粒度写 stage 输出,可再加 + - 但第一版没必要为了观测性重构函数 + +结论: + +- 第一版优先保持 `summarize_selected_articles` 基本不动 +- job 层只把它作为黑盒业务函数调用 + +--- + +## 9. 建议新增代码组织 + +建议新增文件: + +```text +src/summary_mcp/runtime/article_summary_jobs.py +scripts/run_article_summary_job.py +``` + +### 9.1 `article_summary_jobs.py` + +建议包含: + +- `start_article_summary_job(...)` +- `run_article_summary_job(...)` +- `get_article_summary_job_status(...)` +- `get_article_summary_job_result(...)` +- 若干私有 helper: + - job id 生成 + - job dir 解析 + - result artifact 读取 + +### 9.2 `scripts/run_article_summary_job.py` + +作用:作为子进程 runner 入口。 + +例如: + +```bash +python scripts/run_article_summary_job.py --job-id article-summary-... +``` + +runner 只做一件事: + +- 根据 `job_id` 找到 job 目录和 `input.json` +- 真正执行 job + +这样避免在 `Popen("python -c ...")` 里塞长字符串,也方便本地调试。 + +--- + +## 10. 对 `reader-digest-flow` SOP 的影响 + +这块是重点,因为老大的实际痛点就在这。 + +### 10.1 现状 SOP + +当前 skill 已经有硬规则: + +- 优先走 MCP `generate_article_summaries` +- MCP timeout 时,立刻 fallback 到 reader 本地 `.venv` + +这个 fallback 现在是必要的,但它本质是在补“单篇总结没有正式异步接口”。 + +### 10.2 引入异步 job 后的推荐 SOP + +建议调整为: + +1. OpenClaw 在用户确认保留文章后 +2. 调 `start_article_summary_job` +3. 拿到 `job_id` +4. 轮询 `get_article_summary_job_status` +5. 成功后调 `get_article_summary_job_result` +6. 再继续 IMA 格式化与上传 + +### 10.3 对 skill 文档的影响 + +`reader-digest-flow` 需要后续补一条新规则: + +- 单篇总结的正式生产路径从“同步 MCP + 本地 fallback”升级为“异步 job + 状态轮询” +- 本地 `.venv` CLI fallback 仍保留,但降级为: + - job 启动失败 + - job runner 异常 + - reader 服务端出现系统性问题时的应急路径 + +### 10.4 用户体验改善 + +异步 job 后,OpenClaw 可给出更像正式系统的反馈: + +- “已启动单篇总结任务,正在生成” +- “当前状态:generate_markdown” +- “已完成,准备继续沉淀到 IMA” + +而不是现在的: + +- 直接卡住几十秒 +- 然后 tool timeout +- 再走一套 fallback + +--- + +## 11. 验证方案 + +本轮验证不要追求大而全,按 4 层做就够。 + +### 11.1 单元级验证 + +目标:确保 job 状态文件和结果文件行为正确。 + +建议覆盖: + +1. `start_article_summary_job` 能创建 job 目录与 `run-state.json` +2. 输入路径不存在时,启动直接失败 +3. `get_article_summary_job_status` 能正确读取状态 +4. job 成功后 `get_article_summary_job_result` 返回 `written_paths` +5. job 失败后状态为 `failed`,并带错误摘要 + +### 11.2 本地集成验证 + +用一个真实 extracted 文件跑: + +1. 启动 job +2. 轮询 status +3. 成功后检查: + - `result.json` 存在 + - Markdown 文件存在 + - 路径正确 + +### 11.3 OpenClaw 链路验证 + +在实际 `reader-digest-flow` 环节,用一篇已确认保留的文章做: + +1. 启动 async job +2. 等待成功 +3. 继续做 IMA markdown 整理与上传 +4. 确认整个链路不再因为同步 tool timeout 中断 + +### 11.4 异常验证 + +至少测 3 类异常: + +1. `selected_ids` 不存在 +2. LLM 接口失败 / 超时 +3. runner 进程异常退出 + +预期: + +- `run-state.json` 最终为 `failed` +- `error_summary` 可读 +- 调用方能明确知道失败,而不是只看到 transport timeout + +--- + +## 12. 风险与回滚方案 + +## 12.1 风险 + +### 风险 1:后台子进程成功启动,但状态长期卡在 running + +常见原因: + +- runner 进程崩了 +- server 重启时某些路径没写全 +- 子进程命令不对 + +缓解: + +- runner 启动前就先写 `run-state.json` +- runner 一进来先更新 `prepare_job` / `load_input` +- 后续可加“stale running 超时判定”,但第一版先不做自动修复 + +### 风险 2:复用旧函数导致“空输出也算成功” + +当前 `summarize_selected_articles()` 如果没匹配到文章或全部失败,存在返回空列表的可能。 + +缓解: + +- job runner 里把“`written_paths` 为空”视为失败 +- 或者顺手在 `summarize_selected_articles` 里补显式校验 + +### 风险 3:同一时刻大量 job 并发,LLM 调用被打爆 + +第一版不解决系统级并发控制。 + +缓解: + +- SOP 层先按单篇/少量串行用 +- skill 层避免一口气启动很多 job + +### 风险 4:状态目录与结果目录分离,排查时容易迷路 + +缓解: + +- 在 `result.json` 和 artifact 元数据里明确记录 `written_paths` +- 在 `run-state.input.output_dir` 中保留结果目录 + +--- + +## 12.2 回滚方案 + +这个方案很好回滚,因为是“新增,不替换”。 + +### 回滚原则 + +- 保留现有 `generate_article_summaries` +- 保留现有 `scripts/run_article_summaries.py` +- 新增 async job 接口如果不稳定,直接停止在 OpenClaw 层使用即可 + +### 回滚路径 + +1. 停止调用 `start_article_summary_job` +2. 恢复回原 SOP: + - 先尝试同步 MCP `generate_article_summaries` + - 失败则本地 `.venv` fallback +3. async job 相关代码保留但不作为正式入口 + +也就是说,这轮改动不会堵死当前生产路径,风险可控。 + +--- + +## 13. 实施顺序(建议按这个顺序落地) + +### Phase A:先把方案落地成最小代码骨架 + +1. 新增 `src/summary_mcp/runtime/article_summary_jobs.py` +2. 新增 job 目录常量与 helper +3. 新增 `scripts/run_article_summary_job.py` +4. 实现 runner 内部对 `summarize_selected_articles` 的调用 +5. 先本地命令验证 job 可跑通 + +### Phase B:再把 MCP 接口接上 + +6. 在 `server.py` 新增: + - `start_article_summary_job` + - `get_article_summary_job_status` + - `get_article_summary_job_result` +7. 本地 MCP 调用验证 + +### Phase C:最后接 OpenClaw SOP + +8. 更新 `reader-digest-flow` skill,把正式路径切到 async job +9. 保留本地 `.venv` fallback 作为应急方案 +10. 跑一次真实日报保留文章沉淀闭环 + +--- + +## 14. 推荐的最小返回/状态语义 + +为了跟现有 runtime 风格一致,建议继续用: + +- `status`: `running` / `success` / `failed` +- `current_stage`: 当前 stage 名 +- `artifacts`: 注册产物列表 +- `error_summary`: 结构化错误 + +第一版不必引入: + +- `queued` +- `cancelled` +- `retrying` +- `paused` + +避免状态机一开始就复杂化。 + +--- + +## 15. 结论 / 拍板建议 + +拍板建议很直接: + +1. **这件事值得做,而且优先级高**,因为它正好卡在当前日报 SOP 的真实痛点上 +2. **不建议继续先修同步 timeout**,因为同步模型本身就不适合几十秒级、外部 LLM 驱动的任务 +3. **最小真异步应采用“文件状态 + 后台子进程 + start/status/result 三接口”** +4. **业务层严格复用 `summarize_selected_articles`**,不要为异步化重写 summary 核心逻辑 +5. **先把单篇总结 job 做成独立小 runtime 命名空间**,不要急着并入 freshrss 主 run 的 resume 体系 + +如果只允许做一轮最小落地,我建议就做到: + +- `start_article_summary_job` +- `get_article_summary_job_status` +- `get_article_summary_job_result` +- `outputs/freshrss/article_summary_jobs//run-state.json` +- 子进程 runner + +这套已经足够把当前 OpenClaw timeout 问题从“同步等待”改成“正式异步轮询”,并且几乎不碰无关模块。 diff --git a/scripts/run_article_summary_job.py b/scripts/run_article_summary_job.py new file mode 100644 index 0000000..f9da3dd --- /dev/null +++ b/scripts/run_article_summary_job.py @@ -0,0 +1,24 @@ +from __future__ import annotations + +import argparse +import sys +from pathlib import Path + +REPO_ROOT = Path(__file__).resolve().parents[1] +SRC_ROOT = REPO_ROOT / "src" + +if str(SRC_ROOT) not in sys.path: + sys.path.insert(0, str(SRC_ROOT)) + +from summary_mcp.runtime.article_summary_jobs import run_article_summary_job + + +def main() -> None: + parser = argparse.ArgumentParser(description="Run a background article-summary job by job_id.") + parser.add_argument("--job-id", required=True, help="Article summary job id") + args = parser.parse_args() + run_article_summary_job(job_id=args.job_id) + + +if __name__ == "__main__": + main() diff --git a/src/summary_mcp/runtime/article_summary_jobs.py b/src/summary_mcp/runtime/article_summary_jobs.py new file mode 100644 index 0000000..630c3a0 --- /dev/null +++ b/src/summary_mcp/runtime/article_summary_jobs.py @@ -0,0 +1,315 @@ +from __future__ import annotations + +import json +import os +import subprocess +import sys +from datetime import datetime +from pathlib import Path +from typing import Any +from uuid import uuid4 + +from .run_store import RunStore +from .state_models import RunState + +REPO_ROOT = Path(__file__).resolve().parents[3] +OUTPUT_ROOT = REPO_ROOT / "outputs" / "freshrss" +ARTICLE_SUMMARY_JOBS_ROOT = OUTPUT_ROOT / "article_summary_jobs" +WORKFLOW_NAME = "article_summary_job" +RUN_TYPE = "article_summary" +RUN_STATE_FILENAME = "run-state.json" +INPUT_FILENAME = "input.json" +RESULT_FILENAME = "result.json" +JOB_REPORT_FILENAME = "job-report.json" +DEFAULT_STAGES = [ + "prepare_job", + "load_input", + "generate_markdown", + "write_result", +] + + +def _now() -> datetime: + return datetime.now().astimezone() + + +def _new_job_id() -> str: + ts = _now().strftime("%Y%m%d-%H%M%S") + return f"article-summary-{ts}-{uuid4().hex[:8]}" + + +def _job_dir(job_id: str) -> Path: + return ARTICLE_SUMMARY_JOBS_ROOT / job_id + + +def _run_state_path(job_id: str) -> Path: + return _job_dir(job_id) / RUN_STATE_FILENAME + + +def _input_path(job_id: str) -> Path: + return _job_dir(job_id) / INPUT_FILENAME + + +def _result_path(job_id: str) -> Path: + return _job_dir(job_id) / RESULT_FILENAME + + +def _job_report_path(job_id: str) -> Path: + return _job_dir(job_id) / JOB_REPORT_FILENAME + + +def _normalize_repo_path(path: Path) -> str: + try: + return str(path.resolve().relative_to(REPO_ROOT.resolve())) + except ValueError: + return str(path) + + +def _load_run_store(job_id: str) -> RunStore: + return RunStore.load(path=_run_state_path(job_id), repo_root=REPO_ROOT) + + +def _load_result(job_id: str) -> dict[str, Any] | None: + path = _result_path(job_id) + if not path.exists(): + return None + return json.loads(path.read_text(encoding="utf-8-sig")) + + +def start_article_summary_job( + *, + extracted_path: Path, + selected_ids: list[str], + output_dir: Path | None = None, + max_retries: int = 2, + timeout_seconds: float = 120.0, + llm_api_key: str | None = None, + llm_model: str | None = None, + llm_api_url: str | None = None, +) -> dict[str, Any]: + if not extracted_path.exists(): + raise FileNotFoundError(f"extracted_path does not exist: {extracted_path}") + if not selected_ids: + raise ValueError("selected_ids must not be empty") + + job_id = _new_job_id() + job_dir = _job_dir(job_id) + job_dir.mkdir(parents=True, exist_ok=True) + started_at = _now() + + resolved_output_dir = output_dir or (extracted_path.parent / "single_summaries") + input_payload = { + "extracted_path": str(extracted_path), + "selected_ids": selected_ids, + "output_dir": str(resolved_output_dir), + "max_retries": max_retries, + "timeout_seconds": timeout_seconds, + "llm_api_key": llm_api_key, + "llm_model": llm_model, + "llm_api_url": llm_api_url, + "launcher_pid": os.getpid(), + } + + store = RunStore.create( + path=_run_state_path(job_id), + run_id=job_id, + workflow=WORKFLOW_NAME, + run_type=RUN_TYPE, + started_at=started_at, + input_payload=input_payload, + repo_root=REPO_ROOT, + ) + for stage_name in DEFAULT_STAGES: + store._get_or_create_stage(stage_name) + store.save() + + store.start_stage("prepare_job") + + input_file = _input_path(job_id) + input_file.write_text(json.dumps(input_payload, ensure_ascii=False, indent=2), encoding="utf-8") + store.register_artifact(name="job_input", path=input_file, kind="json", stage="prepare_job") + + runner_script = REPO_ROOT / "scripts" / "run_article_summary_job.py" + cmd = [sys.executable, str(runner_script), "--job-id", job_id] + try: + proc = subprocess.Popen( + cmd, + cwd=str(REPO_ROOT), + stdout=subprocess.DEVNULL, + stderr=subprocess.DEVNULL, + start_new_session=True, + ) + except Exception as exc: + store.fail_stage("prepare_job", error=exc) + raise + + store.finish_stage("prepare_job", outputs={"runner_pid": proc.pid, "runner_command": cmd}) + + return { + "job_id": job_id, + "workflow": WORKFLOW_NAME, + "run_type": RUN_TYPE, + "status": "running", + "output_dir": _normalize_repo_path(job_dir), + "message": "Article summary job started successfully. Use get_article_summary_job_status to poll progress.", + } + + +def run_article_summary_job(*, job_id: str) -> dict[str, Any]: + from summary_mcp.workflows.article_summary import ArticleSummaryConfig, summarize_selected_articles + + store = _load_run_store(job_id) + input_payload = json.loads(_input_path(job_id).read_text(encoding="utf-8-sig")) + + try: + store.start_stage("load_input") + extracted_path = Path(input_payload["extracted_path"]) + output_dir = Path(input_payload["output_dir"]) + selected_ids = list(input_payload["selected_ids"]) + if not extracted_path.exists(): + raise FileNotFoundError(f"extracted_path does not exist: {extracted_path}") + if not selected_ids: + raise ValueError("selected_ids must not be empty") + store.finish_stage( + "load_input", + outputs={ + "extracted_path": _normalize_repo_path(extracted_path), + "selected_id_count": len(selected_ids), + "output_dir": _normalize_repo_path(output_dir), + }, + ) + + store.start_stage("generate_markdown") + config = ArticleSummaryConfig( + max_retries=int(input_payload.get("max_retries") or 2), + timeout_seconds=float(input_payload.get("timeout_seconds") or 120.0), + ) + written_paths = summarize_selected_articles( + extracted_path=extracted_path, + selected_ids=selected_ids, + output_dir=output_dir, + config=config, + api_key=input_payload.get("llm_api_key"), + model=input_payload.get("llm_model"), + api_url=input_payload.get("llm_api_url"), + ) + if not written_paths: + raise RuntimeError("Article summary job produced no Markdown outputs.") + normalized_paths = [_normalize_repo_path(Path(p)) for p in written_paths] + store.finish_stage( + "generate_markdown", + outputs={ + "written_count": len(written_paths), + "written_paths": normalized_paths, + }, + ) + for idx, p in enumerate(written_paths, start=1): + store.register_artifact( + name=f"summary_markdown_{idx}", + path=Path(p), + kind="markdown", + stage="generate_markdown", + metadata={"output_type": "single_summary"}, + ) + + store.start_stage("write_result") + result = { + "job_id": job_id, + "selected_ids": selected_ids, + "written_paths": normalized_paths, + "completed_at": _now().isoformat(), + } + result_file = _result_path(job_id) + result_file.write_text(json.dumps(result, ensure_ascii=False, indent=2), encoding="utf-8") + store.register_artifact(name="job_result", path=result_file, kind="json", stage="write_result") + report_file = _job_report_path(job_id) + report_file.write_text( + json.dumps( + { + "job_id": job_id, + "status": "success", + "written_count": len(normalized_paths), + "written_paths": normalized_paths, + }, + ensure_ascii=False, + indent=2, + ), + encoding="utf-8", + ) + store.register_artifact(name="job_report", path=report_file, kind="json", stage="write_result") + store.finish_stage("write_result", outputs={"result_path": _normalize_repo_path(result_file)}) + store.finish_run(status="success") + return result + except Exception as exc: + current_stage = store.state.current_stage or "generate_markdown" + store.fail_stage(current_stage, error=exc) + report_file = _job_report_path(job_id) + report_file.write_text( + json.dumps( + { + "job_id": job_id, + "status": "failed", + "error_type": type(exc).__name__, + "error_message": str(exc), + "failed_stage": current_stage, + }, + ensure_ascii=False, + indent=2, + ), + encoding="utf-8", + ) + try: + store.register_artifact(name="job_report", path=report_file, kind="json", stage=current_stage) + except Exception: + pass + raise + + +def get_article_summary_job_status(*, job_id: str) -> dict[str, Any]: + store = _load_run_store(job_id) + state = store.state + completed_stage_count = sum(1 for s in state.stages if s.status == "success") + running_stage_count = sum(1 for s in state.stages if s.status == "running") + failed_stage_count = sum(1 for s in state.stages if s.status == "failed") + pending_stage_count = sum(1 for s in state.stages if s.status == "pending") + return { + "job_id": state.run_id, + "workflow": state.workflow, + "run_type": state.run_type, + "status": state.status, + "current_stage": state.current_stage, + "started_at": state.started_at.isoformat(), + "updated_at": state.updated_at.isoformat(), + "finished_at": state.finished_at.isoformat() if state.finished_at else None, + "output_dir": _normalize_repo_path(_job_dir(job_id)), + "progress": { + "completed_stage_count": completed_stage_count, + "running_stage_count": running_stage_count, + "failed_stage_count": failed_stage_count, + "pending_stage_count": pending_stage_count, + "total_stage_count": len(state.stages), + }, + "artifacts": [artifact.model_dump(mode="json") for artifact in state.artifacts], + "error_summary": state.error.model_dump(mode="json") if state.error else None, + } + + +def get_article_summary_job_result(*, job_id: str) -> dict[str, Any]: + store = _load_run_store(job_id) + state = store.state + result = _load_result(job_id) + if state.status != "success" or result is None: + return { + "job_id": state.run_id, + "status": state.status, + "message": "Article summary job result is not ready.", + "error_summary": state.error.model_dump(mode="json") if state.error else None, + } + artifact = next((a.model_dump(mode="json") for a in state.artifacts if a.name == "job_result"), None) + return { + "job_id": state.run_id, + "status": state.status, + "written_paths": result.get("written_paths", []), + "artifact": artifact, + "result": result, + } diff --git a/src/summary_mcp/server.py b/src/summary_mcp/server.py index 6e02a6f..1b1a7e6 100644 --- a/src/summary_mcp/server.py +++ b/src/summary_mcp/server.py @@ -20,6 +20,11 @@ from summary_mcp.runtime import get_run_status as load_run_status from summary_mcp.runtime import get_run_report as load_run_report from summary_mcp.runtime import list_run_artifacts as load_run_artifacts from summary_mcp.runtime import list_runs as load_runs +from summary_mcp.runtime.article_summary_jobs import ( + get_article_summary_job_result as load_article_summary_job_result, + get_article_summary_job_status as load_article_summary_job_status, + start_article_summary_job as launch_article_summary_job, +) 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.article_summary import ArticleSummaryConfig, summarize_selected_articles @@ -164,6 +169,42 @@ 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() +def start_article_summary_job( + *, + extracted_path: str, + selected_ids: list[str], + output_dir: str | None = None, + max_retries: int = 2, + timeout_seconds: float = 120.0, + llm_api_key: str | None = None, + llm_model: str | None = None, + llm_api_url: str | None = None, +) -> dict: + """Start an asynchronous article-summary job and return a job_id immediately.""" + return launch_article_summary_job( + extracted_path=Path(extracted_path), + selected_ids=selected_ids, + output_dir=Path(output_dir) if output_dir else None, + max_retries=max_retries, + timeout_seconds=timeout_seconds, + llm_api_key=llm_api_key, + llm_model=llm_model, + llm_api_url=llm_api_url, + ) + + +@mcp.tool() +def get_article_summary_job_status(job_id: str) -> dict: + """Get the current status of an asynchronous article-summary job.""" + return load_article_summary_job_status(job_id=job_id) + + +@mcp.tool() +def get_article_summary_job_result(job_id: str) -> dict: + """Get the final result of an asynchronous article-summary job.""" + return load_article_summary_job_result(job_id=job_id) + @mcp.tool() def generate_article_summaries( diff --git a/src/summary_mcp/workflows/article_summary.py b/src/summary_mcp/workflows/article_summary.py index 98214d5..66b39c3 100644 --- a/src/summary_mcp/workflows/article_summary.py +++ b/src/summary_mcp/workflows/article_summary.py @@ -167,10 +167,12 @@ def summarize_selected_articles( output_dir.mkdir(parents=True, exist_ok=True) written_paths: list[Path] = [] + matched_ids: set[str] = set() from summary_mcp.core.summary_loop import build_summary_input for item_id, entry in _iter_selected_items(payload, selected_ids): + matched_ids.add(item_id) # Prefer the real project format where each entry has ``item`` and # ``extraction.article``; fall back to legacy layout where the # article fields live directly on the element. @@ -288,4 +290,7 @@ def summarize_selected_articles( output_path.write_text("\n".join(lines), encoding="utf-8") written_paths.append(output_path) + if not matched_ids: + raise ValueError(f"No extracted entries matched selected_ids: {list(selected_ids)}") + return written_paths