| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261 |
- 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": {}},
- }
- RULE_NODE = {
- "id": "apply_orders",
- "type": "rule.apply",
- "purpose": "write",
- "idempotency": {
- "strategy": "upsert",
- "key": "orders:2026-07-23",
- },
- "config": {
- "component_binding_id": "01900000-0000-7000-8000-000000000021",
- "rule_version_id": "01900000-0000-7000-8000-000000000022",
- "execution_plan_hash": "a" * 64,
- },
- }
- class Executor:
- def __init__(self):
- self.calls = 0
- self.context = None
- def execute(self, node, parameters, **kwargs):
- self.calls += 1
- self.context = kwargs
- return {
- "node": node["id"],
- "parameters": parameters,
- "output_artifact": "minio://dataops-rules/rules/output.parquet",
- }
- 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",
- deployment_id="01900000-0000-7000-8000-000000000014",
- environment="test",
- 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 first.get_json()["output_artifact"].startswith("minio://")
- assert executor.context["dataflow_uid"] == (
- "01900000-0000-7000-8000-000000000012"
- )
- assert executor.context["workflow_version"] == 7
- assert executor.context["node_id"] == "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",
- deployment_id="01900000-0000-7000-8000-000000000014",
- environment="test",
- 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()
- def _rule_token(issuer):
- return issuer.issue(
- task_uid="01900000-0000-7000-8000-000000000031",
- dataflow_uid="01900000-0000-7000-8000-000000000032",
- deployment_id="01900000-0000-7000-8000-000000000033",
- environment="production",
- workflow_version=7,
- correlation_id="01900000-0000-7000-8000-000000000034",
- node=RULE_NODE,
- write_authorized=True,
- )
- def test_governed_rule_same_jti_replays_durable_response_once():
- issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
- executor = Executor()
- ledger = InMemoryTaskLedger()
- app = create_runner_app(
- verifier=TaskTokenVerifier("x" * 32, clock=lambda: 1_000),
- ledger=ledger,
- registry=NodeRegistry({"rule.apply": executor}),
- )
- payload = {
- "task_token": _rule_token(issuer),
- "node": RULE_NODE,
- "parameters": {},
- }
- 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 == replay.status_code == 200
- assert first.get_json() == replay.get_json()
- assert replay.headers["X-Idempotent-Replay"] == "true"
- assert executor.calls == 1
- def test_governed_rule_recovers_terminal_evidence_after_response_loss():
- class LostFirstFinishLedger(InMemoryTaskLedger):
- def __init__(self):
- super().__init__()
- self.finish_calls = 0
- def finish(self, jti, **outcome):
- self.finish_calls += 1
- if self.finish_calls == 1:
- return
- return super().finish(jti, **outcome)
- class RecoverableExecutor(Executor):
- def replay_task(self, **_context):
- return {
- "node": RULE_NODE["id"],
- "parameters": {},
- "rows_out": 9,
- }
- issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
- executor = RecoverableExecutor()
- ledger = LostFirstFinishLedger()
- app = create_runner_app(
- verifier=TaskTokenVerifier("x" * 32, clock=lambda: 1_000),
- ledger=ledger,
- registry=NodeRegistry({"rule.apply": executor}),
- )
- payload = {
- "task_token": _rule_token(issuer),
- "node": RULE_NODE,
- "parameters": {},
- }
- with app.test_client() as client:
- first = client.post("/v1/tasks/execute", json=payload)
- recovered = client.post("/v1/tasks/execute", json=payload)
- replay = client.post("/v1/tasks/execute", json=payload)
- assert first.status_code == recovered.status_code == replay.status_code == 200
- assert recovered.get_json()["result"]["rows_out"] == 9
- assert replay.get_json() == recovered.get_json()
- assert executor.calls == 1
- def test_running_rule_without_terminal_evidence_returns_retryable_202():
- class RunningExecutor(Executor):
- def replay_task(self, **_context):
- return None
- issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000)
- executor = RunningExecutor()
- ledger = InMemoryTaskLedger()
- token = _rule_token(issuer)
- verifier = TaskTokenVerifier("x" * 32, clock=lambda: 1_000)
- claims = verifier.verify(token, node=RULE_NODE)
- ledger.claim(
- claims.jti,
- {
- "task_uid": claims.task_uid,
- "dataflow_uid": claims.dataflow_uid,
- "deployment_id": claims.deployment_id,
- "environment": claims.environment,
- "workflow_version": claims.workflow_version,
- "correlation_id": claims.correlation_id,
- "node_id": claims.node_id,
- "node_type": claims.node_type,
- "data_source_uid": None,
- "idempotency_key": RULE_NODE["idempotency"]["key"],
- },
- expires_at=claims.expires_at,
- )
- app = create_runner_app(
- verifier=verifier,
- ledger=ledger,
- registry=NodeRegistry({"rule.apply": executor}),
- )
- with app.test_client() as client:
- response = client.post(
- "/v1/tasks/execute",
- json={
- "task_token": token,
- "node": RULE_NODE,
- "parameters": {},
- },
- )
- assert response.status_code == 202
- assert response.headers["Retry-After"] == "2"
- assert executor.calls == 0
|