| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251 |
- #!/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())
|