← Files SavvyARCHIVED FILE

skills/savvy/scripts/workflow/health.py

15.1 KB · Oct 5, 2026 · 18:27 UTC

↓ Download file

#!/usr/bin/env python3
"""Read-only Savant workflow health check.

This helper gathers the first-pass evidence an inspector usually needs before
drilling into a workflow:

- recipe/status/version/run history
- source placeholder and dataset-match health
- optional one serial Analyze preview at the earliest useful checkpoint

It intentionally avoids parallel Analyze calls. The goal is to find the first
likely blocker quickly, then stop with a concise report.
"""

from __future__ import annotations

import argparse
import json
import sys
from pathlib import Path
from typing import Any


SCRIPT_DIR = Path(__file__).resolve().parents[1]
if str(SCRIPT_DIR) not in sys.path:
    sys.path.insert(0, str(SCRIPT_DIR))

from savant_api.cli import (  # noqa: E402
    SavantAppApiError,
    analyze_and_fetch_node,
    discover_session,
    ensure_rns,
    list_recipe_executions,
    parse_flow_url,
    recipe_nodes,
    save_json,
)
from savant_api.recipe_input import assert_flow_id, load_recipe, load_records  # noqa: E402
from workflow.discovery import build_discovery_report  # noqa: E402
from savant_api.fileio import workspace_tmp  # noqa: E402
from savant_api import runmode  # noqa: E402


SOURCE_TYPES = {"source"}
STRUCTURAL_TYPES = {"group", "text"}


def node_label(node: dict[str, Any]) -> str:
    for key in ("name", "label", "id"):
        value = node.get(key)
        if isinstance(value, str) and value.strip():
            return value.strip()
    return "unnamed step"


def source_config_name(node: dict[str, Any]) -> str | None:
    config = node.get("config") if isinstance(node.get("config"), dict) else {}
    for value in (node.get("name"), config.get("name"), config.get("id")):
        if isinstance(value, str) and value.strip():
            return value.strip()
    return None


def source_has_embedded_schema(node: dict[str, Any]) -> bool:
    config = node.get("config") if isinstance(node.get("config"), dict) else {}
    schema = config.get("schema")
    return isinstance(schema, list) and bool(schema)


def source_connector(node: dict[str, Any]) -> str | None:
    config = node.get("config") if isinstance(node.get("config"), dict) else {}
    value = config.get("connector") or config.get("type")
    return value if isinstance(value, str) and value.strip() else None


def source_nodes(recipe: dict[str, Any]) -> list[dict[str, Any]]:
    return [node for node in recipe_nodes(recipe) if isinstance(node, dict) and node.get("type") in SOURCE_TYPES]


def non_structural_nodes(recipe: dict[str, Any]) -> list[dict[str, Any]]:
    return [
        node
        for node in recipe_nodes(recipe)
        if isinstance(node, dict) and node.get("type") not in STRUCTURAL_TYPES
    ]


def map_source_matches(report: dict[str, Any]) -> dict[str, dict[str, Any]]:
    matches = report.get("sources")
    if not isinstance(matches, list):
        return {}
    return {
        str(match.get("node_id")): match
        for match in matches
        if isinstance(match, dict) and match.get("node_id")
    }


def best_candidate(match: dict[str, Any]) -> dict[str, Any] | None:
    search = match.get("dataset_search") if isinstance(match.get("dataset_search"), dict) else {}
    candidates = search.get("candidates")
    if isinstance(candidates, list) and candidates and isinstance(candidates[0], dict):
        return candidates[0]
    return None


def source_health(recipe: dict[str, Any], source_match_report: dict[str, Any]) -> list[dict[str, Any]]:
    matches = map_source_matches(source_match_report)
    rows: list[dict[str, Any]] = []
    for node in source_nodes(recipe):
        node_id = node.get("id")
        match = matches.get(str(node_id), {})
        status = match.get("match_status") or "unknown"
        candidate = best_candidate(match) if isinstance(match, dict) else None
        config_name = source_config_name(node)
        issues: list[str] = []
        if status in {"missing", "ambiguous", "unknown"}:
            issues.append(f"dataset match is {status}")
        if config_name and "." in Path(config_name).name:
            issues.append("source name looks like an imported file placeholder")
        if source_has_embedded_schema(node) and status != "matched":
            issues.append("recipe has embedded schema, but no unique dataset binding was confirmed")
        rows.append(
            {
                "nodeId": node_id,
                "name": node_label(node),
                "connector": source_connector(node),
                "matchStatus": status,
                "matchedDatasetId": match.get("matched_dataset_id") if isinstance(match, dict) else None,
                "matchedDatasetName": match.get("matched_dataset_name") if isinstance(match, dict) else None,
                "bestCandidate": candidate,
                "issues": issues,
            }
        )
    return rows


def choose_preview_node(recipe: dict[str, Any], sources: list[dict[str, Any]]) -> dict[str, str] | None:
    for source in sources:
        if source.get("issues") and isinstance(source.get("nodeId"), str):
            return {
                "nodeId": source["nodeId"],
                "name": str(source.get("name") or source["nodeId"]),
                "reason": "first source with unresolved binding/match issue",
            }
    source_list = source_nodes(recipe)
    if source_list:
        first = source_list[0]
        return {
            "nodeId": str(first.get("id")),
            "name": node_label(first),
            "reason": "first source checkpoint",
        }
    for node in non_structural_nodes(recipe):
        node_id = node.get("id")
        if isinstance(node_id, str):
            return {
                "nodeId": node_id,
                "name": node_label(node),
                "reason": "first non-structural checkpoint",
            }
    return None


def compact_preview_error(exc: Exception) -> dict[str, Any]:
    message = str(exc)
    status_code = None
    for code in ("429", "500", "404", "403", "401"):
        if f" {code} " in message or message.startswith(code) or f": {code} " in message:
            status_code = int(code)
            break
    return {
        "status": "failed",
        "error": message,
        "httpStatus": status_code,
    }


def build_recommendations(report: dict[str, Any]) -> list[str]:
    recommendations: list[str] = []
    problematic_sources = [
        source
        for source in report.get("sources", [])
        if isinstance(source, dict) and source.get("issues")
    ]
    if problematic_sources:
        bind_parts = []
        for source in problematic_sources:
            candidate = source.get("bestCandidate") if isinstance(source.get("bestCandidate"), dict) else None
            if candidate and candidate.get("name"):
                bind_parts.append(f"{source.get('name')} -> {candidate.get('name')}")
        if bind_parts:
            recommendations.append("Bind source datasets: " + "; ".join(bind_parts) + ".")
        else:
            recommendations.append("Resolve missing or ambiguous source datasets before inspecting downstream steps.")
    preview = report.get("preview")
    if isinstance(preview, dict) and preview.get("status") == "failed":
        recommendations.append("Do not debug downstream joins, summaries, or destinations until the first failing checkpoint previews successfully.")
    if not recommendations:
        recommendations.append("Sources did not show an obvious binding issue; inspect the first grain-changing step next.")
    return recommendations


def health_status(report: dict[str, Any]) -> str:
    if any(source.get("issues") for source in report.get("sources", []) if isinstance(source, dict)):
        return "source_attention_required"
    preview = report.get("preview")
    if isinstance(preview, dict) and preview.get("status") == "failed":
        return "preview_failed"
    return "no_first-pass_blocker_found"


def print_summary(report: dict[str, Any]) -> None:
    workflow = report.get("workflow", {})
    print(f"Workflow: {workflow.get('name') or workflow.get('id')}")
    print(f"Status: {workflow.get('status') or 'unknown'}")
    print(f"Version: {workflow.get('version') or 'unknown'}")
    executions = report.get("executionHistory", {}).get("executions") if isinstance(report.get("executionHistory"), dict) else []
    print(f"Run/Test history: {len(executions) if isinstance(executions, list) else 0} execution(s)")
    print(f"Health: {report.get('health')}")
    print("")
    print("Sources:")
    for source in report.get("sources", []):
        if not isinstance(source, dict):
            continue
        candidate = source.get("bestCandidate") if isinstance(source.get("bestCandidate"), dict) else None
        suffix = f"; best match: {candidate.get('name')}" if candidate and candidate.get("name") else ""
        issues = "; ".join(source.get("issues") or []) or "no first-pass issue"
        print(f"- {source.get('name')}: {source.get('matchStatus')}{suffix} ({issues})")
    preview = report.get("preview")
    if isinstance(preview, dict):
        print("")
        print(f"Preview checkpoint: {preview.get('name')} ({preview.get('reason')})")
        if preview.get("status") == "passed":
            output = preview.get("output") if isinstance(preview.get("output"), dict) else {}
            print(f"Preview result: passed, {output.get('rowCount', 0)} row(s)")
        elif preview.get("status") == "skipped":
            print(f"Preview result: skipped ({preview.get('reason')})")
        else:
            print(f"Preview result: failed ({preview.get('error')})")
    print("")
    print("Recommended next step:")
    for item in report.get("recommendations", []):
        print(f"- {item}")


def build_health_report(
    flow_url: str,
    *,
    workflow_json: Path,
    sources_json: Path | None = None,
    sample_tier: str = "1k",
    preview: bool = False,
    preview_timeout_seconds: int = 60,
) -> dict[str, Any]:
    """Health-check a live flow from its MCP-fetched recipe.

    The recipe and the dataset list are inputs now — `workflow_json` from the MCP `fetch` tool on
    savant://workflow/{flowId}, `sources_json` from MCP `search` with `types: ["source"]`. Run
    history and the optional preview still need the API, so this route is not API-free; but the
    recipe and source-matching legs are.

    Omitting `sources_json` skips source matching rather than reporting zero matches: an empty
    dataset list and an unsupplied one are different facts, and conflating them would report every
    bound dataset as unresolvable.
    """
    parsed = parse_flow_url(flow_url)
    context = discover_session(parsed.namespace, origin=parsed.origin)
    recipe = load_recipe(workflow_json)
    assert_flow_id(recipe, parsed.flow_id)
    datasets = load_records(sources_json, flag="--sources-json")
    source_match_report = (
        build_discovery_report(recipe, datasets, limit=10) if sources_json is not None else
        {"skipped": True,
         "reason": "no --sources-json was supplied, so dataset matching was not run. Pass the MCP "
                   "`search` result for types: [\"source\"] to enable it."}
    )
    sources = source_health(recipe, source_match_report)
    execution_types = ["run_now", "scheduled", "test_run"]
    report: dict[str, Any] = {
        "task": "savant_workflow_health_check",
        "flowUrl": ensure_rns(flow_url, context.namespace),
        "workflow": {
            "id": recipe.get("id") or parsed.flow_id,
            "name": recipe.get("name"),
            "status": recipe.get("status"),
            "version": recipe.get("version"),
            "nodeCount": len(recipe_nodes(recipe)),
            "namespace": context.namespace,
            "workspaceId": context.workspace_id,
        },
        "executionHistory": {
            "executionTypes": execution_types,
            "executions": list_recipe_executions(context, parsed.flow_id, execution_types=execution_types),
        },
        "sourceMatchReport": source_match_report,
        "sources": sources,
    }
    checkpoint = choose_preview_node(recipe, sources)
    if not preview:
        report["preview"] = {"status": "skipped", "reason": "preview disabled by caller"}
    elif checkpoint is None:
        report["preview"] = {"status": "skipped", "reason": "no preview checkpoint found"}
    else:
        preview_result: dict[str, Any] = {**checkpoint}
        try:
            inspected = analyze_and_fetch_node(
                context,
                recipe,
                parsed.flow_id,
                checkpoint["nodeId"],
                sample_tier=sample_tier,
                timeout_seconds=preview_timeout_seconds,
            )
            preview_result.update(inspected)
            preview_result["status"] = "passed" if inspected.get("status", {}).get("status") == "Ready" else "not_ready"
        except Exception as exc:  # keep health checks diagnostic, not crash-only
            preview_result.update(compact_preview_error(exc))
        report["preview"] = preview_result
    report["health"] = health_status(report)
    report["recommendations"] = build_recommendations(report)
    return report


def main(argv: list[str] | None = None) -> int:
    parser = argparse.ArgumentParser(description=__doc__)
    parser.add_argument("flow_url", help="Savant flow URL to inspect.")
    parser.add_argument("--workflow-json", type=Path, required=True,
                        help="The flow's recipe, fetched with the MCP `fetch` tool on "
                             "savant://workflow/{flowId}.")
    parser.add_argument("--sources-json", type=Path,
                        help="MCP `search` result for types: [\"source\"]. Omit to skip dataset "
                             "matching rather than report zero matches.")
    parser.add_argument("--output-path", type=Path, help="Where to write health-check JSON.")
    # Health defaults to no-compute: it reads recipe/status/version/run-history/source matches.
    # `--mode interactive|analyze` opts into the single earliest-checkpoint preview.
    runmode.add_run_mode_args(parser, default="cached", legacy_analyze_flag="--analyze-preview")
    parser.add_argument("--preview-timeout-seconds", type=int, default=60)
    parser.add_argument("--skip-preview", action="store_true", help="Compatibility flag; preview is skipped unless a compute mode is set.")
    parser.add_argument("--json-only", action="store_true", help="Do not print the text summary.")
    args = parser.parse_args(argv)

    mode = runmode.resolve_mode(args, default="cached")
    report = build_health_report(
        args.flow_url,
        workflow_json=args.workflow_json,
        sources_json=args.sources_json,
        sample_tier=runmode.sample_tier(mode),
        preview=runmode.should_compute(mode) and not args.skip_preview,
        preview_timeout_seconds=args.preview_timeout_seconds,
    )
    output = args.output_path or workspace_tmp("savant-health-checks", f"{report['workflow']['id']}.health.json")
    save_json(report, output)
    if args.json_only:
        print(output)
    else:
        print_summary(report)
        print("")
        print(f"Wrote health check JSON to {output}")
    return 0


if __name__ == "__main__":
    try:
        raise SystemExit(main())
    except SavantAppApiError as exc:
        print(f"savant_workflow_health_check: {exc}", file=sys.stderr)
        raise SystemExit(2)

SHA-256: 921ef35e872333a704f886834c0bf25285b3aacec84c1a250dfd86a665df4a44