| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467 |
- 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
|