feat: add async resume jobs and doc navigation

This commit is contained in:
root
2026-04-14 15:55:02 +08:00
parent 6705613aa4
commit b8727f1885
26 changed files with 2791 additions and 665 deletions
+9 -2
View File
@@ -1,17 +1,21 @@
from .run_store import RunStore
from .state_models import ArtifactRecord, RecoveryState, RunError, RunState, StageState
from .freshrss_pipeline_jobs import (
get_freshrss_pipeline_job_result,
get_freshrss_pipeline_job_status,
start_freshrss_pipeline_job,
)
from .query_service import get_delivery_payload, get_run_report, get_run_status, list_run_artifacts, list_runs
from .run_store import RunStore
from .state_models import ArtifactRecord, RecoveryState, RunError, RunState, StageState
from .resume_jobs import get_resume_job_result, get_resume_job_status, start_resume_job
from .resume_service import inspect_resume_plan, resume_run
__all__ = [
"ArtifactRecord",
"get_delivery_payload",
"get_freshrss_pipeline_job_result",
"get_freshrss_pipeline_job_status",
"get_resume_job_result",
"get_resume_job_status",
"get_run_report",
"RecoveryState",
"RunError",
@@ -19,7 +23,10 @@ __all__ = [
"RunStore",
"StageState",
"get_run_status",
"inspect_resume_plan",
"list_run_artifacts",
"list_runs",
"resume_run",
"start_resume_job",
"start_freshrss_pipeline_job",
]
+249 -31
View File
@@ -10,6 +10,7 @@ from typing import Any
from uuid import uuid4
from .run_store import RunStore
from .query_service import _resolve_run_record
REPO_ROOT = Path(__file__).resolve().parents[3]
OUTPUT_ROOT = REPO_ROOT / "outputs" / "freshrss"
@@ -26,6 +27,8 @@ DEFAULT_STAGES = [
"run_pipeline",
"write_result",
]
MIN_JOB_STALE_SECONDS = 30 * 60
MAX_JOB_STALE_SECONDS = 6 * 60 * 60
def _now() -> datetime:
@@ -369,33 +372,34 @@ def run_freshrss_pipeline_job(*, job_id: str) -> dict[str, Any]:
def get_freshrss_pipeline_job_status(*, job_id: str) -> dict[str, Any]:
store = _load_run_store(job_id)
state = store.state
completed_stage_count = sum(1 for s in state.stages if s.status == "success")
running_stage_count = sum(1 for s in state.stages if s.status == "running")
failed_stage_count = sum(1 for s in state.stages if s.status == "failed")
pending_stage_count = sum(1 for s in state.stages if s.status == "pending")
effective = _resolve_effective_job_view(job_id=job_id, state=state)
progress = _build_effective_job_progress(state=state, effective_status=effective["status"])
linked_run_id = state.input.get("run_id") if isinstance(state.input, dict) else None
linked_output_dir = _normalize_repo_path_value(state.input.get("output_dir")) if isinstance(state.input, dict) else None
return {
"job_id": state.run_id,
"workflow": state.workflow,
"run_type": state.run_type,
"status": state.status,
"current_stage": state.current_stage,
"status": effective["status"],
"current_stage": effective["current_stage"],
"started_at": state.started_at.isoformat(),
"updated_at": state.updated_at.isoformat(),
"finished_at": state.finished_at.isoformat() if state.finished_at else None,
"updated_at": effective["updated_at"],
"finished_at": effective["finished_at"],
"output_dir": _normalize_repo_path(_job_dir(job_id)),
"linked_run_id": linked_run_id,
"linked_output_dir": linked_output_dir,
"progress": {
"completed_stage_count": completed_stage_count,
"running_stage_count": running_stage_count,
"failed_stage_count": failed_stage_count,
"pending_stage_count": pending_stage_count,
"total_stage_count": len(state.stages),
},
"progress": progress,
"artifacts": [artifact.model_dump(mode="json") for artifact in state.artifacts],
"error_summary": state.error.model_dump(mode="json") if state.error else None,
"status_source": effective["status_source"],
"state_quality": effective["state_quality"],
"state_conflict": effective["state_conflict"],
"state_conflict_reason": effective["state_conflict_reason"],
"status_note": effective["status_note"],
"raw_status": state.status,
"raw_current_stage": state.current_stage,
"linked_run_status": effective["linked_run_status"],
"linked_run_state_conflict": effective["linked_run_state_conflict"],
}
@@ -403,30 +407,244 @@ def get_freshrss_pipeline_job_result(*, job_id: str) -> dict[str, Any]:
store = _load_run_store(job_id)
state = store.state
result = _load_result(job_id)
if state.status != "success" or result is None:
effective = _resolve_effective_job_view(job_id=job_id, state=state, result=result)
synthesized_result = result or effective["result"]
if effective["status"] != "success" or synthesized_result is None:
message = "FreshRSS pipeline job result is not ready."
if effective["status"] == "failed":
message = "FreshRSS pipeline job did not complete successfully, so no terminal job result is available."
return {
"job_id": state.run_id,
"status": state.status,
"message": "FreshRSS pipeline job result is not ready.",
"status": effective["status"],
"message": message,
"linked_run_id": state.input.get("run_id") if isinstance(state.input, dict) else None,
"error_summary": state.error.model_dump(mode="json") if state.error else None,
"status_source": effective["status_source"],
"status_note": effective["status_note"],
}
artifact = next((a.model_dump(mode="json") for a in state.artifacts if a.name == "job_result"), None)
return {
"job_id": state.run_id,
"status": state.status,
"run_id": result.get("run_id"),
"output_dir": result.get("output_dir"),
"raw_output": result.get("raw_output"),
"delivery_output": result.get("delivery_output"),
"digest_brief_output": result.get("digest_brief_output"),
"report_output": result.get("report_output"),
"pulled_count": result.get("pulled_count"),
"delivered_count": result.get("delivered_count"),
"marked_read_count": result.get("marked_read_count"),
"status_counts": result.get("status_counts"),
"keyword_index": result.get("keyword_index"),
"status": effective["status"],
"run_id": synthesized_result.get("run_id"),
"output_dir": synthesized_result.get("output_dir"),
"raw_output": synthesized_result.get("raw_output"),
"delivery_output": synthesized_result.get("delivery_output"),
"digest_brief_output": synthesized_result.get("digest_brief_output"),
"report_output": synthesized_result.get("report_output"),
"pulled_count": synthesized_result.get("pulled_count"),
"delivered_count": synthesized_result.get("delivered_count"),
"marked_read_count": synthesized_result.get("marked_read_count"),
"status_counts": synthesized_result.get("status_counts"),
"keyword_index": synthesized_result.get("keyword_index"),
"artifact": artifact,
"result": result,
"result": synthesized_result,
"status_source": effective["status_source"],
"status_note": effective["status_note"],
"result_source": synthesized_result.get("result_source", "job_result"),
"linked_run_status": effective["linked_run_status"],
}
def _resolve_effective_job_view(
*,
job_id: str,
state: Any,
result: dict[str, Any] | None = None,
) -> dict[str, Any]:
resolved_result = result or _load_result(job_id)
linked_run_record = _load_linked_run_record(state)
raw_status = state.status
raw_current_stage = state.current_stage
effective_status = raw_status
effective_current_stage = raw_current_stage
effective_updated_at = state.updated_at.isoformat()
effective_finished_at = state.finished_at.isoformat() if state.finished_at else None
status_source = "job_run_state"
state_quality = "trusted"
state_conflict = False
state_conflict_reason = None
status_note = None
if resolved_result is not None:
completed_at = resolved_result.get("completed_at")
effective_status = "success"
effective_current_stage = None
effective_updated_at = completed_at or effective_updated_at
effective_finished_at = completed_at or effective_finished_at
if raw_status != "success" or raw_current_stage is not None:
status_source = "job_result_reconciliation"
state_quality = "reconciled"
state_conflict = True
state_conflict_reason = "result.json already exists, but the job run-state did not converge to success."
status_note = "job result exists; the job can be treated as completed."
elif linked_run_record is not None:
if _linked_run_result_available(linked_run_record):
effective_status = "success"
effective_current_stage = None
effective_updated_at = linked_run_record.get("updated_at") or effective_updated_at
effective_finished_at = linked_run_record.get("finished_at") or effective_finished_at
status_source = "linked_run_reconciliation"
state_quality = "reconciled"
state_conflict = raw_status != "success" or raw_current_stage is not None
state_conflict_reason = (
"The linked run already has terminal artifacts, but the outer job run-state did not converge."
)
status_note = "linked run artifacts are complete; the job can be treated as completed."
resolved_result = _build_result_from_linked_run(job_id=job_id, state=state, linked_run_record=linked_run_record)
elif linked_run_record["status"] == "failed" and raw_status == "running":
effective_status = "failed"
effective_current_stage = None
effective_updated_at = linked_run_record.get("updated_at") or effective_updated_at
effective_finished_at = linked_run_record.get("finished_at") or effective_finished_at
status_source = "linked_run_reconciliation"
state_quality = "reconciled"
state_conflict = True
state_conflict_reason = "The linked run is already failed, but the outer job still reports running."
status_note = "linked run failed; the outer job state appears stale."
elif raw_status == "running":
stale_reason = _build_stale_job_reason(state)
if stale_reason is not None:
effective_status = "failed"
effective_current_stage = raw_current_stage
effective_finished_at = effective_updated_at
status_source = "stale_job_state_timeout"
state_quality = "reconciled"
state_conflict = True
state_conflict_reason = stale_reason
status_note = stale_reason
return {
"status": effective_status,
"current_stage": effective_current_stage,
"updated_at": effective_updated_at,
"finished_at": effective_finished_at,
"status_source": status_source,
"state_quality": state_quality,
"state_conflict": state_conflict,
"state_conflict_reason": state_conflict_reason,
"status_note": status_note,
"linked_run_status": linked_run_record["status"] if linked_run_record is not None else None,
"linked_run_state_conflict": linked_run_record.get("state_conflict") if linked_run_record is not None else None,
"result": resolved_result,
}
def _build_effective_job_progress(*, state: Any, effective_status: str) -> dict[str, int]:
total_stage_count = len(state.stages)
if effective_status == "success":
return {
"completed_stage_count": total_stage_count,
"running_stage_count": 0,
"failed_stage_count": 0,
"pending_stage_count": 0,
"total_stage_count": total_stage_count,
}
if effective_status == "failed" and state.status == "running":
completed_stage_count = sum(1 for s in state.stages if s.status == "success")
pending_stage_count = max(total_stage_count - completed_stage_count - 1, 0)
return {
"completed_stage_count": completed_stage_count,
"running_stage_count": 0,
"failed_stage_count": 1,
"pending_stage_count": pending_stage_count,
"total_stage_count": total_stage_count,
}
return {
"completed_stage_count": sum(1 for s in state.stages if s.status == "success"),
"running_stage_count": sum(1 for s in state.stages if s.status == "running"),
"failed_stage_count": sum(1 for s in state.stages if s.status == "failed"),
"pending_stage_count": sum(1 for s in state.stages if s.status == "pending"),
"total_stage_count": total_stage_count,
}
def _load_linked_run_record(state: Any) -> dict[str, Any] | None:
if not isinstance(state.input, dict):
return None
linked_run_id = state.input.get("run_id")
if not isinstance(linked_run_id, str) or not linked_run_id.strip():
return None
try:
return _resolve_run_record(linked_run_id)
except FileNotFoundError:
return None
def _linked_run_result_available(record: dict[str, Any]) -> bool:
artifact_presence = record.get("artifact_presence")
if not isinstance(artifact_presence, dict):
return False
return bool(artifact_presence.get("run_report")) and bool(artifact_presence.get("delivery_payload"))
def _build_result_from_linked_run(
*,
job_id: str,
state: Any,
linked_run_record: dict[str, Any],
) -> dict[str, Any]:
report = linked_run_record.get("report")
if not isinstance(report, dict):
raise RuntimeError("Cannot synthesize job result because linked run-report.json is missing.")
result = {
"job_id": job_id,
"run_id": linked_run_record["run_id"],
"output_dir": _normalize_repo_path_value(str(linked_run_record["run_dir"])),
"raw_output": _normalize_repo_path_value(report.get("raw_output")),
"delivery_output": _normalize_repo_path_value(report.get("delivery_output")),
"digest_brief_output": _normalize_repo_path_value(report.get("digest_brief_output")),
"report_output": _normalize_repo_path(linked_run_record["run_dir"] / "run-report.json"),
"keyword_index": _normalize_keyword_index(report.get("keyword_index")),
"pulled_count": report.get("pulled_count"),
"delivered_count": report.get("delivered_count"),
"marked_read_count": report.get("marked_read_count"),
"status_counts": report.get("status_counts"),
"debug_artifacts": report.get("debug_artifacts"),
"completed_at": report.get("completed_at") or linked_run_record.get("finished_at"),
"result_source": "linked_run_report",
}
if bool(state.input.get("include_item_reports")):
result["items"] = report.get("items", [])
return result
def _build_stale_job_reason(state: Any) -> str | None:
updated_at = state.updated_at
stale_after_seconds = _estimate_job_stale_seconds(state)
age_seconds = (datetime.now(tz=updated_at.tzinfo) - updated_at).total_seconds()
if age_seconds < stale_after_seconds:
return None
current_stage = state.current_stage or "unknown_stage"
age_minutes = int(age_seconds // 60)
stale_after_minutes = int(stale_after_seconds // 60)
return (
f"job run-state has remained in running state at {current_stage} for about {age_minutes} minutes "
f"without result.json; it exceeded the stale threshold of {stale_after_minutes} minutes."
)
def _estimate_job_stale_seconds(state: Any) -> int:
input_payload = state.input if isinstance(state.input, dict) else {}
limit = _safe_int(input_payload.get("limit"), default=5)
timeout_seconds = _safe_float(input_payload.get("timeout_seconds"), default=60.0)
max_retries = _safe_int(input_payload.get("max_retries"), default=2)
estimated = int((limit * max(timeout_seconds, 1.0) * max(max_retries, 1)) + 20 * 60)
return max(MIN_JOB_STALE_SECONDS, min(MAX_JOB_STALE_SECONDS, estimated))
def _safe_int(value: Any, *, default: int) -> int:
try:
return int(value)
except (TypeError, ValueError):
return default
def _safe_float(value: Any, *, default: float) -> float:
try:
return float(value)
except (TypeError, ValueError):
return default
+221 -2
View File
@@ -45,6 +45,8 @@ DISCOVERED_ARTIFACTS = [
"relative_path": Path("candidates/digest-brief.json"),
},
]
MIN_RUNNING_STALE_SECONDS = 30 * 60
MAX_RUNNING_STALE_SECONDS = 6 * 60 * 60
def get_run_status(*, run_id: str) -> dict[str, Any]:
@@ -186,7 +188,7 @@ def _build_run_record(run_dir: Path) -> dict[str, Any]:
run_store = RunStore.load(path=state_path, repo_root=REPO_ROOT)
state = run_store.state
run_id = state.run_id
return {
record = {
"run_id": run_id,
"workflow": state.workflow,
"run_type": state.run_type,
@@ -204,12 +206,13 @@ def _build_run_record(run_dir: Path) -> dict[str, Any]:
"state": state,
"report": _load_json(report_path) if report_path.exists() else None,
}
return _reconcile_run_record(record)
report = _load_json(report_path) if report_path.exists() else None
run_id = str(report.get("run_id")) if isinstance(report, dict) and report.get("run_id") else run_dir.name
inferred_record = _infer_run_record_from_directory(run_dir=run_dir, report=report, run_id=run_id)
inferred_record["aliases"] = {run_id, run_dir.name}
return inferred_record
return _reconcile_run_record(inferred_record)
def _infer_run_record_from_directory(*, run_dir: Path, report: dict[str, Any] | None, run_id: str) -> dict[str, Any]:
@@ -298,6 +301,13 @@ def _build_status_response(record: dict[str, Any]) -> dict[str, Any]:
"artifacts": _collect_artifacts(record),
"recovery": record["recovery"],
"state_source": record["state_source"],
"status_source": record["status_source"],
"state_quality": record["state_quality"],
"state_conflict": record["state_conflict"],
"state_conflict_reason": record["state_conflict_reason"],
"raw_status": record["raw_status"],
"raw_current_stage": record["raw_current_stage"],
"artifact_presence": record["artifact_presence"],
}
@@ -310,6 +320,10 @@ def _build_run_lookup_response(record: dict[str, Any]) -> dict[str, Any]:
"status": record["status"],
"output_dir": _normalize_repo_path(record["run_dir"]),
"state_source": record["state_source"],
"status_source": record["status_source"],
"state_quality": record["state_quality"],
"state_conflict": record["state_conflict"],
"state_conflict_reason": record["state_conflict_reason"],
}
@@ -330,9 +344,191 @@ def _build_list_response(record: dict[str, Any]) -> dict[str, Any]:
"recovery": record["recovery"],
"artifact_count": len(_collect_artifacts(record)),
"state_source": record["state_source"],
"status_source": record["status_source"],
"state_quality": record["state_quality"],
"state_conflict": record["state_conflict"],
"state_conflict_reason": record["state_conflict_reason"],
"raw_status": record["raw_status"],
"raw_current_stage": record["raw_current_stage"],
}
def _reconcile_run_record(record: dict[str, Any]) -> dict[str, Any]:
raw_status = record["status"]
raw_current_stage = record["current_stage"]
raw_stages = record["stages"]
raw_recovery = record["recovery"]
artifact_presence = _build_artifact_presence(record["run_dir"], report=record.get("report"))
reconciled = dict(record)
reconciled["raw_status"] = raw_status
reconciled["raw_current_stage"] = raw_current_stage
reconciled["status_source"] = record["state_source"]
reconciled["state_quality"] = "trusted"
reconciled["state_conflict"] = False
reconciled["state_conflict_reason"] = None
reconciled["artifact_presence"] = artifact_presence
report = record.get("report")
if not isinstance(report, dict):
stale_reason = _build_stale_running_reason(record)
if stale_reason is not None:
reconciled["status"] = "failed"
reconciled["finished_at"] = record["updated_at"]
reconciled["stages"] = _build_stale_failed_stages(raw_stages, raw_current_stage)
reconciled["error"] = {
"type": "StaleRunState",
"message": stale_reason,
"stage": raw_current_stage,
"details": {
"raw_status": raw_status,
"raw_current_stage": raw_current_stage,
},
}
reconciled["status_source"] = "stale_run_state_timeout"
reconciled["state_quality"] = "reconciled"
reconciled["state_conflict"] = True
reconciled["state_conflict_reason"] = stale_reason
return reconciled
effective_status = _infer_status_from_report(report)
report_completed_at = _maybe_iso(report.get("completed_at"))
conflict = (
raw_status != effective_status
or raw_current_stage is not None
or any(stage["status"] in {"running", "failed"} for stage in raw_stages)
or bool(raw_recovery.get("resumable"))
)
if not conflict:
reconciled["artifact_presence"] = artifact_presence
return reconciled
reconciled["status"] = effective_status
reconciled["current_stage"] = None
reconciled["updated_at"] = report_completed_at or record["updated_at"]
reconciled["finished_at"] = report_completed_at or record["finished_at"]
reconciled["stages"] = _build_terminal_success_stages(raw_stages)
reconciled["error"] = None
reconciled["recovery"] = {
"resumable": False,
"resume_from_stage": None,
"last_success_stage": DEFAULT_STAGES[-1],
}
reconciled["status_source"] = "run_report_reconciliation"
reconciled["state_quality"] = "reconciled"
reconciled["state_conflict"] = True
reconciled["state_conflict_reason"] = (
"run-state.json did not converge, but run-report.json already proves the workflow reached a terminal state."
)
return reconciled
def _build_artifact_presence(run_dir: Path, *, report: dict[str, Any] | None) -> dict[str, bool]:
return {
"run_report": isinstance(report, dict) or (run_dir / "run-report.json").exists(),
"delivery_payload": (run_dir / "candidates" / "openclaw-delivery-payload.json").exists(),
"digest_brief": (run_dir / "candidates" / "digest-brief.json").exists(),
}
def _build_terminal_success_stages(stages: list[dict[str, Any]]) -> list[dict[str, Any]]:
stage_by_name = {stage["name"]: stage for stage in stages}
reconciled_stages: list[dict[str, Any]] = []
for stage_name in DEFAULT_STAGES:
existing = stage_by_name.get(stage_name)
if existing is None:
reconciled_stages.append(_stage_dict(name=stage_name, status="success"))
continue
reconciled_stages.append(
{
"name": stage_name,
"status": "success",
"started_at": existing.get("started_at"),
"finished_at": existing.get("finished_at"),
"outputs": existing.get("outputs", {}),
"error": None,
}
)
return reconciled_stages
def _build_stale_failed_stages(stages: list[dict[str, Any]], current_stage: str | None) -> list[dict[str, Any]]:
stage_by_name = {stage["name"]: stage for stage in stages}
reconciled_stages: list[dict[str, Any]] = []
for stage_name in DEFAULT_STAGES:
existing = stage_by_name.get(stage_name)
if existing is None:
reconciled_stages.append(_stage_dict(name=stage_name, status="pending"))
continue
status = existing.get("status")
if status == "running" or (current_stage is not None and stage_name == current_stage):
status = "failed"
reconciled_stages.append(
{
"name": stage_name,
"status": status,
"started_at": existing.get("started_at"),
"finished_at": existing.get("finished_at") or existing.get("started_at"),
"outputs": existing.get("outputs", {}),
"error": existing.get("error"),
}
)
return reconciled_stages
def _build_stale_running_reason(record: dict[str, Any]) -> str | None:
if record["status"] != "running":
return None
updated_at = _parse_iso_datetime(record.get("updated_at"))
if updated_at is None:
return None
stale_after_seconds = _estimate_running_stale_seconds(record)
age_seconds = (datetime.now(tz=updated_at.tzinfo) - updated_at).total_seconds()
if age_seconds < stale_after_seconds:
return None
current_stage = record.get("current_stage") or "unknown_stage"
age_minutes = int(age_seconds // 60)
stale_after_minutes = int(stale_after_seconds // 60)
return (
f"run-state.json has remained in running state at {current_stage} for about {age_minutes} minutes "
f"without terminal artifacts; it exceeded the stale threshold of {stale_after_minutes} minutes."
)
def _estimate_running_stale_seconds(record: dict[str, Any]) -> int:
state = record.get("state")
input_payload = state.input if isinstance(state, RunState) and isinstance(state.input, dict) else {}
limit = _safe_int(input_payload.get("limit"), default=5)
timeout_seconds = _safe_float(input_payload.get("timeout_seconds"), default=60.0)
max_retries = _safe_int(input_payload.get("max_retries"), default=2)
expected_items = limit
current_stage = record.get("current_stage")
if current_stage == "generate_summaries":
expected_items = _stage_output_from_record(record, "generate_summaries", "expected_items") or limit
elif current_stage == "extract_articles":
expected_items = _stage_output_from_record(record, "extract_articles", "expected_items") or limit
estimated = int((expected_items * max(timeout_seconds, 1.0) * max(max_retries, 1)) + 15 * 60)
return max(MIN_RUNNING_STALE_SECONDS, min(MAX_RUNNING_STALE_SECONDS, estimated))
def _stage_output_from_record(record: dict[str, Any], stage_name: str, key: str) -> Any:
for stage in record.get("stages", []):
if stage.get("name") != stage_name:
continue
outputs = stage.get("outputs")
if isinstance(outputs, dict):
return outputs.get(key)
return None
def _collect_artifacts(record: dict[str, Any]) -> list[dict[str, Any]]:
artifacts: list[dict[str, Any]] = []
seen_names: set[str] = set()
@@ -587,6 +783,29 @@ def _maybe_iso(value: Any) -> str | None:
return None
def _parse_iso_datetime(value: Any) -> datetime | None:
if not isinstance(value, str) or not value.strip():
return None
try:
return datetime.fromisoformat(value)
except ValueError:
return None
def _safe_int(value: Any, *, default: int) -> int:
try:
return int(value)
except (TypeError, ValueError):
return default
def _safe_float(value: Any, *, default: float) -> float:
try:
return float(value)
except (TypeError, ValueError):
return default
def _record_sort_key(record: dict[str, Any]) -> tuple[str, str]:
return (record.get("updated_at") or "", record["run_dir"].name)
+595
View File
@@ -0,0 +1,595 @@
from __future__ import annotations
import json
import os
import subprocess
import sys
from datetime import datetime
from pathlib import Path
from typing import Any
from uuid import uuid4
from .query_service import _resolve_run_record
from .resume_service import (
SUPPORTED_RESUME_STAGES,
UNSUPPORTED_RESUME_STAGES,
_build_resume_plan,
_resume_freshrss_run,
_validate_resume_artifacts,
inspect_resume_plan,
)
from .run_store import RunStore
from .state_models import RunState
REPO_ROOT = Path(__file__).resolve().parents[3]
OUTPUT_ROOT = REPO_ROOT / "outputs" / "freshrss"
RESUME_JOBS_ROOT = OUTPUT_ROOT / "resume_jobs"
WORKFLOW_NAME = "freshrss_resume_job"
RUN_TYPE = "resume_job"
RUN_STATE_FILENAME = "run-state.json"
INPUT_FILENAME = "input.json"
RESULT_FILENAME = "result.json"
JOB_REPORT_FILENAME = "job-report.json"
DEFAULT_STAGES = [
"prepare_job",
"validate_resume_plan",
"resume_run",
"write_result",
]
MIN_JOB_STALE_SECONDS = 30 * 60
MAX_JOB_STALE_SECONDS = 6 * 60 * 60
def _now() -> datetime:
return datetime.now().astimezone()
def _new_job_id() -> str:
ts = _now().strftime("%Y%m%d-%H%M%S")
return f"freshrss-resume-job-{ts}-{uuid4().hex[:8]}"
def _job_dir(job_id: str) -> Path:
return RESUME_JOBS_ROOT / job_id
def _run_state_path(job_id: str) -> Path:
return _job_dir(job_id) / RUN_STATE_FILENAME
def _input_path(job_id: str) -> Path:
return _job_dir(job_id) / INPUT_FILENAME
def _result_path(job_id: str) -> Path:
return _job_dir(job_id) / RESULT_FILENAME
def _job_report_path(job_id: str) -> Path:
return _job_dir(job_id) / JOB_REPORT_FILENAME
def _normalize_repo_path(path: Path) -> str:
try:
return str(path.resolve().relative_to(REPO_ROOT.resolve()))
except ValueError:
return str(path)
def _normalize_repo_path_value(path_value: str | None) -> str | None:
if not path_value:
return None
return _normalize_repo_path(Path(path_value))
def _load_run_store(job_id: str) -> RunStore:
return RunStore.load(path=_run_state_path(job_id), repo_root=REPO_ROOT)
def _load_result(job_id: str) -> dict[str, Any] | None:
path = _result_path(job_id)
if not path.exists():
return None
return json.loads(path.read_text(encoding="utf-8-sig"))
def _load_job_report(job_id: str) -> dict[str, Any] | None:
path = _job_report_path(job_id)
if not path.exists():
return None
return json.loads(path.read_text(encoding="utf-8-sig"))
def _write_job_report(*, job_id: str, payload: dict[str, Any]) -> Path:
report_file = _job_report_path(job_id)
report_file.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8")
return report_file
def _build_result_payload(
*,
job_id: str,
run_id: str,
run_dir: Path,
resume_plan: dict[str, Any],
resume_result: dict[str, Any],
) -> dict[str, Any]:
return {
"job_id": job_id,
"run_id": run_id,
"resume_from_stage": resume_plan["effective_resume_from_stage"],
"requested_resume_from_stage": resume_plan["requested_resume_from_stage"],
"resume_decision_source": resume_plan["decision_source"],
"artifact_resume_from_stage": resume_plan["artifact_resume_from_stage"],
"artifact_snapshot": resume_plan["artifact_snapshot"],
"status": resume_result["status"],
"output_dir": _normalize_repo_path(run_dir),
"delivery_output": _normalize_repo_path_value(resume_result.get("delivery_output")),
"report_output": _normalize_repo_path_value(resume_result.get("report_output")),
"digest_brief_output": _normalize_repo_path_value(resume_result.get("digest_brief_output")),
"pulled_count": resume_result.get("pulled_count"),
"delivered_count": resume_result.get("delivered_count"),
"marked_read_count": resume_result.get("marked_read_count"),
"status_counts": resume_result.get("status_counts"),
"keyword_index": resume_result.get("keyword_index"),
"completed_at": _now().isoformat(),
"result_source": "resume_job_result",
}
def _build_result_from_linked_run(*, job_id: str, linked_run_record: dict[str, Any]) -> dict[str, Any] | None:
report = linked_run_record.get("report")
if not isinstance(report, dict):
return None
recovery = linked_run_record.get("recovery")
requested_resume_from_stage = None
if isinstance(recovery, dict):
requested_resume_from_stage = recovery.get("resume_from_stage")
return {
"job_id": job_id,
"run_id": linked_run_record["run_id"],
"resume_from_stage": requested_resume_from_stage,
"requested_resume_from_stage": requested_resume_from_stage,
"resume_decision_source": "linked_run_report",
"artifact_resume_from_stage": None,
"artifact_snapshot": None,
"status": linked_run_record["status"],
"output_dir": _normalize_repo_path(linked_run_record["run_dir"]),
"delivery_output": _normalize_repo_path_value(report.get("delivery_output")),
"report_output": _normalize_repo_path(linked_run_record["run_dir"] / "run-report.json"),
"digest_brief_output": _normalize_repo_path_value(report.get("digest_brief_output")),
"pulled_count": report.get("pulled_count"),
"delivered_count": report.get("delivered_count"),
"marked_read_count": report.get("marked_read_count"),
"status_counts": report.get("status_counts"),
"keyword_index": report.get("keyword_index"),
"completed_at": report.get("completed_at"),
"result_source": "linked_run_report",
}
def _linked_run_result_available(linked_run_record: dict[str, Any]) -> bool:
return linked_run_record["status"] in {"success", "partial"} and isinstance(linked_run_record.get("report"), dict)
def _load_linked_run_record(state: Any) -> dict[str, Any] | None:
if not isinstance(state.input, dict):
return None
run_id = state.input.get("run_id")
if not isinstance(run_id, str) or not run_id.strip():
return None
try:
return _resolve_run_record(run_id)
except FileNotFoundError:
return None
def _build_stale_job_reason(state: Any) -> str | None:
if state.status != "running" or state.current_stage is None or state.finished_at is not None:
return None
now = _now()
age_seconds = max(0.0, (now - state.updated_at).total_seconds())
if age_seconds < MIN_JOB_STALE_SECONDS:
return None
if age_seconds >= MAX_JOB_STALE_SECONDS:
return (
f"Resume job has remained in stage '{state.current_stage}' for more than {int(MAX_JOB_STALE_SECONDS)} seconds "
"without producing a terminal result; treating the job state as stale."
)
return None
def start_resume_job(*, run_id: str) -> dict[str, Any]:
resume_view = inspect_resume_plan(run_id=run_id)
job_id = _new_job_id()
job_dir = _job_dir(job_id)
job_dir.mkdir(parents=True, exist_ok=True)
started_at = _now()
input_payload = {
"run_id": run_id,
"requested_resume_from_stage": resume_view.get("requested_resume_from_stage"),
"resume_from_stage": resume_view.get("resume_from_stage"),
"resume_decision_source": resume_view.get("resume_decision_source"),
"recommended_action": resume_view.get("recommended_action"),
"artifact_resume_from_stage": resume_view.get("artifact_resume_from_stage"),
"artifact_snapshot": resume_view.get("artifact_snapshot"),
"launcher_pid": os.getpid(),
}
store = RunStore.create(
path=_run_state_path(job_id),
run_id=job_id,
workflow=WORKFLOW_NAME,
run_type=RUN_TYPE,
started_at=started_at,
input_payload=input_payload,
repo_root=REPO_ROOT,
)
for stage_name in DEFAULT_STAGES:
store._get_or_create_stage(stage_name)
store.save()
store.start_stage("prepare_job")
input_file = _input_path(job_id)
input_file.write_text(json.dumps(input_payload, ensure_ascii=False, indent=2), encoding="utf-8")
store.register_artifact(name="job_input", path=input_file, kind="json", stage="prepare_job")
if not resume_view.get("can_resume") or resume_view.get("recommended_action") != "resume":
report_file = _write_job_report(
job_id=job_id,
payload={
"job_id": job_id,
"status": "failed",
"error_type": "ResumePreflightRejected",
"error_message": resume_view.get("message"),
"failed_stage": "validate_resume_plan",
"linked_run_id": run_id,
"resume_from_stage": resume_view.get("resume_from_stage"),
"requested_resume_from_stage": resume_view.get("requested_resume_from_stage"),
"recommended_action": resume_view.get("recommended_action"),
},
)
store.register_artifact(name="job_report", path=report_file, kind="json", stage="prepare_job")
store.finish_stage("prepare_job", outputs={"linked_run_id": run_id})
store.fail_stage("validate_resume_plan", error=ValueError(str(resume_view.get("message") or "Resume preflight rejected.")))
return {
"job_id": job_id,
"workflow": WORKFLOW_NAME,
"run_type": RUN_TYPE,
"status": "failed",
"output_dir": _normalize_repo_path(job_dir),
"linked_run_id": run_id,
"resume_from_stage": resume_view.get("resume_from_stage"),
"recommended_action": resume_view.get("recommended_action"),
"message": resume_view.get("message"),
}
runner_script = REPO_ROOT / "scripts" / "run_resume_job.py"
cmd = [sys.executable, str(runner_script), "--job-id", job_id]
try:
proc = subprocess.Popen(
cmd,
cwd=str(REPO_ROOT),
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
start_new_session=True,
)
except Exception as exc:
report_file = _write_job_report(
job_id=job_id,
payload={
"job_id": job_id,
"status": "failed",
"error_type": type(exc).__name__,
"error_message": str(exc),
"failed_stage": "prepare_job",
"linked_run_id": run_id,
"resume_from_stage": resume_view.get("resume_from_stage"),
},
)
store.register_artifact(name="job_report", path=report_file, kind="json", stage="prepare_job")
store.fail_stage("prepare_job", error=exc)
raise
store.finish_stage(
"prepare_job",
outputs={
"runner_pid": proc.pid,
"runner_command": cmd,
"linked_run_id": run_id,
"resume_from_stage": resume_view.get("resume_from_stage"),
},
)
return {
"job_id": job_id,
"workflow": WORKFLOW_NAME,
"run_type": RUN_TYPE,
"status": "running",
"output_dir": _normalize_repo_path(job_dir),
"linked_run_id": run_id,
"resume_from_stage": resume_view.get("resume_from_stage"),
"message": "Resume job started successfully. Use get_resume_job_status to poll progress.",
}
def run_resume_job(*, job_id: str) -> dict[str, Any]:
store = _load_run_store(job_id)
input_payload = json.loads(_input_path(job_id).read_text(encoding="utf-8-sig"))
current_stage = "resume_run"
run_id = str(input_payload["run_id"])
try:
store.start_stage("validate_resume_plan")
record = _resolve_run_record(run_id)
if record["state_source"] != "run_state" or not isinstance(record.get("state"), RunState):
raise RuntimeError("This run cannot be resumed because run-state.json is missing or could not be loaded.")
run_store = RunStore.load(path=record["run_dir"] / "run-state.json", repo_root=REPO_ROOT)
resume_plan = _build_resume_plan(record=record, state=run_store.state)
if resume_plan["decision"] != "resume":
raise RuntimeError(str(resume_plan["message"]))
resume_from_stage = resume_plan["effective_resume_from_stage"]
if resume_from_stage in UNSUPPORTED_RESUME_STAGES:
raise RuntimeError(f"This run cannot be resumed from {resume_from_stage} in the current implementation.")
if resume_from_stage not in SUPPORTED_RESUME_STAGES:
raise RuntimeError(f"This run cannot be resumed because stage '{resume_from_stage}' is not supported.")
missing_artifacts = _validate_resume_artifacts(
record=record,
state=run_store.state,
resume_from_stage=resume_from_stage,
)
if missing_artifacts:
raise RuntimeError(
f"This run cannot be resumed from {resume_from_stage} because required artifacts are missing: {missing_artifacts}"
)
store.finish_stage(
"validate_resume_plan",
outputs={
"linked_run_id": run_id,
"resume_from_stage": resume_from_stage,
"requested_resume_from_stage": resume_plan["requested_resume_from_stage"],
"resume_decision_source": resume_plan["decision_source"],
},
)
current_stage = "resume_run"
store.start_stage("resume_run")
resume_result = _resume_freshrss_run(record=record, run_store=run_store, resume_from_stage=resume_from_stage)
store.finish_stage(
"resume_run",
outputs={
"linked_run_id": run_id,
"resume_from_stage": resume_from_stage,
"linked_run_status": resume_result["status"],
"delivery_output": _normalize_repo_path_value(resume_result.get("delivery_output")),
"report_output": _normalize_repo_path_value(resume_result.get("report_output")),
},
)
store.start_stage("write_result")
result = _build_result_payload(
job_id=job_id,
run_id=run_id,
run_dir=record["run_dir"],
resume_plan=resume_plan,
resume_result=resume_result,
)
result_file = _result_path(job_id)
result_file.write_text(json.dumps(result, ensure_ascii=False, indent=2), encoding="utf-8")
store.register_artifact(name="job_result", path=result_file, kind="json", stage="write_result")
report_file = _write_job_report(
job_id=job_id,
payload={
"job_id": job_id,
"status": "success",
"run_id": run_id,
"resume_from_stage": result["resume_from_stage"],
"delivery_output": result["delivery_output"],
"report_output": result["report_output"],
"digest_brief_output": result["digest_brief_output"],
},
)
store.register_artifact(name="job_report", path=report_file, kind="json", stage="write_result")
store.finish_stage(
"write_result",
outputs={
"result_path": _normalize_repo_path(result_file),
"linked_run_id": run_id,
},
)
store.finish_run(status="success")
return result
except Exception as exc:
current_stage = store.state.current_stage or current_stage
report_file = _write_job_report(
job_id=job_id,
payload={
"job_id": job_id,
"status": "failed",
"error_type": type(exc).__name__,
"error_message": str(exc),
"failed_stage": current_stage,
"linked_run_id": run_id,
},
)
try:
store.register_artifact(name="job_report", path=report_file, kind="json", stage=current_stage)
except Exception:
pass
store.fail_stage(current_stage, error=exc)
raise
def _resolve_effective_job_view(
*,
job_id: str,
state: Any,
result: dict[str, Any] | None = None,
) -> dict[str, Any]:
resolved_result = result or _load_result(job_id)
linked_run_record = _load_linked_run_record(state)
raw_status = state.status
raw_current_stage = state.current_stage
effective_status = raw_status
effective_current_stage = raw_current_stage
effective_updated_at = state.updated_at.isoformat()
effective_finished_at = state.finished_at.isoformat() if state.finished_at else None
status_source = "job_run_state"
state_quality = "trusted"
state_conflict = False
state_conflict_reason = None
status_note = None
if resolved_result is not None:
completed_at = resolved_result.get("completed_at")
effective_status = "success"
effective_current_stage = None
effective_updated_at = completed_at or effective_updated_at
effective_finished_at = completed_at or effective_finished_at
if raw_status != "success" or raw_current_stage is not None:
status_source = "job_result_reconciliation"
state_quality = "reconciled"
state_conflict = True
state_conflict_reason = "result.json already exists, but the resume job run-state did not converge to success."
status_note = "resume job result exists; the job can be treated as completed."
elif linked_run_record is not None:
if _linked_run_result_available(linked_run_record) and raw_status == "running":
effective_status = "success"
effective_current_stage = None
effective_updated_at = linked_run_record.get("updated_at") or effective_updated_at
effective_finished_at = linked_run_record.get("finished_at") or effective_finished_at
status_source = "linked_run_reconciliation"
state_quality = "reconciled"
state_conflict = raw_status != "success" or raw_current_stage is not None
state_conflict_reason = "The linked run already has terminal artifacts, but the resume job run-state did not converge."
status_note = "linked run artifacts are complete; the resume job can be treated as completed."
resolved_result = _build_result_from_linked_run(job_id=job_id, linked_run_record=linked_run_record)
elif raw_status == "running":
stale_reason = _build_stale_job_reason(state)
if stale_reason is not None:
effective_status = "failed"
effective_current_stage = raw_current_stage
effective_finished_at = effective_updated_at
status_source = "stale_job_state_timeout"
state_quality = "reconciled"
state_conflict = True
state_conflict_reason = stale_reason
status_note = stale_reason
return {
"status": effective_status,
"current_stage": effective_current_stage,
"updated_at": effective_updated_at,
"finished_at": effective_finished_at,
"status_source": status_source,
"state_quality": state_quality,
"state_conflict": state_conflict,
"state_conflict_reason": state_conflict_reason,
"status_note": status_note,
"result": resolved_result,
"linked_run_status": linked_run_record["status"] if linked_run_record is not None else None,
"linked_run_state_conflict": linked_run_record.get("state_conflict") if linked_run_record is not None else None,
}
def _build_effective_job_progress(*, state: Any, effective_status: str) -> dict[str, int]:
completed_stage_count = sum(1 for stage in state.stages if stage.status == "success")
running_stage_count = sum(1 for stage in state.stages if stage.status == "running")
failed_stage_count = sum(1 for stage in state.stages if stage.status == "failed")
pending_stage_count = sum(1 for stage in state.stages if stage.status == "pending")
if effective_status == "success" and running_stage_count > 0:
pending_stage_count += running_stage_count
running_stage_count = 0
return {
"completed_stage_count": completed_stage_count,
"running_stage_count": running_stage_count,
"failed_stage_count": failed_stage_count,
"pending_stage_count": pending_stage_count,
"total_stage_count": len(state.stages),
}
def get_resume_job_status(*, job_id: str) -> dict[str, Any]:
store = _load_run_store(job_id)
state = store.state
effective = _resolve_effective_job_view(job_id=job_id, state=state)
linked_run_id = state.input.get("run_id") if isinstance(state.input, dict) else None
return {
"job_id": state.run_id,
"workflow": state.workflow,
"run_type": state.run_type,
"status": effective["status"],
"current_stage": effective["current_stage"],
"started_at": state.started_at.isoformat(),
"updated_at": effective["updated_at"],
"finished_at": effective["finished_at"],
"output_dir": _normalize_repo_path(_job_dir(job_id)),
"linked_run_id": linked_run_id,
"progress": _build_effective_job_progress(state=state, effective_status=effective["status"]),
"artifacts": [artifact.model_dump(mode="json") for artifact in state.artifacts],
"error_summary": state.error.model_dump(mode="json") if state.error else None,
"status_source": effective["status_source"],
"state_quality": effective["state_quality"],
"state_conflict": effective["state_conflict"],
"state_conflict_reason": effective["state_conflict_reason"],
"status_note": effective["status_note"],
"raw_status": state.status,
"raw_current_stage": state.current_stage,
"linked_run_status": effective["linked_run_status"],
"linked_run_state_conflict": effective["linked_run_state_conflict"],
}
def get_resume_job_result(*, job_id: str) -> dict[str, Any]:
store = _load_run_store(job_id)
state = store.state
result = _load_result(job_id)
effective = _resolve_effective_job_view(job_id=job_id, state=state, result=result)
synthesized_result = result or effective["result"]
if effective["status"] != "success" or synthesized_result is None:
message = "Resume job result is not ready."
if effective["status"] == "failed":
message = "Resume job did not complete successfully, so no terminal job result is available."
return {
"job_id": state.run_id,
"status": effective["status"],
"message": message,
"linked_run_id": state.input.get("run_id") if isinstance(state.input, dict) else None,
"error_summary": state.error.model_dump(mode="json") if state.error else None,
"status_source": effective["status_source"],
"status_note": effective["status_note"],
}
artifact = next((a.model_dump(mode="json") for a in state.artifacts if a.name == "job_result"), None)
return {
"job_id": state.run_id,
"status": effective["status"],
"run_id": synthesized_result.get("run_id"),
"resume_from_stage": synthesized_result.get("resume_from_stage"),
"requested_resume_from_stage": synthesized_result.get("requested_resume_from_stage"),
"resume_decision_source": synthesized_result.get("resume_decision_source"),
"artifact_resume_from_stage": synthesized_result.get("artifact_resume_from_stage"),
"artifact_snapshot": synthesized_result.get("artifact_snapshot"),
"output_dir": synthesized_result.get("output_dir"),
"delivery_output": synthesized_result.get("delivery_output"),
"report_output": synthesized_result.get("report_output"),
"digest_brief_output": synthesized_result.get("digest_brief_output"),
"pulled_count": synthesized_result.get("pulled_count"),
"delivered_count": synthesized_result.get("delivered_count"),
"marked_read_count": synthesized_result.get("marked_read_count"),
"status_counts": synthesized_result.get("status_counts"),
"keyword_index": synthesized_result.get("keyword_index"),
"artifact": artifact,
"result": synthesized_result,
"status_source": effective["status_source"],
"status_note": effective["status_note"],
"result_source": synthesized_result.get("result_source", "job_result"),
"linked_run_status": effective["linked_run_status"],
}
+426 -18
View File
@@ -25,6 +25,7 @@ from summary_mcp.models.openclaw_delivery import (
)
from summary_mcp.models.summary_io import ExtractionOutput
from summary_mcp.workflows.freshrss_pipeline import (
CANDIDATE_BATCH_ARTIFACT,
DEFAULT_PROMPT_PATH,
DEFAULT_RULES_PATH,
DEFAULT_TERM_ALIASES_PATH,
@@ -37,13 +38,18 @@ from summary_mcp.workflows.freshrss_pipeline import (
FILTER_STAGE,
REPORT_STAGE,
REPO_ROOT,
SUMMARY_BATCH_ARTIFACT,
SUMMARY_STAGE,
WORKFLOW_NAME,
_build_item_context,
_candidate_batch_output,
_final_run_status,
_load_json,
_load_required_env,
_persist_candidate_batch_artifact,
_persist_summary_batch_artifact,
_save_json,
_summary_batch_output,
)
from .query_service import _normalize_repo_path, _resolve_repo_path, _resolve_run_record
@@ -62,6 +68,12 @@ UNSUPPORTED_RESUME_STAGES = {
}
UTC = timezone.utc
DEFAULT_STREAM_ID = "user/-/state/com.google/reading-list"
RESUME_STAGE_ORDER = {
SUMMARY_STAGE: 1,
FILTER_STAGE: 2,
DELIVERY_STAGE: 3,
REPORT_STAGE: 4,
}
def resume_run(*, run_id: str) -> dict[str, Any]:
@@ -71,37 +83,68 @@ def resume_run(*, run_id: str) -> dict[str, Any]:
"workflow": record["workflow"],
"status": record["status"],
"output_dir": _normalize_repo_path(record["run_dir"]),
"state_source": record["state_source"],
"status_source": record.get("status_source"),
}
if record["status"] in {"success", "partial"} and isinstance(record.get("report"), dict):
return {
**base_response,
"resumed": False,
"resume_from_stage": None,
"requested_resume_from_stage": None,
"resume_decision_source": "artifacts",
"recommended_action": "read_terminal_result",
"message": "This run already has a terminal run-report.json; prefer reading get_run_status/get_run_report instead of resuming.",
"missing_artifacts": [],
"state_conflict": record.get("state_conflict", False),
"state_conflict_reason": record.get("state_conflict_reason"),
}
if record["state_source"] != "run_state" or not isinstance(record.get("state"), RunState):
return {
**base_response,
"resumed": False,
"resume_from_stage": None,
"requested_resume_from_stage": None,
"resume_decision_source": "unavailable",
"recommended_action": "start_new_run",
"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)
resume_plan = _build_resume_plan(record=record, state=state)
requested_resume_from_stage = resume_plan["requested_resume_from_stage"]
resume_from_stage = resume_plan["effective_resume_from_stage"]
if state.workflow != WORKFLOW_NAME:
return {
**base_response,
"resumed": False,
"resume_from_stage": resume_from_stage,
"requested_resume_from_stage": requested_resume_from_stage,
"resume_decision_source": resume_plan["decision_source"],
"recommended_action": "start_new_run",
"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:
if resume_plan["decision"] != "resume":
return {
**base_response,
"resumed": False,
"resume_from_stage": None,
"message": "This run does not expose a recoverable stage in run-state.json.",
"missing_artifacts": [],
"resume_from_stage": resume_from_stage,
"requested_resume_from_stage": requested_resume_from_stage,
"resume_decision_source": resume_plan["decision_source"],
"recommended_action": resume_plan["recommended_action"],
"message": resume_plan["message"],
"missing_artifacts": resume_plan["missing_artifacts"],
"artifact_resume_from_stage": resume_plan["artifact_resume_from_stage"],
"artifact_snapshot": resume_plan["artifact_snapshot"],
"state_conflict": record.get("state_conflict", False),
"state_conflict_reason": record.get("state_conflict_reason"),
}
if resume_from_stage in UNSUPPORTED_RESUME_STAGES:
@@ -109,6 +152,9 @@ def resume_run(*, run_id: str) -> dict[str, Any]:
**base_response,
"resumed": False,
"resume_from_stage": resume_from_stage,
"requested_resume_from_stage": requested_resume_from_stage,
"resume_decision_source": resume_plan["decision_source"],
"recommended_action": "start_new_run",
"message": f"This run cannot be resumed from {resume_from_stage} in the current minimal implementation.",
"missing_artifacts": [],
}
@@ -118,6 +164,9 @@ def resume_run(*, run_id: str) -> dict[str, Any]:
**base_response,
"resumed": False,
"resume_from_stage": resume_from_stage,
"requested_resume_from_stage": requested_resume_from_stage,
"resume_decision_source": resume_plan["decision_source"],
"recommended_action": "start_new_run",
"message": f"This run cannot be resumed because stage '{resume_from_stage}' is not supported.",
"missing_artifacts": [],
}
@@ -128,8 +177,13 @@ def resume_run(*, run_id: str) -> dict[str, Any]:
**base_response,
"resumed": False,
"resume_from_stage": resume_from_stage,
"requested_resume_from_stage": requested_resume_from_stage,
"resume_decision_source": resume_plan["decision_source"],
"recommended_action": "start_new_run",
"message": f"This run cannot be resumed from {resume_from_stage} because required artifacts are missing.",
"missing_artifacts": missing_artifacts,
"artifact_resume_from_stage": resume_plan["artifact_resume_from_stage"],
"artifact_snapshot": resume_plan["artifact_snapshot"],
}
try:
@@ -141,10 +195,15 @@ def resume_run(*, run_id: str) -> dict[str, Any]:
**base_response,
"resumed": True,
"resume_from_stage": resume_from_stage,
"requested_resume_from_stage": requested_resume_from_stage,
"resume_decision_source": resume_plan["decision_source"],
"recommended_action": "inspect_error",
"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,
"artifact_resume_from_stage": resume_plan["artifact_resume_from_stage"],
"artifact_snapshot": resume_plan["artifact_snapshot"],
}
return {
@@ -152,6 +211,9 @@ def resume_run(*, run_id: str) -> dict[str, Any]:
"workflow": run_store.state.workflow,
"resumed": True,
"resume_from_stage": resume_from_stage,
"requested_resume_from_stage": requested_resume_from_stage,
"resume_decision_source": resume_plan["decision_source"],
"recommended_action": "none",
"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']}.",
@@ -166,9 +228,92 @@ def resume_run(*, run_id: str) -> dict[str, Any]:
"marked_read_count": result.get("marked_read_count"),
"status_counts": result.get("status_counts"),
"missing_artifacts": [],
"artifact_resume_from_stage": resume_plan["artifact_resume_from_stage"],
"artifact_snapshot": resume_plan["artifact_snapshot"],
}
def inspect_resume_plan(*, 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"]),
"state_source": record["state_source"],
"status_source": record.get("status_source"),
"state_conflict": record.get("state_conflict", False),
"state_conflict_reason": record.get("state_conflict_reason"),
}
if record["status"] in {"success", "partial"} and isinstance(record.get("report"), dict):
return {
**base_response,
"can_resume": False,
"resume_from_stage": None,
"requested_resume_from_stage": None,
"resume_decision_source": "artifacts",
"recommended_action": "read_terminal_result",
"message": "This run already has a terminal run-report.json; prefer reading get_run_status/get_run_report instead of resuming.",
"missing_artifacts": [],
}
if record["state_source"] != "run_state" or not isinstance(record.get("state"), RunState):
return {
**base_response,
"can_resume": False,
"resume_from_stage": None,
"requested_resume_from_stage": None,
"resume_decision_source": "unavailable",
"recommended_action": "start_new_run",
"message": "This run cannot be resumed because run-state.json is missing or could not be loaded.",
"missing_artifacts": [],
}
state = record["state"]
resume_plan = _build_resume_plan(record=record, state=state)
response = {
**base_response,
"can_resume": False,
"resume_from_stage": resume_plan["effective_resume_from_stage"],
"requested_resume_from_stage": resume_plan["requested_resume_from_stage"],
"resume_decision_source": resume_plan["decision_source"],
"recommended_action": resume_plan["recommended_action"],
"message": resume_plan["message"],
"missing_artifacts": resume_plan["missing_artifacts"],
"artifact_resume_from_stage": resume_plan["artifact_resume_from_stage"],
"artifact_snapshot": resume_plan["artifact_snapshot"],
}
if state.workflow != WORKFLOW_NAME:
response["message"] = (
f"This run cannot be resumed because workflow '{state.workflow}' is not supported by the minimal resume_run implementation."
)
return response
resume_from_stage = resume_plan["effective_resume_from_stage"]
if resume_plan["decision"] != "resume":
return response
if resume_from_stage in UNSUPPORTED_RESUME_STAGES:
response["message"] = f"This run cannot be resumed from {resume_from_stage} in the current minimal implementation."
response["recommended_action"] = "start_new_run"
return response
if resume_from_stage not in SUPPORTED_RESUME_STAGES:
response["message"] = f"This run cannot be resumed because stage '{resume_from_stage}' is not supported."
response["recommended_action"] = "start_new_run"
return response
missing_artifacts = _validate_resume_artifacts(record=record, state=state, resume_from_stage=resume_from_stage)
if missing_artifacts:
response["message"] = f"This run cannot be resumed from {resume_from_stage} because required artifacts are missing."
response["missing_artifacts"] = missing_artifacts
response["recommended_action"] = "start_new_run"
return response
response["can_resume"] = True
return response
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"]
@@ -304,6 +449,75 @@ def _build_resume_config(*, state: RunState, run_dir: Path) -> dict[str, Any]:
}
def _find_summary_batch_path(*, run_dir: Path, state: RunState) -> Path | None:
return _find_artifact_path(
run_dir=run_dir,
state=state,
artifact_name=SUMMARY_BATCH_ARTIFACT,
relative_path=Path("summary/summary-batch.json"),
)
def _find_candidate_batch_path(*, run_dir: Path, state: RunState) -> Path | None:
return _find_artifact_path(
run_dir=run_dir,
state=state,
artifact_name=CANDIDATE_BATCH_ARTIFACT,
relative_path=Path("candidates/candidate-batch.json"),
)
def _load_summary_batch_lookup(path: Path | None) -> tuple[bool, dict[str, dict[str, Any]]]:
if path is None or not path.exists():
return False, {}
try:
payload = _load_json(path)
except Exception:
return False, {}
items = payload.get("items")
if not isinstance(items, list):
return False, {}
summaries_by_item_key: dict[str, dict[str, Any]] = {}
for entry in items:
if not isinstance(entry, dict):
return False, {}
item_key = entry.get("item_key")
summary = entry.get("summary")
if not isinstance(item_key, str) or not item_key.strip() or not isinstance(summary, dict):
return False, {}
summaries_by_item_key[item_key] = summary
return True, summaries_by_item_key
def _load_candidate_batch_lookup(path: Path | None) -> tuple[bool, dict[str, OpenClawCandidateInput]]:
if path is None or not path.exists():
return False, {}
try:
payload = _load_json(path)
except Exception:
return False, {}
items = payload.get("items")
if not isinstance(items, list):
return False, {}
candidates_by_item_key: dict[str, OpenClawCandidateInput] = {}
try:
for entry in items:
if not isinstance(entry, dict):
return False, {}
item_key = entry.get("item_key")
candidate_payload = entry.get("candidate")
if not isinstance(item_key, str) or not item_key.strip() or not isinstance(candidate_payload, dict):
return False, {}
candidates_by_item_key[item_key] = OpenClawCandidateInput.model_validate(candidate_payload)
except Exception:
return False, {}
return True, candidates_by_item_key
def _load_item_contexts(
*,
items: list[Any],
@@ -313,6 +527,8 @@ def _load_item_contexts(
include_candidates: bool = False,
delivery_candidates_by_id: dict[str, OpenClawCandidateInput] | None = None,
) -> list[dict[str, Any]]:
summary_batch_valid, summary_batch_by_item_key = _load_summary_batch_lookup(_summary_batch_output(run_dir))
candidate_batch_valid, candidate_batch_by_item_key = _load_candidate_batch_lookup(_candidate_batch_output(run_dir))
contexts: list[dict[str, Any]] = []
for index, item in enumerate(items, start=1):
context = _build_item_context(
@@ -334,16 +550,24 @@ def _load_item_contexts(
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"
if include_summaries:
summary_payload: dict[str, Any] | None = None
if context["summary_output"] is not None and context["summary_output"].exists():
summary_payload = _load_json(context["summary_output"])
elif summary_batch_valid:
summary_payload = summary_batch_by_item_key.get(context["item_key"])
if summary_payload is not None:
context["summary_payload"] = summary_payload
context["item_report"]["status"] = "summarized"
elif 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 include_candidates and candidate_batch_valid:
candidate = candidate_batch_by_item_key.get(context["item_key"])
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)
@@ -406,6 +630,11 @@ def _run_summary_stage(*, run_store: RunStore, item_contexts: list[dict[str, Any
},
)
summary_batch_output = _persist_summary_batch_artifact(
run_store=run_store,
run_dir=config["run_dir"],
item_contexts=item_contexts,
)
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(
@@ -415,6 +644,7 @@ def _run_summary_stage(*, run_store: RunStore, item_contexts: list[dict[str, Any
"completed_items": summary_success_count + summary_failed_count,
"success_count": summary_success_count,
"failed_count": summary_failed_count,
"summary_batch_output": str(summary_batch_output),
},
)
@@ -500,6 +730,11 @@ def _run_filter_stage(*, run_store: RunStore, item_contexts: list[dict[str, Any]
},
)
candidate_batch_output = _persist_candidate_batch_artifact(
run_store=run_store,
run_dir=config["run_dir"],
item_contexts=item_contexts,
)
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(
@@ -511,6 +746,7 @@ def _run_filter_stage(*, run_store: RunStore, item_contexts: list[dict[str, Any]
"keep_count": keep_count,
"review_count": review_count,
"drop_count": drop_count,
"candidate_batch_output": str(candidate_batch_output),
},
)
@@ -682,6 +918,165 @@ def _build_keyword_index_result(
return result
def _build_resume_plan(*, record: dict[str, Any], state: RunState) -> dict[str, Any]:
requested_resume_from_stage = _resolve_resume_from_stage(state)
artifact_snapshot = _collect_resume_artifact_snapshot(record=record, state=state)
artifact_resume_from_stage = _resolve_resume_stage_from_artifacts(artifact_snapshot)
if artifact_snapshot["has_run_report"]:
return {
"decision": "reject_terminal",
"decision_source": "artifacts",
"requested_resume_from_stage": requested_resume_from_stage,
"effective_resume_from_stage": None,
"artifact_resume_from_stage": artifact_resume_from_stage,
"recommended_action": "read_terminal_result",
"message": "This run already has a terminal run-report.json; prefer reading get_run_status/get_run_report instead of resuming.",
"missing_artifacts": [],
"artifact_snapshot": artifact_snapshot,
}
if artifact_resume_from_stage is None:
return {
"decision": "reject_unrecoverable",
"decision_source": "artifacts",
"requested_resume_from_stage": requested_resume_from_stage,
"effective_resume_from_stage": None,
"artifact_resume_from_stage": None,
"recommended_action": "start_new_run",
"message": "This run does not expose a safe artifact-backed resume point; start a new run instead.",
"missing_artifacts": artifact_snapshot["missing_for_next_resume"],
"artifact_snapshot": artifact_snapshot,
}
decision_source = "artifacts"
message = f"Resume will continue from {artifact_resume_from_stage} based on available artifacts."
if requested_resume_from_stage == artifact_resume_from_stage:
decision_source = "state_and_artifacts"
message = f"Resume stage {artifact_resume_from_stage} was confirmed by both run-state.json and artifacts."
elif requested_resume_from_stage is not None:
requested_rank = RESUME_STAGE_ORDER.get(requested_resume_from_stage, -1)
artifact_rank = RESUME_STAGE_ORDER.get(artifact_resume_from_stage, -1)
if artifact_rank > requested_rank:
message = (
f"Resume stage was advanced from {requested_resume_from_stage} to {artifact_resume_from_stage} "
f"because artifacts prove the run already progressed further."
)
else:
message = (
f"Resume stage was moved back from {requested_resume_from_stage} to {artifact_resume_from_stage} "
f"because later-stage artifacts are not stable enough for a safe resume."
)
return {
"decision": "resume",
"decision_source": decision_source,
"requested_resume_from_stage": requested_resume_from_stage,
"effective_resume_from_stage": artifact_resume_from_stage,
"artifact_resume_from_stage": artifact_resume_from_stage,
"recommended_action": "resume",
"message": message,
"missing_artifacts": [],
"artifact_snapshot": artifact_snapshot,
}
def _collect_resume_artifact_snapshot(*, record: dict[str, Any], state: RunState) -> dict[str, Any]:
run_dir = record["run_dir"]
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"))
summary_batch_output = _find_summary_batch_path(run_dir=run_dir, state=state)
candidate_batch_output = _find_candidate_batch_path(run_dir=run_dir, state=state)
delivery_output = _find_artifact_path(
run_dir=run_dir,
state=state,
artifact_name="delivery_payload",
relative_path=Path("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"),
)
report_output = run_dir / "run-report.json"
items: list[Any] = []
raw_output_valid = False
if raw_output is not None:
try:
items = _load_items(raw_output)
raw_output_valid = True
except Exception:
items = []
raw_output_valid = False
summary_batch_valid, _ = _load_summary_batch_lookup(summary_batch_output)
candidate_batch_valid, _ = _load_candidate_batch_lookup(candidate_batch_output)
item_contexts = _load_item_contexts(
items=items,
run_dir=run_dir,
debug_artifacts=bool(state.input.get("debug_artifacts", False)),
include_summaries=True,
include_candidates=True,
)
extraction_success_count = sum(
1 for context in item_contexts if context.get("extraction") is not None and context["extraction"].success
)
summary_count = sum(1 for context in item_contexts if context.get("summary_payload") is not None)
candidate_count = sum(1 for context in item_contexts if context.get("candidate") is not None)
extracted_complete = raw_output_valid and bool(items) and all(context["extracted_path"].exists() for context in item_contexts)
stable_summary_outputs = summary_batch_valid or (extraction_success_count > 0 and extraction_success_count == summary_count)
stable_candidate_outputs = candidate_batch_valid or (summary_count > 0 and summary_count == candidate_count)
missing_for_next_resume: list[str] = []
if raw_output is None or not raw_output_valid:
missing_for_next_resume.append("raw/freshrss.raw.json")
if extracted_dir is None or not extracted_complete:
missing_for_next_resume.append("extracted/")
if extraction_success_count > 0 and not stable_summary_outputs and not stable_candidate_outputs:
missing_for_next_resume.append("summary/summary-batch.json")
if extraction_success_count > 0 and stable_summary_outputs and not stable_candidate_outputs:
missing_for_next_resume.append("candidates/candidate-batch.json")
return {
"has_raw_output": raw_output is not None,
"raw_output_valid": raw_output_valid,
"has_extracted_dir": extracted_dir is not None,
"has_summary_batch": summary_batch_output is not None,
"summary_batch_valid": summary_batch_valid,
"has_candidate_batch": candidate_batch_output is not None,
"candidate_batch_valid": candidate_batch_valid,
"has_delivery_payload": delivery_output is not None,
"has_digest_brief": digest_brief_output is not None,
"has_run_report": report_output.exists(),
"item_count": len(items),
"extraction_success_count": extraction_success_count,
"summary_count": summary_count,
"candidate_count": candidate_count,
"extracted_complete": extracted_complete,
"stable_summary_outputs": stable_summary_outputs,
"stable_candidate_outputs": stable_candidate_outputs,
"missing_for_next_resume": missing_for_next_resume,
}
def _resolve_resume_stage_from_artifacts(snapshot: dict[str, Any]) -> str | None:
if snapshot["has_run_report"]:
return None
if snapshot["has_delivery_payload"]:
return REPORT_STAGE
if snapshot["stable_candidate_outputs"]:
return DELIVERY_STAGE
if snapshot["has_raw_output"] and snapshot["raw_output_valid"] and snapshot["has_extracted_dir"] and snapshot["extracted_complete"]:
if snapshot["extraction_success_count"] <= 0:
return SUMMARY_STAGE
if snapshot["stable_summary_outputs"]:
return FILTER_STAGE
return SUMMARY_STAGE
return None
def _load_filter_context(config: dict[str, Any]) -> FilterContext:
if config["context"] is not None:
return FilterContext.model_validate(config["context"])
@@ -695,6 +1090,8 @@ def _validate_resume_artifacts(*, record: dict[str, Any], state: RunState, resum
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"))
summary_batch_output = _find_summary_batch_path(run_dir=run_dir, state=state)
candidate_batch_output = _find_candidate_batch_path(run_dir=run_dir, state=state)
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:
@@ -711,6 +1108,8 @@ def _validate_resume_artifacts(*, record: dict[str, Any], state: RunState, resum
include_summaries=resume_from_stage in {FILTER_STAGE, DELIVERY_STAGE},
include_candidates=resume_from_stage == DELIVERY_STAGE,
)
summary_batch_valid, _ = _load_summary_batch_lookup(summary_batch_output)
candidate_batch_valid, _ = _load_candidate_batch_lookup(candidate_batch_output)
if resume_from_stage == SUMMARY_STAGE:
for context in item_contexts:
@@ -718,17 +1117,26 @@ def _validate_resume_artifacts(*, record: dict[str, Any], state: RunState, resum
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"))
)
expected_summary_count = int(_stage_output(state, SUMMARY_STAGE, "success_count") or 0)
actual_summary_count = sum(1 for context in item_contexts if context.get("summary_payload") is not None)
if summary_batch_valid:
if expected_summary_count != actual_summary_count:
missing_artifacts.append(_normalize_repo_path(summary_batch_output or _summary_batch_output(run_dir)))
else:
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:
if candidate_batch_valid:
if expected_candidate_count != actual_candidate_count:
missing_artifacts.append(_normalize_repo_path(candidate_batch_output or _candidate_batch_output(run_dir)))
elif 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(
+28
View File
@@ -18,10 +18,14 @@ from summary_mcp.models.summary_io import ExtractionInput
from summary_mcp.runtime import get_delivery_payload as load_delivery_payload
from summary_mcp.runtime import get_freshrss_pipeline_job_result as load_freshrss_pipeline_job_result
from summary_mcp.runtime import get_freshrss_pipeline_job_status as load_freshrss_pipeline_job_status
from summary_mcp.runtime import get_resume_job_result as load_resume_job_result
from summary_mcp.runtime import get_resume_job_status as load_resume_job_status
from summary_mcp.runtime import get_run_status as load_run_status
from summary_mcp.runtime import get_run_report as load_run_report
from summary_mcp.runtime import inspect_resume_plan as load_resume_plan
from summary_mcp.runtime import list_run_artifacts as load_run_artifacts
from summary_mcp.runtime import list_runs as load_runs
from summary_mcp.runtime import start_resume_job as launch_resume_job
from summary_mcp.runtime import start_freshrss_pipeline_job as launch_freshrss_pipeline_job
from summary_mcp.runtime.article_summary_jobs import (
get_article_summary_job_result as load_article_summary_job_result,
@@ -232,6 +236,30 @@ def resume_run(run_id: str) -> dict:
return resume_existing_run(run_id=run_id)
@mcp.tool()
def inspect_resume_plan(run_id: str) -> dict:
"""Inspect the effective resume plan for a FreshRSS workflow run without executing it."""
return load_resume_plan(run_id=run_id)
@mcp.tool()
def start_resume_job(run_id: str) -> dict:
"""Start an asynchronous resume job for a resumable FreshRSS workflow run."""
return launch_resume_job(run_id=run_id)
@mcp.tool()
def get_resume_job_status(job_id: str) -> dict:
"""Get the current status of an asynchronous resume job."""
return load_resume_job_status(job_id=job_id)
@mcp.tool()
def get_resume_job_result(job_id: str) -> dict:
"""Get the final result of an asynchronous resume job."""
return load_resume_job_result(job_id=job_id)
@mcp.tool()
def start_article_summary_job(
*,
@@ -51,6 +51,10 @@ SUMMARY_STAGE = "generate_summaries"
FILTER_STAGE = "apply_filters"
DELIVERY_STAGE = "build_delivery_payload"
REPORT_STAGE = "write_run_report"
SUMMARY_BATCH_ARTIFACT = "summary_batch"
CANDIDATE_BATCH_ARTIFACT = "candidate_batch"
SUMMARY_BATCH_FILENAME = "summary-batch.json"
CANDIDATE_BATCH_FILENAME = "candidate-batch.json"
def _save_json(path: Path, payload: dict[str, Any] | list[Any]) -> None:
@@ -137,6 +141,77 @@ def _build_item_context(*, index: int, item: Any, resolved_output_dir: Path, deb
}
def _summary_batch_output(run_dir: Path) -> Path:
return run_dir / "summary" / SUMMARY_BATCH_FILENAME
def _candidate_batch_output(run_dir: Path) -> Path:
return run_dir / "candidates" / CANDIDATE_BATCH_FILENAME
def _build_summary_batch_payload(*, run_id: str, item_contexts: list[dict[str, Any]]) -> dict[str, Any]:
items: list[dict[str, Any]] = []
for item_context in item_contexts:
summary_payload = item_context.get("summary_payload")
if summary_payload is None:
continue
item = item_context["item"]
items.append(
{
"item_key": item_context["item_key"],
"item_id": item.item_id,
"summary": summary_payload,
}
)
return {
"run_id": run_id,
"summary_count": len(items),
"items": items,
}
def _build_candidate_batch_payload(*, run_id: str, item_contexts: list[dict[str, Any]]) -> dict[str, Any]:
items: list[dict[str, Any]] = []
for item_context in item_contexts:
candidate = item_context.get("candidate")
if candidate is None:
continue
item = item_context["item"]
items.append(
{
"item_key": item_context["item_key"],
"item_id": item.item_id,
"candidate_id": candidate.candidate_id,
"candidate": candidate.model_dump(mode="json"),
}
)
return {
"run_id": run_id,
"candidate_count": len(items),
"items": items,
}
def _persist_summary_batch_artifact(*, run_store: RunStore, run_dir: Path, item_contexts: list[dict[str, Any]]) -> Path:
output_path = _summary_batch_output(run_dir)
_save_json(
output_path,
_build_summary_batch_payload(run_id=run_store.state.run_id, item_contexts=item_contexts),
)
run_store.register_artifact(name=SUMMARY_BATCH_ARTIFACT, path=output_path, kind="json", stage=SUMMARY_STAGE)
return output_path
def _persist_candidate_batch_artifact(*, run_store: RunStore, run_dir: Path, item_contexts: list[dict[str, Any]]) -> Path:
output_path = _candidate_batch_output(run_dir)
_save_json(
output_path,
_build_candidate_batch_payload(run_id=run_store.state.run_id, item_contexts=item_contexts),
)
run_store.register_artifact(name=CANDIDATE_BATCH_ARTIFACT, path=output_path, kind="json", stage=FILTER_STAGE)
return output_path
def _build_run_report(
*,
resolved_run_id: str,
@@ -418,6 +493,11 @@ def run_freshrss_pipeline(
},
)
summary_batch_output = _persist_summary_batch_artifact(
run_store=run_store,
run_dir=resolved_output_dir,
item_contexts=item_contexts,
)
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(
@@ -427,6 +507,7 @@ def run_freshrss_pipeline(
"completed_items": summary_success_count + summary_failed_count,
"success_count": summary_success_count,
"failed_count": summary_failed_count,
"summary_batch_output": str(summary_batch_output),
},
)
@@ -510,6 +591,11 @@ def run_freshrss_pipeline(
},
)
candidate_batch_output = _persist_candidate_batch_artifact(
run_store=run_store,
run_dir=resolved_output_dir,
item_contexts=item_contexts,
)
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(
@@ -521,6 +607,7 @@ def run_freshrss_pipeline(
"keep_count": keep_count,
"review_count": review_count,
"drop_count": drop_count,
"candidate_batch_output": str(candidate_batch_output),
},
)