| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100 |
- from app.runner.api import create_runner_app
- from app.runner.auth import TaskTokenIssuer, TaskTokenVerifier
- from app.runner.ledger import InMemoryTaskLedger
- from app.runner.nodes import NodeRegistry
- NODE = {
- "id": "read_orders",
- "type": "sql.query",
- "data_source_uid": "01900000-0000-7000-8000-000000000010",
- "purpose": "read",
- "config": {"statement": "SELECT 1", "parameters": {}},
- }
- class Executor:
- def __init__(self):
- self.calls = 0
- def execute(self, node, parameters, **_kwargs):
- self.calls += 1
- return {"node": node["id"], "parameters": parameters}
- def test_runner_accepts_one_signed_task_and_rejects_replay():
- issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
- verifier = TaskTokenVerifier("x" * 32, clock=lambda: 1_000)
- executor = Executor()
- ledger = InMemoryTaskLedger()
- app = create_runner_app(
- verifier=verifier,
- ledger=ledger,
- registry=NodeRegistry({"sql.query": executor}),
- )
- token = issuer.issue(
- task_uid="01900000-0000-7000-8000-000000000011",
- dataflow_uid="01900000-0000-7000-8000-000000000012",
- workflow_version=7,
- correlation_id="01900000-0000-7000-8000-000000000013",
- node=NODE,
- )
- payload = {"task_token": token, "node": NODE, "parameters": {"day": "today"}}
- with app.test_client() as client:
- first = client.post("/v1/tasks/execute", json=payload)
- replay = client.post("/v1/tasks/execute", json=payload)
- assert first.status_code == 200
- assert first.get_json()["result"]["node"] == "read_orders"
- assert replay.status_code == 409
- assert replay.get_json() == {"error": "task token already consumed"}
- assert executor.calls == 1
- def test_runner_health_has_no_secret_or_datasource_details():
- app = create_runner_app(
- verifier=TaskTokenVerifier("x" * 32),
- ledger=InMemoryTaskLedger(),
- registry=NodeRegistry({}),
- )
- with app.test_client() as client:
- response = client.get("/health")
- assert response.status_code == 200
- assert response.get_json() == {"status": "ok"}
- def test_runner_fails_closed_when_the_durable_ledger_is_unavailable():
- class UnavailableLedger:
- def claim(self, *_args, **_kwargs):
- raise RuntimeError("postgres password must never leak")
- issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
- app = create_runner_app(
- verifier=TaskTokenVerifier("x" * 32, clock=lambda: 1_000),
- ledger=UnavailableLedger(),
- registry=NodeRegistry({"sql.query": Executor()}),
- )
- token = issuer.issue(
- task_uid="01900000-0000-7000-8000-000000000011",
- dataflow_uid="01900000-0000-7000-8000-000000000012",
- workflow_version=7,
- correlation_id="01900000-0000-7000-8000-000000000013",
- node=NODE,
- )
- with app.test_client() as client:
- response = client.post(
- "/v1/tasks/execute",
- json={
- "task_token": token,
- "node": NODE,
- "parameters": {},
- },
- )
- assert response.status_code == 503
- assert response.get_json() == {"error": "task ledger is unavailable"}
- assert "password" not in response.get_data(as_text=True).lower()
|