test_api.py 3.2 KB

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