from __future__ import annotations import os from concurrent.futures import ThreadPoolExecutor import pytest from sqlalchemy import create_engine, text from sqlalchemy.exc import IntegrityError from sqlalchemy.orm import Session from app.core.common.identifiers import new_governance_uid from app.core.data_rules.repository import DataRuleRepository def _consume(url, receipt, actor): engine = create_engine(url) try: with Session(engine) as session: try: DataRuleRepository(session).consume_dataflow_draft( receipt, actor_uid=actor ) session.commit() return "accepted" except ValueError: session.rollback() return "rejected" finally: engine.dispose() def test_real_postgres_reservation_fk_expiry_and_concurrent_consume(): url = os.environ.get("DATA_RULE_POSTGRES_ACCEPTANCE_URL") if not url: pytest.skip("real PostgreSQL acceptance URL is not configured") engine = create_engine(url) created_ids = [] actor = new_governance_uid() try: with Session(engine) as session: session.execute( text( "INSERT INTO public.users " "(id, username, display_name, password_hash, status) " "VALUES (CAST(:id AS uuid), :username, " "'Reservation Acceptance', 'not-a-login-hash', 'active')" ), { "id": actor, "username": f"reservation-{actor}", }, ) session.commit() receipt = DataRuleRepository(session).reserve_dataflow_draft( actor_uid=actor ) created_ids.append(receipt["reservation_id"]) session.commit() closed_receipt = { key: receipt[key] for key in ("reservation_id", "dataflow_uid", "nonce") } with ThreadPoolExecutor(max_workers=2) as pool: outcomes = list( pool.map( lambda _index: _consume(url, closed_receipt, actor), range(2), ) ) assert sorted(outcomes) == ["accepted", "rejected"] assert _consume(url, closed_receipt, new_governance_uid()) == "rejected" with Session(engine) as session: expired = DataRuleRepository(session).reserve_dataflow_draft( actor_uid=actor ) created_ids.append(expired["reservation_id"]) session.execute( text( "UPDATE public.dataflow_draft_reservations " "SET created_at = CURRENT_TIMESTAMP - INTERVAL '2 seconds', " "expires_at = CURRENT_TIMESTAMP - INTERVAL '1 second' " "WHERE id = CAST(:id AS uuid)" ), {"id": expired["reservation_id"]}, ) session.commit() expired_receipt = { key: expired[key] for key in ("reservation_id", "dataflow_uid", "nonce") } assert _consume(url, expired_receipt, actor) == "rejected" with Session(engine) as session: with pytest.raises(IntegrityError): DataRuleRepository(session).reserve_dataflow_draft( actor_uid=new_governance_uid() ) session.commit() session.rollback() finally: with Session(engine) as session: session.execute( text( "DELETE FROM public.dataflow_draft_reservations " "WHERE id = ANY(CAST(:ids AS uuid[]))" ), {"ids": created_ids}, ) session.execute( text( "DELETE FROM public.users WHERE id = CAST(:id AS uuid)" ), {"id": actor}, ) session.commit() engine.dispose()