← Files SavvyARCHIVED FILE
skills/savvy/scripts/savant_api/executions.py
16 KB · Oct 2, 2026 · 00:28 UTC
from __future__ import annotations
import time
import urllib.parse
from datetime import datetime, timezone
from typing import Any
from .httpclient import request
from .recipes import recipe_nodes, recipe_parameters
from .models import SavantAppApiError, SavantSessionContext
TERMINAL_NODE_STATUSES = {"Ready", "Failed", "Error", "Canceled", "Skipped"}
FAILED_NODE_STATUSES = {"Failed", "Error", "Canceled"}
EXECUTION_TYPE_ALIASES = {
"run": "run_now",
"runs": "run_now",
"run_now": "run_now",
"on_demand": "run_now",
"test": "test_run",
"tests": "test_run",
"test_run": "test_run",
"scheduled": "scheduled",
"schedule": "scheduled",
}
DEFAULT_EXECUTION_TYPES = ["run_now", "scheduled", "test_run"]
def _normalize_execution_types(values: list[str] | None) -> list[str]:
if not values:
return list(DEFAULT_EXECUTION_TYPES)
normalized: list[str] = []
for value in values:
for part in str(value).split(","):
key = part.strip().lower().replace("-", "_")
if not key:
continue
execution_type = EXECUTION_TYPE_ALIASES.get(key)
if not execution_type:
allowed = ", ".join(sorted(EXECUTION_TYPE_ALIASES))
raise SavantAppApiError(f"Unsupported execution type `{part}`. Use one of: {allowed}.")
if execution_type not in normalized:
normalized.append(execution_type)
return normalized or list(DEFAULT_EXECUTION_TYPES)
def _timestamp_to_iso(value: Any) -> str | None:
if value in (None, ""):
return None
try:
numeric = float(value)
except (TypeError, ValueError):
return str(value)
if numeric <= 0:
return None
if numeric > 10_000_000_000:
numeric = numeric / 1000
return datetime.fromtimestamp(numeric, tz=timezone.utc).isoformat().replace("+00:00", "Z")
def _format_duration_ms(started_at: Any, finished_at: Any) -> str | None:
try:
started = float(started_at)
finished = float(finished_at)
except (TypeError, ValueError):
return None
if started <= 0 or finished <= 0 or finished < started:
return None
seconds = int(round((finished - started) / 1000))
minutes, seconds = divmod(seconds, 60)
hours, minutes = divmod(minutes, 60)
if hours:
return f"{hours} h {minutes} m {seconds} s"
if minutes:
return f"{minutes} m {seconds} s"
return f"{seconds} s"
def normalize_execution(execution: dict[str, Any]) -> dict[str, Any]:
execution_type = execution.get("type")
type_label = {
"run_now": "run",
"scheduled": "run",
"test_run": "test",
}.get(str(execution_type), str(execution_type) if execution_type else None)
started_at = execution.get("startedAt")
finished_at = execution.get("finishedAt")
progress = execution.get("progress")
if isinstance(progress, (int, float)):
progress_value: float | str | None = progress
progress_label = f"{round(progress * 100)}%"
else:
progress_value = progress
progress_label = str(progress) if progress not in (None, "") else None
return {
"type": type_label,
"executionType": execution_type,
"id": execution.get("id"),
"name": execution.get("name"),
"workflowId": execution.get("recipeId") or execution.get("flowId"),
"workflowName": execution.get("recipeName"),
"version": execution.get("recipeVersion"),
"submitter": execution.get("submitterName") or execution.get("submitter"),
"submitterEmail": execution.get("submitter") if "@" in str(execution.get("submitter") or "") else None,
"status": execution.get("phase"),
"phase": execution.get("phase"),
"progress": progress_value,
"progressLabel": progress_label,
"startedAt": _timestamp_to_iso(started_at),
"finishedAt": _timestamp_to_iso(finished_at),
"startedAtEpochMs": started_at,
"finishedAtEpochMs": finished_at,
"duration": execution.get("duration") or _format_duration_ms(started_at, finished_at),
}
def list_recipe_executions(
context: SavantSessionContext,
flow_id: str,
*,
execution_types: list[str] | None = None,
) -> list[dict[str, Any]]:
types = _normalize_execution_types(execution_types)
query = urllib.parse.urlencode({"types": ",".join(types)})
response = request(context, f"/api/recipes/{urllib.parse.quote(flow_id)}/executions?{query}")
executions = response.get("executions") if isinstance(response, dict) else None
if not isinstance(executions, list):
raise SavantAppApiError(f"GET /api/recipes/{flow_id}/executions did not return an execution array.")
return [normalize_execution(item) for item in executions if isinstance(item, dict)]
def get_execution(context: SavantSessionContext, execution_id: str) -> dict[str, Any]:
response = request(context, f"/api/executions/{urllib.parse.quote(execution_id)}")
execution = response.get("execution") if isinstance(response, dict) else None
if not isinstance(execution, dict):
raise SavantAppApiError(f"GET /api/executions/{execution_id} did not return an execution object.")
normalized = normalize_execution(execution)
return {"execution": normalized, "raw": execution}
def _normalize_sample_tier(value: Any) -> str:
"""The backend sampleTier is BINARY: only the literal ``max`` triggers the full
debug-session (Analyze) path; every other value is the 1k interactive path. Collapse
anything that is not ``max`` to ``1k`` so a stray ``10k``/``2k`` does not masquerade as a
larger tier that does not exist."""
return "max" if isinstance(value, str) and value.strip().lower() == "max" else "1k"
def trigger_analysis(
context: SavantSessionContext,
flow_id: str,
nodes: list[dict[str, Any]],
parameters: list[dict[str, Any]],
stopping_nodes: list[str],
*,
sample_tier: str = "1k",
starting_nodes: list[str] | None = None,
action: str | None = None,
) -> Any:
if not stopping_nodes:
raise SavantAppApiError("At least one stopping node is required for analysis.")
body: dict[str, Any] = {
"flowId": flow_id,
"nodes": nodes,
"parameters": parameters,
"stoppingNodes": stopping_nodes,
}
if starting_nodes:
body["startingNodes"] = starting_nodes
# A starting node means "recompute from this edited node forward". APPLY evicts its
# cache first so the recompute is fresh; it is the only backend-honored action (the
# old ANALYZE value mapped to a no-op OTHER).
if action is None:
action = "APPLY"
tier = _normalize_sample_tier(sample_tier)
path = f"/api/interactive/graph-computation?sampleTier={urllib.parse.quote(tier)}"
if action:
path += f"&action={urllib.parse.quote(action)}"
return request(context, path, method="POST", body=body)
def graph_status(context: SavantSessionContext, flow_id: str) -> dict[str, Any]:
status = request(context, f"/api/interactive/graph-computation-status?flowId={urllib.parse.quote(flow_id)}")
if not isinstance(status, dict):
raise SavantAppApiError(f"Graph status for {flow_id} did not return a JSON object.")
return status
def node_status_map(status: dict[str, Any]) -> dict[str, dict[str, Any]]:
details = status.get("nodeStatusDetails")
if not isinstance(details, list):
return {}
return {item.get("nodeId"): item for item in details if isinstance(item, dict) and isinstance(item.get("nodeId"), str)}
def poll_analysis_status(
context: SavantSessionContext,
flow_id: str,
stopping_nodes: list[str],
*,
timeout_seconds: int = 90,
interval_seconds: float = 1.0,
fail_fast_on_error: bool = False,
) -> dict[str, dict[str, Any]]:
deadline = time.time() + timeout_seconds
latest: dict[str, dict[str, Any]] = {}
while time.time() < deadline:
latest = node_status_map(graph_status(context, flow_id))
if fail_fast_on_error and any((latest.get(node_id) or {}).get("status") in FAILED_NODE_STATUSES for node_id in stopping_nodes):
return latest
if stopping_nodes and all((latest.get(node_id) or {}).get("status") in TERMINAL_NODE_STATUSES for node_id in stopping_nodes):
return latest
time.sleep(interval_seconds)
raise SavantAppApiError(f"Timed out waiting for analysis of node(s): {', '.join(stopping_nodes)}")
def fetch_node_output(
context: SavantSessionContext,
flow_id: str,
node_id: str,
*,
sample_tier: str = "1k",
sorts: list[Any] | None = None,
) -> dict[str, Any]:
path = (
f"/api/interactive/graph-computation-sort"
f"?flowId={urllib.parse.quote(flow_id)}"
f"&nodeId={urllib.parse.quote(node_id)}"
f"&sampleTier={urllib.parse.quote(_normalize_sample_tier(sample_tier))}"
)
output = request(context, path, method="POST", body={"sorts": sorts or []})
if not isinstance(output, dict):
raise SavantAppApiError(f"Output for node {node_id} did not return a JSON object.")
return output
def summarize_node_output(output: dict[str, Any], *, row_limit: int = 5) -> dict[str, Any]:
blocks = output.get("outputData")
if not isinstance(blocks, list) or not blocks:
return {"schema": [], "rows": [], "rowCount": 0}
block = blocks[0]
if not isinstance(block, dict):
return {"schema": [], "rows": [], "rowCount": 0}
schema = block.get("schema") if isinstance(block.get("schema"), list) else []
data = block.get("data") if isinstance(block.get("data"), list) else []
names = [column.get("name") for column in schema if isinstance(column, dict)]
rows = [dict(zip(names, row)) for row in data[:row_limit] if isinstance(row, list)]
return {
"schema": schema,
"rows": rows,
"rowCount": block.get("totalRows") or block.get("rowCount") or len(data),
}
def analyze_and_fetch_node(
context: SavantSessionContext,
recipe: dict[str, Any],
flow_id: str,
node_id: str,
*,
sample_tier: str = "1k",
starting_nodes: list[str] | None = None,
timeout_seconds: int = 90,
reuse_ready: bool = True,
fetch_outputs: bool = True,
) -> dict[str, Any]:
"""Fetch one node's output, computing it first only when needed.
With ``reuse_ready`` (default), if the node is already in a terminal ``Ready`` state in the
live graph status, its cached output is read directly — the same data the canvas already shows —
instead of forcing a fresh Analyze and blocking on the poll. Set ``reuse_ready=False`` to force
a recompute (e.g. after an edit that the caller knows invalidated the cache but the status has
not caught up)."""
return analyze_and_fetch_many(
context,
recipe,
flow_id,
[node_id],
sample_tier=sample_tier,
starting_nodes=starting_nodes,
timeout_seconds=timeout_seconds,
reuse_ready=reuse_ready,
fetch_outputs=fetch_outputs,
)[node_id]
def _with_upstream_ids(nodes: list[dict[str, Any]], node_ids: list[str]) -> list[str]:
"""Return node_ids plus every upstream ancestor, deduplicated, sources first.
Used by forced (cache-busting) Analyze so the engine recomputes the full chain
feeding each requested node instead of reusing stale intermediate caches.
Outlet pseudo-nodes (split-blend/filter forks like ``blend_x|1``) are traversed
but NEVER included in the returned ids: the engine's computation-status map does
not report them, so triggering/polling them waits out the full timeout even when
every real node finished — verified live. Their parent real node is what computes.
"""
by_id: dict[str, dict[str, Any]] = {
node["id"]: node for node in nodes if isinstance(node, dict) and isinstance(node.get("id"), str)
}
sources_by_id: dict[str, list[str]] = {}
for node_id, node in by_id.items():
upstream: list[str] = []
for inlet in node.get("inlets") or []:
if not isinstance(inlet, dict):
continue
if isinstance(inlet.get("source"), str) and inlet["source"]:
upstream.append(inlet["source"])
for src in inlet.get("sources") or []:
if isinstance(src, dict) and isinstance(src.get("source"), str) and src["source"]:
upstream.append(src["source"])
sources_by_id[node_id] = upstream
def _is_pseudo(nid: str) -> bool:
node = by_id.get(nid) or {}
return node.get("type") == "outlet" or "|" in nid
ordered: list[str] = []
seen: set[str] = set()
def _visit(nid: str) -> None:
if nid in seen or nid not in sources_by_id:
return
seen.add(nid)
for upstream_id in sources_by_id[nid]:
_visit(upstream_id)
if not _is_pseudo(nid):
ordered.append(nid)
for nid in node_ids:
_visit(nid)
# Callers asked for these ids; keep any pseudo-node the caller explicitly
# requested at the end so its result row still appears in the report.
for nid in node_ids:
if nid not in ordered and _is_pseudo(nid):
ordered.append(nid)
return ordered
def analyze_and_fetch_many(
context: SavantSessionContext,
recipe: dict[str, Any],
flow_id: str,
node_ids: list[str],
*,
sample_tier: str = "1k",
starting_nodes: list[str] | None = None,
timeout_seconds: int = 90,
reuse_ready: bool = True,
fail_fast_on_error: bool = True,
fetch_outputs: bool = True,
) -> dict[str, dict[str, Any]]:
"""Fetch several nodes' outputs in one status-aware pass.
Reading the live ``graph-computation-status`` once, nodes already ``Ready`` are served from
their cached output (no recompute). Only the remaining "cold" nodes are recomputed, and they are
triggered together in a SINGLE Analyze with one shared poll for all of them — rather than a
sequential trigger+poll per node — so verifying N terminals costs one Analyze cycle, not N.
When ``fetch_outputs`` is false, this is a status-only runtime smoke test: it triggers/polls
the requested nodes but never reads preview output. Use that before schema/row inspection so a
failed transform reports immediately instead of burning time on preview endpoints that cannot
succeed.
Returns ``{node_id: {"nodeId", "status", "output"?}}`` for every requested node."""
nodes = recipe_nodes(recipe)
parameters = recipe_parameters(recipe)
valid_ids = {node.get("id") for node in nodes if isinstance(node, dict)}
missing = [nid for nid in node_ids if nid not in valid_ids]
if missing:
raise SavantAppApiError(
f"Workflow {flow_id} does not contain node(s): {', '.join(missing)}."
)
statuses = node_status_map(graph_status(context, flow_id))
if reuse_ready:
cold = [nid for nid in node_ids if (statuses.get(nid) or {}).get("status") != "Ready"]
else:
# Forced refresh: bust the whole upstream chain, not just the requested nodes, so the
# engine cannot serve a requested node from a stale intermediate cache.
cold = _with_upstream_ids(nodes, node_ids)
if cold:
trigger_analysis(
context, flow_id, nodes, parameters, cold,
sample_tier=sample_tier, starting_nodes=starting_nodes,
)
# poll_analysis_status returns the FULL live status map (all nodes), so this also refreshes
# the statuses we use below for the already-warm nodes.
statuses = poll_analysis_status(
context,
flow_id,
cold,
timeout_seconds=timeout_seconds,
fail_fast_on_error=fail_fast_on_error,
)
results: dict[str, dict[str, Any]] = {}
for nid in node_ids:
node_status = statuses.get(nid) or {}
result: dict[str, Any] = {"nodeId": nid, "status": node_status}
if fetch_outputs and node_status.get("status") == "Ready":
raw_output = fetch_node_output(context, flow_id, nid, sample_tier=sample_tier)
result["output"] = summarize_node_output(raw_output)
results[nid] = result
return results
SHA-256: 3dc03cfe9b8185991d76940f21efb4b74ca2691e0dd7135c0f47660ee792d2a3