diff --git a/TODO.md b/TODO.md index 82400a8..e9d3668 100644 --- a/TODO.md +++ b/TODO.md @@ -60,7 +60,7 @@ ## 2. 后续任务队列 -### [TODO][P1] 增加 MCP 状态查询接口 `get_run_status` +### [DONE][P1] 增加 MCP 状态查询接口 `get_run_status` 目标: - 可通过 MCP 查询 run 状态 @@ -69,9 +69,23 @@ - 输入 `run_id` - 返回 status / current_stage / completed_stages / failed_stage / artifacts / recovery +完成情况: +- 已通过 MCP 暴露 `get_run_status` +- 优先读取 `run-state.json`;对无 `run-state.json` 的历史 run 兼容基于现有 run 目录与 `run-report.json` 推断状态 +- 返回补充了 `progress` / `output_dir` / `state_source`,便于 OpenClaw 稳定消费且不必手拼路径 + +改动文件: +- `src/summary_mcp/runtime/query_service.py` +- `src/summary_mcp/runtime/run_store.py` +- `src/summary_mcp/runtime/__init__.py` +- `src/summary_mcp/server.py` + +遗留风险: +- 历史 run 若缺少 `run-state.json`,其阶段状态只能基于现有目录与 `run-report.json` 做保守推断 + --- -### [TODO][P1] 增加 MCP 查询接口 `list_runs` +### [DONE][P1] 增加 MCP 查询接口 `list_runs` 目标: - 查看近期 runs @@ -79,13 +93,39 @@ 要求: - 支持按 workflow / status / latest_n 过滤 +完成情况: +- 已通过 MCP 暴露 `list_runs` +- 支持按 `workflow` / `status` / `latest_n` 过滤近期 runs +- 返回 `run_id`、`status`、`progress`、`recovery`、`output_dir` 与 `state_source` + +改动文件: +- `src/summary_mcp/runtime/query_service.py` +- `src/summary_mcp/runtime/__init__.py` +- `src/summary_mcp/server.py` + +遗留风险: +- 当前按文件系统扫描 `outputs/freshrss/rerun/` 聚合,规模继续增大时可能需要再评估缓存或索引,但本阶段先保持文件系统真相 + --- -### [TODO][P1] 增加 MCP 查询接口 `list_run_artifacts` +### [DONE][P1] 增加 MCP 查询接口 `list_run_artifacts` 目标: - 统一列出 run 下 artifact +完成情况: +- 已通过 MCP 暴露 `list_run_artifacts` +- 对新 run 优先返回 `run-state.json` 已注册 artifacts,并补充 run 目录扫描发现的标准产物 +- 对历史 run 直接基于现有 run 目录发现标准产物,保持兼容 + +改动文件: +- `src/summary_mcp/runtime/query_service.py` +- `src/summary_mcp/runtime/__init__.py` +- `src/summary_mcp/server.py` + +遗留风险: +- 当前仅补充扫描固定的一组标准产物;未注册且不在标准集合内的调试文件不会进入稳定 artifact 列表 + --- ### [TODO][P1] 增加 MCP 结果读取接口 `get_delivery_payload` diff --git a/src/summary_mcp/runtime/__init__.py b/src/summary_mcp/runtime/__init__.py index 6dd253f..f32c948 100644 --- a/src/summary_mcp/runtime/__init__.py +++ b/src/summary_mcp/runtime/__init__.py @@ -1,3 +1,4 @@ +from .query_service import get_run_status, list_run_artifacts, list_runs from .run_store import RunStore from .state_models import ArtifactRecord, RecoveryState, RunError, RunState, StageState @@ -8,4 +9,7 @@ __all__ = [ "RunState", "RunStore", "StageState", + "get_run_status", + "list_run_artifacts", + "list_runs", ] diff --git a/src/summary_mcp/runtime/query_service.py b/src/summary_mcp/runtime/query_service.py new file mode 100644 index 0000000..a544807 --- /dev/null +++ b/src/summary_mcp/runtime/query_service.py @@ -0,0 +1,422 @@ +from __future__ import annotations + +# Compatibility note: +# historical runs may have a timestamp directory name that differs from the recorded run_id, +# so status queries resolve both identifiers without changing the existing output layout. + +import json +from datetime import datetime +from pathlib import Path +from typing import Any + +from .run_store import RunStore +from .state_models import ArtifactRecord, RunError, RunState, StageState + +REPO_ROOT = Path(__file__).resolve().parents[3] +RUNS_ROOT = REPO_ROOT / "outputs" / "freshrss" / "rerun" +DEFAULT_WORKFLOW = "freshrss_daily_digest" +DEFAULT_RUN_TYPE = "daily_digest" +DEFAULT_STAGES = [ + "fetch_feed", + "extract_articles", + "generate_summaries", + "apply_filters", + "build_delivery_payload", + "write_run_report", +] +DISCOVERED_ARTIFACTS = [ + {"name": "run_state", "kind": "json", "stage": "runtime_state", "relative_path": Path("run-state.json")}, + {"name": "run_report", "kind": "json", "stage": "write_run_report", "relative_path": Path("run-report.json")}, + {"name": "raw_output", "kind": "json", "stage": "fetch_feed", "relative_path": Path("raw/freshrss.raw.json")}, + {"name": "items_output", "kind": "json", "stage": "fetch_feed", "relative_path": Path("items/freshrss.items.json")}, + {"name": "extracted_dir", "kind": "directory", "stage": "extract_articles", "relative_path": Path("extracted")}, + {"name": "summary_dir", "kind": "directory", "stage": "generate_summaries", "relative_path": Path("summary")}, + {"name": "candidate_dir", "kind": "directory", "stage": "apply_filters", "relative_path": Path("candidates")}, + { + "name": "delivery_payload", + "kind": "json", + "stage": "build_delivery_payload", + "relative_path": Path("candidates/openclaw-delivery-payload.json"), + }, + { + "name": "digest_brief", + "kind": "json", + "stage": "build_delivery_payload", + "relative_path": Path("candidates/digest-brief.json"), + }, +] + + +def get_run_status(*, run_id: str) -> dict[str, Any]: + record = _resolve_run_record(run_id) + return _build_status_response(record) + + +def list_runs( + *, + workflow: str | None = None, + status: str | None = None, + latest_n: int = 20, +) -> dict[str, Any]: + if latest_n <= 0: + raise ValueError("latest_n must be greater than 0.") + + records = [] + for run_dir in _iter_run_dirs(): + record = _build_run_record(run_dir) + if workflow is not None and record["workflow"] != workflow: + continue + if status is not None and record["status"] != status: + continue + records.append(record) + + records.sort(key=_record_sort_key, reverse=True) + selected_records = records[:latest_n] + return { + "runs": [_build_list_response(record) for record in selected_records], + "count": len(selected_records), + "filters": { + "workflow": workflow, + "status": status, + "latest_n": latest_n, + }, + } + + +def list_run_artifacts(*, run_id: str) -> dict[str, Any]: + record = _resolve_run_record(run_id) + artifacts = _collect_artifacts(record) + return { + "run_id": record["run_id"], + "run_dir": record["run_dir"].name, + "workflow": record["workflow"], + "status": record["status"], + "output_dir": _normalize_repo_path(record["run_dir"]), + "state_source": record["state_source"], + "artifact_count": len(artifacts), + "artifacts": artifacts, + } + + +def _resolve_run_record(run_id: str) -> dict[str, Any]: + for run_dir in _iter_run_dirs(): + record = _build_run_record(run_dir) + if run_id in record["aliases"]: + return record + raise FileNotFoundError(f"Run not found for run_id: {run_id}") + + +def _iter_run_dirs() -> list[Path]: + if not RUNS_ROOT.exists(): + return [] + return sorted((path for path in RUNS_ROOT.iterdir() if path.is_dir()), key=lambda path: path.name, reverse=True) + + +def _build_run_record(run_dir: Path) -> dict[str, Any]: + state_path = run_dir / "run-state.json" + report_path = run_dir / "run-report.json" + + if state_path.exists(): + run_store = RunStore.load(path=state_path, repo_root=REPO_ROOT) + state = run_store.state + run_id = state.run_id + return { + "run_id": 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, + "stages": [_stage_to_dict(stage) for stage in state.stages], + "error": _error_to_dict(state.error), + "recovery": state.recovery.model_dump(mode="json"), + "run_dir": run_dir, + "state_source": "run_state", + "aliases": {run_id, run_dir.name}, + "state": state, + "report": _load_json(report_path) if report_path.exists() else None, + } + + 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 + + +def _infer_run_record_from_directory(*, run_dir: Path, report: dict[str, Any] | None, run_id: str) -> dict[str, Any]: + discovered_artifacts = _discover_artifacts(run_dir) + completed_stages = [artifact["stage"] for artifact in discovered_artifacts if artifact["stage"] in DEFAULT_STAGES] + deduped_completed_stages = [] + for stage_name in DEFAULT_STAGES: + if stage_name in completed_stages: + deduped_completed_stages.append(stage_name) + + if report is not None: + status = _infer_status_from_report(report) + started_at = _maybe_iso(report.get("started_at")) or _parse_run_dir_timestamp(run_dir.name) + updated_at = _maybe_iso(report.get("completed_at")) or started_at + finished_at = _maybe_iso(report.get("completed_at")) or updated_at + stages = [ + _stage_dict(name=stage_name, status="success") + for stage_name in DEFAULT_STAGES + if stage_name in {"fetch_feed", "extract_articles", "generate_summaries", "apply_filters", "build_delivery_payload", "write_run_report"} + ] + error = None + else: + status = "failed" + started_at = _parse_run_dir_timestamp(run_dir.name) + updated_at = started_at + finished_at = updated_at if discovered_artifacts else None + stages = [_stage_dict(name=stage_name, status="success") for stage_name in deduped_completed_stages] + next_stage = _infer_failed_stage(deduped_completed_stages) + if next_stage is not None: + stages.append(_stage_dict(name=next_stage, status="failed")) + error = { + "type": "InferredRunState", + "message": "run-state.json is missing; status inferred from existing run directory.", + "stage": next_stage, + "details": {}, + } + else: + error = { + "type": "InferredRunState", + "message": "run-state.json is missing and no completed stage could be confirmed.", + "stage": None, + "details": {}, + } + + return { + "run_id": run_id, + "workflow": DEFAULT_WORKFLOW, + "run_type": DEFAULT_RUN_TYPE, + "status": status, + "current_stage": None, + "started_at": started_at, + "updated_at": updated_at, + "finished_at": finished_at, + "stages": stages, + "error": error, + "recovery": { + "resumable": False, + "resume_from_stage": None, + "last_success_stage": deduped_completed_stages[-1] if deduped_completed_stages else None, + }, + "run_dir": run_dir, + "state_source": "directory_inference", + "state": None, + "report": report, + } + + +def _build_status_response(record: dict[str, Any]) -> dict[str, Any]: + completed_stages = [stage["name"] for stage in record["stages"] if stage["status"] == "success"] + failed_stage = next((stage["name"] for stage in record["stages"] if stage["status"] == "failed"), None) + return { + "run_id": record["run_id"], + "run_dir": record["run_dir"].name, + "workflow": record["workflow"], + "run_type": record["run_type"], + "status": record["status"], + "current_stage": record["current_stage"], + "started_at": record["started_at"], + "updated_at": record["updated_at"], + "finished_at": record["finished_at"], + "output_dir": _normalize_repo_path(record["run_dir"]), + "progress": _build_progress(record["stages"]), + "completed_stages": completed_stages, + "failed_stage": failed_stage, + "error_summary": record["error"], + "artifacts": _collect_artifacts(record), + "recovery": record["recovery"], + "state_source": record["state_source"], + } + + +def _build_list_response(record: dict[str, Any]) -> dict[str, Any]: + progress = _build_progress(record["stages"]) + return { + "run_id": record["run_id"], + "run_dir": record["run_dir"].name, + "workflow": record["workflow"], + "run_type": record["run_type"], + "status": record["status"], + "current_stage": record["current_stage"], + "started_at": record["started_at"], + "updated_at": record["updated_at"], + "finished_at": record["finished_at"], + "output_dir": _normalize_repo_path(record["run_dir"]), + "progress": progress, + "recovery": record["recovery"], + "artifact_count": len(_collect_artifacts(record)), + "state_source": record["state_source"], + } + + +def _collect_artifacts(record: dict[str, Any]) -> list[dict[str, Any]]: + artifacts: list[dict[str, Any]] = [] + seen_names: set[str] = set() + + state = record.get("state") + if isinstance(state, RunState): + for artifact in state.artifacts: + artifact_dict = _artifact_from_state_record(artifact) + artifacts.append(artifact_dict) + seen_names.add(artifact_dict["name"]) + + for artifact in _discover_artifacts(record["run_dir"]): + if artifact["name"] in seen_names: + continue + artifacts.append(artifact) + seen_names.add(artifact["name"]) + + artifacts.sort(key=lambda artifact: (artifact["stage"], artifact["name"])) + return artifacts + + +def _discover_artifacts(run_dir: Path) -> list[dict[str, Any]]: + artifacts = [] + for spec in DISCOVERED_ARTIFACTS: + path = run_dir / spec["relative_path"] + if not path.exists(): + continue + artifacts.append( + { + "name": spec["name"], + "kind": spec["kind"], + "stage": spec["stage"], + "path": _normalize_repo_path(path), + "exists": True, + "created_at": _iso_from_stat(path), + "metadata": {}, + "source": "directory_scan", + } + ) + return artifacts + + +def _artifact_from_state_record(artifact: ArtifactRecord) -> dict[str, Any]: + return { + "name": artifact.name, + "kind": artifact.kind, + "stage": artifact.stage, + "path": _normalize_repo_path(artifact.path), + "exists": _artifact_exists(artifact.path), + "created_at": artifact.created_at.isoformat(), + "metadata": artifact.metadata, + "source": "run_state", + } + + +def _artifact_exists(path_value: str) -> bool: + path = _resolve_repo_path(path_value) + return path.exists() + + +def _build_progress(stages: list[dict[str, Any]]) -> dict[str, Any]: + expected_stage_names = list(DEFAULT_STAGES) + stage_status_by_name = {stage["name"]: stage["status"] for stage in stages} + for stage_name in stage_status_by_name: + if stage_name not in expected_stage_names: + expected_stage_names.append(stage_name) + + completed_count = sum(1 for stage_name in expected_stage_names if stage_status_by_name.get(stage_name) == "success") + running_count = sum(1 for stage_name in expected_stage_names if stage_status_by_name.get(stage_name) == "running") + failed_count = sum(1 for stage_name in expected_stage_names if stage_status_by_name.get(stage_name) == "failed") + pending_count = sum(1 for stage_name in expected_stage_names if stage_status_by_name.get(stage_name, "pending") == "pending") + return { + "completed_stage_count": completed_count, + "running_stage_count": running_count, + "failed_stage_count": failed_count, + "pending_stage_count": pending_count, + "total_stage_count": len(expected_stage_names), + } + + +def _infer_status_from_report(report: dict[str, Any]) -> str: + status_counts = report.get("status_counts") + if isinstance(status_counts, dict) and any(key in status_counts for key in {"extract_failed", "summary_failed"}): + return "partial" + return "success" + + +def _infer_failed_stage(completed_stages: list[str]) -> str | None: + for stage_name in DEFAULT_STAGES: + if stage_name not in completed_stages: + return stage_name + return None + + +def _stage_to_dict(stage: StageState) -> dict[str, Any]: + return { + "name": stage.name, + "status": stage.status, + "started_at": stage.started_at.isoformat() if stage.started_at else None, + "finished_at": stage.finished_at.isoformat() if stage.finished_at else None, + "outputs": stage.outputs, + "error": _error_to_dict(stage.error), + } + + +def _stage_dict(*, name: str, status: str) -> dict[str, Any]: + return { + "name": name, + "status": status, + "started_at": None, + "finished_at": None, + "outputs": {}, + "error": None, + } + + +def _error_to_dict(error: RunError | dict[str, Any] | None) -> dict[str, Any] | None: + if error is None: + return None + if isinstance(error, RunError): + return error.model_dump(mode="json") + return error + + +def _normalize_repo_path(path_value: str | Path) -> str: + path = _resolve_repo_path(path_value) + try: + return str(path.relative_to(REPO_ROOT)) + except ValueError: + return str(path) + + +def _resolve_repo_path(path_value: str | Path) -> Path: + path = Path(path_value) + if path.is_absolute(): + return path + return REPO_ROOT / path + + +def _iso_from_stat(path: Path) -> str: + return datetime.fromtimestamp(path.stat().st_mtime).astimezone().isoformat() + + +def _parse_run_dir_timestamp(run_dir_name: str) -> str | None: + try: + return datetime.strptime(run_dir_name, "%Y%m%d-%H%M%S").astimezone().isoformat() + except ValueError: + return None + + +def _maybe_iso(value: Any) -> str | None: + if isinstance(value, str): + try: + return datetime.fromisoformat(value).isoformat() + except ValueError: + return value + return None + + +def _record_sort_key(record: dict[str, Any]) -> tuple[str, str]: + return (record.get("updated_at") or "", record["run_dir"].name) + + +def _load_json(path: Path) -> dict[str, Any]: + return json.loads(path.read_text(encoding="utf-8-sig")) diff --git a/src/summary_mcp/runtime/run_store.py b/src/summary_mcp/runtime/run_store.py index 1a1d3c2..ad3cc26 100644 --- a/src/summary_mcp/runtime/run_store.py +++ b/src/summary_mcp/runtime/run_store.py @@ -39,6 +39,11 @@ class RunStore: ) return cls(path=path, state=state, repo_root=repo_root) + @classmethod + def load(cls, *, path: Path, repo_root: Path | None = None) -> "RunStore": + state = RunState.model_validate(json.loads(path.read_text(encoding="utf-8-sig"))) + return cls(path=path, state=state, repo_root=repo_root) + def save(self) -> None: self.state.updated_at = datetime.now(tz=self.state.started_at.tzinfo) self.path.parent.mkdir(parents=True, exist_ok=True) diff --git a/src/summary_mcp/server.py b/src/summary_mcp/server.py index 7982c38..6886c97 100644 --- a/src/summary_mcp/server.py +++ b/src/summary_mcp/server.py @@ -15,6 +15,9 @@ from summary_mcp.models.filtering import FilterContext, FilterInput from summary_mcp.models.item import Item from summary_mcp.models.llm_result import LlmSummaryResult from summary_mcp.models.summary_io import ExtractionInput +from summary_mcp.runtime import get_run_status as load_run_status +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.workflows import run_freshrss_pipeline from summary_mcp.workflows.article_summary import ArticleSummaryConfig, summarize_selected_articles @@ -119,6 +122,28 @@ def run_freshrss_openclaw_pipeline( return result +@mcp.tool() +def get_run_status(run_id: str) -> dict: + """Get the current status of a workflow run by run_id.""" + return load_run_status(run_id=run_id) + + +@mcp.tool() +def list_runs( + workflow: str | None = None, + status: str | None = None, + latest_n: int = 20, +) -> dict: + """List recent workflow runs with optional workflow/status filters.""" + return load_runs(workflow=workflow, status=status, latest_n=latest_n) + + +@mcp.tool() +def list_run_artifacts(run_id: str) -> dict: + """List registered and discovered artifacts for a workflow run.""" + return load_run_artifacts(run_id=run_id) + + @mcp.tool() def generate_article_summaries( *,