| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217 |
- """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")
|