← Files SavvyARCHIVED FILE
skills/savvy/scripts/workflow/create.py
20.1 KB · Oct 2, 2026 · 00:28 UTC
#!/usr/bin/env python3
"""Creator orchestrator — chains the deterministic create-and-verify loop in one call.
Two phases, with the user-confirmation gate between them:
PREFLIGHT (always): validate the workflow JSON, resolve the exact target
folder + session, verify the session namespace against the user-confirmed
destination workspace (--confirmed-namespace), and check for blockers
(validation errors, or an AI-backed node with a missing providerId). Emits a
summary. It does NOT import.
CREATE + VERIFY (only with --confirm-import, and only if preflight is clean): import the JSON
into the folder, then run the shared post-write evidence loop over the created flow: re-fetch
and save the live recipe JSON and inspect terminal/output previews. This is the same evidence
path used by workflow edit after it saves a recipe. Layout quality is judged from `validate
workflow` (which runs the deterministic layout checks) on the persisted JSON — no rendered-canvas inspection.
Creation is folder-bound: the imported workflow must read back with the same folder id the user
specified. A different or missing folder id is a hard failure, not a successful create.
Datasets are already bound in the JSON (each source node carries a resolved dataset id), so this
orchestrator does no dataset discovery, matching, or binding — Creator is dataset-free. Pass
`--confirm-import` only AFTER the user has confirmed creation in chat. INTERNAL ONLY: uses the
authenticated Savant API.
"""
from __future__ import annotations
import argparse
import json
import sys
from pathlib import Path
from savant_api import cli as api
from savant_api import capabilities as savant_capabilities
from savant_api.recipe_input import load_id_set
from contracts.folder_target import is_root_folder_target
from savant_api.fileio import workspace_tmp
from workflow import evidence as workflow_evidence
from validators import workflow as vw
_CREATE_STATUSES = {
"verified",
"created-with-issues",
"imported",
"import-unverified",
}
def preflight(json_path: str, *, folder_id: str, confirmed_namespace: str | None = None,
sources_json: Path | None = None, providers_json: Path | None = None):
wf = json.loads(Path(json_path).read_text(encoding="utf-8"))
nodes = [n for n in (wf.get("nodes") or []) if isinstance(n, dict)]
errors, warnings = vw.validate(Path(json_path), block_documentation_gaps=True)
missing_ai_providers = vw.ai_missing_provider_nodes(nodes)
placeholder_sources = vw.source_placeholder_nodes(nodes)
# A write is coming, so prove the session rather than assuming it: a stale token or an
# unreachable host must fail here, not after the user has confirmed the import.
capability = savant_capabilities.detect(probe_live=True)
api_enabled = bool(capability.get("api_enabled"))
ctx = None
# The target workspace/namespace is the authenticated session's — the credentials' namespace,
# the same one create and fetch all operate in. `home` (or the legacy alias `root`) targets
# the Home folder (namespace root), which the import API represents as a null folder id. The folder id itself is not
# pre-resolved here: `workflow verify --operation create --expect-folder-id` hard-fails unless
# the created workflow reports the requested folder, so a bad/foreign id is caught there.
root_target = is_root_folder_target(folder_id)
folder: dict = {"id": None, "path": "(Home folder)"} if root_target else {"id": folder_id, "path": None}
if api_enabled:
ctx = api.discover_session(None)
# Dataset ids are WORKSPACE-scoped. A real-looking id from another workspace
# does not resolve in the import target and Savant silently DROPS that source
# node (verified live 2026-06-10: four sources discovered via one folder,
# imported via another — all four vanished, twice, before this check existed).
# Resolve every bound id against the target workspace BEFORE import.
unresolved_sources: list[str] = []
# AI providers are WORKSPACE-scoped too: a `providerId` copied from another workspace's
# export looks real but does not resolve here, and Savant silently DROPS the AI node on
# import. Resolve every bound providerId against the target workspace BEFORE import. Trial
# providers (`savant_*`) are platform-stable and always returned, so the builder default
# never false-positives.
unresolved_ai_providers: list[str] = []
if ctx is not None:
# Both id sets are MCP `search` results the caller supplies. None means "not supplied",
# which skips the check — never "the workspace has none", which would block every bound
# dataset and provider.
workspace_ids = load_id_set(sources_json, flag="--sources-json")
if workspace_ids:
unresolved_sources = vw.unresolved_source_dataset_ids(nodes, workspace_ids)
provider_ids = load_id_set(providers_json, flag="--providers-json")
if provider_ids:
unresolved_ai_providers = vw.unresolved_ai_provider_nodes(nodes, provider_ids)
blocked_reasons = []
if not api_enabled:
blocked_reasons.append(
"Savant API is not enabled for this session/package; workflow creation through the API is unavailable."
)
# Workspace-consent guard: the session's namespace must be the one the USER confirmed as the
# destination at intake. This is the deterministic net against silent retargeting — an agent
# that switched workspaces after the user confirmed one (or is creating in a workspace the
# user never named) is caught here at import time, not just at the evidence precheck.
confirmed_ns = (confirmed_namespace or "").strip()
if ctx is not None and confirmed_ns:
session_ns = (ctx.namespace or "").strip()
if not session_ns:
blocked_reasons.append(
"the session namespace could not be determined, so it cannot be verified against "
f"the user-confirmed workspace namespace `{confirmed_ns}` — switch to the confirmed "
"workspace (MCP switch-workspace), re-run `session bind` with --namespace, and retry"
)
elif session_ns != confirmed_ns:
blocked_reasons.append(
f"session namespace `{session_ns}` does not match the user-confirmed workspace "
f"namespace `{confirmed_ns}` — the session is not in the workspace the user "
"confirmed as the destination. Switch back (MCP switch-workspace) or have the "
"user explicitly name the new destination workspace; never retarget on your own."
)
if errors:
blocked_reasons.append(f"{len(errors)} validation error(s)")
if missing_ai_providers:
blocked_reasons.append(f"missing AI provider on: {', '.join(missing_ai_providers)} "
"(resolve the workspace provider or use the Savant Trial default before import)")
if placeholder_sources:
blocked_reasons.append(f"placeholder/missing source dataset id on: {', '.join(placeholder_sources)} "
"(supply a real dataset_id from Planner/dataset discovery before import)")
if unresolved_sources:
blocked_reasons.append(
f"source dataset id(s) do not resolve in the TARGET workspace: {', '.join(unresolved_sources)} "
"(datasets are workspace-scoped; discover or create them in the target workspace — "
"importing anyway silently drops these source nodes)"
)
if unresolved_ai_providers:
blocked_reasons.append(
f"AI provider id(s) do not resolve in the TARGET workspace: {', '.join(unresolved_ai_providers)} "
"(providers are workspace-scoped; use one from MCP `search` with types: [\"ai_provider\"] "
"or the Savant Trial default — importing anyway silently drops these AI nodes)"
)
report = {
"phase": "preflight",
"validation": {"errors": errors, "warnings": warnings},
"folder": {"id": folder.get("id"), "path": folder.get("path")},
"namespace": ctx.namespace if ctx is not None else None,
"confirmedNamespace": confirmed_ns or None,
"capability": capability,
"missingAiProviders": missing_ai_providers,
"sourceDatasetPlaceholders": placeholder_sources,
"sourceDatasetUnresolvedInWorkspace": unresolved_sources,
"aiProvidersUnresolvedInWorkspace": unresolved_ai_providers,
"nodeCount": len([n for n in nodes if n.get("type") not in ("group", "text", "outlet")]),
"blocked": bool(blocked_reasons),
"blockedReasons": blocked_reasons,
}
return report, ctx, folder
def default_output_path(import_json: str | Path) -> Path:
stem = Path(import_json).stem or "workflow"
return workspace_tmp("workflow-create", f"{stem}.create.json")
def _workflow_name(json_path: str | Path) -> str:
wf = json.loads(Path(json_path).read_text(encoding="utf-8"))
return str(wf.get("name") or Path(json_path).stem or "").strip()
def _session_tmp_root() -> Path:
return workspace_tmp()
def duplicate_import_candidates(*, workflow_name: str, folder_id: str, report_path: Path) -> list[dict]:
"""Return same-session live creates for the same workflow name and folder.
Creator must not use repeated imports as an iteration loop. This guard is intentionally scoped
to the current AI session tmp tree because that is where the helper has durable evidence of
flows it created during this run.
"""
candidates: list[dict] = []
root = _session_tmp_root()
if not root.exists():
return candidates
resolved_report = report_path.resolve()
for path in root.rglob("*.json"):
try:
if path.resolve() == resolved_report:
continue
data = json.loads(path.read_text(encoding="utf-8"))
except Exception:
continue
if not isinstance(data, dict):
continue
if data.get("phase") != "create":
continue
if data.get("status") not in _CREATE_STATUSES:
continue
existing_folder = str((data.get("folder") or {}).get("id") or data.get("verifiedFolderId") or "")
existing_name = str(data.get("workflowName") or "").strip()
if not existing_name and isinstance(data.get("importJson"), str):
try:
existing_name = _workflow_name(data["importJson"])
except Exception: # noqa: BLE001 - best-effort compatibility with older reports.
existing_name = Path(data["importJson"]).stem
if existing_folder == str(folder_id) and existing_name == workflow_name and data.get("flowUrl"):
candidates.append({
"reportPath": str(path),
"status": data.get("status"),
"flowId": data.get("flowId"),
"flowUrl": data.get("flowUrl"),
})
return candidates
def after_json_path(flow_id: str) -> Path:
"""Where the caller should write the MCP `fetch` result for the verify/inspect handoff.
A derived, predictable path is what makes the printed follow-up commands copy-pasteable
instead of a template the caller has to fill in.
"""
return workflow_evidence.default_post_write_recipe_path(flow_id, "create")
def main(argv=None) -> int:
p = argparse.ArgumentParser(description=__doc__)
p.add_argument("--import-json", required=True, help="Workflow JSON to create.")
p.add_argument("--folder-id", required=True, help="Target folder id, or `home` for the Home folder (namespace root; `root` is a legacy alias). The workflow is created in the authenticated session's namespace. Find the id via MCP search/fetch on the folder entity.")
p.add_argument("--confirmed-namespace", required=True,
help="Namespace of the workspace the USER confirmed as the destination (from MCP search/whereami, recorded in the handoff's context_confirmation.namespace). Import is blocked if the session is in any other namespace — the destination workspace is user-named, never agent-chosen.")
p.add_argument("--sources-json", type=Path,
help="MCP `search` result for types: [\"source\"], used to catch source dataset "
"ids that do not resolve in the target workspace (they are silently dropped "
"on import). Omit to skip that check.")
p.add_argument("--providers-json", type=Path,
help="MCP `search` result for types: [\"ai_provider\"], used to catch AI nodes "
"bound to a providerId that does not resolve here (the whole node is "
"silently dropped on import). Omit to skip that check.")
p.add_argument("--confirm-import", action="store_true",
help="Proceed past preflight to import + verify. Set ONLY after the user confirmed creation in chat.")
p.add_argument("--expect-columns", help="Comma-separated expected final output columns (for the inspect contract).")
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="Deterministic stage NAME to verify (repeatable).")
p.add_argument("--no-terminal-preview", action="store_true",
help="On --confirm-import, skip the default terminal output preview and inspect only named checkpoints/contracts.")
p.add_argument("--sample-tier", default="1k")
p.add_argument("--timeout-seconds", type=int, default=120)
p.add_argument("--skip-inspect", action="store_true", help="Import only; do not run the verification pass.")
p.add_argument("--allow-duplicate-import", action="store_true",
help="Allow creating another live workflow with the same name in the same folder during this AI session. Use only after explicit user confirmation.")
p.add_argument("--post-save-recipe-path", type=Path,
help="Override where --confirm-import writes the re-fetched created workflow JSON.")
p.add_argument("--no-post-save-recipe", action="store_true",
help="On --confirm-import, skip writing the re-fetched created workflow JSON.")
p.add_argument("--output-path", type=Path, help="Where to write the JSON report.")
args = p.parse_args(argv)
report_path = args.output_path or default_output_path(args.import_json)
report, ctx, folder = preflight(
args.import_json, folder_id=args.folder_id, confirmed_namespace=args.confirmed_namespace,
sources_json=args.sources_json, providers_json=args.providers_json,
)
def emit(rep):
api.save_json(rep, report_path)
print(f"CREATE [{rep['phase']}] folder={rep.get('folder',{}).get('path') or rep.get('folder',{}).get('id')}")
v = rep.get("validation", {})
if rep["phase"] != "create":
print(f" validation: {len(v.get('errors', []))} errors, {len(v.get('warnings', []))} warnings")
for e in v.get("errors", []):
print(f" ERROR: {e}")
for w in v.get("warnings", []):
print(f" warn: {w}")
if rep["blocked"]:
print(" BLOCKED:")
for r in rep["blockedReasons"]:
print(f" - {r}")
# Preflight blocked, or not yet confirmed -> stop before import.
if report["blocked"]:
emit(report)
print(" -> resolve the blockers above, then re-run.")
return 2
workflow_name = _workflow_name(args.import_json)
duplicates = duplicate_import_candidates(
workflow_name=workflow_name,
folder_id=str(folder.get("id") or ""),
report_path=report_path,
)
report["workflowName"] = workflow_name
report["sameSessionExistingCreates"] = duplicates
if not args.confirm_import:
emit(report)
if duplicates:
print(" duplicate guard:")
for item in duplicates:
print(f" [{item['status']}] {item['flowUrl']}")
print(" -> this name/folder already has a same-session create; inspect/edit it instead of creating another copy.")
return 2
print(" -> preflight clean. Confirm creation with the user, then re-run with --confirm-import.")
return 0
if duplicates and not args.allow_duplicate_import:
blocked = {
"phase": "duplicate-import-guard",
"status": "blocked",
"workflowName": workflow_name,
"folder": report["folder"],
"existingCreates": duplicates,
"nextActions": [
"Inspect or edit the already-created workflow instead of creating another copy.",
"If the user explicitly confirms a replacement copy is desired, rerun with --allow-duplicate-import.",
],
}
api.save_json(blocked, report_path)
print(f"CREATE BLOCKED: same-session workflow already exists — {workflow_name}")
for item in duplicates:
print(f" [{item['status']}] {item['flowUrl']}")
print(" -> inspect/edit the existing workflow, or rerun with --allow-duplicate-import after explicit user confirmation.")
return 2
# Create + verify.
import_result = api.create_workflow_from_json(ctx, Path(args.import_json), folder_id=folder.get("id"), poll=True)
flow_url = import_result.get("flowUrl")
flow_id = import_result.get("flowId")
out = {"phase": "create", "workflowName": workflow_name, "folder": report["folder"], "flowId": flow_id, "flowUrl": flow_url,
"importJson": str(Path(args.import_json).resolve()),
"supersededCandidates": duplicates if args.allow_duplicate_import else [],
"verifiedFolderId": import_result.get("verifiedFolderId"), "preflight": report}
if not flow_url:
out["status"] = "import-unverified"
out["note"] = "import promise did not return a flow id/url"
emit(out)
print("CREATE: import did not return a flow id — verify manually.")
return 1
# The import is done; nothing here can confirm it landed correctly, because that needs the
# created recipe read back and this toolchain no longer reads recipes. The caller re-fetches
# through MCP and runs the two commands below. Until `workflow verify` returns ok, this create
# is NOT verified and must not be reported as such.
after_json = after_json_path(str(flow_id))
verify_cmd = [
"python3 savant.py workflow verify --operation create",
f"--workflow-json {after_json}",
f"--expect-flow-id {flow_id}",
f"--expect-folder-id {folder.get('id') or ''}",
f"--source-json {Path(args.import_json).resolve()}",
"--block-documentation-gaps",
]
inspect_cmd = [
"python3 savant.py workflow inspect",
f"{flow_url}",
f"--recipe-json {after_json}",
f"--imported-json {Path(args.import_json).resolve()}",
]
if args.expect_columns:
inspect_cmd.append(f"--expect-columns {args.expect_columns}")
if args.expected_outputs_json:
inspect_cmd.append(f"--expected-outputs-json {args.expected_outputs_json}")
for name in args.checkpoint:
inspect_cmd.append(f"--checkpoint {name!r}")
if args.no_terminal_preview:
inspect_cmd.append("--skip-terminals")
out["status"] = "created-unverified"
out["verified"] = False
out["nextSteps"] = {
"1_fetch": f"MCP `fetch` on savant://workflow/{flow_id}, written to {after_json}",
"2_verify": " ".join(verify_cmd),
"3_inspect": None if args.skip_inspect else " ".join(inspect_cmd),
}
api.save_json(out, report_path)
print(f"CREATE: imported, NOT YET VERIFIED — {flow_url}")
print("")
print("The import returned a flow, but nothing has confirmed it landed correctly yet.")
print("Run these before reporting the workflow as created:")
print(f" 1. Fetch the flow: MCP `fetch` on savant://workflow/{flow_id} -> {after_json}")
print(f" 2. {out['nextSteps']['2_verify']}")
if out["nextSteps"]["3_inspect"]:
print(f" 3. {out['nextSteps']['3_inspect']}")
return 0
if __name__ == "__main__":
raise SystemExit(main())
SHA-256: 3cede4624bb3436fbb9b64dc03034c67ad76d6216278494f700337bb5bbbf684