diff --git a/README.md b/README.md index 6f70f95..137c595 100644 --- a/README.md +++ b/README.md @@ -9,11 +9,13 @@ pip install -e . summary-mcp ``` -服务当前暴露 14 个工具。 +服务当前暴露 17 个工具。 正式 workflow service 相关工具: -- `run_freshrss_openclaw_pipeline` +- `start_freshrss_pipeline_job` +- `get_freshrss_pipeline_job_status` +- `get_freshrss_pipeline_job_result` - `get_run_status` - `list_runs` - `list_run_artifacts` @@ -23,6 +25,7 @@ summary-mcp 单步处理 / 调试相关工具: +- `run_freshrss_openclaw_pipeline`(同步 debug / fallback) - `extract_url_content` - `extract_item_content` - `filter_summary_result` @@ -35,19 +38,24 @@ summary-mcp - 当前正式 workflow 只有 `freshrss_daily_digest` - 每次 FreshRSS 主流水线 run 都会在 `outputs/freshrss/rerun//run-state.json` 落地运行真相 +- FreshRSS 主日报的正式生产启动路径已切到最小异步 job:`start_freshrss_pipeline_job` -> `get_freshrss_pipeline_job_status` -> `get_freshrss_pipeline_job_result` - OpenClaw 正式读取结果应优先使用 `get_delivery_payload` 与 `get_run_report`,而不是自己拼输出目录路径 - `digest-brief.json` 当前会随主流水线产出,但还没有独立的 MCP 读取工具;如需定位它,应通过 `list_run_artifacts` 或 `get_run_report` 返回的信息发现 -- `run_freshrss_openclaw_pipeline` 与 `resume_run` 当前都是同步 MCP 调用;仓库里还没有针对 FreshRSS 主流程的后台队列 / worker / 异步任务管理 +- `run_freshrss_openclaw_pipeline` 仍保留,但定位已降级为同步 debug / fallback 路径,不再是 OpenClaw 的默认生产启动入口 - 单篇总结已补上最小异步 job 形态:`start_article_summary_job` / `get_article_summary_job_status` / `get_article_summary_job_result` - `generate_article_summaries` 仍保留,但定位是同步 debug 路径,而不是 OpenClaw 的正式生产集成入口 ## 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 +1. 调用 `start_freshrss_pipeline_job` 启动正式日报 job,并保存返回的 `job_id` +2. 轮询 `get_freshrss_pipeline_job_status(job_id)`,直到 `status` 变成 `success` 或 `failed` +3. 成功后调用 `get_freshrss_pipeline_job_result(job_id)` 读取 `run_id` 与关键产物路径 +4. 后续所有 run 级状态判断都基于 `get_run_status(run_id)` 或 `list_runs(...)` +5. 需要看产物列表时用 `list_run_artifacts(run_id)`,不要在 OpenClaw 里硬编码 `outputs/freshrss/rerun/...` +6. 需要消费正式结果时优先用 `get_delivery_payload(run_id)` 与 `get_run_report(run_id)` +7. 仅当 `resume_run` 的最小恢复范围满足时,才对失败 run 调用 `resume_run(run_id)`;否则应重启一个新 run + +主日报 async job 的自身状态目录固定在 `outputs/freshrss/pipeline_jobs//`,最少包含 `run-state.json`、`input.json`、`result.json`(成功时)和 `job-report.json`。 ## `resume_run` 当前最小范围 @@ -154,7 +162,7 @@ python scripts/run_freshrss_pipeline.py ^ --mark-read ``` -这条 CLI 与 MCP `run_freshrss_openclaw_pipeline` 共用同一条主流水线逻辑,但正式生产集成应优先走 MCP;CLI 仅用于本地 debug / fallback。默认会写出这些产物: +这条 CLI 与 MCP `run_freshrss_openclaw_pipeline` / `start_freshrss_pipeline_job` 共用同一条主流水线逻辑,但正式生产集成应优先走 async MCP job;CLI 与同步 MCP 入口仅用于本地 debug / fallback。默认会写出这些产物: - `outputs/freshrss/rerun//run-state.json` - `outputs/freshrss/rerun//raw/freshrss.raw.json` @@ -171,7 +179,7 @@ python scripts/run_freshrss_pipeline.py ^ 主流水线默认**不会**产出批量级的 `freshrss.extracted.json`。 如果你需要更多逐条中间产物,例如标准化 items、摘要结果、过滤决策、candidate record、candidate input,可以加 `--debug-artifacts`。 -当 OpenClaw 接入这个 MCP 服务后,应以 `run_id` 作为稳定句柄:先调用 `run_freshrss_openclaw_pipeline`,再通过 `get_run_status` / `get_delivery_payload` / `get_run_report` 读取状态与结果,而不是直接拼接目录路径。 +当 OpenClaw 接入这个 MCP 服务后,应先通过 `start_freshrss_pipeline_job` 启动任务,轮询 `get_freshrss_pipeline_job_status`,再从 `get_freshrss_pipeline_job_result` 读取稳定的 `run_id`。拿到 `run_id` 之后,再通过 `get_run_status` / `get_delivery_payload` / `get_run_report` 读取状态与结果,而不是直接拼接目录路径。`run_freshrss_openclaw_pipeline` 仅保留为同步 debug / fallback 路径。 在排查复杂问题时,也可以把 `debug_artifacts=true` 打开,并结合 `list_run_artifacts` 查看该 run 下实际产物。 ## 对结构化摘要结果执行确定性过滤规则 diff --git a/TODO.md b/TODO.md index 0ae5613..017aa84 100644 --- a/TODO.md +++ b/TODO.md @@ -187,7 +187,7 @@ --- -### [DOING][P1] 单篇总结改为最小真异步 job +### [DONE][P1] 单篇总结改为最小真异步 job 目标: - 解决 `generate_article_summaries` 在 OpenClaw → MCP 同步链路里易 timeout 的问题 @@ -228,6 +228,41 @@ --- +### [DOING][P0] FreshRSS 主日报 run 改为最小真异步 job + +目标: +- 解决 `run_freshrss_openclaw_pipeline` 在正式生产链路里仍为同步 MCP 调用、易超时的问题 +- 将 FreshRSS 主日报启动路径升级为可启动、可轮询、可读取结果的异步 job + +要求: +- 新增最小异步接口: + - `start_freshrss_pipeline_job` + - `get_freshrss_pipeline_job_status` + - `get_freshrss_pipeline_job_result` +- 状态目录固定落到: + - `outputs/freshrss/pipeline_jobs//` +- 至少包含: + - `run-state.json` + - `input.json` + - `result.json`(成功时) + - `job-report.json` +- 执行模型优先使用后台子进程,不使用线程 +- 主业务逻辑继续复用 `run_freshrss_pipeline(...)`,不要重写日报核心逻辑 +- job 成功后结果中必须带回 `run_id` 与关键产物路径 +- README / handoff / OpenClaw 生产建议路径需要同步改成 async start path + +当前进展: +- 2026-04-11:问题定位完成,确认之前异步化的是 article-summary,不是主日报 run +- 2026-04-11:已新增方案文档 `plans/freshrss-pipeline-async-job-plan.md` + +下一步: +- 复用 article-summary job runtime 骨架实现主日报 async job +- 新增后台 runner 脚本 +- 暴露 3 个 MCP tools +- 用真实 MCP 冒烟验证 `start -> status -> result` + +--- + ### [TODO][P2] 评估 `rerun_stage` 是否值得进入第一阶段 目标: diff --git a/docs/openclaw/openclaw-handoff.md b/docs/openclaw/openclaw-handoff.md index acbf0de..9cca94f 100644 --- a/docs/openclaw/openclaw-handoff.md +++ b/docs/openclaw/openclaw-handoff.md @@ -27,21 +27,27 @@ This repository should **not** take over downstream orchestration responsibiliti ## Production Entrypoint -reader 当前正式工作流服务入口是 MCP tool: +reader 当前正式工作流服务启动入口是 MCP tool: -- `run_freshrss_openclaw_pipeline` +- `start_freshrss_pipeline_job` -It starts the only formally supported workflow today: `freshrss_daily_digest`. +OpenClaw 应先拿到 `job_id`,轮询 job 状态,再在成功后读取 `run_id` 作为正式后续句柄。 -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. +`run_freshrss_openclaw_pipeline` 仍保留,但定位是同步 debug / fallback 路径,而不是正式生产启动入口。 + +OpenClaw should treat the returned `run_id` from `get_freshrss_pipeline_job_result` as the only stable handle for follow-up reads. Do not hand-build `outputs/freshrss/rerun/...` paths in OpenClaw. + +Job state is written under `outputs/freshrss/pipeline_jobs//` and will minimally contain `run-state.json`, `input.json`, `result.json` on success, and `job-report.json`. ## Supported MCP Tools -Current MCP tools: 14 total, including the FreshRSS workflow set plus article-summary async job tools. +Current MCP tools: 17 total, including the FreshRSS workflow set plus async job tools for both the main pipeline and article-summary flow. Workflow service tools: -- `run_freshrss_openclaw_pipeline` +- `start_freshrss_pipeline_job` +- `get_freshrss_pipeline_job_status` +- `get_freshrss_pipeline_job_result` - `get_run_status` - `list_runs` - `list_run_artifacts` @@ -57,6 +63,7 @@ Article-summary tools: Single-step / debug tools: +- `run_freshrss_openclaw_pipeline`(同步模式,仅适合 debug / fallback) - `extract_url_content` - `extract_item_content` - `filter_summary_result` @@ -118,12 +125,14 @@ summary-mcp 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 +1. Call `start_freshrss_pipeline_job` and persist the returned `job_id` +2. Poll `get_freshrss_pipeline_job_status(job_id)` until `status` becomes `success` or `failed` +3. On success, call `get_freshrss_pipeline_job_result(job_id)` and persist the returned `run_id` +4. Use `get_run_status(run_id)` as the authoritative run-state read for status, stage, artifacts, and recovery +5. Use `list_runs(...)` when OpenClaw needs recent-run discovery or high-level inspection +6. Use `list_run_artifacts(run_id)` when OpenClaw needs to inspect what this run actually produced +7. Use `get_delivery_payload(run_id)` and `get_run_report(run_id)` as the formal result-reading APIs +8. Use `resume_run(run_id)` only when the run falls inside the minimal supported resume scope OpenClaw should not directly derive or hardcode: @@ -154,8 +163,9 @@ Recommended semantics: - Use `mark_read=false` only for debug, test, or validation runs. - Keep `debug_artifacts=false` for routine production runs. - 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. +- Treat the returned `job_id` as the startup handle, and the later `run_id` from `get_freshrss_pipeline_job_result` as the stable identifier for all follow-up run reads. - If no real `openclaw-delivery-payload.json` was produced, OpenClaw should stop instead of generating a digest from placeholders or examples. +- Use `run_freshrss_openclaw_pipeline` only when a synchronous debug / fallback path is explicitly needed. ## Formal Capability Boundary @@ -163,10 +173,12 @@ reader 当前正式 MCP workflow service 的边界如下: - formal workflow: only `freshrss_daily_digest` - run truth: every FreshRSS run writes `run-state.json` +- main production start path: `start_freshrss_pipeline_job` / `get_freshrss_pipeline_job_status` / `get_freshrss_pipeline_job_result` - 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 for the FreshRSS daily workflow yet +- `run_freshrss_openclaw_pipeline` is still supported, but only as a synchronous debug / fallback path +- the FreshRSS daily workflow now has a minimal background job model backed by a detached runner process, not a full queue / worker system - 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 @@ -175,10 +187,11 @@ 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 Async Job Returns -Primary return fields from `run_freshrss_openclaw_pipeline`: +Primary return fields from `get_freshrss_pipeline_job_result`: +- `job_id` - `run_id` - `output_dir` - `raw_output` @@ -189,7 +202,6 @@ Primary return fields from `run_freshrss_openclaw_pipeline`: - `delivered_count` - `marked_read_count` - `status_counts` -- `delivery_payload` - `keyword_index` Optional: @@ -197,6 +209,8 @@ Optional: - `items` - returned only when `include_item_reports=true` +`run_freshrss_openclaw_pipeline` still returns the same synchronous payload for debug / fallback use. + Follow-up structured reads should use MCP tools rather than re-reading these files directly. ## Minimal Output Files diff --git a/plans/freshrss-pipeline-async-job-plan.md b/plans/freshrss-pipeline-async-job-plan.md new file mode 100644 index 0000000..b766394 --- /dev/null +++ b/plans/freshrss-pipeline-async-job-plan.md @@ -0,0 +1,229 @@ +# FreshRSS 主日报异步 job 方案 + +## 背景 + +当前 `run_freshrss_openclaw_pipeline` 虽然已经作为正式 MCP workflow 入口存在,但执行模型仍是**同步 MCP 调用**。这会带来几个现实问题: + +1. OpenClaw / MCP wrapper 存在超时风险,尤其是 5-10 篇的正式日报批次。 +2. wrapper timeout 与真实 run 是否已落地,容易出现语义分离。 +3. 当前已有 `run-state.json`、`get_run_status`、`get_run_report`、`get_delivery_payload`,但**启动层**仍然是同步调用,不利于正式生产链路稳定运行。 +4. 单篇总结已经验证了“最小 async job + 轮询状态 + 读取结果”模型可行,主日报 run 应收敛到同一套运行模式。 + +## 目标 + +将 FreshRSS 主日报 run 改造成与 article-summary 类似的**最小真异步 job**: + +- 启动即返回 `job_id` +- 真正执行由后台子进程完成 +- 状态可轮询 +- 成功后可读取结构化结果 +- 业务逻辑继续复用既有 `run_freshrss_pipeline(...)` +- 不推翻现有 run-state / result query 能力 + +## 非目标 + +本阶段不做: + +- 分布式任务队列 +- 多 worker 调度 +- 任意 stage 的后台恢复编排 +- 并发控制中心 +- 主流程与 article-summary job 的通用抽象框架一次性大重构 + +先做最小可用。 + +## 设计原则 + +1. **启动层异步化,执行核心不重写** + - `run_freshrss_pipeline(...)` 继续是主业务逻辑真相。 + - async job 只负责启动、状态持久化、结果回读。 + +2. **run truth 与 job truth 分层** + - job truth:这次异步任务有没有启动、运行到哪一步、是否成功。 + - run truth:真正的 freshrss workflow 输出与 `run-state.json`。 + +3. **OpenClaw 正式生产默认改为 async start path** + - 启动走 async job + - 状态和结果优先先看 job + - 真正业务产物仍由现有 run 查询工具承接 + +4. **与 article-summary job 尽量同构** + - 目录结构 + - `run-state.json` / `input.json` / `result.json` / `job-report.json` + - 后台 runner 脚本 + +## 拟新增能力 + +### MCP tools + +新增 3 个工具: + +- `start_freshrss_pipeline_job` +- `get_freshrss_pipeline_job_status` +- `get_freshrss_pipeline_job_result` + +### job 目录 + +固定目录: + +`outputs/freshrss/pipeline_jobs//` + +至少包含: + +- `run-state.json` +- `input.json` +- `result.json`(成功时) +- `job-report.json` + +### 执行模型 + +- `start_freshrss_pipeline_job` 写入 input + 初始化 job state +- 后台 `subprocess.Popen(...)` 启动 runner +- runner 内部调用 `run_freshrss_pipeline(...)` +- 成功后把 `run_id`、核心产物路径、关键计数写入 `result.json` + +## job 输入参数 + +与现有 `run_freshrss_openclaw_pipeline` 尽量对齐: + +- `limit` +- `mark_read` +- `include_read` +- `debug_artifacts` +- `continuation` +- `timeout_seconds` +- `max_retries` +- `stream_id` +- `api_base_url` +- `username` +- `api_password` +- `llm_api_key` +- `llm_model` +- `llm_api_url` +- `context` +- `run_id` +- `date_value` +- `output_dir` +- `include_item_reports` + +## 返回语义 + +### start + +返回: + +- `job_id` +- `workflow` +- `run_type` +- `status=running` +- `output_dir` +- `message` + +### status + +返回: + +- `job_id` +- `status` +- `current_stage` +- `started_at` / `updated_at` / `finished_at` +- `progress` +- `artifacts` +- `error_summary` +- 若主 run 已创建,可附带 `linked_run_id` + +### result + +成功时返回: + +- `job_id` +- `status=success` +- `run_id` +- `delivery_output` +- `report_output` +- `digest_brief_output` +- `pulled_count` +- `delivered_count` +- `marked_read_count` +- `artifact` +- `result` + +## stages 建议 + +最小 job stages: + +1. `prepare_job` +2. `load_input` +3. `run_pipeline` +4. `write_result` + +其中 `run_pipeline` 内部仍由现有 freshrss workflow 自己写它的 run-state。 + +## 与现有同步入口的关系 + +### 保留 + +`run_freshrss_openclaw_pipeline` 暂时保留,作为: + +- debug / light path +- 本地调试工具 +- 向后兼容路径 + +### 正式语义调整 + +文档与 OpenClaw handoff 中,主日报正式生产默认启动入口改为: + +- `start_freshrss_pipeline_job` + +同步入口降级为: + +- debug / fallback +- 小批量验证 + +## OpenClaw 编排建议 + +新的推荐路径: + +1. `start_freshrss_pipeline_job` +2. `get_freshrss_pipeline_job_status` +3. 成功后 `get_freshrss_pipeline_job_result` +4. 后续仍用: + - `get_run_status` + - `get_delivery_payload` + - `get_run_report` + - `list_run_artifacts` + +## 风险点 + +1. **job 成功但 run 部分失败** + - 允许,job 结果应以真实 `run_freshrss_pipeline(...)` 返回为准。 + - `run_id` + `report_output` 仍是最终真相。 + +2. **runner 崩溃但来不及写 result** + - 需保证 `job-report.json` 至少能写下失败摘要。 + +3. **重复状态源导致混淆** + - 文档必须明确: + - job state 管“启动任务” + - run state 管“业务工作流真相” + +4. **同步 / 异步双入口长期漂移** + - 必须要求 async job 内部直接复用 `run_freshrss_pipeline(...)` + - 禁止再实现一套平行主流程 + +## 验收标准 + +1. 能通过 MCP 启动一个主日报 async job 并立即返回 `job_id` +2. 能轮询到 `running -> success/failed` +3. 成功后 `result.json` 含 `run_id` 与核心产物路径 +4. 对应 run 仍能通过既有 `get_run_status` / `get_run_report` / `get_delivery_payload` 正常读取 +5. README / handoff / TODO / plans 同步更新 + +## 建议实施顺序 + +1. 复制 article-summary job 骨架到 freshrss pipeline job +2. 新增 runner 脚本 +3. server.py 暴露 3 个新工具 +4. 补 query/result 读法 +5. 更新 README / handoff +6. 将 TODO 主任务切到“主日报 async job” diff --git a/scripts/run_freshrss_pipeline_job.py b/scripts/run_freshrss_pipeline_job.py new file mode 100644 index 0000000..b53466c --- /dev/null +++ b/scripts/run_freshrss_pipeline_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.freshrss_pipeline_jobs import run_freshrss_pipeline_job + + +def main() -> None: + parser = argparse.ArgumentParser(description="Run a background FreshRSS pipeline job by job_id.") + parser.add_argument("--job-id", required=True, help="FreshRSS pipeline job id") + args = parser.parse_args() + run_freshrss_pipeline_job(job_id=args.job_id) + + +if __name__ == "__main__": + main() diff --git a/src/summary_mcp/runtime/__init__.py b/src/summary_mcp/runtime/__init__.py index 33f2c2e..d6c3766 100644 --- a/src/summary_mcp/runtime/__init__.py +++ b/src/summary_mcp/runtime/__init__.py @@ -1,3 +1,8 @@ +from .freshrss_pipeline_jobs import ( + get_freshrss_pipeline_job_result, + get_freshrss_pipeline_job_status, + start_freshrss_pipeline_job, +) from .query_service import get_delivery_payload, get_run_report, get_run_status, list_run_artifacts, list_runs from .run_store import RunStore from .state_models import ArtifactRecord, RecoveryState, RunError, RunState, StageState @@ -5,6 +10,8 @@ from .state_models import ArtifactRecord, RecoveryState, RunError, RunState, Sta __all__ = [ "ArtifactRecord", "get_delivery_payload", + "get_freshrss_pipeline_job_result", + "get_freshrss_pipeline_job_status", "get_run_report", "RecoveryState", "RunError", @@ -14,4 +21,5 @@ __all__ = [ "get_run_status", "list_run_artifacts", "list_runs", + "start_freshrss_pipeline_job", ] diff --git a/src/summary_mcp/runtime/freshrss_pipeline_jobs.py b/src/summary_mcp/runtime/freshrss_pipeline_jobs.py new file mode 100644 index 0000000..9851390 --- /dev/null +++ b/src/summary_mcp/runtime/freshrss_pipeline_jobs.py @@ -0,0 +1,432 @@ +from __future__ import annotations + +import json +import os +import subprocess +import sys +from datetime import UTC, datetime +from pathlib import Path +from typing import Any +from uuid import uuid4 + +from .run_store import RunStore + +REPO_ROOT = Path(__file__).resolve().parents[3] +OUTPUT_ROOT = REPO_ROOT / "outputs" / "freshrss" +PIPELINE_JOBS_ROOT = OUTPUT_ROOT / "pipeline_jobs" +WORKFLOW_NAME = "freshrss_pipeline_job" +RUN_TYPE = "daily_digest" +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", + "run_pipeline", + "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"freshrss-pipeline-job-{ts}-{uuid4().hex[:8]}" + + +def _job_dir(job_id: str) -> Path: + return PIPELINE_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 _normalize_repo_path_value(path_value: str | None) -> str | None: + if not path_value: + return None + return _normalize_repo_path(Path(path_value)) + + +def _normalize_keyword_index(keyword_index: dict[str, Any] | None) -> dict[str, Any]: + if not isinstance(keyword_index, dict): + return {} + normalized = dict(keyword_index) + for field in ("daily_output", "stats_output"): + normalized[field] = _normalize_repo_path_value(normalized.get(field)) + return normalized + + +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 _build_run_defaults(*, run_id: str | None, output_dir: Path | None) -> tuple[str, Path]: + resolved_output_dir = output_dir + resolved_run_id = run_id + if resolved_run_id is not None and resolved_output_dir is not None: + return resolved_run_id, resolved_output_dir + + run_stamp = datetime.now(tz=UTC).strftime("%Y%m%d-%H%M%S") + if resolved_run_id is None: + resolved_run_id = f"freshrss-pipeline-{run_stamp}" + if resolved_output_dir is None: + resolved_output_dir = OUTPUT_ROOT / "rerun" / run_stamp + return resolved_run_id, resolved_output_dir + + +def _build_result_payload(*, job_id: str, pipeline_result: dict[str, Any], include_item_reports: bool) -> dict[str, Any]: + result = { + "job_id": job_id, + "run_id": pipeline_result["run_id"], + "output_dir": _normalize_repo_path_value(pipeline_result.get("output_dir")), + "raw_output": _normalize_repo_path_value(pipeline_result.get("raw_output")), + "delivery_output": _normalize_repo_path_value(pipeline_result.get("delivery_output")), + "digest_brief_output": _normalize_repo_path_value(pipeline_result.get("digest_brief_output")), + "report_output": _normalize_repo_path_value(pipeline_result.get("report_output")), + "keyword_index": _normalize_keyword_index(pipeline_result.get("keyword_index")), + "pulled_count": pipeline_result.get("pulled_count"), + "delivered_count": pipeline_result.get("delivered_count"), + "marked_read_count": pipeline_result.get("marked_read_count"), + "status_counts": pipeline_result.get("status_counts"), + "debug_artifacts": pipeline_result.get("debug_artifacts"), + "completed_at": _now().isoformat(), + } + if include_item_reports: + result["items"] = pipeline_result.get("items", []) + return result + + +def _write_job_report( + *, + job_id: str, + payload: dict[str, Any], +) -> Path: + report_file = _job_report_path(job_id) + report_file.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8") + return report_file + + +def start_freshrss_pipeline_job( + *, + limit: int = 5, + mark_read: bool = False, + include_read: bool = False, + debug_artifacts: bool = False, + continuation: str | None = None, + timeout_seconds: float = 60.0, + max_retries: int = 2, + stream_id: str = "user/-/state/com.google/reading-list", + api_base_url: str | None = None, + username: str | None = None, + api_password: str | None = None, + llm_api_key: str | None = None, + llm_model: str | None = None, + llm_api_url: str | None = None, + context: dict[str, Any] | None = None, + run_id: str | None = None, + date_value: str | None = None, + output_dir: Path | None = None, + include_item_reports: bool = False, +) -> dict[str, Any]: + job_id = _new_job_id() + job_dir = _job_dir(job_id) + job_dir.mkdir(parents=True, exist_ok=True) + started_at = _now() + resolved_run_id, resolved_output_dir = _build_run_defaults(run_id=run_id, output_dir=output_dir) + + input_payload = { + "limit": limit, + "mark_read": mark_read, + "include_read": include_read, + "debug_artifacts": debug_artifacts, + "continuation": continuation, + "timeout_seconds": timeout_seconds, + "max_retries": max_retries, + "stream_id": stream_id, + "api_base_url": api_base_url, + "username": username, + "api_password": api_password, + "llm_api_key": llm_api_key, + "llm_model": llm_model, + "llm_api_url": llm_api_url, + "context": context, + "run_id": resolved_run_id, + "date_value": date_value, + "output_dir": str(resolved_output_dir), + "include_item_reports": include_item_reports, + "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_freshrss_pipeline_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: + report_file = _write_job_report( + job_id=job_id, + payload={ + "job_id": job_id, + "status": "failed", + "error_type": type(exc).__name__, + "error_message": str(exc), + "failed_stage": "prepare_job", + "linked_run_id": resolved_run_id, + "linked_output_dir": _normalize_repo_path(resolved_output_dir), + }, + ) + store.register_artifact(name="job_report", path=report_file, kind="json", stage="prepare_job") + store.fail_stage("prepare_job", error=exc) + raise + + store.finish_stage( + "prepare_job", + outputs={ + "runner_pid": proc.pid, + "runner_command": cmd, + "linked_run_id": resolved_run_id, + "linked_output_dir": _normalize_repo_path(resolved_output_dir), + }, + ) + + return { + "job_id": job_id, + "workflow": WORKFLOW_NAME, + "run_type": RUN_TYPE, + "status": "running", + "output_dir": _normalize_repo_path(job_dir), + "linked_run_id": resolved_run_id, + "linked_output_dir": _normalize_repo_path(resolved_output_dir), + "message": "FreshRSS pipeline job started successfully. Use get_freshrss_pipeline_job_status to poll progress.", + } + + +def run_freshrss_pipeline_job(*, job_id: str) -> dict[str, Any]: + from datetime import date + + from summary_mcp.workflows import run_freshrss_pipeline + + store = _load_run_store(job_id) + input_payload = json.loads(_input_path(job_id).read_text(encoding="utf-8-sig")) + current_stage = "run_pipeline" + + try: + store.start_stage("load_input") + resolved_output_dir = Path(input_payload["output_dir"]) + resolved_run_id = str(input_payload["run_id"]) + store.finish_stage( + "load_input", + outputs={ + "linked_run_id": resolved_run_id, + "linked_output_dir": _normalize_repo_path(resolved_output_dir), + "limit": int(input_payload.get("limit") or 5), + }, + ) + + store.start_stage("run_pipeline") + pipeline_result = run_freshrss_pipeline( + api_base_url=input_payload.get("api_base_url"), + username=input_payload.get("username"), + api_password=input_payload.get("api_password"), + stream_id=input_payload.get("stream_id") or "user/-/state/com.google/reading-list", + limit=int(input_payload.get("limit") or 5), + continuation=input_payload.get("continuation"), + include_read=bool(input_payload.get("include_read")), + mark_read=bool(input_payload.get("mark_read")), + debug_artifacts=bool(input_payload.get("debug_artifacts")), + context=input_payload.get("context"), + max_retries=int(input_payload.get("max_retries") or 2), + timeout_seconds=float(input_payload.get("timeout_seconds") or 60.0), + llm_api_key=input_payload.get("llm_api_key"), + llm_model=input_payload.get("llm_model"), + llm_api_url=input_payload.get("llm_api_url"), + run_id=resolved_run_id, + delivery_date=date.fromisoformat(input_payload["date_value"]) if input_payload.get("date_value") else None, + output_dir=resolved_output_dir, + ) + store.finish_stage( + "run_pipeline", + outputs={ + "linked_run_id": pipeline_result.get("run_id"), + "linked_output_dir": _normalize_repo_path_value(pipeline_result.get("output_dir")), + "delivery_output": _normalize_repo_path_value(pipeline_result.get("delivery_output")), + "report_output": _normalize_repo_path_value(pipeline_result.get("report_output")), + "delivered_count": pipeline_result.get("delivered_count"), + "marked_read_count": pipeline_result.get("marked_read_count"), + }, + ) + + store.start_stage("write_result") + result = _build_result_payload( + job_id=job_id, + pipeline_result=pipeline_result, + include_item_reports=bool(input_payload.get("include_item_reports")), + ) + 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 = _write_job_report( + job_id=job_id, + payload={ + "job_id": job_id, + "status": "success", + "run_id": result["run_id"], + "delivery_output": result["delivery_output"], + "report_output": result["report_output"], + "digest_brief_output": result["digest_brief_output"], + "pulled_count": result["pulled_count"], + "delivered_count": result["delivered_count"], + "marked_read_count": result["marked_read_count"], + }, + ) + 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), + "linked_run_id": result["run_id"], + }, + ) + store.finish_run(status="success") + return result + except Exception as exc: + current_stage = store.state.current_stage or current_stage + report_file = _write_job_report( + job_id=job_id, + payload={ + "job_id": job_id, + "status": "failed", + "error_type": type(exc).__name__, + "error_message": str(exc), + "failed_stage": current_stage, + "linked_run_id": input_payload.get("run_id"), + "linked_output_dir": _normalize_repo_path_value(input_payload.get("output_dir")), + }, + ) + try: + store.register_artifact(name="job_report", path=report_file, kind="json", stage=current_stage) + except Exception: + pass + store.fail_stage(current_stage, error=exc) + raise + + +def get_freshrss_pipeline_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") + linked_run_id = state.input.get("run_id") if isinstance(state.input, dict) else None + linked_output_dir = _normalize_repo_path_value(state.input.get("output_dir")) if isinstance(state.input, dict) else None + 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)), + "linked_run_id": linked_run_id, + "linked_output_dir": linked_output_dir, + "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_freshrss_pipeline_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": "FreshRSS pipeline job result is not ready.", + "linked_run_id": state.input.get("run_id") if isinstance(state.input, dict) else None, + "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, + "run_id": result.get("run_id"), + "output_dir": result.get("output_dir"), + "raw_output": result.get("raw_output"), + "delivery_output": result.get("delivery_output"), + "digest_brief_output": result.get("digest_brief_output"), + "report_output": result.get("report_output"), + "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"), + "keyword_index": result.get("keyword_index"), + "artifact": artifact, + "result": result, + } diff --git a/src/summary_mcp/server.py b/src/summary_mcp/server.py index 1b1a7e6..2b9d0ed 100644 --- a/src/summary_mcp/server.py +++ b/src/summary_mcp/server.py @@ -1,7 +1,7 @@ from __future__ import annotations # MCP 服务入口:将内容提取、过滤、FreshRSS 全链路管道暴露为 MCP 工具。 -# 生产主入口是 run_freshrss_openclaw_pipeline,其余工具供单步调试使用。 +# 主日报正式启动入口是 start_freshrss_pipeline_job;同步入口仅保留给 debug / fallback。 from datetime import date from pathlib import Path @@ -16,10 +16,13 @@ from summary_mcp.models.item import Item from summary_mcp.models.llm_result import LlmSummaryResult from summary_mcp.models.summary_io import ExtractionInput from summary_mcp.runtime import get_delivery_payload as load_delivery_payload +from summary_mcp.runtime import get_freshrss_pipeline_job_result as load_freshrss_pipeline_job_result +from summary_mcp.runtime import get_freshrss_pipeline_job_status as load_freshrss_pipeline_job_status 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 import start_freshrss_pipeline_job as launch_freshrss_pipeline_job 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, @@ -130,6 +133,64 @@ def run_freshrss_openclaw_pipeline( return result +@mcp.tool() +def start_freshrss_pipeline_job( + limit: int = 5, + mark_read: bool = False, + include_read: bool = False, + debug_artifacts: bool = False, + continuation: str | None = None, + timeout_seconds: float = 60.0, + max_retries: int = 2, + stream_id: str = "user/-/state/com.google/reading-list", + api_base_url: str | None = None, + username: str | None = None, + api_password: str | None = None, + llm_api_key: str | None = None, + llm_model: str | None = None, + llm_api_url: str | None = None, + context: dict | None = None, + run_id: str | None = None, + date_value: str | None = None, + output_dir: str | None = None, + include_item_reports: bool = False, +) -> dict: + """Start an asynchronous FreshRSS pipeline job and return a job_id immediately.""" + return launch_freshrss_pipeline_job( + limit=limit, + mark_read=mark_read, + include_read=include_read, + debug_artifacts=debug_artifacts, + continuation=continuation, + timeout_seconds=timeout_seconds, + max_retries=max_retries, + stream_id=stream_id, + api_base_url=api_base_url, + username=username, + api_password=api_password, + llm_api_key=llm_api_key, + llm_model=llm_model, + llm_api_url=llm_api_url, + context=context, + run_id=run_id, + date_value=date_value, + output_dir=Path(output_dir) if output_dir else None, + include_item_reports=include_item_reports, + ) + + +@mcp.tool() +def get_freshrss_pipeline_job_status(job_id: str) -> dict: + """Get the current status of an asynchronous FreshRSS pipeline job.""" + return load_freshrss_pipeline_job_status(job_id=job_id) + + +@mcp.tool() +def get_freshrss_pipeline_job_result(job_id: str) -> dict: + """Get the final result of an asynchronous FreshRSS pipeline job.""" + return load_freshrss_pipeline_job_result(job_id=job_id) + + @mcp.tool() def get_run_status(run_id: str) -> dict: """Get the current status of a workflow run by run_id.""" @@ -169,6 +230,7 @@ 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( *,