| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169 |
- 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()
|