← Files SavvyARCHIVED FILE

skills/savvy/scripts/workflow/inspection.py

24.9 KB · Oct 4, 2026 · 12:27 UTC

↓ Download file

#!/usr/bin/env python3
"""Live post-import inspection harness, the runtime counterpart to `savant.py validate workflow`.

Given a created Savant flow URL (plus the imported workflow JSON and the expected output
columns), this runs the standard creator verification checks in ONE orchestrated pass and emits
a structured pass/fail report:

  1. persistence    — live node count + node-name set match the imported JSON.
  2. runtime-smoke  — after create/save succeeds, Analyze requested checkpoints status-only and
                      fail fast on the first Failed/Error/Canceled node.
  3. node-ok        — every inspected checkpoint is `Ready`; this is populated from runtime-smoke.
  4. output-contract — only after runtime-smoke passes, read output columns and compare to expected
                      columns, in order (catches drift, leaked columns, wrong rename).
  5. row-sanity     — checkpoints don't silently collapse to 0 rows (WARN, not a hard fail).

It reuses `savant_api` for session/recipe/Analyze, and `--checkpoint` lets the caller name the
deterministic stages to verify (pivot/filter/top-N/etc.) without waiting on slow AI terminals.

INTERNAL ONLY: it uses the authenticated Savant API (Analyze), so it runs only when the live API
is available. It does not remove Analyze latency (each node is a live round-trip; gen_ai nodes
are slow). Layout quality is not inspected here; `savant.py validate workflow` covers the
deterministic layout geometry on the JSON.
"""
from __future__ import annotations

import argparse
from collections import Counter
import json
import sys
from pathlib import Path

from savant_api import cli as api
from savant_api.fileio import workspace_tmp
from savant_api.recipe_input import assert_flow_id, load_recipe
from savant_api import runmode

ERROR, WARN = "error", "warn"
FAILED_STATUSES = {"Failed", "Error", "Canceled"}
READY_STATUS = "Ready"


def _schema_names(schema) -> list:
    names = []
    for col in schema or []:
        names.append(col.get("name") if isinstance(col, dict) else col)
    return names


def _by_name(nodes, name):
    return next((n for n in nodes if n.get("name") == name), None)


def _expected_output_node_name(output: dict) -> str:
    value = output.get("verification_node_name")
    if isinstance(value, str) and value.strip():
        return value.strip()
    value = output.get("output_name")
    return value.strip() if isinstance(value, str) else ""


def _expected_output_columns(output: dict) -> list[str]:
    values = output.get("expected_columns")
    if not isinstance(values, list):
        return []
    return [value.strip() for value in values if isinstance(value, str) and value.strip()]


def load_expected_outputs(path: str | Path | None) -> list[dict] | None:
    if not path:
        return None
    payload = json.loads(Path(path).read_text(encoding="utf-8"))
    outputs = None
    if isinstance(payload, dict):
        if isinstance(payload.get("outputs"), list):
            outputs = payload.get("outputs")
        elif isinstance(payload.get("output_destination_plan"), dict):
            outputs = payload["output_destination_plan"].get("outputs")
        else:
            sections = payload.get("sections")
            if isinstance(sections, dict):
                builder = sections.get("builder_preflight")
                if isinstance(builder, dict) and isinstance(builder.get("output_destination_plan"), dict):
                    outputs = builder["output_destination_plan"].get("outputs")
    else:
        outputs = payload
    if not isinstance(outputs, list):
        raise ValueError(
            "expected outputs JSON must be a list, an object with `outputs`, "
            "a builder_preflight object, or a handoff containing sections.builder_preflight.output_destination_plan.outputs."
        )
    result = []
    for index, output in enumerate(outputs):
        if not isinstance(output, dict):
            raise ValueError(f"expected outputs entry {index} must be an object.")
        result.append(output)
    return result


def _terminals(nodes: list[dict]) -> list[dict]:
    """Data nodes with no outgoing edge (destinations / leaves); skip canvas group/text/outlet."""
    has_downstream = set()
    for n in nodes:
        for outlet in (n.get("outlets") or []):
            if outlet.get("targets"):
                has_downstream.add(n["id"])
                break
    return [n for n in nodes
            if n.get("type") not in ("group", "text", "outlet") and n["id"] not in has_downstream]


def _node_by_id(nodes: list[dict]) -> dict[str, dict]:
    return {node["id"]: node for node in nodes if isinstance(node.get("id"), str)}


def _upstream_map(nodes: list[dict]) -> dict[str, list[str]]:
    upstream: dict[str, list[str]] = {}
    for node in nodes:
        node_id = node.get("id")
        if not isinstance(node_id, str):
            continue
        sources: list[str] = []
        for inlet in node.get("inlets") or []:
            if not isinstance(inlet, dict):
                continue
            source = inlet.get("source")
            if isinstance(source, str) and source:
                sources.append(source)
            for item in inlet.get("sources") or []:
                if isinstance(item, dict) and isinstance(item.get("source"), str):
                    sources.append(item["source"])
        upstream[node_id] = sources
    return upstream


def print_summary_lines(report: dict, *, indent: str = "  ") -> None:
    """Shared CLI summary used by create/edit/inspect — one loop, no variations.

    Prints output smoke, per-check marks, and — critically — the
    `previewSkipped.upstreamBlockers` root-cause lines. A `node-ok: Skipped` line
    alone is not actionable; the blocker line names the Failed node and its engine
    error so no follow-up status-probe call is ever needed.
    """
    output_smoke = report.get("outputSmoke") or []
    if output_smoke:
        print(f"{indent}Outputs:")
        for row in output_smoke:
            columns = row.get("columns") or []
            columns_text = ", ".join(columns) if columns else "n/a"
            print(
                f"{indent}  - {row.get('outputName')}: {row.get('status')} "
                f"({row.get('rowCount')} row(s)); columns: {columns_text}"
            )
    for c in report.get("checks") or []:
        mark = "ok " if c.get("ok") else ("FAIL" if c.get("severity") == ERROR else "warn")
        print(f"{indent}[{mark}] {c.get('check')}: {c.get('detail')}")
    blockers = (report.get("previewSkipped") or {}).get("upstreamBlockers") or {}
    for target_name, blocker in blockers.items():
        print(
            f"{indent}[FAIL] root-cause for {target_name}: "
            f"{blocker.get('name')} [{blocker.get('status')}] — {blocker.get('detail')}"
        )


def _nearest_upstream_non_ready(target_id: str, nodes: list[dict], statuses: dict[str, dict]) -> dict | None:
    """Return the actionable upstream blocker for a non-Ready target.

    Runtime smoke often reports a terminal as `Skipped` while the actual blocker is
    upstream. A `Skipped` node is a SYMPTOM — the engine skipped it because something
    above it failed — so the walk continues THROUGH Skipped nodes and prefers the
    first Failed/Error node, which carries the actionable errorMessage (verified
    live: three Skipped terminals traced to one Failed filter with the real engine
    error). Only when no Failed node exists does the nearest non-Ready node answer.
    """
    by_id = _node_by_id(nodes)
    upstream = _upstream_map(nodes)
    queue = list(upstream.get(target_id, []))
    seen = {target_id}
    nearest_non_ready: dict | None = None
    while queue:
        node_id = queue.pop(0)
        if node_id in seen:
            continue
        seen.add(node_id)
        node = by_id.get(node_id)
        status = statuses.get(node_id) or {}
        status_value = status.get("status")
        if node and node.get("type") not in {"group", "text", "outlet"} and status_value != READY_STATUS:
            entry = {
                "nodeId": node_id,
                "name": node.get("name"),
                "status": status_value,
                "detail": _status_message(status),
            }
            if status_value not in {"Skipped", None}:
                return entry  # the root cause, with its engine error message
            if nearest_non_ready is None:
                nearest_non_ready = entry
        queue.extend(upstream.get(node_id, []))
    return nearest_non_ready


def _status_message(status: dict) -> str:
    value = status.get("status")
    if value == "Ready":
        return "Ready"
    detail = status.get("errorMessage") or status.get("message") or status.get("reason") or ""
    return f"{value} — {detail}" if detail else str(value)


def _runtime_smoke_summary(targets: list[dict], results: dict[str, tuple[dict, dict]]) -> dict:
    failed_nodes = []
    not_ready_nodes = []
    for node in targets:
        node_result = (results.get(node["id"]) or ({}, {}))[1]
        status = node_result.get("status") or {}
        status_value = status.get("status")
        row = {
            "nodeId": node.get("id"),
            "name": node.get("name"),
            "status": status_value,
            "detail": _status_message(status),
        }
        if status_value in FAILED_STATUSES:
            failed_nodes.append(row)
        elif status_value != "Ready":
            not_ready_nodes.append(row)
    return {
        "status": "pass" if not failed_nodes and not not_ready_nodes else "fail",
        "checkedNodes": [{"nodeId": n.get("id"), "name": n.get("name")} for n in targets],
        "failedNodes": failed_nodes,
        "notReadyNodes": not_ready_nodes,
    }


def _output_smoke_rows(targets: list[dict], results: dict[str, tuple[dict, dict]]) -> list[dict]:
    rows = []
    for node in targets:
        if node.get("type") != "destination":
            continue
        result = (results.get(node["id"]) or ({}, {}))[1]
        status = result.get("status") or {}
        output = result.get("output") or {}
        rows.append({
            "outputName": node.get("name") or node.get("id"),
            "nodeId": node.get("id"),
            "status": status.get("status"),
            "rowCount": output.get("rowCount", output.get("numRows")),
            "columns": _schema_names(output.get("schema")),
        })
    return rows


def persistence_check(imported_nodes: list[dict], live_nodes: list[dict]) -> tuple[bool, str]:
    imported_names = [n.get("name") for n in imported_nodes]
    live_names = [n.get("name") for n in live_nodes]
    imported_count = len(imported_nodes)
    live_count = len(live_nodes)
    ok = imported_count == live_count and Counter(imported_names) == Counter(live_names)
    detail = f"source JSON nodes.length={imported_count}, live recipe nodes.length={live_count}"
    if not ok:
        imported_counter = Counter(imported_names)
        live_counter = Counter(live_names)
        imported_only = sorted((imported_counter - live_counter).elements())
        live_only = sorted((live_counter - imported_counter).elements())
        detail += f"; imported-only={imported_only}, live-only={live_only}"
    return ok, detail


def inspect(flow_url: str, *, imported_path: str | None, expect_columns: list[str] | None,
            checkpoints: list[str], sample_tier: str, timeout: int, skip_terminals: bool,
            recipe: dict | None = None, recipe_json: Path | None = None,
            expected_outputs: list[dict] | None = None,
            force_analyze: bool = False) -> dict:
    """Inspect a live flow against its recipe.

    The recipe is supplied, not fetched: pass `recipe` (an already-loaded dict, which is how
    `workflow/evidence.py` hands over the post-write recipe it was given) or `recipe_json` (a path
    to the MCP `fetch` result). Preview/status still use the API.
    """
    checks: list[tuple[str, bool, str, str]] = []  # (check, ok, detail, severity)

    su = api.parse_savant_url(flow_url)
    if su.kind != "flow" or not su.flow_id:
        raise SystemExit("workflow inspect requires a Savant flow URL (.../flow/{flowId}).")
    ctx = api.discover_session(su.namespace, origin=su.origin)
    if recipe is None:
        recipe = load_recipe(recipe_json, flag="--recipe-json")
    assert_flow_id(recipe, su.flow_id, flag="--recipe-json")
    live_nodes = [n for n in api.recipe_nodes(recipe) if isinstance(n, dict)]

    # 1. Persistence vs the imported JSON (import remaps ids but preserves names).
    if imported_path:
        imp = json.loads(Path(imported_path).read_text(encoding="utf-8"))
        imp_nodes = [n for n in (imp.get("nodes") or []) if isinstance(n, dict)]
        ok, detail = persistence_check(imp_nodes, live_nodes)
        checks.append(("persistence", ok, detail, ERROR))

    # 2. Pick checkpoints: terminals (unless skipped) plus any caller-named stages.
    targets: list[dict] = [] if skip_terminals else list(_terminals(live_nodes))
    target_ids = {t["id"] for t in targets}
    for name in checkpoints:
        n = _by_name(live_nodes, name)
        if n is None:
            checks.append((f"checkpoint:{name}", False, "named checkpoint not found in flow", ERROR))
        elif n["id"] not in target_ids:
            targets.append(n); target_ids.add(n["id"])
    for output in expected_outputs or []:
        name = _expected_output_node_name(output)
        if not name:
            checks.append(("output-contract", False, "expected output is missing output_name/verification_node_name", ERROR))
            continue
        n = _by_name(live_nodes, name)
        if n is None:
            checks.append((f"output-contract:{name}", False, "expected output node not found in flow", ERROR))
        elif n["id"] not in target_ids:
            targets.append(n); target_ids.add(n["id"])

    # 3. Runtime smoke first. This is deliberately status-only: after create/save succeeded, check
    #    whether requested checkpoints compute to Ready before touching preview/schema endpoints.
    #    That makes transform/config failures surface immediately instead of being hidden behind a
    #    slow or doomed output preview call. After an API recipe edit, force_analyze avoids stale
    #    Ready previews that have not yet been invalidated by the UI's Apply path.
    results: dict[str, tuple[dict, dict]] = {}
    try:
        fetched = api.analyze_and_fetch_many(
            ctx, recipe, su.flow_id, [n["id"] for n in targets],
            sample_tier=sample_tier, timeout_seconds=timeout, reuse_ready=not force_analyze,
            fetch_outputs=False,
        )
    except Exception as exc:  # noqa: BLE001 — a batch-level failure marks every checkpoint failed
        fetched = {n["id"]: {"status": {"status": "Error", "errorMessage": str(exc)}} for n in targets}
    for n in targets:
        r = fetched.get(n["id"]) or {"status": {"status": "Error", "errorMessage": "no result"}}
        results[n["id"]] = (n, r)
        st = r.get("status") or {}
        status = st.get("status")
        ok = status == "Ready"
        checks.append((f"node-ok:{n.get('name')}", ok,
                       _status_message(st), ERROR))

    runtime_smoke = _runtime_smoke_summary(targets, results)
    failed_errors = [c for c in checks if not c[1] and c[3] == ERROR]
    if failed_errors:
        try:
            status_map = api.node_status_map(api.graph_status(ctx, su.flow_id))
        except Exception:  # noqa: BLE001 - diagnostics should not mask the main failure.
            status_map = {
                node_id: (result.get("status") or {})
                for node_id, (_node, result) in results.items()
            }
        upstream_blockers = {}
        for target in targets:
            target_status = (results.get(target["id"]) or ({}, {}))[1].get("status") or {}
            if target_status.get("status") == READY_STATUS:
                continue
            blocker = _nearest_upstream_non_ready(target["id"], live_nodes, status_map)
            if blocker:
                upstream_blockers[target.get("name") or target["id"]] = blocker
        return {
            "flowId": su.flow_id,
            "overall": "fail",
            "runtimeSmoke": runtime_smoke,
            "outputSmoke": _output_smoke_rows(targets, results),
            "previewSkipped": {
                "reason": "runtime-smoke failed; output previews/contracts were not attempted",
                "upstreamBlockers": upstream_blockers,
            },
            "checks": [{"check": c, "ok": ok, "severity": sev, "detail": d} for c, ok, d, sev in checks],
            "rows": {},
        }

    # 4. Runtime smoke passed; now read output only for Ready checkpoints.
    for n in targets:
        raw_output = api.fetch_node_output(ctx, su.flow_id, n["id"], sample_tier=sample_tier)
        results[n["id"]][1]["output"] = api.summarize_node_output(raw_output)

    # 5. Output contract: terminal/destination schema == expected columns.
    if expected_outputs:
        for output in expected_outputs:
            output_name = str(output.get("output_name") or _expected_output_node_name(output) or "output")
            node_name = _expected_output_node_name(output)
            out_node = _by_name(targets, node_name)
            if out_node is None:
                continue
            output_result = results[out_node["id"]][1].get("output") or {}
            cols = _schema_names(output_result.get("schema"))
            expected = _expected_output_columns(output)
            required_checks = output.get("required_checks")
            checks_to_run = required_checks if isinstance(required_checks, list) and required_checks else ["columns_exact_order", "row_count_nonzero", "no_node_errors"]
            grain = output.get("expected_grain")
            grain_detail = f"; grain={grain}" if isinstance(grain, str) and grain.strip() else ""
            if "columns_exact_order" in checks_to_run:
                ok = cols == expected
                checks.append((f"output-contract:{output_name}:columns_exact_order", ok, f"expected {expected}; got {cols}{grain_detail}", ERROR))
            elif "columns_present" in checks_to_run:
                missing = [column for column in expected if column not in cols]
                checks.append((f"output-contract:{output_name}:columns_present", not missing, f"missing {missing}; got {cols}{grain_detail}", ERROR))
            if "row_count_nonzero" in checks_to_run:
                rc = output_result.get("rowCount", output_result.get("numRows"))
                checks.append((f"output-contract:{output_name}:row_count_nonzero", rc is None or rc > 0, f"rowCount={rc}{grain_detail}", ERROR))
            if "no_node_errors" in checks_to_run:
                st = results[out_node["id"]][1].get("status") or {}
                checks.append((f"output-contract:{output_name}:no_node_errors", st.get("status") == "Ready", _status_message(st), ERROR))
    elif expect_columns:
        dests = [n for n in targets if n.get("type") == "destination"]
        out_node = dests[0] if dests else (targets[0] if len(targets) == 1 else None)
        if out_node is None:
            checks.append(("output-contract", False,
                           "no single terminal/destination to check (name it with --checkpoint)", ERROR))
        else:
            cols = _schema_names((results[out_node["id"]][1].get("output") or {}).get("schema"))
            ok = cols == list(expect_columns)
            checks.append(("output-contract", ok, f"expected {list(expect_columns)}; got {cols}", ERROR))

    # 6. Row-count sanity: a Ready checkpoint with 0 rows is a likely silent collapse (WARN).
    for nid, (n, r) in results.items():
        out = r.get("output") or {}
        rc = out.get("rowCount", out.get("numRows"))
        if (r.get("status") or {}).get("status") == "Ready" and rc == 0:
            checks.append((f"row-sanity:{n.get('name')}", False, "0 rows — possible unintended collapse", WARN))

    failed_errors = [c for c in checks if not c[1] and c[3] == ERROR]
    report = {
        "flowId": su.flow_id,
        "overall": "pass" if not failed_errors else "fail",
        "runtimeSmoke": runtime_smoke,
        "outputSmoke": _output_smoke_rows(targets, results),
        "previewSamples": {
            n.get("name") or n.get("id"): (r.get("output") or {})
            for _, (n, r) in results.items()
        },
        "checks": [{"check": c, "ok": ok, "severity": sev, "detail": d} for c, ok, d, sev in checks],
        "rows": {n.get("name"): (r.get("output") or {}).get("rowCount", (r.get("output") or {}).get("numRows"))
                 for _, (n, r) in results.items()},
    }
    return report


def parse_args(argv=None):
    p = argparse.ArgumentParser(description=__doc__)
    p.add_argument("flow_url", help="Savant flow URL, e.g. https://app.savantlabs.io/en/app/flow/{id}?rns={ns}")
    p.add_argument("--recipe-json", type=Path, required=True,
                   help="The flow's live recipe, fetched with the MCP `fetch` tool on "
                        "savant://workflow/{flowId}.")
    p.add_argument("--imported-json", help="The workflow JSON that was imported (for the persistence check).")
    p.add_argument("--expect-columns", help="Comma-separated expected final output columns, in order.")
    p.add_argument("--expected-outputs-json", help="JSON list/object of output contracts, or a handoff containing builder_preflight.output_destination_plan.outputs.")
    p.add_argument("--checkpoint", action="append", default=[], help="Node NAME to verify (repeatable).")
    p.add_argument("--skip-terminals", action="store_true", help="Only verify named checkpoints (skip auto terminals).")
    # The verify pass always computes, so it defaults to `interactive` (1k). `--mode analyze`
    # validates a checkpoint on full data. `--force-analyze` survives as a hidden alias.
    runmode.add_run_mode_args(p, default="interactive", legacy_analyze_flag=None, legacy_force_analyze=True)
    p.add_argument("--timeout-seconds", type=int, default=120)
    p.add_argument("--post-save-recipe-path", type=Path,
                   help="Override where the re-fetched live workflow JSON is written.")
    p.add_argument("--no-post-save-recipe", action="store_true",
                   help="Skip writing the re-fetched live workflow JSON.")
    p.add_argument("--output-path", type=Path, help="Where to write the JSON report.")
    p.add_argument("--quiet", action="store_true")
    return p.parse_args(argv)


def default_output_path(flow_id: str) -> Path:
    return workspace_tmp("workflow-inspections", f"{flow_id}.inspect.json")


def main(argv=None) -> int:
    args = parse_args(argv)
    expect = [c.strip() for c in args.expect_columns.split(",")] if args.expect_columns else None
    expected_outputs = load_expected_outputs(args.expected_outputs_json)
    mode = runmode.resolve_mode(args, default="interactive")
    from workflow import evidence as workflow_evidence

    su = api.parse_savant_url(args.flow_url)
    if su.kind != "flow" or not su.flow_id:
        raise SystemExit("workflow inspect requires a Savant flow URL (.../flow/{flowId}).")
    ctx = api.discover_session(su.namespace, origin=su.origin)
    live_recipe = load_recipe(args.recipe_json, flag="--recipe-json")
    assert_flow_id(live_recipe, su.flow_id, flag="--recipe-json")
    live_recipe_path = (
        None if args.no_post_save_recipe
        else args.post_save_recipe_path or workflow_evidence.default_post_write_recipe_path(su.flow_id, "inspect")
    )
    evidence = workflow_evidence.collect_post_write_evidence(
        ctx=ctx,
        flow_url=args.flow_url,
        flow_id=su.flow_id,
        operation="inspect",
        recipe=live_recipe,
        imported_path=args.imported_json,
        expect_columns=expect,
        expected_outputs=expected_outputs,
        checkpoints=args.checkpoint,
        sample_tier=runmode.sample_tier(mode),
        timeout=args.timeout_seconds,
        terminal_preview=not args.skip_terminals,
        force_analyze=getattr(args, "_legacy_force_analyze", False),
        post_write_recipe_path=live_recipe_path,
    )
    evidence.pop("refetchedRecipe", None)
    report = evidence.get("inspect") or {"flowId": su.flow_id, "overall": "fail", "checks": []}
    report["savedRecipePath"] = evidence.get("savedRecipePath")
    report["validation"] = evidence.get("validation")
    report["evidence"] = {k: v for k, v in evidence.items() if k not in {"savedRecipePath", "postSaveRecipePath", "inspect", "validation"}}
    api.save_json(report, args.output_path or default_output_path(str(report.get("flowId") or "workflow")))
    if not args.quiet:
        print(f"INSPECT {report['flowId']}: {report['overall'].upper()}")
        if report.get("savedRecipePath"):
            print(f"  [ok ] live-workflow-json: {report['savedRecipePath']}")
        validation = report.get("validation") or {}
        if validation and not validation.get("skipped"):
            mark = "ok " if validation.get("ok") else "FAIL"
            print(
                f"  [{mark}] validation: "
                f"{validation.get('errorCount', 0)} error(s), {validation.get('warningCount', 0)} warning(s)"
            )
        print_summary_lines(report)
    return 0 if evidence.get("ok") else 1


if __name__ == "__main__":
    raise SystemExit(main())

SHA-256: 2512f522ecbf8016a69aebe62ae4bdb9ceedc8a3d341f4ceb8544de3c4e20c3a