#!/usr/bin/env python3 """Compare every node's declared outputSchema against what it actually returns. Why this exists: an outputSchema is what tells a workflow author, and the WebUI's field-path picker, where a value lives in a node's output. It drifts silently - a node's execute() changes and nobody remembers to update the schema next to it - and a wrong schema is worse than no schema, because it looks authoritative while pointing at the wrong path. This reuses the existing tests/nodes/*.json fixtures and the same execution mechanism as scripts/verify-node.py (create a throwaway workflow, run it, read back each node's real output from the execution record) rather than adding a second corpus to maintain. For every node type that completed at least once across the fixtures, it unions the top-level keys of every observed output and checks each one is covered by that node type's declared outputSchema - as a plain property, a patternProperties match (for dynamic-named branches such as switch's case0/case1), or additionalProperties being explicitly true (for nodes whose real shape is inherently dynamic, such as set-fields and workflow-input). A key the schema does not cover is reported as drift. A key the schema declares but no fixture ever produced is not reported - fixtures are a sample, not exhaustive, and many declared keys are genuinely conditional (an error field that only appears on failure, for instance). Usage: ./scripts/check-output-schema-drift.py # all fixtures ./scripts/check-output-schema-drift.py case1 case2 # only these fixtures (by file stem) Exit 0 if every observed key is covered, 1 if any node type shows drift. Requires the webserver and a runner to be up (same as run-node-tests.sh). """ import importlib.util import json import subprocess import sys import time from pathlib import Path ROOT = Path(__file__).resolve().parent.parent TESTS_DIR = ROOT / "tests" / "nodes" def load_verify_node(): spec = importlib.util.spec_from_file_location("verify_node", ROOT / "scripts" / "verify-node.py") mod = importlib.util.module_from_spec(spec) spec.loader.exec_module(mod) return mod def load_schemas(): proc = subprocess.run( ["node", str(ROOT / "scripts" / "extract-output-schemas.js")], capture_output=True, text=True, cwd=ROOT, ) if proc.returncode != 0: raise SystemExit("extract-output-schemas.js failed:\n" + proc.stderr) if proc.stderr: sys.stderr.write(proc.stderr) return json.loads(proc.stdout) # Added by the engine itself to any node's stored result, regardless of node # type - not something an individual node's own outputSchema should have to # declare. See workflow_engine.cpp: applyNodeMarkers (_appliedConfig, # _ignoredConfigKeys, written whenever a Configurator feeds the node) and the # loop body runner (_executedInLoop, written to every node that ran inside a # Loop). A node-specific drift is one the node's own execute() introduced; # these are not that. ENGINE_ADDED_KEYS = {"_executedInLoop", "_appliedConfig", "_ignoredConfigKeys"} def key_covered(key, schema): if key in ENGINE_ADDED_KEYS: return True if key in schema["properties"]: return True for pattern in schema["patternProperties"]: if __import__("re").match(pattern, key): return True return schema["additionalProperties"] is True def run_case(vn, case_path, token): """Runs one fixture like verify-node.py does, but only to observe outputs - assertions in the fixture are ignored here, drift is a separate concern from behavioural correctness, and a fixture that intentionally exercises a failure path (errorContains, expectMissing) is exactly the kind of case that would otherwise make this look like a runner problem.""" case = json.loads(case_path.read_text()) if "http" in case: # These assert on the raw HTTP response, but the workflow underneath # still runs ordinary nodes - run it as a click/manual case instead so # the nodes still execute and still get observed. pass unmet = vn.precondition_unmet(case) if unmet: return None, f"skipped: {unmet}" node_types = {n["id"]: n["type"] for n in case.get("nodes", [])} helper_ids = [] for helper in case.get("helpers", []): made = vn.call("POST", "/workflows", token, { "name": helper["name"], "nodes": helper["nodes"], "connections": helper.get("connections", []), }) helper_id = made.get("id") or made.get("_id") if not helper_id: return None, "no id for helper" helper_ids.append(helper_id) if helper.get("active", True): vn.call("POST", f"/workflows/{helper_id}/activate", token, {}) case = json.loads(json.dumps(case).replace("{{helper:" + helper["key"] + "}}", helper_id)) created = vn.call("POST", "/workflows", token, { "name": case["name"], "nodes": case["nodes"], "connections": case["connections"], "settings": case.get("settings", {}), }) workflow_id = created.get("id") or created.get("_id") if not workflow_id: return None, "no workflow id" observed = {} try: if "http" in case: vn.call("POST", f"/workflows/{workflow_id}/activate", token, {}) try: vn.run_http_case(case, token, workflow_id) except SystemExit: pass # The webhook call already ran the workflow; find its execution. execs = vn.call("GET", f"/executions?pageSize=1&workflowId={workflow_id}", token) executions = execs.get("executions", []) if not executions: return {}, None execution_id = executions[0]["_id"] else: started = vn.call("POST", f"/workflows/{workflow_id}/execute", token, {}) execution_id = started["executionId"] if "resume" in case: try: vn.run_resume_case(case, token, execution_id) except SystemExit as e: return None, f"resume setup failed: {e}" execution = None for _ in range(60): time.sleep(0.5) execution = vn.call("GET", f"/executions/{execution_id}", token) if execution.get("status") in ("completed", "failed", "cancelled", "waiting"): break else: return None, "execution did not finish" for node_exec in execution.get("nodeExecutions", []): if node_exec.get("status") != "completed": continue node_type = node_types.get(node_exec.get("nodeId")) if not node_type: continue output = node_exec.get("output") if not isinstance(output, dict): continue observed.setdefault(node_type, set()).update(output.keys()) return observed, None finally: vn.call("DELETE", f"/workflows/{workflow_id}", token) for helper_id in helper_ids: vn.call("DELETE", f"/workflows/{helper_id}", token) def main(): vn = load_verify_node() schemas = load_schemas() wanted = sys.argv[1:] cases = sorted(TESTS_DIR.glob("*.json")) if wanted: cases = [c for c in cases if c.stem in wanted] if not cases: raise SystemExit(f"no fixtures matched {wanted}") token = vn.login() vn.call("POST", "/nodes/migrate", token, {"nodesPath": "./nodes"}) observed_by_type = {} skipped = [] errors = [] for case_path in cases: try: observed, note = run_case(vn, case_path, token) except Exception as exc: # a broken fixture should not abort the whole scan errors.append(f"{case_path.stem}: {exc}") continue if note and observed is None: skipped.append(f"{case_path.stem}: {note}") continue for node_type, keys in (observed or {}).items(): observed_by_type.setdefault(node_type, set()).update(keys) drift = {} for node_type, keys in sorted(observed_by_type.items()): schema = schemas.get(node_type) if schema is None: drift[node_type] = {"error": "no outputSchema found for this node type", "keys": sorted(keys)} continue uncovered = sorted(k for k in keys if not key_covered(k, schema)) if uncovered: drift[node_type] = {"uncovered": uncovered, "declared": schema["properties"], "file": schema["file"]} all_types = set(schemas.keys()) - {"utils-test"} covered_types = set(observed_by_type.keys()) uncovered_types = sorted(all_types - covered_types) print(f"observed {len(covered_types)}/{len(all_types)} node types across {len(cases) - len(skipped) - len(errors)} fixtures " f"({len(skipped)} skipped, {len(errors)} errored)") if uncovered_types: print("\nno completed execution observed for (not necessarily a problem - just not checked this run):") for t in uncovered_types: print(f" {t}") if errors: print("\nfixtures that errored out (investigate separately, not counted as drift):") for e in errors: print(f" {e}") if not drift: print("\nno drift: every observed output key is covered by its node's outputSchema") return 0 print(f"\nDRIFT in {len(drift)} node type(s):") for node_type, info in drift.items(): if "error" in info: print(f" {node_type}: {info['error']} (keys seen: {', '.join(info['keys'])})") continue print(f" {node_type} ({info['file']}):") print(f" schema declares: {', '.join(info['declared']) or '(none)'}") print(f" undeclared keys actually returned: {', '.join(info['uncovered'])}") return 1 if __name__ == "__main__": sys.exit(main())