|
@@ -0,0 +1,251 @@
|
|
|
|
|
+#!/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())
|