feat: add MCP run status query tools
This commit is contained in:
@@ -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",
|
||||
]
|
||||
|
||||
@@ -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"))
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user