feat(reader): add async article summary jobs

This commit is contained in:
root
2026-04-10 16:37:43 +08:00
parent 87d18e4263
commit 2053cccfef
8 changed files with 1269 additions and 15 deletions
@@ -0,0 +1,315 @@
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 .run_store import RunStore
from .state_models import RunState
REPO_ROOT = Path(__file__).resolve().parents[3]
OUTPUT_ROOT = REPO_ROOT / "outputs" / "freshrss"
ARTICLE_SUMMARY_JOBS_ROOT = OUTPUT_ROOT / "article_summary_jobs"
WORKFLOW_NAME = "article_summary_job"
RUN_TYPE = "article_summary"
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",
"generate_markdown",
"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"article-summary-{ts}-{uuid4().hex[:8]}"
def _job_dir(job_id: str) -> Path:
return ARTICLE_SUMMARY_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 _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 start_article_summary_job(
*,
extracted_path: Path,
selected_ids: list[str],
output_dir: Path | None = None,
max_retries: int = 2,
timeout_seconds: float = 120.0,
llm_api_key: str | None = None,
llm_model: str | None = None,
llm_api_url: str | None = None,
) -> dict[str, Any]:
if not extracted_path.exists():
raise FileNotFoundError(f"extracted_path does not exist: {extracted_path}")
if not selected_ids:
raise ValueError("selected_ids must not be empty")
job_id = _new_job_id()
job_dir = _job_dir(job_id)
job_dir.mkdir(parents=True, exist_ok=True)
started_at = _now()
resolved_output_dir = output_dir or (extracted_path.parent / "single_summaries")
input_payload = {
"extracted_path": str(extracted_path),
"selected_ids": selected_ids,
"output_dir": str(resolved_output_dir),
"max_retries": max_retries,
"timeout_seconds": timeout_seconds,
"llm_api_key": llm_api_key,
"llm_model": llm_model,
"llm_api_url": llm_api_url,
"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_article_summary_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:
store.fail_stage("prepare_job", error=exc)
raise
store.finish_stage("prepare_job", outputs={"runner_pid": proc.pid, "runner_command": cmd})
return {
"job_id": job_id,
"workflow": WORKFLOW_NAME,
"run_type": RUN_TYPE,
"status": "running",
"output_dir": _normalize_repo_path(job_dir),
"message": "Article summary job started successfully. Use get_article_summary_job_status to poll progress.",
}
def run_article_summary_job(*, job_id: str) -> dict[str, Any]:
from summary_mcp.workflows.article_summary import ArticleSummaryConfig, summarize_selected_articles
store = _load_run_store(job_id)
input_payload = json.loads(_input_path(job_id).read_text(encoding="utf-8-sig"))
try:
store.start_stage("load_input")
extracted_path = Path(input_payload["extracted_path"])
output_dir = Path(input_payload["output_dir"])
selected_ids = list(input_payload["selected_ids"])
if not extracted_path.exists():
raise FileNotFoundError(f"extracted_path does not exist: {extracted_path}")
if not selected_ids:
raise ValueError("selected_ids must not be empty")
store.finish_stage(
"load_input",
outputs={
"extracted_path": _normalize_repo_path(extracted_path),
"selected_id_count": len(selected_ids),
"output_dir": _normalize_repo_path(output_dir),
},
)
store.start_stage("generate_markdown")
config = ArticleSummaryConfig(
max_retries=int(input_payload.get("max_retries") or 2),
timeout_seconds=float(input_payload.get("timeout_seconds") or 120.0),
)
written_paths = summarize_selected_articles(
extracted_path=extracted_path,
selected_ids=selected_ids,
output_dir=output_dir,
config=config,
api_key=input_payload.get("llm_api_key"),
model=input_payload.get("llm_model"),
api_url=input_payload.get("llm_api_url"),
)
if not written_paths:
raise RuntimeError("Article summary job produced no Markdown outputs.")
normalized_paths = [_normalize_repo_path(Path(p)) for p in written_paths]
store.finish_stage(
"generate_markdown",
outputs={
"written_count": len(written_paths),
"written_paths": normalized_paths,
},
)
for idx, p in enumerate(written_paths, start=1):
store.register_artifact(
name=f"summary_markdown_{idx}",
path=Path(p),
kind="markdown",
stage="generate_markdown",
metadata={"output_type": "single_summary"},
)
store.start_stage("write_result")
result = {
"job_id": job_id,
"selected_ids": selected_ids,
"written_paths": normalized_paths,
"completed_at": _now().isoformat(),
}
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 = _job_report_path(job_id)
report_file.write_text(
json.dumps(
{
"job_id": job_id,
"status": "success",
"written_count": len(normalized_paths),
"written_paths": normalized_paths,
},
ensure_ascii=False,
indent=2,
),
encoding="utf-8",
)
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)})
store.finish_run(status="success")
return result
except Exception as exc:
current_stage = store.state.current_stage or "generate_markdown"
store.fail_stage(current_stage, error=exc)
report_file = _job_report_path(job_id)
report_file.write_text(
json.dumps(
{
"job_id": job_id,
"status": "failed",
"error_type": type(exc).__name__,
"error_message": str(exc),
"failed_stage": current_stage,
},
ensure_ascii=False,
indent=2,
),
encoding="utf-8",
)
try:
store.register_artifact(name="job_report", path=report_file, kind="json", stage=current_stage)
except Exception:
pass
raise
def get_article_summary_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")
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)),
"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_article_summary_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": "Article summary job result is not ready.",
"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,
"written_paths": result.get("written_paths", []),
"artifact": artifact,
"result": result,
}
+41
View File
@@ -20,6 +20,11 @@ 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.article_summary_jobs import (
get_article_summary_job_result as load_article_summary_job_result,
get_article_summary_job_status as load_article_summary_job_status,
start_article_summary_job as launch_article_summary_job,
)
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
@@ -164,6 +169,42 @@ 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 start_article_summary_job(
*,
extracted_path: str,
selected_ids: list[str],
output_dir: str | None = None,
max_retries: int = 2,
timeout_seconds: float = 120.0,
llm_api_key: str | None = None,
llm_model: str | None = None,
llm_api_url: str | None = None,
) -> dict:
"""Start an asynchronous article-summary job and return a job_id immediately."""
return launch_article_summary_job(
extracted_path=Path(extracted_path),
selected_ids=selected_ids,
output_dir=Path(output_dir) if output_dir else None,
max_retries=max_retries,
timeout_seconds=timeout_seconds,
llm_api_key=llm_api_key,
llm_model=llm_model,
llm_api_url=llm_api_url,
)
@mcp.tool()
def get_article_summary_job_status(job_id: str) -> dict:
"""Get the current status of an asynchronous article-summary job."""
return load_article_summary_job_status(job_id=job_id)
@mcp.tool()
def get_article_summary_job_result(job_id: str) -> dict:
"""Get the final result of an asynchronous article-summary job."""
return load_article_summary_job_result(job_id=job_id)
@mcp.tool()
def generate_article_summaries(
@@ -167,10 +167,12 @@ def summarize_selected_articles(
output_dir.mkdir(parents=True, exist_ok=True)
written_paths: list[Path] = []
matched_ids: set[str] = set()
from summary_mcp.core.summary_loop import build_summary_input
for item_id, entry in _iter_selected_items(payload, selected_ids):
matched_ids.add(item_id)
# Prefer the real project format where each entry has ``item`` and
# ``extraction.article``; fall back to legacy layout where the
# article fields live directly on the element.
@@ -288,4 +290,7 @@ def summarize_selected_articles(
output_path.write_text("\n".join(lines), encoding="utf-8")
written_paths.append(output_path)
if not matched_ids:
raise ValueError(f"No extracted entries matched selected_ids: {list(selected_ids)}")
return written_paths