test_postgres_ledger.py 2.7 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192
  1. from pathlib import Path
  2. from app.runner.ledger import PostgresTaskLedger
  3. class Result:
  4. rowcount = 1
  5. class Connection:
  6. def __init__(self):
  7. self.calls = []
  8. def execute(self, statement, parameters):
  9. self.calls.append((str(statement), parameters))
  10. return Result()
  11. class Begin:
  12. def __init__(self, connection):
  13. self.connection = connection
  14. def __enter__(self):
  15. return self.connection
  16. def __exit__(self, *_args):
  17. return False
  18. class Engine:
  19. def __init__(self):
  20. self.connection = Connection()
  21. def begin(self):
  22. return Begin(self.connection)
  23. def test_postgres_ledger_claim_is_atomic_and_persists_only_bounded_fields():
  24. engine = Engine()
  25. ledger = PostgresTaskLedger(engine)
  26. binding = {
  27. "task_uid": "01900000-0000-7000-8000-000000000011",
  28. "dataflow_uid": "01900000-0000-7000-8000-000000000012",
  29. "deployment_id": "01900000-0000-7000-8000-000000000016",
  30. "environment": "production",
  31. "workflow_version": 7,
  32. "correlation_id": "01900000-0000-7000-8000-000000000013",
  33. "node_id": "write_orders",
  34. "node_type": "sql.execute",
  35. "data_source_uid": "01900000-0000-7000-8000-000000000014",
  36. "idempotency_key": "orders:2026-07-19",
  37. }
  38. assert ledger.claim(
  39. "01900000-0000-7000-8000-000000000015",
  40. binding,
  41. expires_at=2_000,
  42. )
  43. sql, parameters = engine.connection.calls[0]
  44. assert "ON CONFLICT DO NOTHING" in sql
  45. assert "encrypted_payload" not in sql
  46. assert "password" not in sql
  47. assert parameters["idempotency_key"] == "orders:2026-07-19"
  48. assert parameters["deployment_id"] == binding["deployment_id"]
  49. assert parameters["environment"] == "production"
  50. ledger.finish(
  51. "01900000-0000-7000-8000-000000000015",
  52. status="success",
  53. replay_http_status=200,
  54. replay_body={"result": {"rows_out": 3}},
  55. )
  56. finish_sql, finish_parameters = engine.connection.calls[1]
  57. assert "replay_digest" in finish_sql
  58. assert finish_parameters["replay_http_status"] == 200
  59. assert len(finish_parameters["replay_digest"]) == 64
  60. def test_v52_migration_has_single_use_and_commit_outcome_constraints():
  61. root = Path(__file__).resolve().parents[2]
  62. source = (
  63. root
  64. / "migrations/versions/20260719_70_runner_task_ledger.py"
  65. ).read_text(encoding="utf-8")
  66. assert "runner_task_executions" in source
  67. assert "token_jti UUID PRIMARY KEY" in source
  68. assert "UNIQUE (task_uid)" in source
  69. assert "'not_committed', 'committed', 'unknown'" in source
  70. assert "encrypted_payload" not in source
  71. assert "connection_string" not in source