feat: add async job entrypoint for freshrss pipeline

This commit is contained in:
root
2026-04-11 22:12:11 +08:00
parent 10f9cd088f
commit c58d8114cd
8 changed files with 841 additions and 29 deletions
+8
View File
@@ -1,3 +1,8 @@
from .freshrss_pipeline_jobs import (
get_freshrss_pipeline_job_result,
get_freshrss_pipeline_job_status,
start_freshrss_pipeline_job,
)
from .query_service import get_delivery_payload, get_run_report, get_run_status, list_run_artifacts, list_runs
from .run_store import RunStore
from .state_models import ArtifactRecord, RecoveryState, RunError, RunState, StageState
@@ -5,6 +10,8 @@ from .state_models import ArtifactRecord, RecoveryState, RunError, RunState, Sta
__all__ = [
"ArtifactRecord",
"get_delivery_payload",
"get_freshrss_pipeline_job_result",
"get_freshrss_pipeline_job_status",
"get_run_report",
"RecoveryState",
"RunError",
@@ -14,4 +21,5 @@ __all__ = [
"get_run_status",
"list_run_artifacts",
"list_runs",
"start_freshrss_pipeline_job",
]
@@ -0,0 +1,432 @@
from __future__ import annotations
import json
import os
import subprocess
import sys
from datetime import UTC, datetime
from pathlib import Path
from typing import Any
from uuid import uuid4
from .run_store import RunStore
REPO_ROOT = Path(__file__).resolve().parents[3]
OUTPUT_ROOT = REPO_ROOT / "outputs" / "freshrss"
PIPELINE_JOBS_ROOT = OUTPUT_ROOT / "pipeline_jobs"
WORKFLOW_NAME = "freshrss_pipeline_job"
RUN_TYPE = "daily_digest"
RUN_STATE_FILENAME = "run-state.json"
INPUT_FILENAME = "input.json"
RESULT_FILENAME = "result.json"
JOB_REPORT_FILENAME = "job-report.json"
DEFAULT_STAGES = [
"prepare_job",
"load_input",
"run_pipeline",
"write_result",
]
def _now() -> datetime:
return datetime.now().astimezone()
def _new_job_id() -> str:
ts = _now().strftime("%Y%m%d-%H%M%S")
return f"freshrss-pipeline-job-{ts}-{uuid4().hex[:8]}"
def _job_dir(job_id: str) -> Path:
return PIPELINE_JOBS_ROOT / job_id
def _run_state_path(job_id: str) -> Path:
return _job_dir(job_id) / RUN_STATE_FILENAME
def _input_path(job_id: str) -> Path:
return _job_dir(job_id) / INPUT_FILENAME
def _result_path(job_id: str) -> Path:
return _job_dir(job_id) / RESULT_FILENAME
def _job_report_path(job_id: str) -> Path:
return _job_dir(job_id) / JOB_REPORT_FILENAME
def _normalize_repo_path(path: Path) -> str:
try:
return str(path.resolve().relative_to(REPO_ROOT.resolve()))
except ValueError:
return str(path)
def _normalize_repo_path_value(path_value: str | None) -> str | None:
if not path_value:
return None
return _normalize_repo_path(Path(path_value))
def _normalize_keyword_index(keyword_index: dict[str, Any] | None) -> dict[str, Any]:
if not isinstance(keyword_index, dict):
return {}
normalized = dict(keyword_index)
for field in ("daily_output", "stats_output"):
normalized[field] = _normalize_repo_path_value(normalized.get(field))
return normalized
def _load_run_store(job_id: str) -> RunStore:
return RunStore.load(path=_run_state_path(job_id), repo_root=REPO_ROOT)
def _load_result(job_id: str) -> dict[str, Any] | None:
path = _result_path(job_id)
if not path.exists():
return None
return json.loads(path.read_text(encoding="utf-8-sig"))
def _build_run_defaults(*, run_id: str | None, output_dir: Path | None) -> tuple[str, Path]:
resolved_output_dir = output_dir
resolved_run_id = run_id
if resolved_run_id is not None and resolved_output_dir is not None:
return resolved_run_id, resolved_output_dir
run_stamp = datetime.now(tz=UTC).strftime("%Y%m%d-%H%M%S")
if resolved_run_id is None:
resolved_run_id = f"freshrss-pipeline-{run_stamp}"
if resolved_output_dir is None:
resolved_output_dir = OUTPUT_ROOT / "rerun" / run_stamp
return resolved_run_id, resolved_output_dir
def _build_result_payload(*, job_id: str, pipeline_result: dict[str, Any], include_item_reports: bool) -> dict[str, Any]:
result = {
"job_id": job_id,
"run_id": pipeline_result["run_id"],
"output_dir": _normalize_repo_path_value(pipeline_result.get("output_dir")),
"raw_output": _normalize_repo_path_value(pipeline_result.get("raw_output")),
"delivery_output": _normalize_repo_path_value(pipeline_result.get("delivery_output")),
"digest_brief_output": _normalize_repo_path_value(pipeline_result.get("digest_brief_output")),
"report_output": _normalize_repo_path_value(pipeline_result.get("report_output")),
"keyword_index": _normalize_keyword_index(pipeline_result.get("keyword_index")),
"pulled_count": pipeline_result.get("pulled_count"),
"delivered_count": pipeline_result.get("delivered_count"),
"marked_read_count": pipeline_result.get("marked_read_count"),
"status_counts": pipeline_result.get("status_counts"),
"debug_artifacts": pipeline_result.get("debug_artifacts"),
"completed_at": _now().isoformat(),
}
if include_item_reports:
result["items"] = pipeline_result.get("items", [])
return result
def _write_job_report(
*,
job_id: str,
payload: dict[str, Any],
) -> Path:
report_file = _job_report_path(job_id)
report_file.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8")
return report_file
def start_freshrss_pipeline_job(
*,
limit: int = 5,
mark_read: bool = False,
include_read: bool = False,
debug_artifacts: bool = False,
continuation: str | None = None,
timeout_seconds: float = 60.0,
max_retries: int = 2,
stream_id: str = "user/-/state/com.google/reading-list",
api_base_url: str | None = None,
username: str | None = None,
api_password: str | None = None,
llm_api_key: str | None = None,
llm_model: str | None = None,
llm_api_url: str | None = None,
context: dict[str, Any] | None = None,
run_id: str | None = None,
date_value: str | None = None,
output_dir: Path | None = None,
include_item_reports: bool = False,
) -> dict[str, Any]:
job_id = _new_job_id()
job_dir = _job_dir(job_id)
job_dir.mkdir(parents=True, exist_ok=True)
started_at = _now()
resolved_run_id, resolved_output_dir = _build_run_defaults(run_id=run_id, output_dir=output_dir)
input_payload = {
"limit": limit,
"mark_read": mark_read,
"include_read": include_read,
"debug_artifacts": debug_artifacts,
"continuation": continuation,
"timeout_seconds": timeout_seconds,
"max_retries": max_retries,
"stream_id": stream_id,
"api_base_url": api_base_url,
"username": username,
"api_password": api_password,
"llm_api_key": llm_api_key,
"llm_model": llm_model,
"llm_api_url": llm_api_url,
"context": context,
"run_id": resolved_run_id,
"date_value": date_value,
"output_dir": str(resolved_output_dir),
"include_item_reports": include_item_reports,
"launcher_pid": os.getpid(),
}
store = RunStore.create(
path=_run_state_path(job_id),
run_id=job_id,
workflow=WORKFLOW_NAME,
run_type=RUN_TYPE,
started_at=started_at,
input_payload=input_payload,
repo_root=REPO_ROOT,
)
for stage_name in DEFAULT_STAGES:
store._get_or_create_stage(stage_name)
store.save()
store.start_stage("prepare_job")
input_file = _input_path(job_id)
input_file.write_text(json.dumps(input_payload, ensure_ascii=False, indent=2), encoding="utf-8")
store.register_artifact(name="job_input", path=input_file, kind="json", stage="prepare_job")
runner_script = REPO_ROOT / "scripts" / "run_freshrss_pipeline_job.py"
cmd = [sys.executable, str(runner_script), "--job-id", job_id]
try:
proc = subprocess.Popen(
cmd,
cwd=str(REPO_ROOT),
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
start_new_session=True,
)
except Exception as exc:
report_file = _write_job_report(
job_id=job_id,
payload={
"job_id": job_id,
"status": "failed",
"error_type": type(exc).__name__,
"error_message": str(exc),
"failed_stage": "prepare_job",
"linked_run_id": resolved_run_id,
"linked_output_dir": _normalize_repo_path(resolved_output_dir),
},
)
store.register_artifact(name="job_report", path=report_file, kind="json", stage="prepare_job")
store.fail_stage("prepare_job", error=exc)
raise
store.finish_stage(
"prepare_job",
outputs={
"runner_pid": proc.pid,
"runner_command": cmd,
"linked_run_id": resolved_run_id,
"linked_output_dir": _normalize_repo_path(resolved_output_dir),
},
)
return {
"job_id": job_id,
"workflow": WORKFLOW_NAME,
"run_type": RUN_TYPE,
"status": "running",
"output_dir": _normalize_repo_path(job_dir),
"linked_run_id": resolved_run_id,
"linked_output_dir": _normalize_repo_path(resolved_output_dir),
"message": "FreshRSS pipeline job started successfully. Use get_freshrss_pipeline_job_status to poll progress.",
}
def run_freshrss_pipeline_job(*, job_id: str) -> dict[str, Any]:
from datetime import date
from summary_mcp.workflows import run_freshrss_pipeline
store = _load_run_store(job_id)
input_payload = json.loads(_input_path(job_id).read_text(encoding="utf-8-sig"))
current_stage = "run_pipeline"
try:
store.start_stage("load_input")
resolved_output_dir = Path(input_payload["output_dir"])
resolved_run_id = str(input_payload["run_id"])
store.finish_stage(
"load_input",
outputs={
"linked_run_id": resolved_run_id,
"linked_output_dir": _normalize_repo_path(resolved_output_dir),
"limit": int(input_payload.get("limit") or 5),
},
)
store.start_stage("run_pipeline")
pipeline_result = run_freshrss_pipeline(
api_base_url=input_payload.get("api_base_url"),
username=input_payload.get("username"),
api_password=input_payload.get("api_password"),
stream_id=input_payload.get("stream_id") or "user/-/state/com.google/reading-list",
limit=int(input_payload.get("limit") or 5),
continuation=input_payload.get("continuation"),
include_read=bool(input_payload.get("include_read")),
mark_read=bool(input_payload.get("mark_read")),
debug_artifacts=bool(input_payload.get("debug_artifacts")),
context=input_payload.get("context"),
max_retries=int(input_payload.get("max_retries") or 2),
timeout_seconds=float(input_payload.get("timeout_seconds") or 60.0),
llm_api_key=input_payload.get("llm_api_key"),
llm_model=input_payload.get("llm_model"),
llm_api_url=input_payload.get("llm_api_url"),
run_id=resolved_run_id,
delivery_date=date.fromisoformat(input_payload["date_value"]) if input_payload.get("date_value") else None,
output_dir=resolved_output_dir,
)
store.finish_stage(
"run_pipeline",
outputs={
"linked_run_id": pipeline_result.get("run_id"),
"linked_output_dir": _normalize_repo_path_value(pipeline_result.get("output_dir")),
"delivery_output": _normalize_repo_path_value(pipeline_result.get("delivery_output")),
"report_output": _normalize_repo_path_value(pipeline_result.get("report_output")),
"delivered_count": pipeline_result.get("delivered_count"),
"marked_read_count": pipeline_result.get("marked_read_count"),
},
)
store.start_stage("write_result")
result = _build_result_payload(
job_id=job_id,
pipeline_result=pipeline_result,
include_item_reports=bool(input_payload.get("include_item_reports")),
)
result_file = _result_path(job_id)
result_file.write_text(json.dumps(result, ensure_ascii=False, indent=2), encoding="utf-8")
store.register_artifact(name="job_result", path=result_file, kind="json", stage="write_result")
report_file = _write_job_report(
job_id=job_id,
payload={
"job_id": job_id,
"status": "success",
"run_id": result["run_id"],
"delivery_output": result["delivery_output"],
"report_output": result["report_output"],
"digest_brief_output": result["digest_brief_output"],
"pulled_count": result["pulled_count"],
"delivered_count": result["delivered_count"],
"marked_read_count": result["marked_read_count"],
},
)
store.register_artifact(name="job_report", path=report_file, kind="json", stage="write_result")
store.finish_stage(
"write_result",
outputs={
"result_path": _normalize_repo_path(result_file),
"linked_run_id": result["run_id"],
},
)
store.finish_run(status="success")
return result
except Exception as exc:
current_stage = store.state.current_stage or current_stage
report_file = _write_job_report(
job_id=job_id,
payload={
"job_id": job_id,
"status": "failed",
"error_type": type(exc).__name__,
"error_message": str(exc),
"failed_stage": current_stage,
"linked_run_id": input_payload.get("run_id"),
"linked_output_dir": _normalize_repo_path_value(input_payload.get("output_dir")),
},
)
try:
store.register_artifact(name="job_report", path=report_file, kind="json", stage=current_stage)
except Exception:
pass
store.fail_stage(current_stage, error=exc)
raise
def get_freshrss_pipeline_job_status(*, job_id: str) -> dict[str, Any]:
store = _load_run_store(job_id)
state = store.state
completed_stage_count = sum(1 for s in state.stages if s.status == "success")
running_stage_count = sum(1 for s in state.stages if s.status == "running")
failed_stage_count = sum(1 for s in state.stages if s.status == "failed")
pending_stage_count = sum(1 for s in state.stages if s.status == "pending")
linked_run_id = state.input.get("run_id") if isinstance(state.input, dict) else None
linked_output_dir = _normalize_repo_path_value(state.input.get("output_dir")) if isinstance(state.input, dict) else None
return {
"job_id": state.run_id,
"workflow": state.workflow,
"run_type": state.run_type,
"status": state.status,
"current_stage": state.current_stage,
"started_at": state.started_at.isoformat(),
"updated_at": state.updated_at.isoformat(),
"finished_at": state.finished_at.isoformat() if state.finished_at else None,
"output_dir": _normalize_repo_path(_job_dir(job_id)),
"linked_run_id": linked_run_id,
"linked_output_dir": linked_output_dir,
"progress": {
"completed_stage_count": completed_stage_count,
"running_stage_count": running_stage_count,
"failed_stage_count": failed_stage_count,
"pending_stage_count": pending_stage_count,
"total_stage_count": len(state.stages),
},
"artifacts": [artifact.model_dump(mode="json") for artifact in state.artifacts],
"error_summary": state.error.model_dump(mode="json") if state.error else None,
}
def get_freshrss_pipeline_job_result(*, job_id: str) -> dict[str, Any]:
store = _load_run_store(job_id)
state = store.state
result = _load_result(job_id)
if state.status != "success" or result is None:
return {
"job_id": state.run_id,
"status": state.status,
"message": "FreshRSS pipeline job result is not ready.",
"linked_run_id": state.input.get("run_id") if isinstance(state.input, dict) else None,
"error_summary": state.error.model_dump(mode="json") if state.error else None,
}
artifact = next((a.model_dump(mode="json") for a in state.artifacts if a.name == "job_result"), None)
return {
"job_id": state.run_id,
"status": state.status,
"run_id": result.get("run_id"),
"output_dir": result.get("output_dir"),
"raw_output": result.get("raw_output"),
"delivery_output": result.get("delivery_output"),
"digest_brief_output": result.get("digest_brief_output"),
"report_output": result.get("report_output"),
"pulled_count": result.get("pulled_count"),
"delivered_count": result.get("delivered_count"),
"marked_read_count": result.get("marked_read_count"),
"status_counts": result.get("status_counts"),
"keyword_index": result.get("keyword_index"),
"artifact": artifact,
"result": result,
}