check-output-schema-drift.py 9.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251
  1. #!/usr/bin/env python3
  2. """Compare every node's declared outputSchema against what it actually returns.
  3. Why this exists: an outputSchema is what tells a workflow author, and the
  4. WebUI's field-path picker, where a value lives in a node's output. It drifts
  5. silently - a node's execute() changes and nobody remembers to update the
  6. schema next to it - and a wrong schema is worse than no schema, because it
  7. looks authoritative while pointing at the wrong path.
  8. This reuses the existing tests/nodes/*.json fixtures and the same execution
  9. mechanism as scripts/verify-node.py (create a throwaway workflow, run it,
  10. read back each node's real output from the execution record) rather than
  11. adding a second corpus to maintain. For every node type that completed at
  12. least once across the fixtures, it unions the top-level keys of every
  13. observed output and checks each one is covered by that node type's declared
  14. outputSchema - as a plain property, a patternProperties match (for
  15. dynamic-named branches such as switch's case0/case1), or additionalProperties
  16. being explicitly true (for nodes whose real shape is inherently dynamic, such
  17. as set-fields and workflow-input).
  18. A key the schema does not cover is reported as drift. A key the schema
  19. declares but no fixture ever produced is not reported - fixtures are a
  20. sample, not exhaustive, and many declared keys are genuinely conditional
  21. (an error field that only appears on failure, for instance).
  22. Usage:
  23. ./scripts/check-output-schema-drift.py # all fixtures
  24. ./scripts/check-output-schema-drift.py case1 case2 # only these fixtures (by file stem)
  25. Exit 0 if every observed key is covered, 1 if any node type shows drift.
  26. Requires the webserver and a runner to be up (same as run-node-tests.sh).
  27. """
  28. import importlib.util
  29. import json
  30. import subprocess
  31. import sys
  32. import time
  33. from pathlib import Path
  34. ROOT = Path(__file__).resolve().parent.parent
  35. TESTS_DIR = ROOT / "tests" / "nodes"
  36. def load_verify_node():
  37. spec = importlib.util.spec_from_file_location("verify_node", ROOT / "scripts" / "verify-node.py")
  38. mod = importlib.util.module_from_spec(spec)
  39. spec.loader.exec_module(mod)
  40. return mod
  41. def load_schemas():
  42. proc = subprocess.run(
  43. ["node", str(ROOT / "scripts" / "extract-output-schemas.js")],
  44. capture_output=True, text=True, cwd=ROOT,
  45. )
  46. if proc.returncode != 0:
  47. raise SystemExit("extract-output-schemas.js failed:\n" + proc.stderr)
  48. if proc.stderr:
  49. sys.stderr.write(proc.stderr)
  50. return json.loads(proc.stdout)
  51. # Added by the engine itself to any node's stored result, regardless of node
  52. # type - not something an individual node's own outputSchema should have to
  53. # declare. See workflow_engine.cpp: applyNodeMarkers (_appliedConfig,
  54. # _ignoredConfigKeys, written whenever a Configurator feeds the node) and the
  55. # loop body runner (_executedInLoop, written to every node that ran inside a
  56. # Loop). A node-specific drift is one the node's own execute() introduced;
  57. # these are not that.
  58. ENGINE_ADDED_KEYS = {"_executedInLoop", "_appliedConfig", "_ignoredConfigKeys"}
  59. def key_covered(key, schema):
  60. if key in ENGINE_ADDED_KEYS:
  61. return True
  62. if key in schema["properties"]:
  63. return True
  64. for pattern in schema["patternProperties"]:
  65. if __import__("re").match(pattern, key):
  66. return True
  67. return schema["additionalProperties"] is True
  68. def run_case(vn, case_path, token):
  69. """Runs one fixture like verify-node.py does, but only to observe outputs -
  70. assertions in the fixture are ignored here, drift is a separate concern
  71. from behavioural correctness, and a fixture that intentionally exercises a
  72. failure path (errorContains, expectMissing) is exactly the kind of case
  73. that would otherwise make this look like a runner problem."""
  74. case = json.loads(case_path.read_text())
  75. if "http" in case:
  76. # These assert on the raw HTTP response, but the workflow underneath
  77. # still runs ordinary nodes - run it as a click/manual case instead so
  78. # the nodes still execute and still get observed.
  79. pass
  80. unmet = vn.precondition_unmet(case)
  81. if unmet:
  82. return None, f"skipped: {unmet}"
  83. node_types = {n["id"]: n["type"] for n in case.get("nodes", [])}
  84. helper_ids = []
  85. for helper in case.get("helpers", []):
  86. made = vn.call("POST", "/workflows", token, {
  87. "name": helper["name"],
  88. "nodes": helper["nodes"],
  89. "connections": helper.get("connections", []),
  90. })
  91. helper_id = made.get("id") or made.get("_id")
  92. if not helper_id:
  93. return None, "no id for helper"
  94. helper_ids.append(helper_id)
  95. if helper.get("active", True):
  96. vn.call("POST", f"/workflows/{helper_id}/activate", token, {})
  97. case = json.loads(json.dumps(case).replace("{{helper:" + helper["key"] + "}}", helper_id))
  98. created = vn.call("POST", "/workflows", token, {
  99. "name": case["name"],
  100. "nodes": case["nodes"],
  101. "connections": case["connections"],
  102. "settings": case.get("settings", {}),
  103. })
  104. workflow_id = created.get("id") or created.get("_id")
  105. if not workflow_id:
  106. return None, "no workflow id"
  107. observed = {}
  108. try:
  109. if "http" in case:
  110. vn.call("POST", f"/workflows/{workflow_id}/activate", token, {})
  111. try:
  112. vn.run_http_case(case, token, workflow_id)
  113. except SystemExit:
  114. pass
  115. # The webhook call already ran the workflow; find its execution.
  116. execs = vn.call("GET", f"/executions?pageSize=1&workflowId={workflow_id}", token)
  117. executions = execs.get("executions", [])
  118. if not executions:
  119. return {}, None
  120. execution_id = executions[0]["_id"]
  121. else:
  122. started = vn.call("POST", f"/workflows/{workflow_id}/execute", token, {})
  123. execution_id = started["executionId"]
  124. if "resume" in case:
  125. try:
  126. vn.run_resume_case(case, token, execution_id)
  127. except SystemExit as e:
  128. return None, f"resume setup failed: {e}"
  129. execution = None
  130. for _ in range(60):
  131. time.sleep(0.5)
  132. execution = vn.call("GET", f"/executions/{execution_id}", token)
  133. if execution.get("status") in ("completed", "failed", "cancelled", "waiting"):
  134. break
  135. else:
  136. return None, "execution did not finish"
  137. for node_exec in execution.get("nodeExecutions", []):
  138. if node_exec.get("status") != "completed":
  139. continue
  140. node_type = node_types.get(node_exec.get("nodeId"))
  141. if not node_type:
  142. continue
  143. output = node_exec.get("output")
  144. if not isinstance(output, dict):
  145. continue
  146. observed.setdefault(node_type, set()).update(output.keys())
  147. return observed, None
  148. finally:
  149. vn.call("DELETE", f"/workflows/{workflow_id}", token)
  150. for helper_id in helper_ids:
  151. vn.call("DELETE", f"/workflows/{helper_id}", token)
  152. def main():
  153. vn = load_verify_node()
  154. schemas = load_schemas()
  155. wanted = sys.argv[1:]
  156. cases = sorted(TESTS_DIR.glob("*.json"))
  157. if wanted:
  158. cases = [c for c in cases if c.stem in wanted]
  159. if not cases:
  160. raise SystemExit(f"no fixtures matched {wanted}")
  161. token = vn.login()
  162. vn.call("POST", "/nodes/migrate", token, {"nodesPath": "./nodes"})
  163. observed_by_type = {}
  164. skipped = []
  165. errors = []
  166. for case_path in cases:
  167. try:
  168. observed, note = run_case(vn, case_path, token)
  169. except Exception as exc: # a broken fixture should not abort the whole scan
  170. errors.append(f"{case_path.stem}: {exc}")
  171. continue
  172. if note and observed is None:
  173. skipped.append(f"{case_path.stem}: {note}")
  174. continue
  175. for node_type, keys in (observed or {}).items():
  176. observed_by_type.setdefault(node_type, set()).update(keys)
  177. drift = {}
  178. for node_type, keys in sorted(observed_by_type.items()):
  179. schema = schemas.get(node_type)
  180. if schema is None:
  181. drift[node_type] = {"error": "no outputSchema found for this node type", "keys": sorted(keys)}
  182. continue
  183. uncovered = sorted(k for k in keys if not key_covered(k, schema))
  184. if uncovered:
  185. drift[node_type] = {"uncovered": uncovered, "declared": schema["properties"], "file": schema["file"]}
  186. all_types = set(schemas.keys()) - {"utils-test"}
  187. covered_types = set(observed_by_type.keys())
  188. uncovered_types = sorted(all_types - covered_types)
  189. print(f"observed {len(covered_types)}/{len(all_types)} node types across {len(cases) - len(skipped) - len(errors)} fixtures "
  190. f"({len(skipped)} skipped, {len(errors)} errored)")
  191. if uncovered_types:
  192. print("\nno completed execution observed for (not necessarily a problem - just not checked this run):")
  193. for t in uncovered_types:
  194. print(f" {t}")
  195. if errors:
  196. print("\nfixtures that errored out (investigate separately, not counted as drift):")
  197. for e in errors:
  198. print(f" {e}")
  199. if not drift:
  200. print("\nno drift: every observed output key is covered by its node's outputSchema")
  201. return 0
  202. print(f"\nDRIFT in {len(drift)} node type(s):")
  203. for node_type, info in drift.items():
  204. if "error" in info:
  205. print(f" {node_type}: {info['error']} (keys seen: {', '.join(info['keys'])})")
  206. continue
  207. print(f" {node_type} ({info['file']}):")
  208. print(f" schema declares: {', '.join(info['declared']) or '(none)'}")
  209. print(f" undeclared keys actually returned: {', '.join(info['uncovered'])}")
  210. return 1
  211. if __name__ == "__main__":
  212. sys.exit(main())