verify-node.py 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415
  1. #!/usr/bin/env python3
  2. """Run one node verification case against the live SmartBotic services.
  3. Usage: scripts/verify-node.py tests/nodes/<case>.json
  4. Exit 0 if every assertion holds, 1 if one does not, 2 if the case declares a
  5. precondition this machine does not meet - reported as a skip, never as a pass.
  6. """
  7. import base64
  8. import json
  9. import os
  10. import sys
  11. import time
  12. import urllib.error
  13. import urllib.parse
  14. import urllib.request
  15. # Where the API is. Overridable because the harness does not have to run on the
  16. # webserver's host - and, more importantly, because a few cases ask a NODE to
  17. # fetch a URL, and the node runs wherever the runner is. Those cases used to
  18. # hardcode http://localhost:8090; the day the runner moved to a different
  19. # machine from the webserver they failed with "Could not connect to server",
  20. # which looks like a broken http-request node rather than a test that assumed
  21. # the two were co-located.
  22. API_ROOT = os.environ.get("SMARTBOTIC_API", "http://localhost:8090").rstrip("/")
  23. BASE = API_ROOT + "/api/v1"
  24. # What a node should use to reach the API. Same thing by default; set it when
  25. # the runner reaches the webserver by a different name than the harness does.
  26. RUNNER_API_ROOT = os.environ.get("SMARTBOTIC_RUNNER_API", API_ROOT).rstrip("/")
  27. def call(method, path, token=None, body=None):
  28. data = json.dumps(body).encode() if body is not None else None
  29. req = urllib.request.Request(BASE + path, data=data, method=method)
  30. req.add_header("Content-Type", "application/json")
  31. if token:
  32. req.add_header("Authorization", "Bearer " + token)
  33. try:
  34. with urllib.request.urlopen(req, timeout=30) as res:
  35. raw = res.read()
  36. return json.loads(raw) if raw else {}
  37. except urllib.error.HTTPError as e:
  38. raise SystemExit(f"{method} {path} failed: {e.code} {e.read().decode()[:400]}")
  39. def login():
  40. return call("POST", "/auth/login", body={"username": "admin", "password": "admin"})["accessToken"]
  41. def subset_matches(expected, actual, path):
  42. """Every key in expected must be present and equal in actual. Returns list of failures."""
  43. fails = []
  44. if isinstance(expected, dict):
  45. if not isinstance(actual, dict):
  46. return [f"{path}: expected an object, got {type(actual).__name__}"]
  47. for key, want in expected.items():
  48. if key not in actual:
  49. fails.append(f"{path}.{key}: missing")
  50. else:
  51. fails += subset_matches(want, actual[key], f"{path}.{key}")
  52. elif isinstance(expected, list):
  53. if not isinstance(actual, list):
  54. return [f"{path}: expected a list, got {type(actual).__name__}"]
  55. if len(expected) != len(actual):
  56. fails.append(f"{path}: expected {len(expected)} items, got {len(actual)}")
  57. else:
  58. for i, want in enumerate(expected):
  59. fails += subset_matches(want, actual[i], f"{path}[{i}]")
  60. elif expected != actual:
  61. fails.append(f"{path}: expected {expected!r}, got {actual!r}")
  62. return fails
  63. def run_http_case(case, token, workflow_id):
  64. """A webhook case asserts on the HTTP response, not on node outputs."""
  65. spec = case["http"]
  66. url = BASE.replace("/api/v1", "") + spec["path"].replace("{workflowId}", workflow_id)
  67. method = spec.get("method", "POST")
  68. multipart = spec.get("multipart")
  69. form_body = spec.get("formBody")
  70. if form_body is not None:
  71. # application/x-www-form-urlencoded, not JSON - cpp-httplib only
  72. # fills req.params (what handleWebhook's password check reads via
  73. # req.has_param) from a urlencoded body. A JSON body never reaches
  74. # that check, and multipart puts non-file parts in req.files instead
  75. # of req.params, so neither can exercise the password gate's
  76. # success path - only this can.
  77. payload = urllib.parse.urlencode(form_body).encode()
  78. req = urllib.request.Request(url, data=payload, method="POST")
  79. req.add_header("Content-Type", "application/x-www-form-urlencoded")
  80. elif multipart:
  81. boundary = "----smartboticverify"
  82. parts = []
  83. for name, value in multipart.get("fields", {}).items():
  84. parts.append(
  85. '--%s\r\nContent-Disposition: form-data; name="%s"\r\n\r\n%s\r\n' % (boundary, name, value))
  86. for name, f in multipart.get("files", {}).items():
  87. # A tiny real PNG, so the payload is genuinely binary rather than text
  88. # that happens to be labelled as an image.
  89. blob = base64.b64decode(
  90. "iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAYAAAAfFcSJAAAADUlEQVR42mP8z8BQDwAEhQGAhKmMIQAAAABJRU5ErkJggg==")
  91. head = ('--%s\r\nContent-Disposition: form-data; name="%s"; filename="%s"\r\n'
  92. 'Content-Type: %s\r\n\r\n' % (boundary, name, f.get("filename", "upload.png"),
  93. f.get("contentType", "image/png")))
  94. parts.append(head.encode() + blob + b"\r\n")
  95. payload = b"".join(p.encode() if isinstance(p, str) else p for p in parts)
  96. payload += ("--%s--\r\n" % boundary).encode()
  97. req = urllib.request.Request(url, data=payload, method="POST")
  98. req.add_header("Content-Type", "multipart/form-data; boundary=%s" % boundary)
  99. else:
  100. data = None
  101. if method in ("POST", "PUT", "PATCH"):
  102. body = spec.get("body", {})
  103. pad = spec.get("bodyPadBytes", 0)
  104. if pad:
  105. body = dict(body)
  106. body["_pad"] = "x" * pad
  107. data = json.dumps(body).encode()
  108. req = urllib.request.Request(url, data=data, method=method)
  109. if data is not None:
  110. req.add_header("Content-Type", "application/json")
  111. failures = []
  112. try:
  113. with urllib.request.urlopen(req, timeout=45) as res:
  114. status, headers, raw = res.status, dict(res.headers), res.read().decode()
  115. except urllib.error.HTTPError as e:
  116. status, headers, raw = e.code, dict(e.headers), e.read().decode()
  117. print("http: %s %s -> %s" % (spec.get("method", "POST"), spec["path"], status))
  118. print(" body: %s" % raw[:200])
  119. if "expectStatus" in spec and status != spec["expectStatus"]:
  120. failures.append("status: expected %s, got %s" % (spec["expectStatus"], status))
  121. for name, want in spec.get("expectHeaders", {}).items():
  122. got = headers.get(name)
  123. if got is None:
  124. failures.append("header %s: missing" % name)
  125. elif want not in got:
  126. failures.append("header %s: expected to contain %r, got %r" % (name, want, got))
  127. for fragment in spec.get("expectBodyContains", []):
  128. if fragment not in raw:
  129. failures.append("body: expected to contain %r" % fragment)
  130. for fragment in spec.get("expectBodyExcludes", []):
  131. if fragment in raw:
  132. failures.append("body: expected NOT to contain %r" % fragment)
  133. return failures
  134. def run_resume_case(case, token, execution_id):
  135. """Wait for the execution to pause, answer it, then let it finish."""
  136. spec = case["resume"]
  137. want_status = spec.get("waitForStatus", "waiting")
  138. execution = {}
  139. for _ in range(60):
  140. time.sleep(0.5)
  141. execution = call("GET", "/executions/%s" % execution_id, token)
  142. if execution.get("status") == want_status:
  143. break
  144. if execution.get("status") in ("completed", "failed", "cancelled"):
  145. raise SystemExit(
  146. "execution %s reached %s without pausing" % (execution_id, execution.get("status")))
  147. else:
  148. raise SystemExit("execution %s never reached %s" % (execution_id, want_status))
  149. by_id = {n["nodeId"]: n for n in execution.get("nodeExecutions", [])}
  150. source = by_id.get(spec["tokenFrom"])
  151. if not source:
  152. raise SystemExit("resume: node %s did not run" % spec["tokenFrom"])
  153. resume_token = (source.get("output") or {}).get("token")
  154. if not resume_token:
  155. raise SystemExit("resume: node %s produced no token" % spec["tokenFrom"])
  156. body = dict(spec.get("payload", {}))
  157. body["token"] = resume_token
  158. call("POST", "/executions/%s/resume" % execution_id, token, body)
  159. print("resumed %s with token %s" % (execution_id, resume_token))
  160. def precondition_unmet(case):
  161. """Why this case cannot run here, or None if it can.
  162. A couple of cases need a model loaded on the SD.cpp server, which unloads
  163. itself when idle. Without this they fail for a reason that has nothing to do
  164. with the code, and a failure everyone learns to ignore is worse than no
  165. test. Skipping is only honest because it is reported as a skip - never as a
  166. pass.
  167. """
  168. needs = case.get("requires")
  169. if not needs:
  170. return None
  171. probe = needs.get("http")
  172. if probe:
  173. try:
  174. with urllib.request.urlopen(probe, timeout=10) as response:
  175. body = json.loads(response.read().decode())
  176. except Exception as exc:
  177. return f"{probe} could not be read ({exc})"
  178. for path, expected in (needs.get("expect") or {}).items():
  179. actual = body
  180. for part in path.split("."):
  181. actual = (actual or {}).get(part) if isinstance(actual, dict) else None
  182. if actual != expected:
  183. return f"{probe} reports {path}={actual!r}, this case needs {expected!r}"
  184. return None
  185. def project_id(token):
  186. """The project new rows belong to. The first one, which is what every other
  187. part of the harness implicitly uses."""
  188. return call("GET", "/projects", token)["projects"][0]["_id"]
  189. def main():
  190. if len(sys.argv) != 2:
  191. raise SystemExit("usage: verify-node.py <case.json>")
  192. raw_case = open(sys.argv[1]).read()
  193. # {{API}} is what a case writes when it needs a URL the RUNNER can fetch.
  194. raw_case = raw_case.replace("{{API}}", RUNNER_API_ROOT)
  195. case = json.loads(raw_case)
  196. unmet = precondition_unmet(case)
  197. if unmet:
  198. print(f"SKIP {case['name']}: {unmet}")
  199. return 2
  200. token = login()
  201. call("POST", "/nodes/migrate", token, {"nodesPath": "./nodes"})
  202. # Workflows the case needs to exist before it runs - a sub-workflow it
  203. # calls. Created here rather than named by id, so a case is not silently
  204. # tied to a row someone made by hand once and can delete.
  205. helper_ids = []
  206. for helper in case.get("helpers", []):
  207. made = call("POST", "/workflows", token, {
  208. "name": helper["name"],
  209. "nodes": helper["nodes"],
  210. "connections": helper.get("connections", []),
  211. })
  212. helper_id = made.get("id") or made.get("_id")
  213. if not helper_id:
  214. raise SystemExit(f"no id for helper {helper['key']}")
  215. helper_ids.append(helper_id)
  216. if helper.get("active", True):
  217. call("POST", f"/workflows/{helper_id}/activate", token, {})
  218. # Substituted everywhere, so the case refers to it by name.
  219. case = json.loads(json.dumps(case).replace("{{helper:" + helper["key"] + "}}", helper_id))
  220. # Credentials the case needs. Created here and deleted afterwards, so a case
  221. # is portable: sdcpp-model-load-assert names a credential by id, which ties
  222. # it to one installation and fails everywhere else.
  223. credential_ids = []
  224. for cred in case.get("credentials", []):
  225. made = call("POST", "/credentials", token, {
  226. "name": cred.get("name", "zz test: " + cred["key"]),
  227. "type": cred["type"],
  228. "projectId": project_id(token),
  229. "data": cred.get("data", {}),
  230. })
  231. cred_id = made.get("id") or made.get("_id")
  232. if not cred_id:
  233. raise SystemExit(f"no id for credential {cred['key']}")
  234. credential_ids.append(cred_id)
  235. case = json.loads(json.dumps(case).replace("{{credential:" + cred["key"] + "}}", cred_id))
  236. created = call("POST", "/workflows", token, {
  237. "name": case["name"],
  238. "nodes": case["nodes"],
  239. "connections": case["connections"],
  240. "settings": case.get("settings", {}),
  241. })
  242. workflow_id = created.get("id") or created.get("_id")
  243. if not workflow_id:
  244. raise SystemExit(f"no workflow id in create response: {json.dumps(created)[:300]}")
  245. try:
  246. if "http" in case:
  247. call("POST", f"/workflows/{workflow_id}/activate", token, {})
  248. failures = run_http_case(case, token, workflow_id)
  249. print("\nFAIL" if failures else "\nPASS")
  250. for f in failures:
  251. print(" " + f)
  252. return 1 if failures else 0
  253. started = call("POST", f"/workflows/{workflow_id}/execute", token, case.get("trigger", {}))
  254. execution_id = started["executionId"]
  255. if "resume" in case:
  256. run_resume_case(case, token, execution_id)
  257. execution = None
  258. for _ in range(60):
  259. time.sleep(0.5)
  260. execution = call("GET", f"/executions/{execution_id}", token)
  261. if execution.get("status") in ("completed", "failed", "cancelled"):
  262. break
  263. else:
  264. raise SystemExit(f"execution {execution_id} did not finish in 30s")
  265. by_id = {n["nodeId"]: n for n in execution.get("nodeExecutions", [])}
  266. failures = []
  267. for node_id, want in case.get("expect", {}).items():
  268. got = by_id.get(node_id)
  269. if got is None:
  270. failures.append(f"{node_id}: did not run")
  271. continue
  272. if "status" in want and got.get("status") != want["status"]:
  273. failures.append(
  274. f"{node_id}: status {got.get('status')!r}, expected {want['status']!r}"
  275. + (f" (error: {got.get('error')})" if got.get("error") else "")
  276. )
  277. if "errorContains" in want:
  278. actual_error = got.get("error") or ""
  279. if want["errorContains"] not in actual_error:
  280. failures.append(
  281. f"{node_id}: error does not contain {want['errorContains']!r}, "
  282. f"got {actual_error!r}"
  283. )
  284. if "output" in want:
  285. failures += subset_matches(want["output"], got.get("output"), node_id)
  286. if "loopNodeId" in want and got.get("loopNodeId") != want["loopNodeId"]:
  287. failures.append(
  288. f"{node_id}: loopNodeId {got.get('loopNodeId')!r}, expected {want['loopNodeId']!r}"
  289. )
  290. if "loopIteration" in want and got.get("loopIteration") != want["loopIteration"]:
  291. failures.append(
  292. f"{node_id}: loopIteration {got.get('loopIteration')!r}, "
  293. f"expected {want['loopIteration']!r}"
  294. )
  295. if want.get("noLoopTag") and (
  296. "loopNodeId" in got or "loopIteration" in got
  297. ):
  298. failures.append(
  299. f"{node_id}: expected no loop tag, got loopNodeId={got.get('loopNodeId')!r} "
  300. f"loopIteration={got.get('loopIteration')!r}"
  301. )
  302. # How many raw nodeExecutions entries carry a given node id - the
  303. # count itself is the assertion for things like "a disabled node
  304. # inside a loop body is recorded once, not once per iteration" and
  305. # "a loop body node's per-iteration records are not shadowed by a
  306. # duplicate mirror entry".
  307. all_node_executions = execution.get("nodeExecutions", [])
  308. for node_id, want_count in case.get("expectCounts", {}).items():
  309. got_count = sum(1 for n in all_node_executions if n.get("nodeId") == node_id)
  310. if got_count != want_count:
  311. failures.append(
  312. f"{node_id}: {got_count} nodeExecutions entries, expected {want_count}"
  313. )
  314. # nodeExecutions has to come back sorted by startedAt - the array
  315. # order is the contract, not something callers are meant to re-sort.
  316. if case.get("expectChronological"):
  317. starts = [n.get("startedAt", 0) for n in all_node_executions]
  318. if starts != sorted(starts):
  319. failures.append(
  320. f"nodeExecutions is not in startedAt order: {starts}"
  321. )
  322. # What a watcher is told when the run fails. "Loop iteration failed" on
  323. # its own passed every assertion that looked at nodes, because the nodes
  324. # were all correct - the message was the broken part.
  325. wanted_error = case.get("expectExecutionError")
  326. if wanted_error and wanted_error not in (execution.get("error") or ""):
  327. failures.append(
  328. f"execution error does not contain {wanted_error!r}, "
  329. f"got {execution.get('error')!r}"
  330. )
  331. # Assertions against the execution record itself, not a node inside it -
  332. # e.g. "singleNodeTarget" naming which node a 'test this node' run
  333. # targeted, absent entirely on an ordinary run.
  334. failures += subset_matches(case.get("expectExecution", {}), execution, "execution")
  335. for key in case.get("expectExecutionAbsent", []):
  336. if key in execution:
  337. failures.append(f"execution.{key}: expected absent, got {execution[key]!r}")
  338. for node_id in case.get("expectMissing", []):
  339. got = by_id.get(node_id)
  340. if got is not None and got.get("status") != "skipped":
  341. failures.append(
  342. f"{node_id}: ran, but should not have (status: {got.get('status')!r})"
  343. )
  344. print(f"case: {case['name']} execution: {execution_id} status: {execution.get('status')}")
  345. for node_id, node in sorted(by_id.items()):
  346. print(f" {node_id:24} {node.get('status'):10} {json.dumps(node.get('output'))[:120]}")
  347. if failures:
  348. print("\nFAIL")
  349. for f in failures:
  350. print(" " + f)
  351. return 1
  352. print("\nPASS")
  353. return 0
  354. finally:
  355. call("DELETE", f"/workflows/{workflow_id}", token)
  356. for helper_id in helper_ids:
  357. call("DELETE", f"/workflows/{helper_id}", token)
  358. for cred_id in credential_ids:
  359. call("DELETE", f"/credentials/{cred_id}", token)
  360. if __name__ == "__main__":
  361. sys.exit(main())