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 _claim(url, receipt, actor): engine = create_engine(url) try: with Session(engine) as session: try: value = DataRuleRepository(session).begin_dataflow_create( receipt, actor_uid=actor ) session.commit() return value except ValueError as exc: session.rollback() return {"status": str(exc)} finally: engine.dispose() def test_real_postgres_saga_fk_expiry_concurrent_claim_replay_and_lease_recovery(): 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, " "'Saga Acceptance', 'not-a-login-hash', 'active')" ), {"id": actor, "username": f"saga-{actor}"}, ) session.commit() receipt = DataRuleRepository(session).reserve_dataflow_draft( actor_uid=actor ) created_ids.append(receipt["reservation_id"]) session.commit() closed = { key: receipt[key] for key in ("reservation_id", "dataflow_uid", "nonce") } with ThreadPoolExecutor(max_workers=2) as pool: outcomes = list( pool.map(lambda _index: _claim(url, closed, actor), range(2)) ) claimed = [value for value in outcomes if value["status"] == "claimed"] assert len(claimed) == 1 assert sorted(value["status"] for value in outcomes) == [ "claimed", "dataflow_create_in_progress", ] result = {"id": 73, "uid": receipt["dataflow_uid"]} with Session(engine) as session: repository = DataRuleRepository(session) assert repository.complete_dataflow_create( reservation_id=receipt["reservation_id"], dataflow_uid=receipt["dataflow_uid"], lease_token=claimed[0]["lease_token"], dataflow_node_id=73, result=result, ) == result session.commit() replay = _claim(url, closed, actor) assert replay["status"] == "completed" assert replay["result"] == result 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_closed = { key: expired[key] for key in ("reservation_id", "dataflow_uid", "nonce") } assert _claim(url, expired_closed, actor)["status"] == ( "draft_reservation_refresh_required" ) with Session(engine) as session: recoverable = DataRuleRepository(session).reserve_dataflow_draft( actor_uid=actor ) created_ids.append(recoverable["reservation_id"]) session.commit() first = DataRuleRepository(session).begin_dataflow_create( { key: recoverable[key] for key in ("reservation_id", "dataflow_uid", "nonce") }, actor_uid=actor, lease_seconds=10, ) session.commit() session.execute( text( "UPDATE public.dataflow_draft_reservations " "SET lease_expires_at = CURRENT_TIMESTAMP - INTERVAL '1 second' " "WHERE id = CAST(:id AS uuid)" ), {"id": recoverable["reservation_id"]}, ) session.commit() recovered = _claim( url, { key: recoverable[key] for key in ("reservation_id", "dataflow_uid", "nonce") }, actor, ) assert first["attempt"] == 1 assert recovered["status"] == "claimed" assert recovered["attempt"] == 2 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()