"""Opt-in Docker outage, recovery and credential-rotation acceptance test.""" import os import subprocess import threading from dataclasses import replace from pathlib import Path import pytest from sqlalchemy import create_engine, text from sqlalchemy.engine import make_url from app.core.data_source.adapters import adapter_for from app.core.data_source.errors import ( DataSourceCircuitOpen, DataSourceConnectionFailed, ) from app.core.data_source.manager import DataSourceConnectionManager from app.core.data_source.models import ( DataSourceCredential, DataSourceDefinition, ) from app.core.data_source.pool_registry import PoolRegistry ROOT = Path(__file__).resolve().parents[2] COMPOSE = ROOT / "deploy/docker/docker-compose.yml" POSTGRES_UID = "01900000-0000-7000-8000-000000000201" MYSQL_UID = "01900000-0000-7000-8000-000000000202" pytestmark = pytest.mark.integration class FakeClock: def __init__(self): self._now = 1000.0 self._lock = threading.Lock() def __call__(self): with self._lock: return self._now def advance(self, seconds): with self._lock: self._now += seconds class Definitions: def __init__(self, values): self.values = {item.uid: item for item in values} def get(self, uid): return self.values.get(uid) def rotate(self, uid): self.values[uid] = replace( self.values[uid], credential_version=self.values[uid].credential_version + 1, ) class Credentials: def __init__(self, values): self.values = values def get_active(self, _session, uid, _version): return self.values[uid] def _compose(*arguments): subprocess.run( ["docker", "compose", "-f", str(COMPOSE), *arguments], cwd=ROOT, check=True, text=True, capture_output=True, ) def _definition(uid, database_type, raw_url): url = make_url(raw_url) return DataSourceDefinition( uid=uid, name_en=f"outage-{database_type}", database_type=database_type, host=url.host, port=url.port, database=url.database, credential_ref=uid, credential_version=1, pool_size=2, max_overflow=3, ) def _credentials(raw_url): url = make_url(raw_url) return DataSourceCredential(url.username, url.password) def _build_manager(): postgres_url = os.environ.get( "TEST_SOURCE_POSTGRES_URL", "postgresql://source_reader:source-test-password" "@127.0.0.1:25432/acceptance", ) mysql_url = os.environ.get( "TEST_SOURCE_MYSQL_URL", "mysql+pymysql://source_reader:source-test-password" "@127.0.0.1:23306/acceptance", ) definitions = Definitions( [ _definition(POSTGRES_UID, "postgresql", postgres_url), _definition(MYSQL_UID, "mysql", mysql_url), ] ) credentials = Credentials( { POSTGRES_UID: _credentials(postgres_url), MYSQL_UID: _credentials(mysql_url), } ) clock = FakeClock() settings = { "pool_size": 2, "max_overflow": 3, "pool_timeout": 1, "pool_recycle": 1800, "idle_ttl": 900, "max_idle_pools": 20, "drain_timeout": 30, "query_timeout": 5, } registry = PoolRegistry(settings, clock=clock) manager = DataSourceConnectionManager( definitions=definitions, credentials=credentials, platform_session=lambda: object(), adapter_resolver=adapter_for, registry=registry, settings_resolver=lambda overrides: { **settings, **overrides, }, clock=clock, ) return definitions, clock, registry, manager def _query_count(manager, uid): with manager.connect(uid, purpose="dataflow_read") as connection: return connection.execute( text("SELECT COUNT(*) FROM acceptance_customers") ).scalar_one() def _platform_database_is_healthy(): url = os.environ.get( "TEST_DATABASE_URL", "postgresql://dataops:dataops-test-password" "@127.0.0.1:15432/dataops", ) engine = create_engine(url, pool_pre_ping=True) try: with engine.connect() as connection: return connection.execute(text("SELECT 1")).scalar_one() == 1 finally: engine.dispose() def test_mysql_outage_is_isolated_and_recovers_before_credential_rotation(): if os.environ.get("RUN_DATASOURCE_OUTAGE_TEST") != "1": pytest.skip("set RUN_DATASOURCE_OUTAGE_TEST=1 to control Docker") _compose("up", "-d", "--wait", "source-postgres", "source-mysql") definitions, clock, registry, manager = _build_manager() try: assert _query_count(manager, POSTGRES_UID) == 2 assert _query_count(manager, MYSQL_UID) == 2 old_mysql_engine = next( entry.engine for entry in registry._entries.values() if entry.key.data_source_uid == MYSQL_UID ) _compose("stop", "source-mysql") for _attempt in range(3): with pytest.raises(DataSourceConnectionFailed): _query_count(manager, MYSQL_UID) with pytest.raises(DataSourceCircuitOpen): _query_count(manager, MYSQL_UID) assert _query_count(manager, POSTGRES_UID) == 2 assert _platform_database_is_healthy() is True _compose("up", "-d", "--wait", "source-mysql") clock.advance(30) assert _query_count(manager, MYSQL_UID) == 2 mysql_status = next( item for item in registry.snapshot() if item.data_source_uid == MYSQL_UID ) assert mysql_status.pool_state == "healthy" definitions.rotate(MYSQL_UID) assert _query_count(manager, MYSQL_UID) == 2 assert old_mysql_engine._dataops_disposed is True rotated_status = next( item for item in registry.snapshot() if item.data_source_uid == MYSQL_UID ) assert rotated_status.credential_version == 2 finally: manager.close() _compose("up", "-d", "--wait", "source-mysql")