test_database_migrations.py 7.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248
  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. "ingestion_sources",
  35. "source_artifacts",
  36. "ingestion_jobs",
  37. "evidence_fragments",
  38. "extraction_candidates",
  39. "data_elements",
  40. "data_element_versions",
  41. "candidate_decisions",
  42. "ontologies",
  43. "ontology_versions",
  44. "ontology_domain_links",
  45. "ontology_change_sets",
  46. "ontology_publish_runs",
  47. }
  48. def test_data_research_ingestion_migration_is_additive_and_constrained():
  49. migration = (
  50. ROOT
  51. / "migrations"
  52. / "versions"
  53. / "20260722_100_data_research_ingestion.py"
  54. ).read_text(encoding="utf-8")
  55. assert 'revision = "20260722_100"' in migration
  56. assert 'down_revision = "20260719_90"' in migration
  57. for table in (
  58. "ingestion_sources",
  59. "source_artifacts",
  60. "ingestion_jobs",
  61. "evidence_fragments",
  62. "extraction_candidates",
  63. ):
  64. assert f"CREATE TABLE IF NOT EXISTS public.{table}" in migration
  65. assert "idempotency_key" in migration
  66. assert "parser_version" in migration
  67. assert "locator JSONB" in migration
  68. assert "DROP TABLE" not in migration.upper()
  69. def test_data_element_migration_adds_versioned_governance_tables():
  70. migration = (
  71. ROOT
  72. / "migrations"
  73. / "versions"
  74. / "20260722_105_data_research_elements.py"
  75. ).read_text(encoding="utf-8")
  76. assert 'revision = "20260722_105"' in migration
  77. assert 'down_revision = "20260722_100"' in migration
  78. for table in ("data_elements", "data_element_versions", "candidate_decisions"):
  79. assert f"CREATE TABLE IF NOT EXISTS public.{table}" in migration
  80. assert "UNIQUE (data_element_uid, version)" in migration
  81. assert "DROP TABLE" not in migration.upper()
  82. def test_ontology_migration_adds_versioned_control_plane_tables():
  83. migration = (
  84. ROOT / "migrations" / "versions" / "20260722_110_data_research_ontology.py"
  85. ).read_text(encoding="utf-8")
  86. assert 'revision = "20260722_110"' in migration
  87. assert 'down_revision = "20260722_105"' in migration
  88. for table in (
  89. "ontologies",
  90. "ontology_versions",
  91. "ontology_domain_links",
  92. "ontology_change_sets",
  93. "ontology_publish_runs",
  94. ):
  95. assert f"CREATE TABLE IF NOT EXISTS public.{table}" in migration
  96. assert "UNIQUE (ontology_uid, version)" in migration
  97. assert "DROP TABLE" not in migration.upper()
  98. def test_alembic_configuration_is_environment_only():
  99. ini = (ROOT / "alembic.ini").read_text(encoding="utf-8")
  100. env = (ROOT / "migrations" / "env.py").read_text(encoding="utf-8")
  101. assert "sqlalchemy.url" not in ini
  102. assert "SQLALCHEMY_DATABASE_URI" in env
  103. assert "DATABASE_URL" in env
  104. assert "password" not in env.lower()
  105. def test_baseline_migration_is_non_destructive():
  106. baseline = (
  107. ROOT / "migrations" / "versions" / "20260716_01_baseline.py"
  108. ).read_text(encoding="utf-8")
  109. assert "def upgrade" in baseline
  110. assert "def downgrade" in baseline
  111. assert "drop_table" not in baseline
  112. for table in EXPECTED_BASELINE_TABLES:
  113. assert table in baseline
  114. @pytest.mark.integration
  115. def test_alembic_upgrade_is_repeatable_and_downgrade_preserves_tables():
  116. admin_url = os.environ.get("TEST_POSTGRES_ADMIN_URL")
  117. if not admin_url:
  118. pytest.skip("TEST_POSTGRES_ADMIN_URL is not configured")
  119. parsed = make_url(admin_url)
  120. database_name = f"dataops_migration_{uuid.uuid4().hex[:12]}"
  121. target_url = parsed.set(database=database_name).render_as_string(
  122. hide_password=False
  123. )
  124. connection = psycopg2.connect(admin_url)
  125. connection.autocommit = True
  126. try:
  127. with connection.cursor() as cursor:
  128. cursor.execute(f'CREATE DATABASE "{database_name}"')
  129. env = os.environ.copy()
  130. env["SQLALCHEMY_DATABASE_URI"] = target_url
  131. command = [str(ROOT / ".venv" / "bin" / "alembic"), "-c", "alembic.ini"]
  132. subprocess.run(command + ["upgrade", "head"], cwd=ROOT, env=env, check=True)
  133. subprocess.run(command + ["upgrade", "head"], cwd=ROOT, env=env, check=True)
  134. engine = create_engine(target_url)
  135. try:
  136. tables = set(inspect(engine).get_table_names(schema="public"))
  137. assert EXPECTED_BASELINE_TABLES | {"alembic_version"} <= tables
  138. assert EXPECTED_UPGRADED_TABLES <= tables
  139. finally:
  140. engine.dispose()
  141. subprocess.run(command + ["downgrade", "-1"], cwd=ROOT, env=env, check=True)
  142. engine = create_engine(target_url)
  143. try:
  144. tables = set(inspect(engine).get_table_names(schema="public"))
  145. assert EXPECTED_BASELINE_TABLES <= tables
  146. finally:
  147. engine.dispose()
  148. finally:
  149. with connection.cursor() as cursor:
  150. cursor.execute(
  151. "SELECT pg_terminate_backend(pid) FROM pg_stat_activity "
  152. "WHERE datname = %s AND pid <> pg_backend_pid()",
  153. (database_name,),
  154. )
  155. cursor.execute(f'DROP DATABASE IF EXISTS "{database_name}"')
  156. connection.close()
  157. @pytest.mark.integration
  158. def test_concurrent_alembic_upgrades_are_serialized():
  159. admin_url = os.environ.get("TEST_POSTGRES_ADMIN_URL")
  160. if not admin_url:
  161. pytest.skip("TEST_POSTGRES_ADMIN_URL is not configured")
  162. parsed = make_url(admin_url)
  163. database_name = f"dataops_concurrent_{uuid.uuid4().hex[:12]}"
  164. target_url = parsed.set(database=database_name).render_as_string(
  165. hide_password=False
  166. )
  167. connection = psycopg2.connect(admin_url)
  168. connection.autocommit = True
  169. try:
  170. with connection.cursor() as cursor:
  171. cursor.execute(f'CREATE DATABASE "{database_name}"')
  172. env = os.environ.copy()
  173. env["SQLALCHEMY_DATABASE_URI"] = target_url
  174. command = [
  175. str(ROOT / ".venv" / "bin" / "alembic"),
  176. "-c",
  177. "alembic.ini",
  178. "upgrade",
  179. "head",
  180. ]
  181. processes = [
  182. subprocess.Popen(
  183. command,
  184. cwd=ROOT,
  185. env=env,
  186. stdout=subprocess.PIPE,
  187. stderr=subprocess.STDOUT,
  188. text=True,
  189. )
  190. for _ in range(2)
  191. ]
  192. results = [process.communicate(timeout=120) for process in processes]
  193. failures = [
  194. output
  195. for process, (output, _) in zip(processes, results)
  196. if process.returncode != 0
  197. ]
  198. assert failures == []
  199. engine = create_engine(target_url)
  200. try:
  201. assert set(inspect(engine).get_table_names(schema="public")) >= (
  202. EXPECTED_UPGRADED_TABLES | {"alembic_version"}
  203. )
  204. finally:
  205. engine.dispose()
  206. finally:
  207. with connection.cursor() as cursor:
  208. cursor.execute(
  209. "SELECT pg_terminate_backend(pid) FROM pg_stat_activity "
  210. "WHERE datname = %s AND pid <> pg_backend_pid()",
  211. (database_name,),
  212. )
  213. cursor.execute(f'DROP DATABASE IF EXISTS "{database_name}"')
  214. connection.close()