test_database_migrations.py 7.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222
  1. from __future__ import annotations
  2. import os
  3. import subprocess
  4. import uuid
  5. from pathlib import Path
  6. import psycopg2
  7. import pytest
  8. from sqlalchemy import create_engine, inspect
  9. from sqlalchemy.engine import make_url
  10. ROOT = Path(__file__).resolve().parents[1]
  11. EXPECTED_BASELINE_TABLES = {
  12. "data_orders",
  13. "data_products",
  14. "metadata_review_records",
  15. "metadata_version_history",
  16. "task_list",
  17. "users",
  18. }
  19. EXPECTED_UPGRADED_TABLES = {
  20. "datasource_credentials",
  21. "datasource_credential_audit_events",
  22. "workflow_schedules",
  23. "workflow_runs",
  24. "workflow_task_runs",
  25. "workflow_engine_bindings",
  26. "workflow_plan_audits",
  27. "runner_task_executions",
  28. "workflow_canary_evidence",
  29. "workflow_gateway_operations",
  30. "workflow_migration_states",
  31. "workflow_dual_runs",
  32. "workflow_reconciliation_reports",
  33. "workflow_cutover_operations",
  34. "rule_sql_staging_receipts",
  35. }
  36. def test_alembic_configuration_is_environment_only():
  37. ini = (ROOT / "alembic.ini").read_text(encoding="utf-8")
  38. env = (ROOT / "migrations" / "env.py").read_text(encoding="utf-8")
  39. assert "sqlalchemy.url" not in ini
  40. assert "SQLALCHEMY_DATABASE_URI" in env
  41. assert "DATABASE_URL" in env
  42. assert "password" not in env.lower()
  43. def test_baseline_migration_is_non_destructive():
  44. baseline = (
  45. ROOT / "migrations" / "versions" / "20260716_01_baseline.py"
  46. ).read_text(encoding="utf-8")
  47. assert "def upgrade" in baseline
  48. assert "def downgrade" in baseline
  49. assert "drop_table" not in baseline
  50. for table in EXPECTED_BASELINE_TABLES:
  51. assert table in baseline
  52. @pytest.mark.integration
  53. def test_alembic_upgrade_is_repeatable_and_downgrade_preserves_tables():
  54. admin_url = os.environ.get("TEST_POSTGRES_ADMIN_URL")
  55. if not admin_url:
  56. pytest.skip("TEST_POSTGRES_ADMIN_URL is not configured")
  57. parsed = make_url(admin_url)
  58. database_name = f"dataops_migration_{uuid.uuid4().hex[:12]}"
  59. target_url = parsed.set(database=database_name).render_as_string(
  60. hide_password=False
  61. )
  62. connection = psycopg2.connect(admin_url)
  63. connection.autocommit = True
  64. try:
  65. with connection.cursor() as cursor:
  66. cursor.execute(f'CREATE DATABASE "{database_name}"')
  67. env = os.environ.copy()
  68. env["SQLALCHEMY_DATABASE_URI"] = target_url
  69. command = [str(ROOT / ".venv" / "bin" / "alembic"), "-c", "alembic.ini"]
  70. subprocess.run(command + ["upgrade", "head"], cwd=ROOT, env=env, check=True)
  71. subprocess.run(command + ["upgrade", "head"], cwd=ROOT, env=env, check=True)
  72. engine = create_engine(target_url)
  73. try:
  74. tables = set(inspect(engine).get_table_names(schema="public"))
  75. assert EXPECTED_BASELINE_TABLES | {"alembic_version"} <= tables
  76. assert tables >= EXPECTED_UPGRADED_TABLES
  77. finally:
  78. engine.dispose()
  79. subprocess.run(command + ["downgrade", "-1"], cwd=ROOT, env=env, check=True)
  80. engine = create_engine(target_url)
  81. try:
  82. tables = set(inspect(engine).get_table_names(schema="public"))
  83. assert tables >= EXPECTED_BASELINE_TABLES
  84. finally:
  85. engine.dispose()
  86. finally:
  87. with connection.cursor() as cursor:
  88. cursor.execute(
  89. "SELECT pg_terminate_backend(pid) FROM pg_stat_activity "
  90. "WHERE datname = %s AND pid <> pg_backend_pid()",
  91. (database_name,),
  92. )
  93. cursor.execute(f'DROP DATABASE IF EXISTS "{database_name}"')
  94. connection.close()
  95. @pytest.mark.integration
  96. @pytest.mark.parametrize(
  97. ("legacy_rows", "diagnostic"),
  98. (
  99. ((1, 1), "duplicate rule_run_id"),
  100. ((101,), "rows above 100"),
  101. ),
  102. )
  103. def test_upgrade_from_150_rejects_malformed_legacy_samples(
  104. legacy_rows,
  105. diagnostic,
  106. ):
  107. admin_url = os.environ.get("TEST_POSTGRES_ADMIN_URL")
  108. if not admin_url:
  109. pytest.skip("TEST_POSTGRES_ADMIN_URL is not configured")
  110. parsed = make_url(admin_url)
  111. database_name = f"dataops_evidence_{uuid.uuid4().hex[:12]}"
  112. target_url = parsed.set(database=database_name).render_as_string(
  113. hide_password=False
  114. )
  115. admin = psycopg2.connect(admin_url)
  116. admin.autocommit = True
  117. try:
  118. with admin.cursor() as cursor:
  119. cursor.execute(f'CREATE DATABASE "{database_name}"')
  120. env = os.environ.copy()
  121. env["SQLALCHEMY_DATABASE_URI"] = target_url
  122. command = [
  123. str(ROOT / ".venv" / "bin" / "alembic"),
  124. "-c",
  125. "alembic.ini",
  126. ]
  127. subprocess.run(
  128. command + ["upgrade", "20260723_150"],
  129. cwd=ROOT,
  130. env=env,
  131. check=True,
  132. )
  133. database = psycopg2.connect(target_url)
  134. try:
  135. database.autocommit = True
  136. with database.cursor() as cursor:
  137. cursor.execute("SET session_replication_role = replica")
  138. run_id = str(uuid.uuid4())
  139. cursor.execute(
  140. """
  141. INSERT INTO public.rule_runs (
  142. id, deployment_id, component_binding_id,
  143. rule_version_id, plan_hash, status, correlation_id
  144. ) VALUES (%s, %s, %s, %s, %s, 'failed', %s)
  145. """,
  146. (
  147. run_id,
  148. str(uuid.uuid4()),
  149. str(uuid.uuid4()),
  150. str(uuid.uuid4()),
  151. "a" * 64,
  152. str(uuid.uuid4()),
  153. ),
  154. )
  155. for count in legacy_rows:
  156. cursor.execute(
  157. """
  158. INSERT INTO public.rule_violation_samples (
  159. id, rule_run_id, artifact_ref, sample_count,
  160. redaction_policy, expires_at
  161. ) VALUES (%s, %s, %s, %s, 'legacy-v1', NOW())
  162. """,
  163. (
  164. str(uuid.uuid4()),
  165. run_id,
  166. f"minio://legacy/{uuid.uuid4()}",
  167. count,
  168. ),
  169. )
  170. cursor.execute("SET session_replication_role = origin")
  171. finally:
  172. database.close()
  173. failed = subprocess.run(
  174. command + ["upgrade", "head"],
  175. cwd=ROOT,
  176. env=env,
  177. text=True,
  178. capture_output=True,
  179. )
  180. assert failed.returncode != 0
  181. assert diagnostic in (failed.stdout + failed.stderr)
  182. database = psycopg2.connect(target_url)
  183. try:
  184. with database.cursor() as cursor:
  185. cursor.execute(
  186. "SELECT version_num FROM public.alembic_version"
  187. )
  188. assert cursor.fetchone()[0] == "20260723_150"
  189. finally:
  190. database.close()
  191. finally:
  192. with admin.cursor() as cursor:
  193. cursor.execute(
  194. "SELECT pg_terminate_backend(pid) FROM pg_stat_activity "
  195. "WHERE datname = %s AND pid <> pg_backend_pid()",
  196. (database_name,),
  197. )
  198. cursor.execute(f'DROP DATABASE IF EXISTS "{database_name}"')
  199. admin.close()