#!/usr/bin/env python3 """Run one node verification case against the live SmartBotic services. Usage: scripts/verify-node.py tests/nodes/.json Exit 0 if every assertion holds, 1 if one does not, 2 if the case declares a precondition this machine does not meet - reported as a skip, never as a pass. """ import base64 import json import os import sys import time import urllib.error import urllib.parse import urllib.request # Where the API is. Overridable because the harness does not have to run on the # webserver's host - and, more importantly, because a few cases ask a NODE to # fetch a URL, and the node runs wherever the runner is. Those cases used to # hardcode http://localhost:8090; the day the runner moved to a different # machine from the webserver they failed with "Could not connect to server", # which looks like a broken http-request node rather than a test that assumed # the two were co-located. API_ROOT = os.environ.get("SMARTBOTIC_API", "http://localhost:8090").rstrip("/") BASE = API_ROOT + "/api/v1" # What a node should use to reach the API. Same thing by default; set it when # the runner reaches the webserver by a different name than the harness does. RUNNER_API_ROOT = os.environ.get("SMARTBOTIC_RUNNER_API", API_ROOT).rstrip("/") def call(method, path, token=None, body=None): data = json.dumps(body).encode() if body is not None else None req = urllib.request.Request(BASE + path, data=data, method=method) req.add_header("Content-Type", "application/json") if token: req.add_header("Authorization", "Bearer " + token) try: with urllib.request.urlopen(req, timeout=30) as res: raw = res.read() return json.loads(raw) if raw else {} except urllib.error.HTTPError as e: raise SystemExit(f"{method} {path} failed: {e.code} {e.read().decode()[:400]}") def login(): return call("POST", "/auth/login", body={"username": "admin", "password": "admin"})["accessToken"] def subset_matches(expected, actual, path): """Every key in expected must be present and equal in actual. Returns list of failures.""" fails = [] if isinstance(expected, dict): if not isinstance(actual, dict): return [f"{path}: expected an object, got {type(actual).__name__}"] for key, want in expected.items(): if key not in actual: fails.append(f"{path}.{key}: missing") else: fails += subset_matches(want, actual[key], f"{path}.{key}") elif isinstance(expected, list): if not isinstance(actual, list): return [f"{path}: expected a list, got {type(actual).__name__}"] if len(expected) != len(actual): fails.append(f"{path}: expected {len(expected)} items, got {len(actual)}") else: for i, want in enumerate(expected): fails += subset_matches(want, actual[i], f"{path}[{i}]") elif expected != actual: fails.append(f"{path}: expected {expected!r}, got {actual!r}") return fails def run_http_case(case, token, workflow_id): """A webhook case asserts on the HTTP response, not on node outputs.""" spec = case["http"] url = BASE.replace("/api/v1", "") + spec["path"].replace("{workflowId}", workflow_id) method = spec.get("method", "POST") multipart = spec.get("multipart") form_body = spec.get("formBody") if form_body is not None: # application/x-www-form-urlencoded, not JSON - cpp-httplib only # fills req.params (what handleWebhook's password check reads via # req.has_param) from a urlencoded body. A JSON body never reaches # that check, and multipart puts non-file parts in req.files instead # of req.params, so neither can exercise the password gate's # success path - only this can. payload = urllib.parse.urlencode(form_body).encode() req = urllib.request.Request(url, data=payload, method="POST") req.add_header("Content-Type", "application/x-www-form-urlencoded") elif multipart: boundary = "----smartboticverify" parts = [] for name, value in multipart.get("fields", {}).items(): parts.append( '--%s\r\nContent-Disposition: form-data; name="%s"\r\n\r\n%s\r\n' % (boundary, name, value)) for name, f in multipart.get("files", {}).items(): # A tiny real PNG, so the payload is genuinely binary rather than text # that happens to be labelled as an image. blob = base64.b64decode( "iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAYAAAAfFcSJAAAADUlEQVR42mP8z8BQDwAEhQGAhKmMIQAAAABJRU5ErkJggg==") head = ('--%s\r\nContent-Disposition: form-data; name="%s"; filename="%s"\r\n' 'Content-Type: %s\r\n\r\n' % (boundary, name, f.get("filename", "upload.png"), f.get("contentType", "image/png"))) parts.append(head.encode() + blob + b"\r\n") payload = b"".join(p.encode() if isinstance(p, str) else p for p in parts) payload += ("--%s--\r\n" % boundary).encode() req = urllib.request.Request(url, data=payload, method="POST") req.add_header("Content-Type", "multipart/form-data; boundary=%s" % boundary) else: data = None if method in ("POST", "PUT", "PATCH"): body = spec.get("body", {}) pad = spec.get("bodyPadBytes", 0) if pad: body = dict(body) body["_pad"] = "x" * pad data = json.dumps(body).encode() req = urllib.request.Request(url, data=data, method=method) if data is not None: req.add_header("Content-Type", "application/json") failures = [] try: with urllib.request.urlopen(req, timeout=45) as res: status, headers, raw = res.status, dict(res.headers), res.read().decode() except urllib.error.HTTPError as e: status, headers, raw = e.code, dict(e.headers), e.read().decode() print("http: %s %s -> %s" % (spec.get("method", "POST"), spec["path"], status)) print(" body: %s" % raw[:200]) if "expectStatus" in spec and status != spec["expectStatus"]: failures.append("status: expected %s, got %s" % (spec["expectStatus"], status)) for name, want in spec.get("expectHeaders", {}).items(): got = headers.get(name) if got is None: failures.append("header %s: missing" % name) elif want not in got: failures.append("header %s: expected to contain %r, got %r" % (name, want, got)) for fragment in spec.get("expectBodyContains", []): if fragment not in raw: failures.append("body: expected to contain %r" % fragment) for fragment in spec.get("expectBodyExcludes", []): if fragment in raw: failures.append("body: expected NOT to contain %r" % fragment) return failures def run_resume_case(case, token, execution_id): """Wait for the execution to pause, answer it, then let it finish.""" spec = case["resume"] want_status = spec.get("waitForStatus", "waiting") execution = {} for _ in range(60): time.sleep(0.5) execution = call("GET", "/executions/%s" % execution_id, token) if execution.get("status") == want_status: break if execution.get("status") in ("completed", "failed", "cancelled"): raise SystemExit( "execution %s reached %s without pausing" % (execution_id, execution.get("status"))) else: raise SystemExit("execution %s never reached %s" % (execution_id, want_status)) by_id = {n["nodeId"]: n for n in execution.get("nodeExecutions", [])} source = by_id.get(spec["tokenFrom"]) if not source: raise SystemExit("resume: node %s did not run" % spec["tokenFrom"]) resume_token = (source.get("output") or {}).get("token") if not resume_token: raise SystemExit("resume: node %s produced no token" % spec["tokenFrom"]) body = dict(spec.get("payload", {})) body["token"] = resume_token call("POST", "/executions/%s/resume" % execution_id, token, body) print("resumed %s with token %s" % (execution_id, resume_token)) def precondition_unmet(case): """Why this case cannot run here, or None if it can. A couple of cases need a model loaded on the SD.cpp server, which unloads itself when idle. Without this they fail for a reason that has nothing to do with the code, and a failure everyone learns to ignore is worse than no test. Skipping is only honest because it is reported as a skip - never as a pass. """ needs = case.get("requires") if not needs: return None probe = needs.get("http") if probe: try: with urllib.request.urlopen(probe, timeout=10) as response: body = json.loads(response.read().decode()) except Exception as exc: return f"{probe} could not be read ({exc})" for path, expected in (needs.get("expect") or {}).items(): actual = body for part in path.split("."): actual = (actual or {}).get(part) if isinstance(actual, dict) else None if actual != expected: return f"{probe} reports {path}={actual!r}, this case needs {expected!r}" return None def main(): if len(sys.argv) != 2: raise SystemExit("usage: verify-node.py ") raw_case = open(sys.argv[1]).read() # {{API}} is what a case writes when it needs a URL the RUNNER can fetch. raw_case = raw_case.replace("{{API}}", RUNNER_API_ROOT) case = json.loads(raw_case) unmet = precondition_unmet(case) if unmet: print(f"SKIP {case['name']}: {unmet}") return 2 token = login() call("POST", "/nodes/migrate", token, {"nodesPath": "./nodes"}) # Workflows the case needs to exist before it runs - a sub-workflow it # calls. Created here rather than named by id, so a case is not silently # tied to a row someone made by hand once and can delete. helper_ids = [] for helper in case.get("helpers", []): made = 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: raise SystemExit(f"no id for helper {helper['key']}") helper_ids.append(helper_id) if helper.get("active", True): call("POST", f"/workflows/{helper_id}/activate", token, {}) # Substituted everywhere, so the case refers to it by name. case = json.loads(json.dumps(case).replace("{{helper:" + helper["key"] + "}}", helper_id)) created = 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: raise SystemExit(f"no workflow id in create response: {json.dumps(created)[:300]}") try: if "http" in case: call("POST", f"/workflows/{workflow_id}/activate", token, {}) failures = run_http_case(case, token, workflow_id) print("\nFAIL" if failures else "\nPASS") for f in failures: print(" " + f) return 1 if failures else 0 started = call("POST", f"/workflows/{workflow_id}/execute", token, case.get("trigger", {})) execution_id = started["executionId"] if "resume" in case: run_resume_case(case, token, execution_id) execution = None for _ in range(60): time.sleep(0.5) execution = call("GET", f"/executions/{execution_id}", token) if execution.get("status") in ("completed", "failed", "cancelled"): break else: raise SystemExit(f"execution {execution_id} did not finish in 30s") by_id = {n["nodeId"]: n for n in execution.get("nodeExecutions", [])} failures = [] for node_id, want in case.get("expect", {}).items(): got = by_id.get(node_id) if got is None: failures.append(f"{node_id}: did not run") continue if "status" in want and got.get("status") != want["status"]: failures.append( f"{node_id}: status {got.get('status')!r}, expected {want['status']!r}" + (f" (error: {got.get('error')})" if got.get("error") else "") ) if "errorContains" in want: actual_error = got.get("error") or "" if want["errorContains"] not in actual_error: failures.append( f"{node_id}: error does not contain {want['errorContains']!r}, " f"got {actual_error!r}" ) if "output" in want: failures += subset_matches(want["output"], got.get("output"), node_id) if "loopNodeId" in want and got.get("loopNodeId") != want["loopNodeId"]: failures.append( f"{node_id}: loopNodeId {got.get('loopNodeId')!r}, expected {want['loopNodeId']!r}" ) if "loopIteration" in want and got.get("loopIteration") != want["loopIteration"]: failures.append( f"{node_id}: loopIteration {got.get('loopIteration')!r}, " f"expected {want['loopIteration']!r}" ) if want.get("noLoopTag") and ( "loopNodeId" in got or "loopIteration" in got ): failures.append( f"{node_id}: expected no loop tag, got loopNodeId={got.get('loopNodeId')!r} " f"loopIteration={got.get('loopIteration')!r}" ) # How many raw nodeExecutions entries carry a given node id - the # count itself is the assertion for things like "a disabled node # inside a loop body is recorded once, not once per iteration" and # "a loop body node's per-iteration records are not shadowed by a # duplicate mirror entry". all_node_executions = execution.get("nodeExecutions", []) for node_id, want_count in case.get("expectCounts", {}).items(): got_count = sum(1 for n in all_node_executions if n.get("nodeId") == node_id) if got_count != want_count: failures.append( f"{node_id}: {got_count} nodeExecutions entries, expected {want_count}" ) # nodeExecutions has to come back sorted by startedAt - the array # order is the contract, not something callers are meant to re-sort. if case.get("expectChronological"): starts = [n.get("startedAt", 0) for n in all_node_executions] if starts != sorted(starts): failures.append( f"nodeExecutions is not in startedAt order: {starts}" ) # What a watcher is told when the run fails. "Loop iteration failed" on # its own passed every assertion that looked at nodes, because the nodes # were all correct - the message was the broken part. wanted_error = case.get("expectExecutionError") if wanted_error and wanted_error not in (execution.get("error") or ""): failures.append( f"execution error does not contain {wanted_error!r}, " f"got {execution.get('error')!r}" ) # Assertions against the execution record itself, not a node inside it - # e.g. "singleNodeTarget" naming which node a 'test this node' run # targeted, absent entirely on an ordinary run. failures += subset_matches(case.get("expectExecution", {}), execution, "execution") for key in case.get("expectExecutionAbsent", []): if key in execution: failures.append(f"execution.{key}: expected absent, got {execution[key]!r}") for node_id in case.get("expectMissing", []): got = by_id.get(node_id) if got is not None and got.get("status") != "skipped": failures.append( f"{node_id}: ran, but should not have (status: {got.get('status')!r})" ) print(f"case: {case['name']} execution: {execution_id} status: {execution.get('status')}") for node_id, node in sorted(by_id.items()): print(f" {node_id:24} {node.get('status'):10} {json.dumps(node.get('output'))[:120]}") if failures: print("\nFAIL") for f in failures: print(" " + f) return 1 print("\nPASS") return 0 finally: call("DELETE", f"/workflows/{workflow_id}", token) for helper_id in helper_ids: call("DELETE", f"/workflows/{helper_id}", token) if __name__ == "__main__": sys.exit(main())