test_api.py 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467
  1. from app.runner.api import create_runner_app
  2. from app.runner.auth import TaskTokenIssuer, TaskTokenVerifier
  3. from app.runner.ledger import InMemoryTaskLedger
  4. from app.runner.nodes import NodeRegistry
  5. NODE = {
  6. "id": "read_orders",
  7. "type": "sql.query",
  8. "data_source_uid": "01900000-0000-7000-8000-000000000010",
  9. "purpose": "read",
  10. "config": {"statement": "SELECT 1", "parameters": {}},
  11. }
  12. RULE_NODE = {
  13. "id": "apply_orders",
  14. "type": "rule.apply",
  15. "purpose": "write",
  16. "idempotency": {
  17. "strategy": "upsert",
  18. "key": "orders:2026-07-23",
  19. },
  20. "config": {
  21. "component_binding_id": "01900000-0000-7000-8000-000000000021",
  22. "rule_version_id": "01900000-0000-7000-8000-000000000022",
  23. "execution_plan_hash": "a" * 64,
  24. },
  25. }
  26. class Executor:
  27. def __init__(self):
  28. self.calls = 0
  29. self.context = None
  30. def execute(self, node, parameters, **kwargs):
  31. self.calls += 1
  32. self.context = kwargs
  33. return {
  34. "node": node["id"],
  35. "parameters": parameters,
  36. "output_artifact": "minio://dataops-rules/rules/output.parquet",
  37. }
  38. def test_runner_accepts_one_signed_task_and_rejects_replay():
  39. issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
  40. verifier = TaskTokenVerifier("x" * 32, clock=lambda: 1_000)
  41. executor = Executor()
  42. ledger = InMemoryTaskLedger()
  43. app = create_runner_app(
  44. verifier=verifier,
  45. ledger=ledger,
  46. registry=NodeRegistry({"sql.query": executor}),
  47. )
  48. token = issuer.issue(
  49. task_uid="01900000-0000-7000-8000-000000000011",
  50. dataflow_uid="01900000-0000-7000-8000-000000000012",
  51. deployment_id="01900000-0000-7000-8000-000000000014",
  52. environment="test",
  53. workflow_version=7,
  54. correlation_id="01900000-0000-7000-8000-000000000013",
  55. node=NODE,
  56. )
  57. payload = {"task_token": token, "node": NODE, "parameters": {"day": "today"}}
  58. with app.test_client() as client:
  59. first = client.post("/v1/tasks/execute", json=payload)
  60. replay = client.post("/v1/tasks/execute", json=payload)
  61. assert first.status_code == 200
  62. assert first.get_json()["result"]["node"] == "read_orders"
  63. assert first.get_json()["output_artifact"].startswith("minio://")
  64. assert executor.context["dataflow_uid"] == (
  65. "01900000-0000-7000-8000-000000000012"
  66. )
  67. assert executor.context["workflow_version"] == 7
  68. assert executor.context["node_id"] == "read_orders"
  69. assert replay.status_code == 409
  70. assert replay.get_json() == {"error": "task token already consumed"}
  71. assert executor.calls == 1
  72. def test_runner_health_has_no_secret_or_datasource_details():
  73. app = create_runner_app(
  74. verifier=TaskTokenVerifier("x" * 32),
  75. ledger=InMemoryTaskLedger(),
  76. registry=NodeRegistry({}),
  77. )
  78. with app.test_client() as client:
  79. response = client.get("/health")
  80. assert response.status_code == 200
  81. assert response.get_json() == {"status": "ok"}
  82. def test_runner_fails_closed_when_the_durable_ledger_is_unavailable():
  83. class UnavailableLedger:
  84. def claim(self, *_args, **_kwargs):
  85. raise RuntimeError("postgres password must never leak")
  86. issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
  87. app = create_runner_app(
  88. verifier=TaskTokenVerifier("x" * 32, clock=lambda: 1_000),
  89. ledger=UnavailableLedger(),
  90. registry=NodeRegistry({"sql.query": Executor()}),
  91. )
  92. token = issuer.issue(
  93. task_uid="01900000-0000-7000-8000-000000000011",
  94. dataflow_uid="01900000-0000-7000-8000-000000000012",
  95. deployment_id="01900000-0000-7000-8000-000000000014",
  96. environment="test",
  97. workflow_version=7,
  98. correlation_id="01900000-0000-7000-8000-000000000013",
  99. node=NODE,
  100. )
  101. with app.test_client() as client:
  102. response = client.post(
  103. "/v1/tasks/execute",
  104. json={
  105. "task_token": token,
  106. "node": NODE,
  107. "parameters": {},
  108. },
  109. )
  110. assert response.status_code == 503
  111. assert response.get_json() == {"error": "task ledger is unavailable"}
  112. assert "password" not in response.get_data(as_text=True).lower()
  113. def _rule_token(issuer):
  114. return issuer.issue(
  115. task_uid="01900000-0000-7000-8000-000000000031",
  116. dataflow_uid="01900000-0000-7000-8000-000000000032",
  117. deployment_id="01900000-0000-7000-8000-000000000033",
  118. environment="production",
  119. workflow_version=7,
  120. correlation_id="01900000-0000-7000-8000-000000000034",
  121. node=RULE_NODE,
  122. write_authorized=True,
  123. )
  124. def test_governed_rule_same_jti_replays_durable_response_once():
  125. issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
  126. executor = Executor()
  127. ledger = InMemoryTaskLedger()
  128. app = create_runner_app(
  129. verifier=TaskTokenVerifier("x" * 32, clock=lambda: 1_000),
  130. ledger=ledger,
  131. registry=NodeRegistry({"rule.apply": executor}),
  132. )
  133. payload = {
  134. "task_token": _rule_token(issuer),
  135. "node": RULE_NODE,
  136. "parameters": {},
  137. }
  138. with app.test_client() as client:
  139. first = client.post("/v1/tasks/execute", json=payload)
  140. replay = client.post("/v1/tasks/execute", json=payload)
  141. assert first.status_code == replay.status_code == 200
  142. assert first.get_json() == replay.get_json()
  143. assert replay.headers["X-Idempotent-Replay"] == "true"
  144. assert executor.calls == 1
  145. def test_governed_rule_recovers_terminal_evidence_after_response_loss():
  146. class LostFirstFinishLedger(InMemoryTaskLedger):
  147. def __init__(self):
  148. super().__init__()
  149. self.finish_calls = 0
  150. def finish(self, jti, **outcome):
  151. self.finish_calls += 1
  152. if self.finish_calls == 1:
  153. return
  154. return super().finish(jti, **outcome)
  155. class RecoverableExecutor(Executor):
  156. def replay_task(self, **_context):
  157. return {
  158. "node": RULE_NODE["id"],
  159. "parameters": {},
  160. "rows_out": 9,
  161. }
  162. issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
  163. executor = RecoverableExecutor()
  164. ledger = LostFirstFinishLedger()
  165. app = create_runner_app(
  166. verifier=TaskTokenVerifier("x" * 32, clock=lambda: 1_000),
  167. ledger=ledger,
  168. registry=NodeRegistry({"rule.apply": executor}),
  169. )
  170. payload = {
  171. "task_token": _rule_token(issuer),
  172. "node": RULE_NODE,
  173. "parameters": {},
  174. }
  175. with app.test_client() as client:
  176. first = client.post("/v1/tasks/execute", json=payload)
  177. recovered = client.post("/v1/tasks/execute", json=payload)
  178. replay = client.post("/v1/tasks/execute", json=payload)
  179. assert first.status_code == recovered.status_code == replay.status_code == 200
  180. assert recovered.get_json()["result"]["rows_out"] == 9
  181. assert replay.get_json() == recovered.get_json()
  182. assert executor.calls == 1
  183. def test_running_rule_without_terminal_evidence_returns_retryable_202():
  184. class RunningExecutor(Executor):
  185. def replay_task(self, **_context):
  186. return None
  187. issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
  188. executor = RunningExecutor()
  189. ledger = InMemoryTaskLedger()
  190. token = _rule_token(issuer)
  191. verifier = TaskTokenVerifier("x" * 32, clock=lambda: 1_000)
  192. claims = verifier.verify(token, node=RULE_NODE)
  193. ledger.claim(
  194. claims.jti,
  195. {
  196. "task_uid": claims.task_uid,
  197. "dataflow_uid": claims.dataflow_uid,
  198. "deployment_id": claims.deployment_id,
  199. "environment": claims.environment,
  200. "workflow_version": claims.workflow_version,
  201. "correlation_id": claims.correlation_id,
  202. "node_id": claims.node_id,
  203. "node_type": claims.node_type,
  204. "data_source_uid": None,
  205. "idempotency_key": RULE_NODE["idempotency"]["key"],
  206. },
  207. expires_at=claims.expires_at,
  208. )
  209. app = create_runner_app(
  210. verifier=verifier,
  211. ledger=ledger,
  212. registry=NodeRegistry({"rule.apply": executor}),
  213. )
  214. with app.test_client() as client:
  215. response = client.post(
  216. "/v1/tasks/execute",
  217. json={
  218. "task_token": token,
  219. "node": RULE_NODE,
  220. "parameters": {},
  221. },
  222. )
  223. assert response.status_code == 202
  224. assert response.headers["Retry-After"] == "2"
  225. assert executor.calls == 0
  226. def test_running_ledger_without_rule_run_expires_to_unknown():
  227. class StaleLedger(InMemoryTaskLedger):
  228. def reconcile_running(self, jti, binding):
  229. record = self.get(jti)
  230. assert record.binding == binding
  231. record.status = "unknown"
  232. record.commit_outcome = "unknown"
  233. return record
  234. class MissingEvidenceExecutor(Executor):
  235. def replay_task(self, **_context):
  236. return {"state": "missing"}
  237. issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
  238. verifier = TaskTokenVerifier("x" * 32, clock=lambda: 1_000)
  239. token = _rule_token(issuer)
  240. claims = verifier.verify(token, node=RULE_NODE)
  241. ledger = StaleLedger()
  242. binding = {
  243. "task_uid": claims.task_uid,
  244. "dataflow_uid": claims.dataflow_uid,
  245. "deployment_id": claims.deployment_id,
  246. "environment": claims.environment,
  247. "workflow_version": claims.workflow_version,
  248. "correlation_id": claims.correlation_id,
  249. "node_id": claims.node_id,
  250. "node_type": claims.node_type,
  251. "data_source_uid": None,
  252. "idempotency_key": RULE_NODE["idempotency"]["key"],
  253. }
  254. ledger.claim(claims.jti, binding, expires_at=claims.expires_at)
  255. app = create_runner_app(
  256. verifier=verifier,
  257. ledger=ledger,
  258. registry=NodeRegistry(
  259. {"rule.apply": MissingEvidenceExecutor()}
  260. ),
  261. )
  262. with app.test_client() as client:
  263. response = client.post(
  264. "/v1/tasks/execute",
  265. json={
  266. "task_token": token,
  267. "node": RULE_NODE,
  268. "parameters": {},
  269. },
  270. )
  271. assert response.status_code == 409
  272. assert response.get_json() == {
  273. "error": "task execution outcome is unknown"
  274. }
  275. assert ledger.get(claims.jti).status == "unknown"
  276. def test_expired_rule_run_is_reconciled_and_ledger_finalized_unknown():
  277. class ExpiredEvidenceExecutor(Executor):
  278. def replay_task(self, **_context):
  279. return {
  280. "state": "terminal",
  281. "status": "unknown",
  282. "commit_outcome": "unknown",
  283. }
  284. issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
  285. verifier = TaskTokenVerifier("x" * 32, clock=lambda: 1_000)
  286. token = _rule_token(issuer)
  287. claims = verifier.verify(token, node=RULE_NODE)
  288. ledger = InMemoryTaskLedger()
  289. ledger.claim(
  290. claims.jti,
  291. {
  292. "task_uid": claims.task_uid,
  293. "dataflow_uid": claims.dataflow_uid,
  294. "deployment_id": claims.deployment_id,
  295. "environment": claims.environment,
  296. "workflow_version": claims.workflow_version,
  297. "correlation_id": claims.correlation_id,
  298. "node_id": claims.node_id,
  299. "node_type": claims.node_type,
  300. "data_source_uid": None,
  301. "idempotency_key": RULE_NODE["idempotency"]["key"],
  302. },
  303. expires_at=claims.expires_at,
  304. )
  305. app = create_runner_app(
  306. verifier=verifier,
  307. ledger=ledger,
  308. registry=NodeRegistry(
  309. {"rule.apply": ExpiredEvidenceExecutor()}
  310. ),
  311. )
  312. with app.test_client() as client:
  313. response = client.post(
  314. "/v1/tasks/execute",
  315. json={
  316. "task_token": token,
  317. "node": RULE_NODE,
  318. "parameters": {},
  319. },
  320. )
  321. assert response.status_code == 409
  322. assert ledger.get(claims.jti).status == "unknown"
  323. def test_expired_token_is_replay_only_and_never_starts_execution():
  324. now = [1_000]
  325. issuer = TaskTokenIssuer(
  326. "x" * 32,
  327. clock=lambda: now[0],
  328. ttl_seconds=10,
  329. )
  330. verifier = TaskTokenVerifier("x" * 32, clock=lambda: now[0])
  331. token = _rule_token(issuer)
  332. ledger = InMemoryTaskLedger()
  333. executor = Executor()
  334. app = create_runner_app(
  335. verifier=verifier,
  336. ledger=ledger,
  337. registry=NodeRegistry({"rule.apply": executor}),
  338. )
  339. payload = {
  340. "task_token": token,
  341. "node": RULE_NODE,
  342. "parameters": {},
  343. }
  344. with app.test_client() as client:
  345. first = client.post("/v1/tasks/execute", json=payload)
  346. now[0] = 1_011
  347. replay = client.post("/v1/tasks/execute", json=payload)
  348. assert first.status_code == replay.status_code == 200
  349. assert replay.headers["X-Idempotent-Replay"] == "true"
  350. assert executor.calls == 1
  351. fresh_token = _rule_token(issuer)
  352. now[0] = 1_022
  353. with app.test_client() as client:
  354. rejected = client.post(
  355. "/v1/tasks/execute",
  356. json={
  357. "task_token": fresh_token,
  358. "node": RULE_NODE,
  359. "parameters": {},
  360. },
  361. )
  362. assert rejected.status_code == 401
  363. assert executor.calls == 1
  364. now[0] = 1_030
  365. running_token = _rule_token(issuer)
  366. running_claims = verifier.verify(
  367. running_token,
  368. node=RULE_NODE,
  369. )
  370. ledger.claim(
  371. running_claims.jti,
  372. {
  373. "task_uid": running_claims.task_uid,
  374. "dataflow_uid": running_claims.dataflow_uid,
  375. "deployment_id": running_claims.deployment_id,
  376. "environment": running_claims.environment,
  377. "workflow_version": running_claims.workflow_version,
  378. "correlation_id": running_claims.correlation_id,
  379. "node_id": running_claims.node_id,
  380. "node_type": running_claims.node_type,
  381. "data_source_uid": None,
  382. "idempotency_key": RULE_NODE["idempotency"]["key"],
  383. },
  384. expires_at=running_claims.expires_at,
  385. )
  386. now[0] = 1_041
  387. header, body, signature = running_token.split(".")
  388. changed_signature = (
  389. ("A" if signature[0] != "A" else "B") + signature[1:]
  390. )
  391. with app.test_client() as client:
  392. running = client.post(
  393. "/v1/tasks/execute",
  394. json={
  395. "task_token": running_token,
  396. "node": RULE_NODE,
  397. "parameters": {},
  398. },
  399. )
  400. tampered = client.post(
  401. "/v1/tasks/execute",
  402. json={
  403. "task_token": ".".join(
  404. (header, body, changed_signature)
  405. ),
  406. "node": RULE_NODE,
  407. "parameters": {},
  408. },
  409. )
  410. assert running.status_code == 401
  411. assert running.get_json() == {
  412. "error": "expired task is not replayable"
  413. }
  414. assert tampered.status_code == 401
  415. assert executor.calls == 1