Add MCP delivery payload and run report queries
This commit is contained in:
@@ -1,9 +1,11 @@
|
||||
from .query_service import get_run_status, list_run_artifacts, list_runs
|
||||
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
|
||||
|
||||
__all__ = [
|
||||
"ArtifactRecord",
|
||||
"get_delivery_payload",
|
||||
"get_run_report",
|
||||
"RecoveryState",
|
||||
"RunError",
|
||||
"RunState",
|
||||
|
||||
@@ -98,6 +98,72 @@ def list_run_artifacts(*, run_id: str) -> dict[str, Any]:
|
||||
}
|
||||
|
||||
|
||||
def get_delivery_payload(*, run_id: str) -> dict[str, Any]:
|
||||
record = _resolve_run_record(run_id)
|
||||
artifact = _resolve_artifact(
|
||||
record,
|
||||
artifact_name="delivery_payload",
|
||||
relative_path=Path("candidates/openclaw-delivery-payload.json"),
|
||||
fallback_stage="build_delivery_payload",
|
||||
fallback_kind="json",
|
||||
report_field="delivery_output",
|
||||
)
|
||||
payload = _load_json(_resolve_repo_path(artifact["path"]))
|
||||
candidates = payload.get("candidates")
|
||||
candidate_count = len(candidates) if isinstance(candidates, list) else 0
|
||||
stats = payload.get("stats") if isinstance(payload.get("stats"), dict) else None
|
||||
response = _build_run_lookup_response(record)
|
||||
response.update(
|
||||
{
|
||||
"artifact": artifact,
|
||||
"payload_schema_version": payload.get("schema_version"),
|
||||
"generated_at": payload.get("generated_at"),
|
||||
"delivery_date": payload.get("date"),
|
||||
"candidate_count": candidate_count,
|
||||
"stats": stats,
|
||||
"payload": payload,
|
||||
}
|
||||
)
|
||||
return response
|
||||
|
||||
|
||||
def get_run_report(*, run_id: str) -> dict[str, Any]:
|
||||
record = _resolve_run_record(run_id)
|
||||
artifact = _resolve_artifact(
|
||||
record,
|
||||
artifact_name="run_report",
|
||||
relative_path=Path("run-report.json"),
|
||||
fallback_stage="write_run_report",
|
||||
fallback_kind="json",
|
||||
report_field="report_output",
|
||||
)
|
||||
report = _load_json(_resolve_repo_path(artifact["path"]))
|
||||
items = report.get("items")
|
||||
item_count = len(items) if isinstance(items, list) else 0
|
||||
response = _build_run_lookup_response(record)
|
||||
response.update(
|
||||
{
|
||||
"artifact": artifact,
|
||||
"started_at": report.get("started_at"),
|
||||
"completed_at": report.get("completed_at"),
|
||||
"requested_limit": report.get("requested_limit"),
|
||||
"pulled_count": report.get("pulled_count"),
|
||||
"delivered_count": report.get("delivered_count"),
|
||||
"marked_read_count": report.get("marked_read_count"),
|
||||
"mark_read_requested": report.get("mark_read_requested"),
|
||||
"debug_artifacts": report.get("debug_artifacts"),
|
||||
"raw_output": _normalize_optional_repo_path(report.get("raw_output")),
|
||||
"delivery_output": _normalize_optional_repo_path(report.get("delivery_output")),
|
||||
"digest_brief_output": _normalize_optional_repo_path(report.get("digest_brief_output")),
|
||||
"status_counts": report.get("status_counts"),
|
||||
"keyword_index": _normalize_keyword_index(report.get("keyword_index")),
|
||||
"item_count": item_count,
|
||||
"report": report,
|
||||
}
|
||||
)
|
||||
return response
|
||||
|
||||
|
||||
def _resolve_run_record(run_id: str) -> dict[str, Any]:
|
||||
for run_dir in _iter_run_dirs():
|
||||
record = _build_run_record(run_dir)
|
||||
@@ -235,6 +301,18 @@ def _build_status_response(record: dict[str, Any]) -> dict[str, Any]:
|
||||
}
|
||||
|
||||
|
||||
def _build_run_lookup_response(record: dict[str, Any]) -> dict[str, Any]:
|
||||
return {
|
||||
"run_id": record["run_id"],
|
||||
"run_dir": record["run_dir"].name,
|
||||
"workflow": record["workflow"],
|
||||
"run_type": record["run_type"],
|
||||
"status": record["status"],
|
||||
"output_dir": _normalize_repo_path(record["run_dir"]),
|
||||
"state_source": record["state_source"],
|
||||
}
|
||||
|
||||
|
||||
def _build_list_response(record: dict[str, Any]) -> dict[str, Any]:
|
||||
progress = _build_progress(record["stages"])
|
||||
return {
|
||||
@@ -276,6 +354,85 @@ def _collect_artifacts(record: dict[str, Any]) -> list[dict[str, Any]]:
|
||||
return artifacts
|
||||
|
||||
|
||||
def _resolve_artifact(
|
||||
record: dict[str, Any],
|
||||
*,
|
||||
artifact_name: str,
|
||||
relative_path: Path,
|
||||
fallback_stage: str,
|
||||
fallback_kind: str,
|
||||
report_field: str | None = None,
|
||||
) -> dict[str, Any]:
|
||||
state = record.get("state")
|
||||
if isinstance(state, RunState):
|
||||
for artifact in state.artifacts:
|
||||
if artifact.name != artifact_name:
|
||||
continue
|
||||
if _artifact_exists(artifact.path):
|
||||
return _artifact_from_state_record(artifact)
|
||||
|
||||
standard_path = record["run_dir"] / relative_path
|
||||
if standard_path.exists():
|
||||
return _build_resolved_artifact(
|
||||
name=artifact_name,
|
||||
path=standard_path,
|
||||
kind=fallback_kind,
|
||||
stage=fallback_stage,
|
||||
source="standard_path",
|
||||
)
|
||||
|
||||
report = record.get("report")
|
||||
if report_field and isinstance(report, dict):
|
||||
report_path_value = report.get(report_field)
|
||||
if isinstance(report_path_value, str) and report_path_value.strip():
|
||||
report_path = _resolve_repo_path(report_path_value)
|
||||
if report_path.exists():
|
||||
return _build_resolved_artifact(
|
||||
name=artifact_name,
|
||||
path=report_path,
|
||||
kind=fallback_kind,
|
||||
stage=fallback_stage,
|
||||
source="run_report_reference",
|
||||
)
|
||||
|
||||
for path in sorted(record["run_dir"].rglob(relative_path.name)):
|
||||
if path.is_file():
|
||||
return _build_resolved_artifact(
|
||||
name=artifact_name,
|
||||
path=path,
|
||||
kind=fallback_kind,
|
||||
stage=fallback_stage,
|
||||
source="directory_scan",
|
||||
)
|
||||
|
||||
available_artifacts = [artifact["name"] for artifact in _collect_artifacts(record)]
|
||||
available_summary = ", ".join(available_artifacts) if available_artifacts else "none"
|
||||
raise FileNotFoundError(
|
||||
f"Artifact '{artifact_name}' not found for run_id: {record['run_id']} "
|
||||
f"(status={record['status']}, run_dir={record['run_dir'].name}, available_artifacts={available_summary})"
|
||||
)
|
||||
|
||||
|
||||
def _build_resolved_artifact(
|
||||
*,
|
||||
name: str,
|
||||
path: Path,
|
||||
kind: str,
|
||||
stage: str,
|
||||
source: str,
|
||||
) -> dict[str, Any]:
|
||||
return {
|
||||
"name": name,
|
||||
"kind": kind,
|
||||
"stage": stage,
|
||||
"path": _normalize_repo_path(path),
|
||||
"exists": True,
|
||||
"created_at": _iso_from_stat(path),
|
||||
"metadata": {},
|
||||
"source": source,
|
||||
}
|
||||
|
||||
|
||||
def _discover_artifacts(run_dir: Path) -> list[dict[str, Any]]:
|
||||
artifacts = []
|
||||
for spec in DISCOVERED_ARTIFACTS:
|
||||
@@ -394,6 +551,22 @@ def _resolve_repo_path(path_value: str | Path) -> Path:
|
||||
return REPO_ROOT / path
|
||||
|
||||
|
||||
def _normalize_optional_repo_path(path_value: Any) -> str | None:
|
||||
if not isinstance(path_value, str) or not path_value.strip():
|
||||
return None
|
||||
return _normalize_repo_path(path_value)
|
||||
|
||||
|
||||
def _normalize_keyword_index(keyword_index: Any) -> dict[str, Any] | None:
|
||||
if not isinstance(keyword_index, dict):
|
||||
return None
|
||||
|
||||
normalized = dict(keyword_index)
|
||||
normalized["daily_output"] = _normalize_optional_repo_path(keyword_index.get("daily_output"))
|
||||
normalized["stats_output"] = _normalize_optional_repo_path(keyword_index.get("stats_output"))
|
||||
return normalized
|
||||
|
||||
|
||||
def _iso_from_stat(path: Path) -> str:
|
||||
return datetime.fromtimestamp(path.stat().st_mtime).astimezone().isoformat()
|
||||
|
||||
|
||||
Reference in New Issue
Block a user