From 72a6853c0311c88eef34bb543b8729d06a3f98fe Mon Sep 17 00:00:00 2001 From: root Date: Tue, 7 Apr 2026 10:48:00 +0800 Subject: [PATCH] Add run-state persistence for FreshRSS pipeline --- TODO.md | 186 +++-- src/summary_mcp/runtime/__init__.py | 11 + src/summary_mcp/runtime/run_store.py | 185 +++++ src/summary_mcp/runtime/state_models.py | 58 ++ .../workflows/freshrss_pipeline.py | 670 +++++++++++------- 5 files changed, 802 insertions(+), 308 deletions(-) create mode 100644 src/summary_mcp/runtime/__init__.py create mode 100644 src/summary_mcp/runtime/run_store.py create mode 100644 src/summary_mcp/runtime/state_models.py diff --git a/TODO.md b/TODO.md index b49f133..82400a8 100644 --- a/TODO.md +++ b/TODO.md @@ -1,83 +1,151 @@ -# TODO +# TODO - Reader MCP 正式化 -## 当前状态 +> 本文件用于架构与 Codex 协作同步。 +> +> 规则: +> - `TODO` = 未开始 +> - `DOING` = 正在进行 +> - `DONE` = 已完成 +> - 每次只允许一个最高优先级主任务处于 `DOING` -项目当前已经进入“可交付给 OpenClaw 调用”的阶段。 +## 0. 协作约束 -当前主链路: +开始编码前必须阅读: -`FreshRSS 未读 -> RSS 内容提取 -> LLM 总结 -> 规则过滤 -> OpenClaw delivery payload` - -当前已经完成: - -- [x] FreshRSS `greader` API 接入 -- [x] `entry -> item` 标准化映射 -- [x] RSS-first 内容提取策略 -- [x] LLM 摘要与校验闭环 -- [x] 第一版规则引擎 -- [x] `ArticleCandidateRecord` / `OpenClawCandidateInput` 分层 -- [x] `OpenClawDeliveryPayload` 批量投递结构 -- [x] 最终 payload 成功后才标记 FreshRSS 已读 -- [x] MCP 工具 `run_freshrss_openclaw_pipeline` -- [x] 默认精简输出模式 -- [x] 日报级 `keywords` 词元库与周期性词元清洗 skill 设计完成 -- [x] 日报级 `keywords` 词元库与全局词频统计实现完成 -- [x] `keyword-cleanup-review` skill 骨架与 review bundle 脚本实现完成 -- [x] 词元清洗低复杂治理层落地:`term_cleanup_policy` / `term_watchlist` / `term_change_log` -- [x] 已支持人工确认采纳建议并写入 `term_watchlist` / `term_change_log` +1. `plans/reader-mcp-architecture-design.md` +2. `plans/reader-mcp-implementation-plan.md` +3. 本文件 +4. `plans/issues/2026-04-06-reader-digest-sigterm.md` +5. `docs/openclaw/openclaw-handoff.md` --- -## P0 - 交接前后最优先 +## 1. 当前主任务 -- [x] 为 OpenClaw 补齐交接文档 -- [x] 将 MCP 工具作为统一生产入口 -- [x] 将默认输出收敛为最小必要文件 -- [ ] 设计 OpenClaw webhook / delivery payload 的主动推送方式 -- [ ] 明确 OpenClaw 侧如何注册和启动本 MCP 服务 +### [DONE][P0] 建立 run-state 运行态基础设施 + +目标: +- 给 freshrss pipeline 引入正式 run state +- 即使失败或中断,也能留下明确运行真相 + +要求: +- 新增 `RunState / StageState / ArtifactRecord` 模型 +- 在 `outputs/freshrss/rerun//run-state.json` 持久化 +- 至少覆盖以下 stages: + - `fetch_feed` + - `extract_articles` + - `generate_summaries` + - `apply_filters` + - `build_delivery_payload` + - `write_run_report` +- 失败时写入失败阶段与错误摘要 +- 不破坏现有输出目录兼容性 + +建议文件: +- `src/summary_mcp/runtime/state_models.py` +- `src/summary_mcp/runtime/run_store.py` +- `src/summary_mcp/workflows/...` + +完成标准: +- 跑一次 pipeline 后,无论成功失败,都存在 `run-state.json` +- 文件中可看出当前/最后阶段、整体状态、关键 artifacts + +进展备注: +- 2026-04-07:架构设计文档已建立;开始进入实现阶段。 +- 2026-04-07:已新增 `runtime` 包骨架,落地 `RunState / StageState / ArtifactRecord` 与文件存储接口。 +- 2026-04-07:已将 `run-state.json` 接入 `freshrss` 主流程,按阶段持续写入状态与关键 artifacts。 +- 2026-04-07:已完成成功/失败路径自检,确认 `run-state.json` 在两类路径下都保留且不改变既有对外返回字段。 --- -## P1 - 下一阶段推进 +## 2. 后续任务队列 -- [ ] 设计“人工确认后再沉淀知识库”的状态流转 -- [ ] 收敛 `paywall` 误判规则,降低中文文本误报 -- [ ] 细化过滤规则并引入更多个性化上下文 -- [ ] 将 `keyword-cleanup-review` skill 接入周期性执行流程,产出别名/停用词/兴趣词建议 -- [ ] 增加清洗前后效果对比报告,验证配置调整是否真的改善过滤质量 -- [ ] 增加批量 run 的保留策略与历史清理策略 -- [ ] 为 OpenClaw 补一份更正式的 MCP 调用示例和接线说明 +### [TODO][P1] 增加 MCP 状态查询接口 `get_run_status` + +目标: +- 可通过 MCP 查询 run 状态 + +要求: +- 输入 `run_id` +- 返回 status / current_stage / completed_stages / failed_stage / artifacts / recovery --- -## P2 - 后续增强 +### [TODO][P1] 增加 MCP 查询接口 `list_runs` -- [ ] 将 `Markdown sink` 进一步降级为 debug / fallback 能力 -- [ ] 增加按天聚合 `ArticleCandidateRecord` 的批处理能力 -- [ ] 让 OpenClaw 聚合候选内容并生成日级摘要 -- [ ] 将日级摘要写入知识库,并同步生成面向用户的日报消息 -- [ ] 支持更多 `content_kind` -- [ ] 增加提取缓存、重试和更细粒度日志 -- [ ] 整理历史 rerun 目录与调试产物保留策略 +目标: +- 查看近期 runs + +要求: +- 支持按 workflow / status / latest_n 过滤 --- -## 当前建议的下一步 +### [TODO][P1] 增加 MCP 查询接口 `list_run_artifacts` -优先做这三件事: - -1. 将 `keyword-cleanup-review` skill 接入周期性执行流程 -2. 设计“人工确认后再沉淀知识库”的状态流转 -3. 收敛规则误判,尤其是 `paywall` 相关启发式 +目标: +- 统一列出 run 下 artifact --- -## 交接时优先阅读 +### [TODO][P1] 增加 MCP 结果读取接口 `get_delivery_payload` -- `README.md` -- `docs/openclaw/openclaw-handoff.md` -- `docs/openclaw/openclaw-candidate-input-field-spec.md` -- `docs/openclaw/openclaw-delivery-payload-spec.md` -- `docs/design/daily-keyword-index-design.md` -- `skills/keyword-cleanup-review/SKILL.md` -- `docs/current/context-reset-brief.md` \ No newline at end of file +目标: +- 按 run_id 读取 delivery payload + +--- + +### [TODO][P1] 增加 MCP 结果读取接口 `get_run_report` + +目标: +- 按 run_id 读取 run report + +--- + +### [TODO][P2] 设计并实现 `resume_run` + +目标: +- 基于 `run-state.json` 和现有中间产物继续执行 + +说明: +- 先做最小可用恢复 +- 暂不追求任意 stage 任意重入 + +--- + +### [TODO][P2] 评估 `rerun_stage` 是否值得进入第一阶段 + +目标: +- 在 `resume_run` 之后评估是否继续增加更细粒度补跑能力 + +--- + +### [TODO][P2] 调整 `run_freshrss_openclaw_pipeline` 内部实现以复用 runtime + +目标: +- 保持外部兼容 +- 内部不再是黑箱长函数 + +--- + +### [TODO][P3] 更新 README / handoff / docs,明确 MCP 为正式入口 + +目标: +- 把生产建议从 CLI 迁移到 MCP +- CLI 明确降级为 debug / fallback + +--- + +## 3. 记录区 + +### 已完成记录 + +- 2026-04-07:新增架构设计文档 `plans/reader-mcp-architecture-design.md` +- 2026-04-07:新增实施计划文档 `plans/reader-mcp-implementation-plan.md` +- 2026-04-07:完成 `freshrss` pipeline 的 run-state 基础设施,新增 `runtime` 包并覆盖关键 stages 状态持久化。 + +### 风险提醒 + +- 不要在第一阶段引入复杂任务队列 +- 不要让 CLI 和 MCP 背后变成两套独立逻辑 +- 若实现偏离架构,先更新 `plans/` 再改代码 diff --git a/src/summary_mcp/runtime/__init__.py b/src/summary_mcp/runtime/__init__.py new file mode 100644 index 0000000..6dd253f --- /dev/null +++ b/src/summary_mcp/runtime/__init__.py @@ -0,0 +1,11 @@ +from .run_store import RunStore +from .state_models import ArtifactRecord, RecoveryState, RunError, RunState, StageState + +__all__ = [ + "ArtifactRecord", + "RecoveryState", + "RunError", + "RunState", + "RunStore", + "StageState", +] diff --git a/src/summary_mcp/runtime/run_store.py b/src/summary_mcp/runtime/run_store.py new file mode 100644 index 0000000..1a1d3c2 --- /dev/null +++ b/src/summary_mcp/runtime/run_store.py @@ -0,0 +1,185 @@ +from __future__ import annotations + +import json +from datetime import datetime +from pathlib import Path +from typing import Any + +from .state_models import ArtifactRecord, RecoveryState, RunError, RunState, StageState + + +class RunStore: + def __init__(self, *, path: Path, state: RunState, repo_root: Path | None = None) -> None: + self.path = path + self.state = state + self.repo_root = repo_root + + @classmethod + def create( + cls, + *, + path: Path, + run_id: str, + workflow: str, + run_type: str, + started_at: datetime, + input_payload: dict[str, Any] | None = None, + repo_root: Path | None = None, + ) -> "RunStore": + state = RunState( + run_id=run_id, + workflow=workflow, + run_type=run_type, + status="running", + current_stage=None, + started_at=started_at, + updated_at=started_at, + input=input_payload or {}, + recovery=RecoveryState(), + ) + return cls(path=path, state=state, repo_root=repo_root) + + def save(self) -> None: + self.state.updated_at = datetime.now(tz=self.state.started_at.tzinfo) + self.path.parent.mkdir(parents=True, exist_ok=True) + self.path.write_text( + json.dumps(self.state.model_dump(mode="json"), ensure_ascii=False, indent=2), + encoding="utf-8", + ) + + def start_stage(self, name: str, *, outputs: dict[str, Any] | None = None) -> None: + stage = self._get_or_create_stage(name) + now = datetime.now(tz=self.state.started_at.tzinfo) + stage.status = "running" + stage.started_at = stage.started_at or now + stage.finished_at = None + if outputs: + stage.outputs.update(outputs) + stage.error = None + self.state.current_stage = name + self.state.status = "running" + self.state.error = None + self._refresh_recovery() + self.save() + + def update_stage(self, name: str, *, outputs: dict[str, Any] | None = None) -> None: + stage = self._get_or_create_stage(name) + if outputs: + stage.outputs.update(outputs) + self.state.current_stage = name + self._refresh_recovery() + self.save() + + def finish_stage(self, name: str, *, outputs: dict[str, Any] | None = None) -> None: + stage = self._get_or_create_stage(name) + now = datetime.now(tz=self.state.started_at.tzinfo) + stage.status = "success" + stage.started_at = stage.started_at or now + stage.finished_at = now + if outputs: + stage.outputs.update(outputs) + stage.error = None + if self.state.current_stage == name: + self.state.current_stage = None + self._refresh_recovery() + self.save() + + def fail_stage( + self, + name: str, + *, + error: BaseException | RunError, + outputs: dict[str, Any] | None = None, + ) -> None: + stage = self._get_or_create_stage(name) + now = datetime.now(tz=self.state.started_at.tzinfo) + stage.status = "failed" + stage.started_at = stage.started_at or now + stage.finished_at = now + if outputs: + stage.outputs.update(outputs) + stage_error = error if isinstance(error, RunError) else self._build_error(error, stage=name) + stage.error = stage_error + self.state.current_stage = name + self.state.status = "failed" + self.state.error = stage_error + self.state.finished_at = now + self._refresh_recovery() + self.save() + + def register_artifact( + self, + *, + name: str, + path: Path, + kind: str, + stage: str, + metadata: dict[str, Any] | None = None, + ) -> None: + artifact = ArtifactRecord( + name=name, + path=self._normalize_path(path), + kind=kind, + stage=stage, + exists=path.exists(), + created_at=datetime.now(tz=self.state.started_at.tzinfo), + metadata=metadata or {}, + ) + existing = next((item for item in self.state.artifacts if item.name == name), None) + if existing is None: + self.state.artifacts.append(artifact) + else: + existing.path = artifact.path + existing.kind = artifact.kind + existing.stage = artifact.stage + existing.exists = artifact.exists + existing.created_at = artifact.created_at + existing.metadata = artifact.metadata + self.save() + + def finish_run(self, *, status: str) -> None: + now = datetime.now(tz=self.state.started_at.tzinfo) + self.state.status = status + self.state.current_stage = None + self.state.finished_at = now + self.state.error = None + self._refresh_recovery() + self.save() + + def _get_or_create_stage(self, name: str) -> StageState: + for stage in self.state.stages: + if stage.name == name: + return stage + + stage = StageState(name=name) + self.state.stages.append(stage) + return stage + + def _refresh_recovery(self) -> None: + successful_stages = [stage.name for stage in self.state.stages if stage.status == "success"] + last_success_stage = successful_stages[-1] if successful_stages else None + resume_from_stage = None + + if self.state.status == "running": + resume_from_stage = self.state.current_stage + elif self.state.status == "failed": + failed_stage = next((stage.name for stage in self.state.stages if stage.status == "failed"), None) + resume_from_stage = failed_stage or self.state.current_stage + + self.state.recovery = RecoveryState( + resumable=self.state.status in {"running", "failed"} and resume_from_stage is not None, + resume_from_stage=resume_from_stage, + last_success_stage=last_success_stage, + ) + + def _build_error(self, error: BaseException, *, stage: str | None = None) -> RunError: + return RunError(type=type(error).__name__, message=str(error), stage=stage) + + def _normalize_path(self, path: Path) -> str: + resolved_path = path.resolve() + if self.repo_root is not None: + try: + return str(resolved_path.relative_to(self.repo_root.resolve())) + except ValueError: + return str(path) + return str(path) diff --git a/src/summary_mcp/runtime/state_models.py b/src/summary_mcp/runtime/state_models.py new file mode 100644 index 0000000..1b88e82 --- /dev/null +++ b/src/summary_mcp/runtime/state_models.py @@ -0,0 +1,58 @@ +from __future__ import annotations + +from datetime import datetime +from typing import Any, Literal + +from pydantic import BaseModel, Field + + +RunStatus = Literal["running", "success", "partial", "failed"] +StageStatus = Literal["pending", "running", "success", "failed"] + + +class RunError(BaseModel): + type: str + message: str + stage: str | None = None + details: dict[str, Any] = Field(default_factory=dict) + + +class ArtifactRecord(BaseModel): + name: str + path: str + kind: str + stage: str + exists: bool = True + created_at: datetime + metadata: dict[str, Any] = Field(default_factory=dict) + + +class StageState(BaseModel): + name: str + status: StageStatus = "pending" + started_at: datetime | None = None + finished_at: datetime | None = None + outputs: dict[str, Any] = Field(default_factory=dict) + error: RunError | None = None + + +class RecoveryState(BaseModel): + resumable: bool = False + resume_from_stage: str | None = None + last_success_stage: str | None = None + + +class RunState(BaseModel): + run_id: str + workflow: str + run_type: str + status: RunStatus = "running" + current_stage: str | None = None + started_at: datetime + updated_at: datetime + finished_at: datetime | None = None + input: dict[str, Any] = Field(default_factory=dict) + stages: list[StageState] = Field(default_factory=list) + artifacts: list[ArtifactRecord] = Field(default_factory=list) + error: RunError | None = None + recovery: RecoveryState = Field(default_factory=RecoveryState) diff --git a/src/summary_mcp/workflows/freshrss_pipeline.py b/src/summary_mcp/workflows/freshrss_pipeline.py index 1ec410d..0de8810 100644 --- a/src/summary_mcp/workflows/freshrss_pipeline.py +++ b/src/summary_mcp/workflows/freshrss_pipeline.py @@ -1,14 +1,8 @@ from __future__ import annotations -# FreshRSS 全链路管道:拉取未读条目 -> 内容提取 -> LLM 摘要 -> 规则过滤 -> -# 构建 OpenClaw delivery payload -> 写盘 -> 词元统计 -> 标记已读。 -# 生产入口:run_freshrss_pipeline(),由 MCP 工具 run_freshrss_openclaw_pipeline 调用。 - import json import os from datetime import date, datetime, timezone - -UTC = timezone.utc from pathlib import Path from typing import Any @@ -32,8 +26,10 @@ from summary_mcp.models.openclaw_delivery import ( build_openclaw_digest_brief, ) from summary_mcp.models.summary_io import ExtractionInput +from summary_mcp.runtime import RunStore +UTC = timezone.utc REPO_ROOT = Path(__file__).resolve().parents[3] OUTPUT_ROOT = REPO_ROOT / "outputs" FRESHRSS_OUTPUT_ROOT = OUTPUT_ROOT / "freshrss" @@ -44,6 +40,14 @@ DEFAULT_TERM_ALIASES_PATH = REPO_ROOT / "configs" / "term_aliases.json" DEFAULT_TERM_STOPWORDS_PATH = REPO_ROOT / "configs" / "term_stopwords.json" DEFAULT_TERM_DAILY_DIR = DATA_ROOT / "daily" DEFAULT_TERM_STATS_PATH = DATA_ROOT / "term_stats.json" +WORKFLOW_NAME = "freshrss_daily_digest" +RUN_TYPE = "daily_digest" +FETCH_STAGE = "fetch_feed" +EXTRACT_STAGE = "extract_articles" +SUMMARY_STAGE = "generate_summaries" +FILTER_STAGE = "apply_filters" +DELIVERY_STAGE = "build_delivery_payload" +REPORT_STAGE = "write_run_report" def _save_json(path: Path, payload: dict[str, Any] | list[Any]) -> None: @@ -56,7 +60,6 @@ def _load_json(path: Path) -> dict[str, Any]: def _load_required_env(name: str, value: str | None) -> str: - # 优先使用传入的 value,否则读取同名环境变量;两者均缺失时抛出 RuntimeError if value: return value env_value = os.environ.get(name) @@ -75,25 +78,7 @@ def _maybe_path(enabled: bool, path: Path) -> Path | None: return path if enabled else None -def _process_item( - *, - index: int, - item: Any, - resolved_output_dir: Path, - resolved_prompt_path: Path, - resolved_run_id: str, - debug_artifacts: bool, - loaded_rules: list, - filter_context: Any, - max_retries: int, - timeout_seconds: float, - resolved_llm_api_key: str, - resolved_llm_model: str, - resolved_llm_api_url: str, -) -> dict[str, Any]: - # 处理单条 item:提取 -> LLM 摘要 -> 规则过滤 -> 候选构建 - # 返回 item_report dict;delivered 时额外携带 _candidate/_external_id 供调用方解包 - # 步骤 1:确定各中间文件路径(debug_artifacts=False 时大部分路径为 None,不写盘) +def _build_item_context(*, index: int, item: Any, resolved_output_dir: Path, debug_artifacts: bool) -> dict[str, Any]: item_key = f"item-{index:02d}" item_path = _maybe_path(debug_artifacts, resolved_output_dir / "items" / f"{item_key}.item.json") extracted_path = resolved_output_dir / "extracted" / f"{item_key}.extracted.json" @@ -102,10 +87,6 @@ def _process_item( record_path = _maybe_path(debug_artifacts, resolved_output_dir / "candidates" / f"{item_key}.article-candidate-record.json") openclaw_path = _maybe_path(debug_artifacts, resolved_output_dir / "candidates" / f"{item_key}.openclaw-candidate-input.json") - if item_path is not None: - _save_json(item_path, item.model_dump(mode="json")) - - # 步骤 2:初始化 item_report,记录基础元信息;debug 模式下附加各中间文件路径 item_report: dict[str, Any] = { "item_key": item_key, "item_id": item.item_id, @@ -117,115 +98,27 @@ def _process_item( if debug_artifacts: item_report["paths"] = { "item": str(item_path) if item_path else None, - "extracted": str(extracted_path) if extracted_path else None, + "extracted": str(extracted_path), "summary": str(summary_output) if summary_output else None, "filter": str(filter_path) if filter_path else None, "article_candidate": str(record_path) if record_path else None, "openclaw_candidate": str(openclaw_path) if openclaw_path else None, } - # 步骤 3:内容提取(RSS 内联内容 or 回源抓取);extracted.json 始终写盘 - extraction = extract_content(ExtractionInput(item=item)) - extracted_payload = extraction.model_dump(mode="json") - if extracted_path is not None: - _save_json(extracted_path, extracted_payload) - if not extraction.success or extraction.article is None: - item_report["status"] = "extract_failed" - item_report["error"] = extraction.error.model_dump(mode="json") if extraction.error else None - return item_report - - # 步骤 4:LLM 摘要循环,失败时最多重试 max_retries 次 - item_report["status"] = "extracted" - summary_exit_code, summary_payload, summary_report = run_loop_payload( - extracted_payload=extracted_payload, - prompt_path=resolved_prompt_path, - output_path=summary_output, - max_retries=max_retries, - timeout_seconds=timeout_seconds, - api_key=resolved_llm_api_key, - model=resolved_llm_model, - api_url=resolved_llm_api_url, - ) - if summary_exit_code != 0 or summary_payload is None: - item_report["status"] = "summary_failed" - if summary_report is not None: - item_report["summary_errors"] = summary_report.errors - return item_report - - # 步骤 5:规则引擎过滤,产出 keep/review/drop 决策及 digest_rank - summary = LlmSummaryResult.model_validate(summary_payload) - decision = evaluate_filter_rules( - FilterInput(item=item, article=extraction.article, summary=summary, context=filter_context), - loaded_rules, - ) - if filter_path is not None: - _save_json(filter_path, decision.model_dump(mode="json")) - - # 步骤 6:构建 ArticleCandidateRecord 和 OpenClawCandidateInput - record = build_article_candidate_record( - summary=summary, - article=extraction.article, - filter_result=decision, - item=item, - source_refs=CandidateSourceRefs( - item_path=str(item_path) if item_path else None, - extracted_path=str(extracted_path) if extracted_path else None, - summary_path=str(summary_output) if summary_output else None, - filter_path=str(filter_path) if filter_path else None, - ), - metadata=CandidateMetadata( - generated_at=datetime.now(tz=UTC), - producer="run_freshrss_pipeline", - run_id=resolved_run_id, - ), - ) - openclaw_input = build_openclaw_candidate_input(record) - - if record_path is not None: - _save_json(record_path, record.model_dump(mode="json")) - if openclaw_path is not None: - _save_json(openclaw_path, openclaw_input.model_dump(mode="json")) - - # 步骤 7:标记 delivered,附加临时键 _candidate/_external_id 供主函数解包后加入投递列表 - item_report["status"] = "delivered" - item_report["selection_decision"] = decision.decision - item_report["candidate_id"] = openclaw_input.candidate_id - item_report["_candidate"] = openclaw_input - item_report["_external_id"] = item.external_id - return item_report - - -def _build_and_persist_delivery( - *, - delivered_candidates: list[OpenClawCandidateInput], - resolved_run_id: str, - resolved_delivery_date: Any, - delivery_output: Path, -) -> tuple[OpenClawDeliveryPayload, Path, dict[str, Any]]: - # 按 digest_rank 排序,构建 delivery payload;同时派生一个更轻量的 digest-brief.json 供日报生成使用。 - delivered_candidates.sort(key=lambda candidate: candidate.digest_rank, reverse=True) - delivery_payload = build_openclaw_delivery_payload( - delivered_candidates, - run_id=resolved_run_id, - for_date=resolved_delivery_date, - ) - _save_json(delivery_output, delivery_payload.model_dump(mode="json")) - - digest_brief = build_openclaw_digest_brief(delivery_payload) - digest_brief_output = delivery_output.with_name("digest-brief.json") - _save_json(digest_brief_output, digest_brief.model_dump(mode="json")) - - keyword_index_result = persist_keyword_indexes( - delivery_payload.candidates, - for_date=delivery_payload.date, - digest_id=delivery_payload.run_id, - source="openclaw_delivery_payload", - daily_dir=DEFAULT_TERM_DAILY_DIR, - stats_path=DEFAULT_TERM_STATS_PATH, - aliases_path=DEFAULT_TERM_ALIASES_PATH, - stopwords_path=DEFAULT_TERM_STOPWORDS_PATH, - ) - return delivery_payload, digest_brief_output, keyword_index_result + return { + "item": item, + "item_key": item_key, + "item_path": item_path, + "extracted_path": extracted_path, + "summary_output": summary_output, + "filter_path": filter_path, + "record_path": record_path, + "openclaw_path": openclaw_path, + "item_report": item_report, + "extraction": None, + "extracted_payload": None, + "summary_payload": None, + } def _build_run_report( @@ -233,8 +126,8 @@ def _build_run_report( resolved_run_id: str, started_at: datetime, limit: int, - items: list, - delivered_candidates: list, + items: list[Any], + delivered_candidates: list[OpenClawCandidateInput], marked_count: int, mark_read: bool, debug_artifacts: bool, @@ -244,7 +137,6 @@ def _build_run_report( keyword_index_result: dict[str, Any], item_reports: list[dict[str, Any]], ) -> dict[str, Any]: - # 统计各状态计数,组装 run report dict status_counts: dict[str, int] = {} for item_report in item_reports: status = str(item_report["status"]) @@ -269,6 +161,12 @@ def _build_run_report( } +def _final_run_status(item_reports: list[dict[str, Any]]) -> str: + if any(item_report.get("status") in {"extract_failed", "summary_failed"} for item_report in item_reports): + return "partial" + return "success" + + def run_freshrss_pipeline( *, api_base_url: str | None = None, @@ -293,139 +191,413 @@ def run_freshrss_pipeline( delivery_date: date | None = None, output_dir: Path | None = None, ) -> dict[str, Any]: - # --- 阶段 1:初始化 run_id、输出路径、凭证 --- started_at = datetime.now(tz=UTC) resolved_output_dir = output_dir or default_output_dir() run_stamp = started_at.strftime("%Y%m%d-%H%M%S") resolved_run_id = run_id or f"freshrss-pipeline-{run_stamp}" resolved_delivery_date = delivery_date or datetime.now(tz=UTC).date() - resolved_api_base_url = _load_required_env("FRESHRSS_API_BASE_URL", api_base_url) - resolved_username = _load_required_env("FRESHRSS_USERNAME", username) - resolved_api_password = _load_required_env("FRESHRSS_API_PASSWORD", api_password) - resolved_llm_api_key, resolved_llm_model, resolved_llm_api_url = resolve_llm_settings( - api_key=llm_api_key, - model=llm_model, - api_url=llm_api_url, - ) - resolved_prompt_path = prompt or DEFAULT_PROMPT_PATH resolved_rules_path = rules or DEFAULT_RULES_PATH raw_output = resolved_output_dir / "raw" / "freshrss.raw.json" delivery_output = resolved_output_dir / "candidates" / "openclaw-delivery-payload.json" + digest_brief_output = delivery_output.with_name("digest-brief.json") report_output = resolved_output_dir / "run-report.json" + run_state_output = resolved_output_dir / "run-state.json" items_list_output = _maybe_path(debug_artifacts, resolved_output_dir / "items" / "freshrss.items.json") - # --- 阶段 2:登录 FreshRSS,拉取未读条目,写原始 payload --- - client = FreshRSSClient( - api_base_url=resolved_api_base_url, - username=resolved_username, - api_password=resolved_api_password, - timeout_seconds=timeout_seconds, + run_store = RunStore.create( + path=run_state_output, + run_id=resolved_run_id, + workflow=WORKFLOW_NAME, + run_type=RUN_TYPE, + started_at=started_at, + input_payload={ + "limit": limit, + "mark_read": mark_read, + "include_read": include_read, + "debug_artifacts": debug_artifacts, + "continuation": continuation, + "stream_id": stream_id, + }, + repo_root=REPO_ROOT, ) - auth_token = client.client_login() - payload = client.fetch_stream_contents( - auth_token=auth_token, - stream_id=stream_id, - limit=limit, - continuation=continuation, - exclude_targets=[] if include_read else [READ_TAG], - ) - entries = payload.get("items") - if not isinstance(entries, list): - raise RuntimeError("FreshRSS stream response does not contain an items array.") + run_store.save() - _save_json(raw_output, payload) - - items = [map_entry_to_item(entry) for entry in entries] - if items_list_output is not None: - _save_json(items_list_output, [item.model_dump(mode="json") for item in items]) - - loaded_rules = load_filter_rules(resolved_rules_path) - # context 优先使用直接传入的 dict,其次读取 context_path 文件,两者均缺失则使用空 context - if context is not None: - filter_context = FilterContext.model_validate(context) - elif context_path is not None: - filter_context = FilterContext.model_validate(_load_json(context_path)) - else: - filter_context = FilterContext() - - # --- 阶段 3:逐条处理(提取 -> LLM 摘要 -> 规则过滤 -> 候选构建) --- + client: FreshRSSClient | None = None + auth_token: str | None = None + items: list[Any] = [] + item_contexts: list[dict[str, Any]] = [] + item_reports: list[dict[str, Any]] = [] delivered_candidates: list[OpenClawCandidateInput] = [] delivered_item_ids: list[str] = [] - item_reports: list[dict[str, Any]] = [] - - for index, item in enumerate(items, start=1): - result = _process_item( - index=index, - item=item, - resolved_output_dir=resolved_output_dir, - resolved_prompt_path=resolved_prompt_path, - resolved_run_id=resolved_run_id, - debug_artifacts=debug_artifacts, - loaded_rules=loaded_rules, - filter_context=filter_context, - max_retries=max_retries, - timeout_seconds=timeout_seconds, - resolved_llm_api_key=resolved_llm_api_key, - resolved_llm_model=resolved_llm_model, - resolved_llm_api_url=resolved_llm_api_url, - ) - if result.get("status") == "delivered": - delivered_candidates.append(result.pop("_candidate")) - external_id = result.pop("_external_id", None) - if external_id: - delivered_item_ids.append(external_id) - item_reports.append(result) - - # --- 阶段 4+5:构建 delivery payload 并持久化词元索引 --- - delivery_payload, digest_brief_output, keyword_index_result = _build_and_persist_delivery( - delivered_candidates=delivered_candidates, - resolved_run_id=resolved_run_id, - resolved_delivery_date=resolved_delivery_date, - delivery_output=delivery_output, - ) - - # --- 阶段 6:标记已读,汇总报告,返回结果 --- + keyword_index_result: dict[str, Any] = {} marked_count = 0 - if mark_read and delivered_item_ids: - # 仅标记成功投递(delivered)的 item;drop/review 的 item 保持未读状态 - client.mark_items_as_read(auth_token=auth_token, item_ids=delivered_item_ids) - marked_count = len({item_id for item_id in delivered_item_ids if item_id}) - report = _build_run_report( - resolved_run_id=resolved_run_id, - started_at=started_at, - limit=limit, - items=items, - delivered_candidates=delivered_candidates, - marked_count=marked_count, - mark_read=mark_read, - debug_artifacts=debug_artifacts, - raw_output=raw_output, - delivery_output=delivery_output, - digest_brief_output=digest_brief_output, - keyword_index_result=keyword_index_result, - item_reports=item_reports, - ) - _save_json(report_output, report) + try: + run_store.start_stage(FETCH_STAGE, outputs={"output_dir": str(resolved_output_dir)}) + resolved_api_base_url = _load_required_env("FRESHRSS_API_BASE_URL", api_base_url) + resolved_username = _load_required_env("FRESHRSS_USERNAME", username) + resolved_api_password = _load_required_env("FRESHRSS_API_PASSWORD", api_password) + resolved_llm_api_key, resolved_llm_model, resolved_llm_api_url = resolve_llm_settings( + api_key=llm_api_key, + model=llm_model, + api_url=llm_api_url, + ) - return { - "run_id": resolved_run_id, - "output_dir": str(resolved_output_dir), - "raw_output": str(raw_output), - "delivery_output": str(delivery_output), - "digest_brief_output": str(digest_brief_output), - "report_output": str(report_output), - "keyword_index": keyword_index_result, - "pulled_count": len(items), - "delivered_count": len(delivered_candidates), - "marked_read_count": marked_count, - "status_counts": report["status_counts"], - "debug_artifacts": debug_artifacts, - "delivery_payload": delivery_payload.model_dump(mode="json"), - "items": item_reports, - } + client = FreshRSSClient( + api_base_url=resolved_api_base_url, + username=resolved_username, + api_password=resolved_api_password, + timeout_seconds=timeout_seconds, + ) + auth_token = client.client_login() + payload = client.fetch_stream_contents( + auth_token=auth_token, + stream_id=stream_id, + limit=limit, + continuation=continuation, + exclude_targets=[] if include_read else [READ_TAG], + ) + entries = payload.get("items") + if not isinstance(entries, list): + raise RuntimeError("FreshRSS stream response does not contain an items array.") + + _save_json(raw_output, payload) + run_store.register_artifact(name="raw_output", path=raw_output, kind="json", stage=FETCH_STAGE) + + items = [map_entry_to_item(entry) for entry in entries] + if items_list_output is not None: + _save_json(items_list_output, [item.model_dump(mode="json") for item in items]) + run_store.register_artifact(name="items_output", path=items_list_output, kind="json", stage=FETCH_STAGE) + + loaded_rules = load_filter_rules(resolved_rules_path) + if context is not None: + filter_context = FilterContext.model_validate(context) + elif context_path is not None: + filter_context = FilterContext.model_validate(_load_json(context_path)) + else: + filter_context = FilterContext() + + run_store.finish_stage( + FETCH_STAGE, + outputs={ + "pulled_count": len(items), + "raw_output": str(raw_output), + "items_output": str(items_list_output) if items_list_output else None, + }, + ) + + run_store.start_stage( + EXTRACT_STAGE, + outputs={ + "expected_items": len(items), + "completed_items": 0, + "success_count": 0, + "failed_count": 0, + }, + ) + extracted_success_count = 0 + extracted_failed_count = 0 + for index, item in enumerate(items, start=1): + item_context = _build_item_context( + index=index, + item=item, + resolved_output_dir=resolved_output_dir, + debug_artifacts=debug_artifacts, + ) + item_contexts.append(item_context) + item_report = item_context["item_report"] + item_reports.append(item_report) + + item_path = item_context["item_path"] + if item_path is not None: + _save_json(item_path, item.model_dump(mode="json")) + + extraction = extract_content(ExtractionInput(item=item)) + extracted_payload = extraction.model_dump(mode="json") + item_context["extraction"] = extraction + item_context["extracted_payload"] = extracted_payload + _save_json(item_context["extracted_path"], extracted_payload) + run_store.register_artifact(name="extracted_dir", path=resolved_output_dir / "extracted", kind="directory", stage=EXTRACT_STAGE) + + if not extraction.success or extraction.article is None: + item_report["status"] = "extract_failed" + item_report["error"] = extraction.error.model_dump(mode="json") if extraction.error else None + extracted_failed_count += 1 + else: + item_report["status"] = "extracted" + extracted_success_count += 1 + + run_store.update_stage( + EXTRACT_STAGE, + outputs={ + "expected_items": len(items), + "completed_items": extracted_success_count + extracted_failed_count, + "success_count": extracted_success_count, + "failed_count": extracted_failed_count, + }, + ) + + run_store.finish_stage( + EXTRACT_STAGE, + outputs={ + "expected_items": len(items), + "completed_items": extracted_success_count + extracted_failed_count, + "success_count": extracted_success_count, + "failed_count": extracted_failed_count, + "extracted_dir": str(resolved_output_dir / "extracted"), + }, + ) + + run_store.start_stage( + SUMMARY_STAGE, + outputs={ + "expected_items": extracted_success_count, + "completed_items": 0, + "success_count": 0, + "failed_count": 0, + }, + ) + summary_success_count = 0 + summary_failed_count = 0 + summary_candidates = [ctx for ctx in item_contexts if ctx["extraction"] is not None and ctx["extraction"].success] + for item_context in summary_candidates: + item_report = item_context["item_report"] + summary_exit_code, summary_payload, summary_report = run_loop_payload( + extracted_payload=item_context["extracted_payload"], + prompt_path=resolved_prompt_path, + output_path=item_context["summary_output"], + max_retries=max_retries, + timeout_seconds=timeout_seconds, + api_key=resolved_llm_api_key, + model=resolved_llm_model, + api_url=resolved_llm_api_url, + ) + if summary_exit_code != 0 or summary_payload is None: + item_report["status"] = "summary_failed" + if summary_report is not None: + item_report["summary_errors"] = summary_report.errors + summary_failed_count += 1 + else: + item_context["summary_payload"] = summary_payload + item_report["status"] = "summarized" + summary_success_count += 1 + + run_store.update_stage( + SUMMARY_STAGE, + outputs={ + "expected_items": extracted_success_count, + "completed_items": summary_success_count + summary_failed_count, + "success_count": summary_success_count, + "failed_count": summary_failed_count, + }, + ) + + if debug_artifacts and (resolved_output_dir / "summary").exists(): + run_store.register_artifact(name="summary_dir", path=resolved_output_dir / "summary", kind="directory", stage=SUMMARY_STAGE) + run_store.finish_stage( + SUMMARY_STAGE, + outputs={ + "expected_items": extracted_success_count, + "completed_items": summary_success_count + summary_failed_count, + "success_count": summary_success_count, + "failed_count": summary_failed_count, + }, + ) + + run_store.start_stage( + FILTER_STAGE, + outputs={ + "expected_items": summary_success_count, + "completed_items": 0, + "candidate_count": 0, + "keep_count": 0, + "review_count": 0, + "drop_count": 0, + }, + ) + filter_completed_count = 0 + keep_count = 0 + review_count = 0 + drop_count = 0 + for item_context in [ctx for ctx in item_contexts if ctx["summary_payload"] is not None]: + item = item_context["item"] + item_report = item_context["item_report"] + extraction = item_context["extraction"] + summary = LlmSummaryResult.model_validate(item_context["summary_payload"]) + decision = evaluate_filter_rules( + FilterInput(item=item, article=extraction.article, summary=summary, context=filter_context), + loaded_rules, + ) + filter_path = item_context["filter_path"] + if filter_path is not None: + _save_json(filter_path, decision.model_dump(mode="json")) + + record = build_article_candidate_record( + summary=summary, + article=extraction.article, + filter_result=decision, + item=item, + source_refs=CandidateSourceRefs( + item_path=str(item_context["item_path"]) if item_context["item_path"] else None, + extracted_path=str(item_context["extracted_path"]), + summary_path=str(item_context["summary_output"]) if item_context["summary_output"] else None, + filter_path=str(filter_path) if filter_path else None, + ), + metadata=CandidateMetadata( + generated_at=datetime.now(tz=UTC), + producer="run_freshrss_pipeline", + run_id=resolved_run_id, + ), + ) + openclaw_input = build_openclaw_candidate_input(record) + item_context["candidate"] = openclaw_input + + if item_context["record_path"] is not None: + _save_json(item_context["record_path"], record.model_dump(mode="json")) + if item_context["openclaw_path"] is not None: + _save_json(item_context["openclaw_path"], openclaw_input.model_dump(mode="json")) + + item_report["status"] = "delivered" + item_report["selection_decision"] = decision.decision + item_report["candidate_id"] = openclaw_input.candidate_id + delivered_candidates.append(openclaw_input) + if item.external_id: + delivered_item_ids.append(item.external_id) + + if decision.decision == "keep": + keep_count += 1 + elif decision.decision == "review": + review_count += 1 + elif decision.decision == "drop": + drop_count += 1 + + filter_completed_count += 1 + run_store.update_stage( + FILTER_STAGE, + outputs={ + "expected_items": summary_success_count, + "completed_items": filter_completed_count, + "candidate_count": len(delivered_candidates), + "keep_count": keep_count, + "review_count": review_count, + "drop_count": drop_count, + }, + ) + + if debug_artifacts and (resolved_output_dir / "candidates").exists(): + run_store.register_artifact(name="candidate_dir", path=resolved_output_dir / "candidates", kind="directory", stage=FILTER_STAGE) + run_store.finish_stage( + FILTER_STAGE, + outputs={ + "expected_items": summary_success_count, + "completed_items": filter_completed_count, + "candidate_count": len(delivered_candidates), + "keep_count": keep_count, + "review_count": review_count, + "drop_count": drop_count, + }, + ) + + run_store.start_stage(DELIVERY_STAGE, outputs={"candidate_count": len(delivered_candidates)}) + delivered_candidates.sort(key=lambda candidate: candidate.digest_rank, reverse=True) + delivery_payload = build_openclaw_delivery_payload( + delivered_candidates, + run_id=resolved_run_id, + for_date=resolved_delivery_date, + ) + _save_json(delivery_output, delivery_payload.model_dump(mode="json")) + run_store.register_artifact(name="delivery_payload", path=delivery_output, kind="json", stage=DELIVERY_STAGE) + + digest_brief = build_openclaw_digest_brief(delivery_payload) + _save_json(digest_brief_output, digest_brief.model_dump(mode="json")) + run_store.register_artifact(name="digest_brief", path=digest_brief_output, kind="json", stage=DELIVERY_STAGE) + + keyword_index_result = persist_keyword_indexes( + delivery_payload.candidates, + for_date=delivery_payload.date, + digest_id=delivery_payload.run_id, + source="openclaw_delivery_payload", + daily_dir=DEFAULT_TERM_DAILY_DIR, + stats_path=DEFAULT_TERM_STATS_PATH, + aliases_path=DEFAULT_TERM_ALIASES_PATH, + stopwords_path=DEFAULT_TERM_STOPWORDS_PATH, + ) + run_store.register_artifact( + name="keyword_daily_index", + path=Path(str(keyword_index_result["daily_output"])), + kind="json", + stage=DELIVERY_STAGE, + ) + run_store.register_artifact( + name="keyword_stats_index", + path=Path(str(keyword_index_result["stats_output"])), + kind="json", + stage=DELIVERY_STAGE, + ) + run_store.finish_stage( + DELIVERY_STAGE, + outputs={ + "candidate_count": len(delivered_candidates), + "delivery_output": str(delivery_output), + "digest_brief_output": str(digest_brief_output), + "keyword_daily_output": str(keyword_index_result["daily_output"]), + "keyword_stats_output": str(keyword_index_result["stats_output"]), + }, + ) + + run_store.start_stage(REPORT_STAGE, outputs={"mark_read_requested": mark_read}) + if mark_read and delivered_item_ids: + client.mark_items_as_read(auth_token=auth_token, item_ids=delivered_item_ids) + marked_count = len({item_id for item_id in delivered_item_ids if item_id}) + + report = _build_run_report( + resolved_run_id=resolved_run_id, + started_at=started_at, + limit=limit, + items=items, + delivered_candidates=delivered_candidates, + marked_count=marked_count, + mark_read=mark_read, + debug_artifacts=debug_artifacts, + raw_output=raw_output, + delivery_output=delivery_output, + digest_brief_output=digest_brief_output, + keyword_index_result=keyword_index_result, + item_reports=item_reports, + ) + _save_json(report_output, report) + run_store.register_artifact(name="run_report", path=report_output, kind="json", stage=REPORT_STAGE) + run_store.finish_stage( + REPORT_STAGE, + outputs={ + "marked_read_count": marked_count, + "report_output": str(report_output), + }, + ) + + run_store.finish_run(status=_final_run_status(item_reports)) + return { + "run_id": resolved_run_id, + "output_dir": str(resolved_output_dir), + "raw_output": str(raw_output), + "delivery_output": str(delivery_output), + "digest_brief_output": str(digest_brief_output), + "report_output": str(report_output), + "keyword_index": keyword_index_result, + "pulled_count": len(items), + "delivered_count": len(delivered_candidates), + "marked_read_count": marked_count, + "status_counts": report["status_counts"], + "debug_artifacts": debug_artifacts, + "delivery_payload": delivery_payload.model_dump(mode="json"), + "items": item_reports, + } + except Exception as error: + failed_stage = run_store.state.current_stage or FETCH_STAGE + run_store.fail_stage(failed_stage, error=error) + raise def read_delivery_payload(path: Path) -> OpenClawDeliveryPayload: