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 def test_running_ledger_without_rule_run_expires_to_unknown(): class StaleLedger(InMemoryTaskLedger): def reconcile_running(self, jti, binding): record = self.get(jti) assert record.binding == binding record.status = "unknown" record.commit_outcome = "unknown" return record class MissingEvidenceExecutor(Executor): def replay_task(self, **_context): return {"state": "missing"} issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000) verifier = TaskTokenVerifier("x" * 32, clock=lambda: 1_000) token = _rule_token(issuer) claims = verifier.verify(token, node=RULE_NODE) ledger = StaleLedger() binding = { "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"], } ledger.claim(claims.jti, binding, expires_at=claims.expires_at) app = create_runner_app( verifier=verifier, ledger=ledger, registry=NodeRegistry( {"rule.apply": MissingEvidenceExecutor()} ), ) with app.test_client() as client: response = client.post( "/v1/tasks/execute", json={ "task_token": token, "node": RULE_NODE, "parameters": {}, }, ) assert response.status_code == 409 assert response.get_json() == { "error": "task execution outcome is unknown" } assert ledger.get(claims.jti).status == "unknown" def test_expired_rule_run_is_reconciled_and_ledger_finalized_unknown(): class ExpiredEvidenceExecutor(Executor): def replay_task(self, **_context): return { "state": "terminal", "status": "unknown", "commit_outcome": "unknown", } issuer = TaskTokenIssuer("x" * 32, clock=lambda: 1_000) verifier = TaskTokenVerifier("x" * 32, clock=lambda: 1_000) token = _rule_token(issuer) claims = verifier.verify(token, node=RULE_NODE) ledger = InMemoryTaskLedger() 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": ExpiredEvidenceExecutor()} ), ) with app.test_client() as client: response = client.post( "/v1/tasks/execute", json={ "task_token": token, "node": RULE_NODE, "parameters": {}, }, ) assert response.status_code == 409 assert ledger.get(claims.jti).status == "unknown" def test_expired_token_is_replay_only_and_never_starts_execution(): now = [1_000] issuer = TaskTokenIssuer( "x" * 32, clock=lambda: now[0], ttl_seconds=10, ) verifier = TaskTokenVerifier("x" * 32, clock=lambda: now[0]) token = _rule_token(issuer) ledger = InMemoryTaskLedger() executor = Executor() app = create_runner_app( verifier=verifier, ledger=ledger, registry=NodeRegistry({"rule.apply": executor}), ) payload = { "task_token": token, "node": RULE_NODE, "parameters": {}, } with app.test_client() as client: first = client.post("/v1/tasks/execute", json=payload) now[0] = 1_011 replay = client.post("/v1/tasks/execute", json=payload) assert first.status_code == replay.status_code == 200 assert replay.headers["X-Idempotent-Replay"] == "true" assert executor.calls == 1 fresh_token = _rule_token(issuer) now[0] = 1_022 with app.test_client() as client: rejected = client.post( "/v1/tasks/execute", json={ "task_token": fresh_token, "node": RULE_NODE, "parameters": {}, }, ) assert rejected.status_code == 401 assert executor.calls == 1 now[0] = 1_030 running_token = _rule_token(issuer) running_claims = verifier.verify( running_token, node=RULE_NODE, ) ledger.claim( running_claims.jti, { "task_uid": running_claims.task_uid, "dataflow_uid": running_claims.dataflow_uid, "deployment_id": running_claims.deployment_id, "environment": running_claims.environment, "workflow_version": running_claims.workflow_version, "correlation_id": running_claims.correlation_id, "node_id": running_claims.node_id, "node_type": running_claims.node_type, "data_source_uid": None, "idempotency_key": RULE_NODE["idempotency"]["key"], }, expires_at=running_claims.expires_at, ) now[0] = 1_041 header, body, signature = running_token.split(".") changed_signature = ( ("A" if signature[0] != "A" else "B") + signature[1:] ) with app.test_client() as client: running = client.post( "/v1/tasks/execute", json={ "task_token": running_token, "node": RULE_NODE, "parameters": {}, }, ) tampered = client.post( "/v1/tasks/execute", json={ "task_token": ".".join( (header, body, changed_signature) ), "node": RULE_NODE, "parameters": {}, }, ) assert running.status_code == 401 assert running.get_json() == { "error": "expired task is not replayable" } assert tampered.status_code == 401 assert executor.calls == 1