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", "rule_sql_staging_receipts", } 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 tables >= EXPECTED_UPGRADED_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 tables >= EXPECTED_BASELINE_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() @pytest.mark.integration @pytest.mark.parametrize( ("legacy_rows", "diagnostic"), ( ((1, 1), "duplicate rule_run_id"), ((101,), "rows above 100"), ), ) def test_upgrade_from_150_rejects_malformed_legacy_samples( legacy_rows, diagnostic, ): 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_evidence_{uuid.uuid4().hex[:12]}" target_url = parsed.set(database=database_name).render_as_string( hide_password=False ) admin = psycopg2.connect(admin_url) admin.autocommit = True try: with admin.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", "20260723_150"], cwd=ROOT, env=env, check=True, ) database = psycopg2.connect(target_url) try: database.autocommit = True with database.cursor() as cursor: cursor.execute("SET session_replication_role = replica") run_id = str(uuid.uuid4()) cursor.execute( """ INSERT INTO public.rule_runs ( id, deployment_id, component_binding_id, rule_version_id, plan_hash, status, correlation_id ) VALUES (%s, %s, %s, %s, %s, 'failed', %s) """, ( run_id, str(uuid.uuid4()), str(uuid.uuid4()), str(uuid.uuid4()), "a" * 64, str(uuid.uuid4()), ), ) for count in legacy_rows: cursor.execute( """ INSERT INTO public.rule_violation_samples ( id, rule_run_id, artifact_ref, sample_count, redaction_policy, expires_at ) VALUES (%s, %s, %s, %s, 'legacy-v1', NOW()) """, ( str(uuid.uuid4()), run_id, f"minio://legacy/{uuid.uuid4()}", count, ), ) cursor.execute("SET session_replication_role = origin") finally: database.close() failed = subprocess.run( command + ["upgrade", "head"], cwd=ROOT, env=env, text=True, capture_output=True, ) assert failed.returncode != 0 assert diagnostic in (failed.stdout + failed.stderr) database = psycopg2.connect(target_url) try: with database.cursor() as cursor: cursor.execute( "SELECT version_num FROM public.alembic_version" ) assert cursor.fetchone()[0] == "20260723_150" finally: database.close() finally: with admin.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}"') admin.close()