Implement minimal resume_run for freshrss runs

This commit is contained in:
root
2026-04-07 15:19:44 +08:00
parent d91cdbc6c9
commit c622bc6247
5 changed files with 819 additions and 2 deletions
+6 -1
View File
@@ -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` 是否值得进入第一阶段
+798
View File
@@ -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]
+1 -1
View File
@@ -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(
+7
View File
@@ -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(
*,
@@ -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,
)