test_dataflow_draft_reservation_postgres.py 4.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118
  1. from __future__ import annotations
  2. import os
  3. from concurrent.futures import ThreadPoolExecutor
  4. import pytest
  5. from sqlalchemy import create_engine, text
  6. from sqlalchemy.exc import IntegrityError
  7. from sqlalchemy.orm import Session
  8. from app.core.common.identifiers import new_governance_uid
  9. from app.core.data_rules.repository import DataRuleRepository
  10. def _consume(url, receipt, actor):
  11. engine = create_engine(url)
  12. try:
  13. with Session(engine) as session:
  14. try:
  15. DataRuleRepository(session).consume_dataflow_draft(
  16. receipt, actor_uid=actor
  17. )
  18. session.commit()
  19. return "accepted"
  20. except ValueError:
  21. session.rollback()
  22. return "rejected"
  23. finally:
  24. engine.dispose()
  25. def test_real_postgres_reservation_fk_expiry_and_concurrent_consume():
  26. url = os.environ.get("DATA_RULE_POSTGRES_ACCEPTANCE_URL")
  27. if not url:
  28. pytest.skip("real PostgreSQL acceptance URL is not configured")
  29. engine = create_engine(url)
  30. created_ids = []
  31. actor = new_governance_uid()
  32. try:
  33. with Session(engine) as session:
  34. session.execute(
  35. text(
  36. "INSERT INTO public.users "
  37. "(id, username, display_name, password_hash, status) "
  38. "VALUES (CAST(:id AS uuid), :username, "
  39. "'Reservation Acceptance', 'not-a-login-hash', 'active')"
  40. ),
  41. {
  42. "id": actor,
  43. "username": f"reservation-{actor}",
  44. },
  45. )
  46. session.commit()
  47. receipt = DataRuleRepository(session).reserve_dataflow_draft(
  48. actor_uid=actor
  49. )
  50. created_ids.append(receipt["reservation_id"])
  51. session.commit()
  52. closed_receipt = {
  53. key: receipt[key]
  54. for key in ("reservation_id", "dataflow_uid", "nonce")
  55. }
  56. with ThreadPoolExecutor(max_workers=2) as pool:
  57. outcomes = list(
  58. pool.map(
  59. lambda _index: _consume(url, closed_receipt, actor),
  60. range(2),
  61. )
  62. )
  63. assert sorted(outcomes) == ["accepted", "rejected"]
  64. assert _consume(url, closed_receipt, new_governance_uid()) == "rejected"
  65. with Session(engine) as session:
  66. expired = DataRuleRepository(session).reserve_dataflow_draft(
  67. actor_uid=actor
  68. )
  69. created_ids.append(expired["reservation_id"])
  70. session.execute(
  71. text(
  72. "UPDATE public.dataflow_draft_reservations "
  73. "SET created_at = CURRENT_TIMESTAMP - INTERVAL '2 seconds', "
  74. "expires_at = CURRENT_TIMESTAMP - INTERVAL '1 second' "
  75. "WHERE id = CAST(:id AS uuid)"
  76. ),
  77. {"id": expired["reservation_id"]},
  78. )
  79. session.commit()
  80. expired_receipt = {
  81. key: expired[key]
  82. for key in ("reservation_id", "dataflow_uid", "nonce")
  83. }
  84. assert _consume(url, expired_receipt, actor) == "rejected"
  85. with Session(engine) as session:
  86. with pytest.raises(IntegrityError):
  87. DataRuleRepository(session).reserve_dataflow_draft(
  88. actor_uid=new_governance_uid()
  89. )
  90. session.commit()
  91. session.rollback()
  92. finally:
  93. with Session(engine) as session:
  94. session.execute(
  95. text(
  96. "DELETE FROM public.dataflow_draft_reservations "
  97. "WHERE id = ANY(CAST(:ids AS uuid[]))"
  98. ),
  99. {"ids": created_ids},
  100. )
  101. session.execute(
  102. text(
  103. "DELETE FROM public.users WHERE id = CAST(:id AS uuid)"
  104. ),
  105. {"id": actor},
  106. )
  107. session.commit()
  108. engine.dispose()