test_api.py 21 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660
  1. import hashlib
  2. import json
  3. from app.runner.api import create_runner_app
  4. from app.runner.auth import TaskTokenIssuer, TaskTokenVerifier
  5. from app.runner.ledger import InMemoryTaskLedger
  6. from app.runner.nodes import NodeRegistry
  7. NODE = {
  8. "id": "read_orders",
  9. "type": "sql.query",
  10. "data_source_uid": "01900000-0000-7000-8000-000000000010",
  11. "purpose": "read",
  12. "config": {"statement": "SELECT 1", "parameters": {}},
  13. }
  14. RULE_NODE = {
  15. "id": "apply_orders",
  16. "type": "rule.apply",
  17. "purpose": "write",
  18. "idempotency": {
  19. "strategy": "upsert",
  20. "key": "orders:2026-07-23",
  21. },
  22. "config": {
  23. "component_binding_id": "01900000-0000-7000-8000-000000000021",
  24. "rule_version_id": "01900000-0000-7000-8000-000000000022",
  25. "execution_plan_hash": "a" * 64,
  26. },
  27. }
  28. class Executor:
  29. def __init__(self):
  30. self.calls = 0
  31. self.context = None
  32. def execute(self, node, parameters, **kwargs):
  33. self.calls += 1
  34. self.context = kwargs
  35. return {
  36. "node": node["id"],
  37. "parameters": parameters,
  38. "output_artifact": "minio://dataops-rules/rules/output.parquet",
  39. }
  40. def _digest(value):
  41. encoded = json.dumps(
  42. value,
  43. sort_keys=True,
  44. separators=(",", ":"),
  45. ensure_ascii=False,
  46. ).encode("utf-8")
  47. return hashlib.sha256(encoded).hexdigest()
  48. def test_runner_accepts_one_signed_task_and_rejects_replay():
  49. issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
  50. verifier = TaskTokenVerifier("x" * 32, clock=lambda: 1_000)
  51. executor = Executor()
  52. ledger = InMemoryTaskLedger()
  53. app = create_runner_app(
  54. verifier=verifier,
  55. ledger=ledger,
  56. registry=NodeRegistry({"sql.query": executor}),
  57. )
  58. token = issuer.issue(
  59. task_uid="01900000-0000-7000-8000-000000000011",
  60. dataflow_uid="01900000-0000-7000-8000-000000000012",
  61. deployment_id="01900000-0000-7000-8000-000000000014",
  62. environment="test",
  63. workflow_version=7,
  64. correlation_id="01900000-0000-7000-8000-000000000013",
  65. node=NODE,
  66. )
  67. payload = {"task_token": token, "node": NODE, "parameters": {"day": "today"}}
  68. with app.test_client() as client:
  69. first = client.post("/v1/tasks/execute", json=payload)
  70. replay = client.post("/v1/tasks/execute", json=payload)
  71. assert first.status_code == 200
  72. assert first.get_json()["result"]["node"] == "read_orders"
  73. assert first.get_json()["output_artifact"].startswith("minio://")
  74. assert executor.context["dataflow_uid"] == (
  75. "01900000-0000-7000-8000-000000000012"
  76. )
  77. assert executor.context["workflow_version"] == 7
  78. assert executor.context["node_id"] == "read_orders"
  79. assert replay.status_code == 409
  80. assert replay.get_json() == {"error": "task token already consumed"}
  81. assert executor.calls == 1
  82. def test_runner_health_has_no_secret_or_datasource_details():
  83. app = create_runner_app(
  84. verifier=TaskTokenVerifier("x" * 32),
  85. ledger=InMemoryTaskLedger(),
  86. registry=NodeRegistry({}),
  87. )
  88. with app.test_client() as client:
  89. response = client.get("/health")
  90. assert response.status_code == 200
  91. assert response.get_json() == {"status": "ok"}
  92. def test_runner_fails_closed_when_the_durable_ledger_is_unavailable():
  93. class UnavailableLedger:
  94. def claim(self, *_args, **_kwargs):
  95. raise RuntimeError("postgres password must never leak")
  96. issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
  97. app = create_runner_app(
  98. verifier=TaskTokenVerifier("x" * 32, clock=lambda: 1_000),
  99. ledger=UnavailableLedger(),
  100. registry=NodeRegistry({"sql.query": Executor()}),
  101. )
  102. token = issuer.issue(
  103. task_uid="01900000-0000-7000-8000-000000000011",
  104. dataflow_uid="01900000-0000-7000-8000-000000000012",
  105. deployment_id="01900000-0000-7000-8000-000000000014",
  106. environment="test",
  107. workflow_version=7,
  108. correlation_id="01900000-0000-7000-8000-000000000013",
  109. node=NODE,
  110. )
  111. with app.test_client() as client:
  112. response = client.post(
  113. "/v1/tasks/execute",
  114. json={
  115. "task_token": token,
  116. "node": NODE,
  117. "parameters": {},
  118. },
  119. )
  120. assert response.status_code == 503
  121. assert response.get_json() == {"error": "task ledger is unavailable"}
  122. assert "password" not in response.get_data(as_text=True).lower()
  123. def _rule_token(issuer):
  124. return issuer.issue(
  125. task_uid="01900000-0000-7000-8000-000000000031",
  126. dataflow_uid="01900000-0000-7000-8000-000000000032",
  127. deployment_id="01900000-0000-7000-8000-000000000033",
  128. environment="production",
  129. workflow_version=7,
  130. correlation_id="01900000-0000-7000-8000-000000000034",
  131. node=RULE_NODE,
  132. write_authorized=True,
  133. )
  134. def test_governed_rule_same_jti_replays_durable_response_once():
  135. issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
  136. executor = Executor()
  137. ledger = InMemoryTaskLedger()
  138. app = create_runner_app(
  139. verifier=TaskTokenVerifier("x" * 32, clock=lambda: 1_000),
  140. ledger=ledger,
  141. registry=NodeRegistry({"rule.apply": executor}),
  142. )
  143. payload = {
  144. "task_token": _rule_token(issuer),
  145. "node": RULE_NODE,
  146. "parameters": {},
  147. }
  148. with app.test_client() as client:
  149. first = client.post("/v1/tasks/execute", json=payload)
  150. replay = client.post("/v1/tasks/execute", json=payload)
  151. assert first.status_code == replay.status_code == 200
  152. assert first.get_json() == replay.get_json()
  153. assert replay.headers["X-Idempotent-Replay"] == "true"
  154. assert executor.calls == 1
  155. def test_governed_rule_recovers_terminal_evidence_after_response_loss():
  156. class LostFirstFinishLedger(InMemoryTaskLedger):
  157. def __init__(self):
  158. super().__init__()
  159. self.finish_calls = 0
  160. def finish(self, jti, **outcome):
  161. self.finish_calls += 1
  162. if self.finish_calls == 1:
  163. return
  164. return super().finish(jti, **outcome)
  165. class RecoverableExecutor(Executor):
  166. def replay_task(self, **_context):
  167. result = {
  168. "node": RULE_NODE["id"],
  169. "parameters": {},
  170. "rows_out": 9,
  171. "commit_outcome": "committed",
  172. }
  173. return {
  174. "state": "terminal",
  175. "status": "success",
  176. "commit_outcome": "committed",
  177. "result": result,
  178. "result_digest": _digest(result),
  179. "evidence_digest": "b" * 64,
  180. }
  181. issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
  182. executor = RecoverableExecutor()
  183. ledger = LostFirstFinishLedger()
  184. app = create_runner_app(
  185. verifier=TaskTokenVerifier("x" * 32, clock=lambda: 1_000),
  186. ledger=ledger,
  187. registry=NodeRegistry({"rule.apply": executor}),
  188. )
  189. payload = {
  190. "task_token": _rule_token(issuer),
  191. "node": RULE_NODE,
  192. "parameters": {},
  193. }
  194. with app.test_client() as client:
  195. first = client.post("/v1/tasks/execute", json=payload)
  196. recovered = client.post("/v1/tasks/execute", json=payload)
  197. replay = client.post("/v1/tasks/execute", json=payload)
  198. assert first.status_code == recovered.status_code == replay.status_code == 200
  199. assert recovered.get_json()["result"]["rows_out"] == 9
  200. assert replay.get_json() == recovered.get_json()
  201. assert executor.calls == 1
  202. def test_running_rule_without_terminal_evidence_returns_retryable_202():
  203. class RunningExecutor(Executor):
  204. def replay_task(self, **_context):
  205. return None
  206. issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
  207. executor = RunningExecutor()
  208. ledger = InMemoryTaskLedger()
  209. token = _rule_token(issuer)
  210. verifier = TaskTokenVerifier("x" * 32, clock=lambda: 1_000)
  211. claims = verifier.verify(token, node=RULE_NODE)
  212. ledger.claim(
  213. claims.jti,
  214. {
  215. "task_uid": claims.task_uid,
  216. "dataflow_uid": claims.dataflow_uid,
  217. "deployment_id": claims.deployment_id,
  218. "environment": claims.environment,
  219. "workflow_version": claims.workflow_version,
  220. "correlation_id": claims.correlation_id,
  221. "node_id": claims.node_id,
  222. "node_type": claims.node_type,
  223. "data_source_uid": None,
  224. "idempotency_key": RULE_NODE["idempotency"]["key"],
  225. },
  226. expires_at=claims.expires_at,
  227. )
  228. app = create_runner_app(
  229. verifier=verifier,
  230. ledger=ledger,
  231. registry=NodeRegistry({"rule.apply": executor}),
  232. )
  233. with app.test_client() as client:
  234. response = client.post(
  235. "/v1/tasks/execute",
  236. json={
  237. "task_token": token,
  238. "node": RULE_NODE,
  239. "parameters": {},
  240. },
  241. )
  242. assert response.status_code == 202
  243. assert response.headers["Retry-After"] == "2"
  244. assert executor.calls == 0
  245. def test_explicit_running_evidence_never_finalizes_ledger_success():
  246. class RunningExecutor(Executor):
  247. def replay_task(self, **_context):
  248. return {"state": "running"}
  249. issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
  250. verifier = TaskTokenVerifier("x" * 32, clock=lambda: 1_000)
  251. token = _rule_token(issuer)
  252. claims = verifier.verify(token, node=RULE_NODE)
  253. ledger = InMemoryTaskLedger()
  254. binding = {
  255. "task_uid": claims.task_uid,
  256. "dataflow_uid": claims.dataflow_uid,
  257. "deployment_id": claims.deployment_id,
  258. "environment": claims.environment,
  259. "workflow_version": claims.workflow_version,
  260. "correlation_id": claims.correlation_id,
  261. "node_id": claims.node_id,
  262. "node_type": claims.node_type,
  263. "data_source_uid": None,
  264. "idempotency_key": RULE_NODE["idempotency"]["key"],
  265. }
  266. ledger.claim(claims.jti, binding, expires_at=claims.expires_at)
  267. executor = RunningExecutor()
  268. app = create_runner_app(
  269. verifier=verifier,
  270. ledger=ledger,
  271. registry=NodeRegistry({"rule.apply": executor}),
  272. )
  273. with app.test_client() as client:
  274. response = client.post(
  275. "/v1/tasks/execute",
  276. json={
  277. "task_token": token,
  278. "node": RULE_NODE,
  279. "parameters": {},
  280. },
  281. )
  282. assert response.status_code == 202
  283. assert ledger.get(claims.jti).status == "running"
  284. assert ledger.get(claims.jti).replay_body is None
  285. assert executor.calls == 0
  286. def test_terminal_success_requires_matching_result_digest():
  287. class MismatchedExecutor(Executor):
  288. def replay_task(self, **_context):
  289. result = {
  290. "rows_out": 9,
  291. "commit_outcome": "committed",
  292. }
  293. return {
  294. "state": "terminal",
  295. "status": "success",
  296. "commit_outcome": "committed",
  297. "result": result,
  298. "result_digest": "0" * 64,
  299. "evidence_digest": "b" * 64,
  300. }
  301. issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
  302. verifier = TaskTokenVerifier("x" * 32, clock=lambda: 1_000)
  303. token = _rule_token(issuer)
  304. claims = verifier.verify(token, node=RULE_NODE)
  305. ledger = InMemoryTaskLedger()
  306. binding = {
  307. "task_uid": claims.task_uid,
  308. "dataflow_uid": claims.dataflow_uid,
  309. "deployment_id": claims.deployment_id,
  310. "environment": claims.environment,
  311. "workflow_version": claims.workflow_version,
  312. "correlation_id": claims.correlation_id,
  313. "node_id": claims.node_id,
  314. "node_type": claims.node_type,
  315. "data_source_uid": None,
  316. "idempotency_key": RULE_NODE["idempotency"]["key"],
  317. }
  318. ledger.claim(claims.jti, binding, expires_at=claims.expires_at)
  319. app = create_runner_app(
  320. verifier=verifier,
  321. ledger=ledger,
  322. registry=NodeRegistry(
  323. {"rule.apply": MismatchedExecutor()}
  324. ),
  325. )
  326. with app.test_client() as client:
  327. response = client.post(
  328. "/v1/tasks/execute",
  329. json={
  330. "task_token": token,
  331. "node": RULE_NODE,
  332. "parameters": {},
  333. },
  334. )
  335. assert response.status_code == 503
  336. assert ledger.get(claims.jti).status == "running"
  337. assert ledger.get(claims.jti).replay_body is None
  338. def test_terminal_replay_rejects_unclosed_state_fields():
  339. class UnclosedExecutor(Executor):
  340. def replay_task(self, **_context):
  341. result = {
  342. "rows_out": 9,
  343. "commit_outcome": "committed",
  344. }
  345. return {
  346. "state": "terminal",
  347. "status": "success",
  348. "commit_outcome": "committed",
  349. "result": result,
  350. "result_digest": _digest(result),
  351. "evidence_digest": "b" * 64,
  352. "raw_internal_row": {"secret": "must-not-pass"},
  353. }
  354. issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
  355. verifier = TaskTokenVerifier("x" * 32, clock=lambda: 1_000)
  356. token = _rule_token(issuer)
  357. claims = verifier.verify(token, node=RULE_NODE)
  358. ledger = InMemoryTaskLedger()
  359. ledger.claim(
  360. claims.jti,
  361. {
  362. "task_uid": claims.task_uid,
  363. "dataflow_uid": claims.dataflow_uid,
  364. "deployment_id": claims.deployment_id,
  365. "environment": claims.environment,
  366. "workflow_version": claims.workflow_version,
  367. "correlation_id": claims.correlation_id,
  368. "node_id": claims.node_id,
  369. "node_type": claims.node_type,
  370. "data_source_uid": None,
  371. "idempotency_key": RULE_NODE["idempotency"]["key"],
  372. },
  373. expires_at=claims.expires_at,
  374. )
  375. app = create_runner_app(
  376. verifier=verifier,
  377. ledger=ledger,
  378. registry=NodeRegistry(
  379. {"rule.apply": UnclosedExecutor()}
  380. ),
  381. )
  382. with app.test_client() as client:
  383. response = client.post(
  384. "/v1/tasks/execute",
  385. json={
  386. "task_token": token,
  387. "node": RULE_NODE,
  388. "parameters": {},
  389. },
  390. )
  391. assert response.status_code == 503
  392. assert ledger.get(claims.jti).status == "running"
  393. def test_running_ledger_without_rule_run_expires_to_unknown():
  394. class StaleLedger(InMemoryTaskLedger):
  395. def reconcile_running(self, jti, binding):
  396. record = self.get(jti)
  397. assert record.binding == binding
  398. record.status = "unknown"
  399. record.commit_outcome = "unknown"
  400. return record
  401. class MissingEvidenceExecutor(Executor):
  402. def replay_task(self, **_context):
  403. return {"state": "missing"}
  404. issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
  405. verifier = TaskTokenVerifier("x" * 32, clock=lambda: 1_000)
  406. token = _rule_token(issuer)
  407. claims = verifier.verify(token, node=RULE_NODE)
  408. ledger = StaleLedger()
  409. binding = {
  410. "task_uid": claims.task_uid,
  411. "dataflow_uid": claims.dataflow_uid,
  412. "deployment_id": claims.deployment_id,
  413. "environment": claims.environment,
  414. "workflow_version": claims.workflow_version,
  415. "correlation_id": claims.correlation_id,
  416. "node_id": claims.node_id,
  417. "node_type": claims.node_type,
  418. "data_source_uid": None,
  419. "idempotency_key": RULE_NODE["idempotency"]["key"],
  420. }
  421. ledger.claim(claims.jti, binding, expires_at=claims.expires_at)
  422. app = create_runner_app(
  423. verifier=verifier,
  424. ledger=ledger,
  425. registry=NodeRegistry(
  426. {"rule.apply": MissingEvidenceExecutor()}
  427. ),
  428. )
  429. with app.test_client() as client:
  430. response = client.post(
  431. "/v1/tasks/execute",
  432. json={
  433. "task_token": token,
  434. "node": RULE_NODE,
  435. "parameters": {},
  436. },
  437. )
  438. assert response.status_code == 409
  439. assert response.get_json() == {
  440. "error": "task execution outcome is unknown"
  441. }
  442. assert ledger.get(claims.jti).status == "unknown"
  443. def test_expired_rule_run_is_reconciled_and_ledger_finalized_unknown():
  444. class ExpiredEvidenceExecutor(Executor):
  445. def replay_task(self, **_context):
  446. result = {"commit_outcome": "unknown"}
  447. return {
  448. "state": "terminal",
  449. "status": "unknown",
  450. "commit_outcome": "unknown",
  451. "result": result,
  452. "result_digest": _digest(result),
  453. "evidence_digest": "",
  454. }
  455. issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
  456. verifier = TaskTokenVerifier("x" * 32, clock=lambda: 1_000)
  457. token = _rule_token(issuer)
  458. claims = verifier.verify(token, node=RULE_NODE)
  459. ledger = InMemoryTaskLedger()
  460. ledger.claim(
  461. claims.jti,
  462. {
  463. "task_uid": claims.task_uid,
  464. "dataflow_uid": claims.dataflow_uid,
  465. "deployment_id": claims.deployment_id,
  466. "environment": claims.environment,
  467. "workflow_version": claims.workflow_version,
  468. "correlation_id": claims.correlation_id,
  469. "node_id": claims.node_id,
  470. "node_type": claims.node_type,
  471. "data_source_uid": None,
  472. "idempotency_key": RULE_NODE["idempotency"]["key"],
  473. },
  474. expires_at=claims.expires_at,
  475. )
  476. app = create_runner_app(
  477. verifier=verifier,
  478. ledger=ledger,
  479. registry=NodeRegistry(
  480. {"rule.apply": ExpiredEvidenceExecutor()}
  481. ),
  482. )
  483. with app.test_client() as client:
  484. response = client.post(
  485. "/v1/tasks/execute",
  486. json={
  487. "task_token": token,
  488. "node": RULE_NODE,
  489. "parameters": {},
  490. },
  491. )
  492. assert response.status_code == 409
  493. assert ledger.get(claims.jti).status == "unknown"
  494. def test_expired_token_is_replay_only_and_never_starts_execution():
  495. class RunningReplayExecutor(Executor):
  496. def replay_task(self, **_context):
  497. return {"state": "running"}
  498. now = [1_000]
  499. issuer = TaskTokenIssuer(
  500. "x" * 32,
  501. clock=lambda: now[0],
  502. ttl_seconds=10,
  503. )
  504. verifier = TaskTokenVerifier("x" * 32, clock=lambda: now[0])
  505. token = _rule_token(issuer)
  506. ledger = InMemoryTaskLedger()
  507. executor = RunningReplayExecutor()
  508. app = create_runner_app(
  509. verifier=verifier,
  510. ledger=ledger,
  511. registry=NodeRegistry({"rule.apply": executor}),
  512. )
  513. payload = {
  514. "task_token": token,
  515. "node": RULE_NODE,
  516. "parameters": {},
  517. }
  518. with app.test_client() as client:
  519. first = client.post("/v1/tasks/execute", json=payload)
  520. now[0] = 1_011
  521. replay = client.post("/v1/tasks/execute", json=payload)
  522. assert first.status_code == replay.status_code == 200
  523. assert replay.headers["X-Idempotent-Replay"] == "true"
  524. assert executor.calls == 1
  525. fresh_token = _rule_token(issuer)
  526. now[0] = 1_022
  527. with app.test_client() as client:
  528. rejected = client.post(
  529. "/v1/tasks/execute",
  530. json={
  531. "task_token": fresh_token,
  532. "node": RULE_NODE,
  533. "parameters": {},
  534. },
  535. )
  536. assert rejected.status_code == 401
  537. assert executor.calls == 1
  538. now[0] = 1_030
  539. running_token = _rule_token(issuer)
  540. running_claims = verifier.verify(
  541. running_token,
  542. node=RULE_NODE,
  543. )
  544. ledger.claim(
  545. running_claims.jti,
  546. {
  547. "task_uid": running_claims.task_uid,
  548. "dataflow_uid": running_claims.dataflow_uid,
  549. "deployment_id": running_claims.deployment_id,
  550. "environment": running_claims.environment,
  551. "workflow_version": running_claims.workflow_version,
  552. "correlation_id": running_claims.correlation_id,
  553. "node_id": running_claims.node_id,
  554. "node_type": running_claims.node_type,
  555. "data_source_uid": None,
  556. "idempotency_key": RULE_NODE["idempotency"]["key"],
  557. },
  558. expires_at=running_claims.expires_at,
  559. )
  560. now[0] = 1_041
  561. header, body, signature = running_token.split(".")
  562. changed_signature = (
  563. ("A" if signature[0] != "A" else "B") + signature[1:]
  564. )
  565. with app.test_client() as client:
  566. running = client.post(
  567. "/v1/tasks/execute",
  568. json={
  569. "task_token": running_token,
  570. "node": RULE_NODE,
  571. "parameters": {},
  572. },
  573. )
  574. tampered = client.post(
  575. "/v1/tasks/execute",
  576. json={
  577. "task_token": ".".join(
  578. (header, body, changed_signature)
  579. ),
  580. "node": RULE_NODE,
  581. "parameters": {},
  582. },
  583. )
  584. assert running.status_code == 401
  585. assert running.get_json() == {
  586. "error": "expired task is not replayable"
  587. }
  588. assert tampered.status_code == 401
  589. assert executor.calls == 1