diff --git a/TODO.md b/TODO.md index c00481a..2c35214 100644 --- a/TODO.md +++ b/TODO.md @@ -168,7 +168,7 @@ --- -### [TODO][P2] 设计并实现 `resume_run` +### [DONE][P2] 设计并实现 `resume_run` 目标: - 基于 `run-state.json` 和现有中间产物继续执行 @@ -180,6 +180,11 @@ - 第一版只支持 freshrss workflow 且仅支持有 `run-state.json` 的 run - 第一版仅考虑从最近可恢复点继续;`fetch_feed` / `extract_articles` 暂不支持恢复 +进展备注: +- 2026-04-07:已完成 `resume_run` minimal design 与现有 runtime/workflow/server 代码对齐分析,开始实现最小恢复链路。 +- 2026-04-07:已完成 `resume_run` 最小实现编码,新增 runtime 恢复服务并接入 MCP server;当前进入设计对齐与本地自检。 +- 2026-04-07:已完成 `resume_run` 架构对齐与本地自检;已验证 `write_run_report` 可恢复,且 `extract_articles` 会被明确拒绝恢复。 + --- ### [TODO][P2] 评估 `rerun_stage` 是否值得进入第一阶段 diff --git a/src/summary_mcp/runtime/resume_service.py b/src/summary_mcp/runtime/resume_service.py new file mode 100644 index 0000000..006f689 --- /dev/null +++ b/src/summary_mcp/runtime/resume_service.py @@ -0,0 +1,798 @@ +from __future__ import annotations + +from datetime import date, datetime, timezone +from pathlib import Path +from typing import Any + +from summary_mcp.core.keyword_index import persist_keyword_indexes +from summary_mcp.core.summary_loop import resolve_llm_settings, run_loop_payload +from summary_mcp.filters.engine import evaluate_filter_rules, load_filter_rules +from summary_mcp.integrations.freshrss import FreshRSSClient, map_entry_to_item +from summary_mcp.models.article_candidate import ( + CandidateMetadata, + CandidateSourceRefs, + OpenClawCandidateInput, + build_article_candidate_record, + build_openclaw_candidate_input, + candidate_id_for, +) +from summary_mcp.models.filtering import FilterContext, FilterInput +from summary_mcp.models.llm_result import LlmSummaryResult +from summary_mcp.models.openclaw_delivery import ( + OpenClawDeliveryPayload, + build_openclaw_delivery_payload, + build_openclaw_digest_brief, +) +from summary_mcp.models.summary_io import ExtractionOutput +from summary_mcp.workflows.freshrss_pipeline import ( + DEFAULT_PROMPT_PATH, + DEFAULT_RULES_PATH, + DEFAULT_TERM_ALIASES_PATH, + DEFAULT_TERM_DAILY_DIR, + DEFAULT_TERM_STATS_PATH, + DEFAULT_TERM_STOPWORDS_PATH, + DELIVERY_STAGE, + EXTRACT_STAGE, + FETCH_STAGE, + FILTER_STAGE, + REPORT_STAGE, + REPO_ROOT, + SUMMARY_STAGE, + WORKFLOW_NAME, + _build_item_context, + _final_run_status, + _load_json, + _load_required_env, + _save_json, +) + +from .query_service import _normalize_repo_path, _resolve_repo_path, _resolve_run_record +from .run_store import RunStore +from .state_models import RunState, StageState + +SUPPORTED_RESUME_STAGES = { + SUMMARY_STAGE, + FILTER_STAGE, + DELIVERY_STAGE, + REPORT_STAGE, +} +UNSUPPORTED_RESUME_STAGES = { + FETCH_STAGE, + EXTRACT_STAGE, +} +UTC = timezone.utc +DEFAULT_STREAM_ID = "user/-/state/com.google/reading-list" + + +def resume_run(*, run_id: str) -> dict[str, Any]: + record = _resolve_run_record(run_id) + base_response = { + "run_id": record["run_id"], + "workflow": record["workflow"], + "status": record["status"], + "output_dir": _normalize_repo_path(record["run_dir"]), + } + + if record["state_source"] != "run_state" or not isinstance(record.get("state"), RunState): + return { + **base_response, + "resumed": False, + "resume_from_stage": None, + "message": "This run cannot be resumed because run-state.json is missing or could not be loaded.", + "missing_artifacts": [], + } + + run_store = RunStore.load(path=record["run_dir"] / "run-state.json", repo_root=REPO_ROOT) + state = run_store.state + resume_from_stage = _resolve_resume_from_stage(state) + + if state.workflow != WORKFLOW_NAME: + return { + **base_response, + "resumed": False, + "resume_from_stage": resume_from_stage, + "message": f"This run cannot be resumed because workflow '{state.workflow}' is not supported by the minimal resume_run implementation.", + "missing_artifacts": [], + } + + if resume_from_stage is None: + return { + **base_response, + "resumed": False, + "resume_from_stage": None, + "message": "This run does not expose a recoverable stage in run-state.json.", + "missing_artifacts": [], + } + + if resume_from_stage in UNSUPPORTED_RESUME_STAGES: + return { + **base_response, + "resumed": False, + "resume_from_stage": resume_from_stage, + "message": f"This run cannot be resumed from {resume_from_stage} in the current minimal implementation.", + "missing_artifacts": [], + } + + if resume_from_stage not in SUPPORTED_RESUME_STAGES: + return { + **base_response, + "resumed": False, + "resume_from_stage": resume_from_stage, + "message": f"This run cannot be resumed because stage '{resume_from_stage}' is not supported.", + "missing_artifacts": [], + } + + missing_artifacts = _validate_resume_artifacts(record=record, state=state, resume_from_stage=resume_from_stage) + if missing_artifacts: + return { + **base_response, + "resumed": False, + "resume_from_stage": resume_from_stage, + "message": f"This run cannot be resumed from {resume_from_stage} because required artifacts are missing.", + "missing_artifacts": missing_artifacts, + } + + try: + result = _resume_freshrss_run(record=record, run_store=run_store, resume_from_stage=resume_from_stage) + except Exception as error: + failed_stage = run_store.state.current_stage or resume_from_stage + run_store.fail_stage(failed_stage, error=error) + return { + **base_response, + "resumed": True, + "resume_from_stage": resume_from_stage, + "status": run_store.state.status, + "message": f"Run resumed from {resume_from_stage} but failed again at {failed_stage}: {error}", + "missing_artifacts": [], + "error_summary": run_store.state.error.model_dump(mode="json") if run_store.state.error else None, + } + + return { + "run_id": run_store.state.run_id, + "workflow": run_store.state.workflow, + "resumed": True, + "resume_from_stage": resume_from_stage, + "status": result["status"], + "output_dir": _normalize_repo_path(record["run_dir"]), + "message": f"Run resumed from {resume_from_stage} and completed with status {result['status']}.", + "delivery_payload": result.get("delivery_payload"), + "run_report": result.get("run_report"), + "delivery_output": result.get("delivery_output"), + "report_output": result.get("report_output"), + "digest_brief_output": result.get("digest_brief_output"), + "keyword_index": result.get("keyword_index"), + "pulled_count": result.get("pulled_count"), + "delivered_count": result.get("delivered_count"), + "marked_read_count": result.get("marked_read_count"), + "status_counts": result.get("status_counts"), + "missing_artifacts": [], + } + + +def _resume_freshrss_run(*, record: dict[str, Any], run_store: RunStore, resume_from_stage: str) -> dict[str, Any]: + state = run_store.state + run_dir = record["run_dir"] + config = _build_resume_config(state=state, run_dir=run_dir) + + items = _load_items(config["raw_output"]) if config["raw_output"] is not None else [] + item_contexts = _load_item_contexts( + items=items, + run_dir=run_dir, + debug_artifacts=config["debug_artifacts"], + include_summaries=resume_from_stage in {FILTER_STAGE, DELIVERY_STAGE}, + include_candidates=resume_from_stage == DELIVERY_STAGE, + ) + item_reports = [context["item_report"] for context in item_contexts] + delivered_candidates = [context["candidate"] for context in item_contexts if context.get("candidate") is not None] + + if resume_from_stage == SUMMARY_STAGE: + _run_summary_stage(run_store=run_store, item_contexts=item_contexts, config=config) + _run_filter_stage(run_store=run_store, item_contexts=item_contexts, config=config) + delivered_candidates = [context["candidate"] for context in item_contexts if context.get("candidate") is not None] + delivery_payload, keyword_index_result = _run_delivery_stage( + run_store=run_store, + delivered_candidates=delivered_candidates, + config=config, + ) + elif resume_from_stage == FILTER_STAGE: + _run_filter_stage(run_store=run_store, item_contexts=item_contexts, config=config) + delivered_candidates = [context["candidate"] for context in item_contexts if context.get("candidate") is not None] + delivery_payload, keyword_index_result = _run_delivery_stage( + run_store=run_store, + delivered_candidates=delivered_candidates, + config=config, + ) + elif resume_from_stage == DELIVERY_STAGE: + delivery_payload, keyword_index_result = _run_delivery_stage( + run_store=run_store, + delivered_candidates=delivered_candidates, + config=config, + ) + else: + delivery_payload = OpenClawDeliveryPayload.model_validate(_load_json(config["delivery_output"])) + candidate_by_id = {candidate.candidate_id: candidate for candidate in delivery_payload.candidates} + item_contexts = _load_item_contexts( + items=items, + run_dir=run_dir, + debug_artifacts=config["debug_artifacts"], + delivery_candidates_by_id=candidate_by_id, + ) + item_reports = [context["item_report"] for context in item_contexts] + delivered_candidates = list(delivery_payload.candidates) + keyword_index_result = _build_keyword_index_result( + delivery_payload=delivery_payload, + keyword_daily_output=_stage_output(state, DELIVERY_STAGE, "keyword_daily_output"), + keyword_stats_output=_stage_output(state, DELIVERY_STAGE, "keyword_stats_output"), + ) + + item_reports = [context["item_report"] for context in item_contexts] + delivered_item_ids = [ + context["item"].external_id + for context in item_contexts + if context.get("candidate") is not None and context["item"].external_id + ] + report, marked_count = _run_report_stage( + run_store=run_store, + items=items, + item_reports=item_reports, + delivered_candidates=delivered_candidates, + delivered_item_ids=delivered_item_ids, + delivery_payload=delivery_payload, + keyword_index_result=keyword_index_result, + config=config, + ) + final_status = _final_run_status(item_reports) + run_store.finish_run(status=final_status) + return { + "status": final_status, + "delivery_payload": delivery_payload.model_dump(mode="json"), + "run_report": report, + "delivery_output": _normalize_repo_path(config["delivery_output"]), + "report_output": _normalize_repo_path(config["report_output"]), + "digest_brief_output": _normalize_repo_path(config["digest_brief_output"]), + "keyword_index": keyword_index_result, + "pulled_count": report["pulled_count"], + "delivered_count": report["delivered_count"], + "marked_read_count": marked_count, + "status_counts": report["status_counts"], + } + + +def _build_resume_config(*, state: RunState, run_dir: Path) -> dict[str, Any]: + input_payload = state.input if isinstance(state.input, dict) else {} + prompt_path = Path(input_payload.get("prompt")) if input_payload.get("prompt") else DEFAULT_PROMPT_PATH + rules_path = Path(input_payload.get("rules")) if input_payload.get("rules") else DEFAULT_RULES_PATH + delivery_date_value = input_payload.get("delivery_date") + resolved_delivery_date = date.fromisoformat(delivery_date_value) if isinstance(delivery_date_value, str) else datetime.now(tz=UTC).date() + context = input_payload.get("context") if isinstance(input_payload.get("context"), dict) else None + context_path_value = input_payload.get("context_path") + context_path = Path(context_path_value) if isinstance(context_path_value, str) and context_path_value.strip() else None + + raw_output = _find_artifact_path(run_dir=run_dir, state=state, artifact_name="raw_output", relative_path=Path("raw/freshrss.raw.json")) + delivery_output = _find_artifact_path( + run_dir=run_dir, + state=state, + artifact_name="delivery_payload", + relative_path=Path("candidates/openclaw-delivery-payload.json"), + ) or (run_dir / "candidates" / "openclaw-delivery-payload.json") + digest_brief_output = _find_artifact_path( + run_dir=run_dir, + state=state, + artifact_name="digest_brief", + relative_path=Path("candidates/digest-brief.json"), + ) or (run_dir / "candidates" / "digest-brief.json") + + return { + "run_dir": run_dir, + "raw_output": raw_output, + "delivery_output": delivery_output, + "digest_brief_output": digest_brief_output, + "report_output": run_dir / "run-report.json", + "prompt_path": prompt_path, + "rules_path": rules_path, + "context": context, + "context_path": context_path, + "limit": int(input_payload.get("limit") or _stage_output(state, FETCH_STAGE, "pulled_count") or 0), + "mark_read": bool(input_payload.get("mark_read", False)), + "include_read": bool(input_payload.get("include_read", False)), + "debug_artifacts": bool(input_payload.get("debug_artifacts", False)), + "continuation": input_payload.get("continuation"), + "stream_id": str(input_payload.get("stream_id") or DEFAULT_STREAM_ID), + "timeout_seconds": float(input_payload.get("timeout_seconds") or 60.0), + "max_retries": int(input_payload.get("max_retries") or 2), + "delivery_date": resolved_delivery_date, + } + + +def _load_item_contexts( + *, + items: list[Any], + run_dir: Path, + debug_artifacts: bool, + include_summaries: bool = False, + include_candidates: bool = False, + delivery_candidates_by_id: dict[str, OpenClawCandidateInput] | None = None, +) -> list[dict[str, Any]]: + contexts: list[dict[str, Any]] = [] + for index, item in enumerate(items, start=1): + context = _build_item_context( + index=index, + item=item, + resolved_output_dir=run_dir, + debug_artifacts=debug_artifacts, + ) + context["candidate"] = None + extraction_path = context["extracted_path"] + if extraction_path.exists(): + extracted_payload = _load_json(extraction_path) + extraction = ExtractionOutput.model_validate(extracted_payload) + context["extraction"] = extraction + context["extracted_payload"] = extracted_payload + if not extraction.success or extraction.article is None: + context["item_report"]["status"] = "extract_failed" + context["item_report"]["error"] = extraction.error.model_dump(mode="json") if extraction.error else None + else: + context["item_report"]["status"] = "extracted" + + if include_summaries and context["summary_output"] is not None and context["summary_output"].exists(): + summary_payload = _load_json(context["summary_output"]) + context["summary_payload"] = summary_payload + context["item_report"]["status"] = "summarized" + elif include_summaries and context.get("extraction") is not None and context["extraction"].success: + context["item_report"]["status"] = "summary_failed" + + candidate: OpenClawCandidateInput | None = None + if include_candidates and context["openclaw_path"] is not None and context["openclaw_path"].exists(): + candidate = OpenClawCandidateInput.model_validate(_load_json(context["openclaw_path"])) + elif delivery_candidates_by_id is not None and context.get("extraction") is not None and context["extraction"].success: + candidate_id = candidate_id_for(item, context["extraction"].article) + candidate = delivery_candidates_by_id.get(candidate_id) + if candidate is None and context["item_report"]["status"] == "extracted": + context["item_report"]["status"] = "summary_failed" + + if candidate is not None: + context["candidate"] = candidate + context["item_report"]["status"] = "delivered" + context["item_report"]["selection_decision"] = candidate.selection_decision + context["item_report"]["candidate_id"] = candidate.candidate_id + + contexts.append(context) + return contexts + + +def _run_summary_stage(*, run_store: RunStore, item_contexts: list[dict[str, Any]], config: dict[str, Any]) -> None: + extraction_success_count = sum(1 for context in item_contexts if context.get("extraction") is not None and context["extraction"].success) + run_store.start_stage( + SUMMARY_STAGE, + outputs={ + "expected_items": extraction_success_count, + "completed_items": 0, + "success_count": 0, + "failed_count": 0, + }, + ) + summary_success_count = 0 + summary_failed_count = 0 + resolved_llm_api_key, resolved_llm_model, resolved_llm_api_url = resolve_llm_settings() + + for item_context in [context for context in item_contexts if context.get("extraction") is not None and context["extraction"].success]: + summary_exit_code, summary_payload, summary_report = run_loop_payload( + extracted_payload=item_context["extracted_payload"], + prompt_path=config["prompt_path"], + output_path=item_context["summary_output"], + max_retries=config["max_retries"], + timeout_seconds=config["timeout_seconds"], + api_key=resolved_llm_api_key, + model=resolved_llm_model, + api_url=resolved_llm_api_url, + ) + if summary_exit_code != 0 or summary_payload is None: + item_context["item_report"]["status"] = "summary_failed" + if summary_report is not None: + item_context["item_report"]["summary_errors"] = summary_report.errors + summary_failed_count += 1 + else: + item_context["summary_payload"] = summary_payload + item_context["item_report"]["status"] = "summarized" + summary_success_count += 1 + + run_store.update_stage( + SUMMARY_STAGE, + outputs={ + "expected_items": extraction_success_count, + "completed_items": summary_success_count + summary_failed_count, + "success_count": summary_success_count, + "failed_count": summary_failed_count, + }, + ) + + if config["debug_artifacts"] and (config["run_dir"] / "summary").exists(): + run_store.register_artifact(name="summary_dir", path=config["run_dir"] / "summary", kind="directory", stage=SUMMARY_STAGE) + run_store.finish_stage( + SUMMARY_STAGE, + outputs={ + "expected_items": extraction_success_count, + "completed_items": summary_success_count + summary_failed_count, + "success_count": summary_success_count, + "failed_count": summary_failed_count, + }, + ) + + +def _run_filter_stage(*, run_store: RunStore, item_contexts: list[dict[str, Any]], config: dict[str, Any]) -> None: + loaded_rules = load_filter_rules(config["rules_path"]) + filter_context = _load_filter_context(config) + summary_success_count = sum(1 for context in item_contexts if context.get("summary_payload") is not None) + run_store.start_stage( + FILTER_STAGE, + outputs={ + "expected_items": summary_success_count, + "completed_items": 0, + "candidate_count": 0, + "keep_count": 0, + "review_count": 0, + "drop_count": 0, + }, + ) + filter_completed_count = 0 + keep_count = 0 + review_count = 0 + drop_count = 0 + + for item_context in [context for context in item_contexts if context.get("summary_payload") is not None]: + item = item_context["item"] + extraction = item_context["extraction"] + summary = LlmSummaryResult.model_validate(item_context["summary_payload"]) + decision = evaluate_filter_rules( + FilterInput(item=item, article=extraction.article, summary=summary, context=filter_context), + loaded_rules, + ) + if item_context["filter_path"] is not None: + _save_json(item_context["filter_path"], decision.model_dump(mode="json")) + + record = build_article_candidate_record( + summary=summary, + article=extraction.article, + filter_result=decision, + item=item, + source_refs=CandidateSourceRefs( + item_path=str(item_context["item_path"]) if item_context["item_path"] else None, + extracted_path=str(item_context["extracted_path"]), + summary_path=str(item_context["summary_output"]) if item_context["summary_output"] else None, + filter_path=str(item_context["filter_path"]) if item_context["filter_path"] else None, + ), + metadata=CandidateMetadata( + generated_at=datetime.now(tz=UTC), + producer="resume_run", + run_id=run_store.state.run_id, + ), + ) + candidate = build_openclaw_candidate_input(record) + item_context["candidate"] = candidate + + if item_context["record_path"] is not None: + _save_json(item_context["record_path"], record.model_dump(mode="json")) + if item_context["openclaw_path"] is not None: + _save_json(item_context["openclaw_path"], candidate.model_dump(mode="json")) + + item_context["item_report"]["status"] = "delivered" + item_context["item_report"]["selection_decision"] = decision.decision + item_context["item_report"]["candidate_id"] = candidate.candidate_id + + if decision.decision == "keep": + keep_count += 1 + elif decision.decision == "review": + review_count += 1 + elif decision.decision == "drop": + drop_count += 1 + + filter_completed_count += 1 + candidate_count = sum(1 for context in item_contexts if context.get("candidate") is not None) + run_store.update_stage( + FILTER_STAGE, + outputs={ + "expected_items": summary_success_count, + "completed_items": filter_completed_count, + "candidate_count": candidate_count, + "keep_count": keep_count, + "review_count": review_count, + "drop_count": drop_count, + }, + ) + + if config["debug_artifacts"] and (config["run_dir"] / "candidates").exists(): + run_store.register_artifact(name="candidate_dir", path=config["run_dir"] / "candidates", kind="directory", stage=FILTER_STAGE) + run_store.finish_stage( + FILTER_STAGE, + outputs={ + "expected_items": summary_success_count, + "completed_items": filter_completed_count, + "candidate_count": sum(1 for context in item_contexts if context.get("candidate") is not None), + "keep_count": keep_count, + "review_count": review_count, + "drop_count": drop_count, + }, + ) + + +def _run_delivery_stage( + *, + run_store: RunStore, + delivered_candidates: list[OpenClawCandidateInput], + config: dict[str, Any], +) -> tuple[OpenClawDeliveryPayload, dict[str, Any]]: + run_store.start_stage(DELIVERY_STAGE, outputs={"candidate_count": len(delivered_candidates)}) + delivered_candidates.sort(key=lambda candidate: candidate.digest_rank, reverse=True) + delivery_payload = build_openclaw_delivery_payload( + delivered_candidates, + run_id=run_store.state.run_id, + for_date=config["delivery_date"], + ) + _save_json(config["delivery_output"], delivery_payload.model_dump(mode="json")) + run_store.register_artifact(name="delivery_payload", path=config["delivery_output"], kind="json", stage=DELIVERY_STAGE) + + digest_brief = build_openclaw_digest_brief(delivery_payload) + _save_json(config["digest_brief_output"], digest_brief.model_dump(mode="json")) + run_store.register_artifact(name="digest_brief", path=config["digest_brief_output"], kind="json", stage=DELIVERY_STAGE) + + keyword_index_result = persist_keyword_indexes( + delivery_payload.candidates, + for_date=delivery_payload.date, + digest_id=delivery_payload.run_id, + source="openclaw_delivery_payload", + daily_dir=DEFAULT_TERM_DAILY_DIR, + stats_path=DEFAULT_TERM_STATS_PATH, + aliases_path=DEFAULT_TERM_ALIASES_PATH, + stopwords_path=DEFAULT_TERM_STOPWORDS_PATH, + ) + run_store.register_artifact( + name="keyword_daily_index", + path=Path(str(keyword_index_result["daily_output"])), + kind="json", + stage=DELIVERY_STAGE, + ) + run_store.register_artifact( + name="keyword_stats_index", + path=Path(str(keyword_index_result["stats_output"])), + kind="json", + stage=DELIVERY_STAGE, + ) + run_store.finish_stage( + DELIVERY_STAGE, + outputs={ + "candidate_count": len(delivered_candidates), + "delivery_output": str(config["delivery_output"]), + "digest_brief_output": str(config["digest_brief_output"]), + "keyword_daily_output": str(keyword_index_result["daily_output"]), + "keyword_stats_output": str(keyword_index_result["stats_output"]), + }, + ) + return delivery_payload, keyword_index_result + + +def _run_report_stage( + *, + run_store: RunStore, + items: list[Any], + item_reports: list[dict[str, Any]], + delivered_candidates: list[OpenClawCandidateInput], + delivered_item_ids: list[str], + delivery_payload: OpenClawDeliveryPayload, + keyword_index_result: dict[str, Any], + config: dict[str, Any], +) -> tuple[dict[str, Any], int]: + run_store.start_stage(REPORT_STAGE, outputs={"mark_read_requested": config["mark_read"]}) + marked_count = 0 + if config["mark_read"] and delivered_item_ids: + resolved_api_base_url = _load_required_env("FRESHRSS_API_BASE_URL", None) + resolved_username = _load_required_env("FRESHRSS_USERNAME", None) + resolved_api_password = _load_required_env("FRESHRSS_API_PASSWORD", None) + client = FreshRSSClient( + api_base_url=resolved_api_base_url, + username=resolved_username, + api_password=resolved_api_password, + timeout_seconds=config["timeout_seconds"], + ) + auth_token = client.client_login() + client.mark_items_as_read(auth_token=auth_token, item_ids=delivered_item_ids) + marked_count = len({item_id for item_id in delivered_item_ids if item_id}) + + report = _build_resume_run_report( + state=run_store.state, + items=items, + item_reports=item_reports, + delivered_candidates=delivered_candidates, + marked_count=marked_count, + delivery_payload=delivery_payload, + keyword_index_result=keyword_index_result, + config=config, + ) + _save_json(config["report_output"], report) + run_store.register_artifact(name="run_report", path=config["report_output"], kind="json", stage=REPORT_STAGE) + run_store.finish_stage( + REPORT_STAGE, + outputs={ + "marked_read_count": marked_count, + "report_output": str(config["report_output"]), + }, + ) + return report, marked_count + + +def _build_resume_run_report( + *, + state: RunState, + items: list[Any], + item_reports: list[dict[str, Any]], + delivered_candidates: list[OpenClawCandidateInput], + marked_count: int, + delivery_payload: OpenClawDeliveryPayload, + keyword_index_result: dict[str, Any], + config: dict[str, Any], +) -> dict[str, Any]: + status_counts: dict[str, int] = {} + for item_report in item_reports: + status = str(item_report["status"]) + status_counts[status] = status_counts.get(status, 0) + 1 + + pulled_count = len(items) + if not items: + pulled_count = int(_stage_output(state, FETCH_STAGE, "pulled_count") or len(item_reports)) + if not status_counts: + failed_count = int(_stage_output(state, SUMMARY_STAGE, "failed_count") or 0) + if failed_count: + status_counts["summary_failed"] = failed_count + if delivered_candidates: + status_counts["delivered"] = len(delivered_candidates) + + report: dict[str, Any] = { + "run_id": state.run_id, + "started_at": state.started_at.isoformat(), + "completed_at": datetime.now(tz=UTC).isoformat(), + "requested_limit": config["limit"], + "pulled_count": pulled_count, + "delivered_count": len(delivered_candidates), + "marked_read_count": marked_count, + "mark_read_requested": config["mark_read"], + "debug_artifacts": config["debug_artifacts"], + "raw_output": str(config["raw_output"]) if config["raw_output"] is not None else None, + "delivery_output": str(config["delivery_output"]), + "digest_brief_output": str(config["digest_brief_output"]), + "keyword_index": keyword_index_result, + "status_counts": status_counts, + "items": item_reports, + } + return report + + +def _build_keyword_index_result( + *, + delivery_payload: OpenClawDeliveryPayload, + keyword_daily_output: str | None, + keyword_stats_output: str | None, +) -> dict[str, Any]: + result: dict[str, Any] = { + "source": "openclaw_delivery_payload", + "candidate_count": len(delivery_payload.candidates), + } + if isinstance(keyword_daily_output, str) and keyword_daily_output.strip(): + result["daily_output"] = keyword_daily_output + if isinstance(keyword_stats_output, str) and keyword_stats_output.strip(): + result["stats_output"] = keyword_stats_output + return result + + +def _load_filter_context(config: dict[str, Any]) -> FilterContext: + if config["context"] is not None: + return FilterContext.model_validate(config["context"]) + if config["context_path"] is not None and config["context_path"].exists(): + return FilterContext.model_validate(_load_json(config["context_path"])) + return FilterContext() + + +def _validate_resume_artifacts(*, record: dict[str, Any], state: RunState, resume_from_stage: str) -> list[str]: + run_dir = record["run_dir"] + missing_artifacts: list[str] = [] + raw_output = _find_artifact_path(run_dir=run_dir, state=state, artifact_name="raw_output", relative_path=Path("raw/freshrss.raw.json")) + extracted_dir = _find_artifact_path(run_dir=run_dir, state=state, artifact_name="extracted_dir", relative_path=Path("extracted")) + if resume_from_stage in {SUMMARY_STAGE, FILTER_STAGE, DELIVERY_STAGE} and raw_output is None: + missing_artifacts.append("raw/freshrss.raw.json") + if resume_from_stage in {SUMMARY_STAGE, FILTER_STAGE} and extracted_dir is None: + missing_artifacts.append("extracted/") + + if missing_artifacts: + return missing_artifacts + + items = _load_items(raw_output) if raw_output is not None else [] + item_contexts = _load_item_contexts( + items=items, + run_dir=run_dir, + debug_artifacts=bool(state.input.get("debug_artifacts", False)), + include_summaries=resume_from_stage in {FILTER_STAGE, DELIVERY_STAGE}, + include_candidates=resume_from_stage == DELIVERY_STAGE, + ) + + if resume_from_stage == SUMMARY_STAGE: + for context in item_contexts: + if not context["extracted_path"].exists(): + missing_artifacts.append(_normalize_repo_path(context["extracted_path"])) + + if resume_from_stage == FILTER_STAGE: + for context in item_contexts: + if context.get("extraction") is not None and context["extraction"].success: + if context["summary_output"] is None or not context["summary_output"].exists(): + missing_artifacts.append( + _normalize_repo_path(context["summary_output"] or (run_dir / "summary" / context["item_key"] / "result.loop.json")) + ) + + if resume_from_stage == DELIVERY_STAGE: + expected_candidate_count = int(_stage_output(state, FILTER_STAGE, "candidate_count") or 0) + actual_candidate_count = sum(1 for context in item_contexts if context.get("candidate") is not None) + if expected_candidate_count != actual_candidate_count: + for context in item_contexts: + if context.get("summary_payload") is not None and (context["openclaw_path"] is None or not context["openclaw_path"].exists()): + missing_artifacts.append( + _normalize_repo_path(context["openclaw_path"] or (run_dir / "candidates" / f"{context['item_key']}.openclaw-candidate-input.json")) + ) + + if resume_from_stage == REPORT_STAGE: + delivery_output = _find_artifact_path( + run_dir=run_dir, + state=state, + artifact_name="delivery_payload", + relative_path=Path("candidates/openclaw-delivery-payload.json"), + ) + if delivery_output is None: + missing_artifacts.append("candidates/openclaw-delivery-payload.json") + elif bool(state.input.get("mark_read", False)) and raw_output is None: + missing_artifacts.append("raw/freshrss.raw.json") + + deduped_missing_artifacts: list[str] = [] + for artifact in missing_artifacts: + if artifact not in deduped_missing_artifacts: + deduped_missing_artifacts.append(artifact) + return deduped_missing_artifacts + + +def _resolve_resume_from_stage(state: RunState) -> str | None: + if state.recovery.resume_from_stage: + return state.recovery.resume_from_stage + if state.status == "failed": + failed_stage = next((stage.name for stage in reversed(state.stages) if stage.status == "failed"), None) + return failed_stage + return None + + +def _stage_output(state: RunState, stage_name: str, key: str) -> Any: + stage_state = _find_stage_state(state, stage_name) + if stage_state is None: + return None + return stage_state.outputs.get(key) + + +def _find_stage_state(state: RunState, stage_name: str) -> StageState | None: + for stage in state.stages: + if stage.name == stage_name: + return stage + return None + + +def _find_artifact_path(*, run_dir: Path, state: RunState, artifact_name: str, relative_path: Path) -> Path | None: + for artifact in state.artifacts: + if artifact.name != artifact_name: + continue + resolved = _resolve_repo_path(artifact.path) + if resolved.exists(): + return resolved + path = run_dir / relative_path + if path.exists(): + return path + return None + + +def _load_items(raw_output: Path) -> list[Any]: + payload = _load_json(raw_output) + entries = payload.get("items") + if not isinstance(entries, list): + raise RuntimeError("FreshRSS raw output does not contain an items array.") + return [map_entry_to_item(entry) for entry in entries] diff --git a/src/summary_mcp/runtime/run_store.py b/src/summary_mcp/runtime/run_store.py index ad3cc26..e636a0a 100644 --- a/src/summary_mcp/runtime/run_store.py +++ b/src/summary_mcp/runtime/run_store.py @@ -168,7 +168,7 @@ class RunStore: 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) + failed_stage = next((stage.name for stage in reversed(self.state.stages) if stage.status == "failed"), None) resume_from_stage = failed_stage or self.state.current_stage self.state.recovery = RecoveryState( diff --git a/src/summary_mcp/server.py b/src/summary_mcp/server.py index 0d8df67..05c2f16 100644 --- a/src/summary_mcp/server.py +++ b/src/summary_mcp/server.py @@ -20,6 +20,7 @@ from summary_mcp.runtime import get_run_status as load_run_status from summary_mcp.runtime import get_run_report as load_run_report from summary_mcp.runtime import list_run_artifacts as load_run_artifacts from summary_mcp.runtime import list_runs as load_runs +from summary_mcp.runtime.resume_service import resume_run as resume_existing_run from summary_mcp.workflows import run_freshrss_pipeline from summary_mcp.workflows.article_summary import ArticleSummaryConfig, summarize_selected_articles @@ -158,6 +159,12 @@ def get_run_report(run_id: str) -> dict: return load_run_report(run_id=run_id) +@mcp.tool() +def resume_run(run_id: str) -> dict: + """Resume a failed or interrupted FreshRSS workflow run from its latest supported recovery point.""" + return resume_existing_run(run_id=run_id) + + @mcp.tool() def generate_article_summaries( *, diff --git a/src/summary_mcp/workflows/freshrss_pipeline.py b/src/summary_mcp/workflows/freshrss_pipeline.py index 0de8810..145a6e0 100644 --- a/src/summary_mcp/workflows/freshrss_pipeline.py +++ b/src/summary_mcp/workflows/freshrss_pipeline.py @@ -219,6 +219,13 @@ def run_freshrss_pipeline( "debug_artifacts": debug_artifacts, "continuation": continuation, "stream_id": stream_id, + "timeout_seconds": timeout_seconds, + "max_retries": max_retries, + "delivery_date": resolved_delivery_date.isoformat(), + "prompt": str(resolved_prompt_path), + "rules": str(resolved_rules_path), + "context": context, + "context_path": str(context_path) if context_path else None, }, repo_root=REPO_ROOT, )