| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179 |
- from __future__ import annotations
- import os
- import subprocess
- import uuid
- from pathlib import Path
- import psycopg2
- import pytest
- from sqlalchemy import create_engine, inspect
- from sqlalchemy.engine import make_url
- ROOT = Path(__file__).resolve().parents[1]
- EXPECTED_BASELINE_TABLES = {
- "data_orders",
- "data_products",
- "metadata_review_records",
- "metadata_version_history",
- "task_list",
- "users",
- }
- EXPECTED_UPGRADED_TABLES = {
- "datasource_credentials",
- "datasource_credential_audit_events",
- "workflow_schedules",
- "workflow_runs",
- "workflow_task_runs",
- "workflow_engine_bindings",
- "workflow_plan_audits",
- "runner_task_executions",
- "workflow_canary_evidence",
- "workflow_gateway_operations",
- "workflow_migration_states",
- "workflow_dual_runs",
- "workflow_reconciliation_reports",
- "workflow_cutover_operations",
- "ingestion_sources",
- "source_artifacts",
- "ingestion_jobs",
- "evidence_fragments",
- "extraction_candidates",
- "data_elements",
- "data_element_versions",
- "candidate_decisions",
- }
- def test_data_research_ingestion_migration_is_additive_and_constrained():
- migration = (
- ROOT
- / "migrations"
- / "versions"
- / "20260722_100_data_research_ingestion.py"
- ).read_text(encoding="utf-8")
- assert 'revision = "20260722_100"' in migration
- assert 'down_revision = "20260719_90"' in migration
- for table in (
- "ingestion_sources",
- "source_artifacts",
- "ingestion_jobs",
- "evidence_fragments",
- "extraction_candidates",
- ):
- assert f"CREATE TABLE IF NOT EXISTS public.{table}" in migration
- assert "idempotency_key" in migration
- assert "parser_version" in migration
- assert "locator JSONB" in migration
- assert "DROP TABLE" not in migration.upper()
- def test_data_element_migration_adds_versioned_governance_tables():
- migration = (
- ROOT
- / "migrations"
- / "versions"
- / "20260722_105_data_research_elements.py"
- ).read_text(encoding="utf-8")
- assert 'revision = "20260722_105"' in migration
- assert 'down_revision = "20260722_100"' in migration
- for table in ("data_elements", "data_element_versions", "candidate_decisions"):
- assert f"CREATE TABLE IF NOT EXISTS public.{table}" in migration
- assert "UNIQUE (data_element_uid, version)" in migration
- assert "DROP TABLE" not in migration.upper()
- def test_ontology_migration_adds_versioned_control_plane_tables():
- migration = (
- ROOT / "migrations" / "versions" / "20260722_110_data_research_ontology.py"
- ).read_text(encoding="utf-8")
- assert 'revision = "20260722_110"' in migration
- assert 'down_revision = "20260722_105"' in migration
- for table in (
- "ontologies",
- "ontology_versions",
- "ontology_domain_links",
- "ontology_change_sets",
- "ontology_publish_runs",
- ):
- assert f"CREATE TABLE IF NOT EXISTS public.{table}" in migration
- assert "UNIQUE (ontology_uid, version)" in migration
- assert "DROP TABLE" not in migration.upper()
- def test_alembic_configuration_is_environment_only():
- ini = (ROOT / "alembic.ini").read_text(encoding="utf-8")
- env = (ROOT / "migrations" / "env.py").read_text(encoding="utf-8")
- assert "sqlalchemy.url" not in ini
- assert "SQLALCHEMY_DATABASE_URI" in env
- assert "DATABASE_URL" in env
- assert "password" not in env.lower()
- def test_baseline_migration_is_non_destructive():
- baseline = (
- ROOT / "migrations" / "versions" / "20260716_01_baseline.py"
- ).read_text(encoding="utf-8")
- assert "def upgrade" in baseline
- assert "def downgrade" in baseline
- assert "drop_table" not in baseline
- for table in EXPECTED_BASELINE_TABLES:
- assert table in baseline
- @pytest.mark.integration
- def test_alembic_upgrade_is_repeatable_and_downgrade_preserves_tables():
- admin_url = os.environ.get("TEST_POSTGRES_ADMIN_URL")
- if not admin_url:
- pytest.skip("TEST_POSTGRES_ADMIN_URL is not configured")
- parsed = make_url(admin_url)
- database_name = f"dataops_migration_{uuid.uuid4().hex[:12]}"
- target_url = parsed.set(database=database_name).render_as_string(
- hide_password=False
- )
- connection = psycopg2.connect(admin_url)
- connection.autocommit = True
- try:
- with connection.cursor() as cursor:
- cursor.execute(f'CREATE DATABASE "{database_name}"')
- env = os.environ.copy()
- env["SQLALCHEMY_DATABASE_URI"] = target_url
- command = [str(ROOT / ".venv" / "bin" / "alembic"), "-c", "alembic.ini"]
- subprocess.run(command + ["upgrade", "head"], cwd=ROOT, env=env, check=True)
- subprocess.run(command + ["upgrade", "head"], cwd=ROOT, env=env, check=True)
- engine = create_engine(target_url)
- try:
- tables = set(inspect(engine).get_table_names(schema="public"))
- assert EXPECTED_BASELINE_TABLES | {"alembic_version"} <= tables
- assert EXPECTED_UPGRADED_TABLES <= tables
- finally:
- engine.dispose()
- subprocess.run(command + ["downgrade", "-1"], cwd=ROOT, env=env, check=True)
- engine = create_engine(target_url)
- try:
- tables = set(inspect(engine).get_table_names(schema="public"))
- assert EXPECTED_BASELINE_TABLES <= tables
- finally:
- engine.dispose()
- finally:
- with connection.cursor() as cursor:
- cursor.execute(
- "SELECT pg_terminate_backend(pid) FROM pg_stat_activity "
- "WHERE datname = %s AND pid <> pg_backend_pid()",
- (database_name,),
- )
- cursor.execute(f'DROP DATABASE IF EXISTS "{database_name}"')
- connection.close()
|