test_api.py 3.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111
  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. class Executor:
  13. def __init__(self):
  14. self.calls = 0
  15. self.context = None
  16. def execute(self, node, parameters, **kwargs):
  17. self.calls += 1
  18. self.context = kwargs
  19. return {
  20. "node": node["id"],
  21. "parameters": parameters,
  22. "output_artifact": "minio://dataops-rules/rules/output.parquet",
  23. }
  24. def test_runner_accepts_one_signed_task_and_rejects_replay():
  25. issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
  26. verifier = TaskTokenVerifier("x" * 32, clock=lambda: 1_000)
  27. executor = Executor()
  28. ledger = InMemoryTaskLedger()
  29. app = create_runner_app(
  30. verifier=verifier,
  31. ledger=ledger,
  32. registry=NodeRegistry({"sql.query": executor}),
  33. )
  34. token = issuer.issue(
  35. task_uid="01900000-0000-7000-8000-000000000011",
  36. dataflow_uid="01900000-0000-7000-8000-000000000012",
  37. workflow_version=7,
  38. correlation_id="01900000-0000-7000-8000-000000000013",
  39. node=NODE,
  40. )
  41. payload = {"task_token": token, "node": NODE, "parameters": {"day": "today"}}
  42. with app.test_client() as client:
  43. first = client.post("/v1/tasks/execute", json=payload)
  44. replay = client.post("/v1/tasks/execute", json=payload)
  45. assert first.status_code == 200
  46. assert first.get_json()["result"]["node"] == "read_orders"
  47. assert first.get_json()["output_artifact"].startswith("minio://")
  48. assert executor.context["dataflow_uid"] == (
  49. "01900000-0000-7000-8000-000000000012"
  50. )
  51. assert executor.context["workflow_version"] == 7
  52. assert executor.context["node_id"] == "read_orders"
  53. assert replay.status_code == 409
  54. assert replay.get_json() == {"error": "task token already consumed"}
  55. assert executor.calls == 1
  56. def test_runner_health_has_no_secret_or_datasource_details():
  57. app = create_runner_app(
  58. verifier=TaskTokenVerifier("x" * 32),
  59. ledger=InMemoryTaskLedger(),
  60. registry=NodeRegistry({}),
  61. )
  62. with app.test_client() as client:
  63. response = client.get("/health")
  64. assert response.status_code == 200
  65. assert response.get_json() == {"status": "ok"}
  66. def test_runner_fails_closed_when_the_durable_ledger_is_unavailable():
  67. class UnavailableLedger:
  68. def claim(self, *_args, **_kwargs):
  69. raise RuntimeError("postgres password must never leak")
  70. issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
  71. app = create_runner_app(
  72. verifier=TaskTokenVerifier("x" * 32, clock=lambda: 1_000),
  73. ledger=UnavailableLedger(),
  74. registry=NodeRegistry({"sql.query": Executor()}),
  75. )
  76. token = issuer.issue(
  77. task_uid="01900000-0000-7000-8000-000000000011",
  78. dataflow_uid="01900000-0000-7000-8000-000000000012",
  79. workflow_version=7,
  80. correlation_id="01900000-0000-7000-8000-000000000013",
  81. node=NODE,
  82. )
  83. with app.test_client() as client:
  84. response = client.post(
  85. "/v1/tasks/execute",
  86. json={
  87. "task_token": token,
  88. "node": NODE,
  89. "parameters": {},
  90. },
  91. )
  92. assert response.status_code == 503
  93. assert response.get_json() == {"error": "task ledger is unavailable"}
  94. assert "password" not in response.get_data(as_text=True).lower()