|
@@ -55,6 +55,72 @@ def subset_matches(expected, actual, path):
|
|
|
return fails
|
|
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)
|
|
|
|
|
+ data = json.dumps(spec.get("body", {})).encode()
|
|
|
|
|
+ req = urllib.request.Request(url, data=data, method=spec.get("method", "POST"))
|
|
|
|
|
+ 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)
|
|
|
|
|
+
|
|
|
|
|
+ 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 main():
|
|
def main():
|
|
|
if len(sys.argv) != 2:
|
|
if len(sys.argv) != 2:
|
|
|
raise SystemExit("usage: verify-node.py <case.json>")
|
|
raise SystemExit("usage: verify-node.py <case.json>")
|
|
@@ -74,9 +140,19 @@ def main():
|
|
|
raise SystemExit(f"no workflow id in create response: {json.dumps(created)[:300]}")
|
|
raise SystemExit(f"no workflow id in create response: {json.dumps(created)[:300]}")
|
|
|
|
|
|
|
|
try:
|
|
try:
|
|
|
|
|
+ if "http" in case:
|
|
|
|
|
+ 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, {})
|
|
started = call("POST", f"/workflows/{workflow_id}/execute", token, {})
|
|
|
execution_id = started["executionId"]
|
|
execution_id = started["executionId"]
|
|
|
|
|
|
|
|
|
|
+ if "resume" in case:
|
|
|
|
|
+ run_resume_case(case, token, execution_id)
|
|
|
|
|
+
|
|
|
execution = None
|
|
execution = None
|
|
|
for _ in range(60):
|
|
for _ in range(60):
|
|
|
time.sleep(0.5)
|
|
time.sleep(0.5)
|