test_api.py 8.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261
  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