test_database_migrations.py 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361
  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. "ingestion_sources",
  36. "source_artifacts",
  37. "ingestion_jobs",
  38. "evidence_fragments",
  39. "extraction_candidates",
  40. "data_elements",
  41. "data_element_versions",
  42. "candidate_decisions",
  43. "ontologies",
  44. "ontology_versions",
  45. "ontology_domain_links",
  46. "ontology_change_sets",
  47. "ontology_publish_runs",
  48. "governance_responsibility_scopes",
  49. "governance_responsibility_assignments",
  50. "governance_responsibility_audit_events",
  51. }
  52. def test_data_research_ingestion_migration_is_additive_and_constrained():
  53. migration = (
  54. ROOT
  55. / "migrations"
  56. / "versions"
  57. / "20260722_100_data_research_ingestion.py"
  58. ).read_text(encoding="utf-8")
  59. assert 'revision = "20260722_100"' in migration
  60. assert 'down_revision = "20260720_100"' in migration
  61. for table in (
  62. "ingestion_sources",
  63. "source_artifacts",
  64. "ingestion_jobs",
  65. "evidence_fragments",
  66. "extraction_candidates",
  67. ):
  68. assert f"CREATE TABLE IF NOT EXISTS public.{table}" in migration
  69. assert "idempotency_key" in migration
  70. assert "parser_version" in migration
  71. assert "locator JSONB" in migration
  72. assert "DROP TABLE" not in migration.upper()
  73. def test_data_element_migration_adds_versioned_governance_tables():
  74. migration = (
  75. ROOT
  76. / "migrations"
  77. / "versions"
  78. / "20260722_105_data_research_elements.py"
  79. ).read_text(encoding="utf-8")
  80. assert 'revision = "20260722_105"' in migration
  81. assert 'down_revision = "20260722_100"' in migration
  82. for table in ("data_elements", "data_element_versions", "candidate_decisions"):
  83. assert f"CREATE TABLE IF NOT EXISTS public.{table}" in migration
  84. assert "UNIQUE (data_element_uid, version)" in migration
  85. assert "DROP TABLE" not in migration.upper()
  86. def test_ontology_migration_adds_versioned_control_plane_tables():
  87. migration = (
  88. ROOT / "migrations" / "versions" / "20260722_110_data_research_ontology.py"
  89. ).read_text(encoding="utf-8")
  90. assert 'revision = "20260722_110"' in migration
  91. assert 'down_revision = "20260722_105"' in migration
  92. for table in (
  93. "ontologies",
  94. "ontology_versions",
  95. "ontology_domain_links",
  96. "ontology_change_sets",
  97. "ontology_publish_runs",
  98. ):
  99. assert f"CREATE TABLE IF NOT EXISTS public.{table}" in migration
  100. assert "UNIQUE (ontology_uid, version)" in migration
  101. assert "DROP TABLE" not in migration.upper()
  102. def test_alembic_configuration_is_environment_only():
  103. ini = (ROOT / "alembic.ini").read_text(encoding="utf-8")
  104. env = (ROOT / "migrations" / "env.py").read_text(encoding="utf-8")
  105. assert "sqlalchemy.url" not in ini
  106. assert "SQLALCHEMY_DATABASE_URI" in env
  107. assert "DATABASE_URL" in env
  108. assert "password" not in env.lower()
  109. def test_baseline_migration_is_non_destructive():
  110. baseline = (
  111. ROOT / "migrations" / "versions" / "20260716_01_baseline.py"
  112. ).read_text(encoding="utf-8")
  113. assert "def upgrade" in baseline
  114. assert "def downgrade" in baseline
  115. assert "drop_table" not in baseline
  116. for table in EXPECTED_BASELINE_TABLES:
  117. assert table in baseline
  118. @pytest.mark.integration
  119. def test_alembic_upgrade_is_repeatable_and_downgrade_preserves_tables():
  120. admin_url = os.environ.get("TEST_POSTGRES_ADMIN_URL")
  121. if not admin_url:
  122. pytest.skip("TEST_POSTGRES_ADMIN_URL is not configured")
  123. parsed = make_url(admin_url)
  124. database_name = f"dataops_migration_{uuid.uuid4().hex[:12]}"
  125. target_url = parsed.set(database=database_name).render_as_string(
  126. hide_password=False
  127. )
  128. connection = psycopg2.connect(admin_url)
  129. connection.autocommit = True
  130. try:
  131. with connection.cursor() as cursor:
  132. cursor.execute(f'CREATE DATABASE "{database_name}"')
  133. env = os.environ.copy()
  134. env["SQLALCHEMY_DATABASE_URI"] = target_url
  135. command = [str(ROOT / ".venv" / "bin" / "alembic"), "-c", "alembic.ini"]
  136. subprocess.run(command + ["upgrade", "head"], cwd=ROOT, env=env, check=True)
  137. subprocess.run(command + ["upgrade", "head"], cwd=ROOT, env=env, check=True)
  138. engine = create_engine(target_url)
  139. try:
  140. tables = set(inspect(engine).get_table_names(schema="public"))
  141. assert EXPECTED_BASELINE_TABLES | {"alembic_version"} <= tables
  142. assert tables >= EXPECTED_UPGRADED_TABLES
  143. finally:
  144. engine.dispose()
  145. subprocess.run(command + ["downgrade", "-1"], cwd=ROOT, env=env, check=True)
  146. engine = create_engine(target_url)
  147. try:
  148. tables = set(inspect(engine).get_table_names(schema="public"))
  149. assert tables >= EXPECTED_BASELINE_TABLES
  150. finally:
  151. engine.dispose()
  152. finally:
  153. with connection.cursor() as cursor:
  154. cursor.execute(
  155. "SELECT pg_terminate_backend(pid) FROM pg_stat_activity "
  156. "WHERE datname = %s AND pid <> pg_backend_pid()",
  157. (database_name,),
  158. )
  159. cursor.execute(f'DROP DATABASE IF EXISTS "{database_name}"')
  160. connection.close()
  161. @pytest.mark.integration
  162. @pytest.mark.parametrize(
  163. ("legacy_rows", "diagnostic"),
  164. (
  165. ((1, 1), "duplicate rule_run_id"),
  166. ((101,), "rows above 100"),
  167. ),
  168. )
  169. def test_upgrade_from_150_rejects_malformed_legacy_samples(
  170. legacy_rows,
  171. diagnostic,
  172. ):
  173. admin_url = os.environ.get("TEST_POSTGRES_ADMIN_URL")
  174. if not admin_url:
  175. pytest.skip("TEST_POSTGRES_ADMIN_URL is not configured")
  176. parsed = make_url(admin_url)
  177. database_name = f"dataops_evidence_{uuid.uuid4().hex[:12]}"
  178. target_url = parsed.set(database=database_name).render_as_string(
  179. hide_password=False
  180. )
  181. admin = psycopg2.connect(admin_url)
  182. admin.autocommit = True
  183. try:
  184. with admin.cursor() as cursor:
  185. cursor.execute(f'CREATE DATABASE "{database_name}"')
  186. env = os.environ.copy()
  187. env["SQLALCHEMY_DATABASE_URI"] = target_url
  188. command = [
  189. str(ROOT / ".venv" / "bin" / "alembic"),
  190. "-c",
  191. "alembic.ini",
  192. ]
  193. subprocess.run(
  194. command + ["upgrade", "20260723_150"],
  195. cwd=ROOT,
  196. env=env,
  197. check=True,
  198. )
  199. database = psycopg2.connect(target_url)
  200. try:
  201. database.autocommit = True
  202. with database.cursor() as cursor:
  203. cursor.execute("SET session_replication_role = replica")
  204. run_id = str(uuid.uuid4())
  205. cursor.execute(
  206. """
  207. INSERT INTO public.rule_runs (
  208. id, deployment_id, component_binding_id,
  209. rule_version_id, plan_hash, status, correlation_id
  210. ) VALUES (%s, %s, %s, %s, %s, 'failed', %s)
  211. """,
  212. (
  213. run_id,
  214. str(uuid.uuid4()),
  215. str(uuid.uuid4()),
  216. str(uuid.uuid4()),
  217. "a" * 64,
  218. str(uuid.uuid4()),
  219. ),
  220. )
  221. for count in legacy_rows:
  222. cursor.execute(
  223. """
  224. INSERT INTO public.rule_violation_samples (
  225. id, rule_run_id, artifact_ref, sample_count,
  226. redaction_policy, expires_at
  227. ) VALUES (%s, %s, %s, %s, 'legacy-v1', NOW())
  228. """,
  229. (
  230. str(uuid.uuid4()),
  231. run_id,
  232. f"minio://legacy/{uuid.uuid4()}",
  233. count,
  234. ),
  235. )
  236. cursor.execute("SET session_replication_role = origin")
  237. finally:
  238. database.close()
  239. failed = subprocess.run(
  240. command + ["upgrade", "head"],
  241. cwd=ROOT,
  242. env=env,
  243. text=True,
  244. capture_output=True,
  245. )
  246. assert failed.returncode != 0
  247. assert diagnostic in (failed.stdout + failed.stderr)
  248. database = psycopg2.connect(target_url)
  249. try:
  250. with database.cursor() as cursor:
  251. cursor.execute(
  252. "SELECT version_num FROM public.alembic_version"
  253. )
  254. assert cursor.fetchone()[0] == "20260723_150"
  255. finally:
  256. database.close()
  257. finally:
  258. with admin.cursor() as cursor:
  259. cursor.execute(
  260. "SELECT pg_terminate_backend(pid) FROM pg_stat_activity "
  261. "WHERE datname = %s AND pid <> pg_backend_pid()",
  262. (database_name,),
  263. )
  264. cursor.execute(f'DROP DATABASE IF EXISTS "{database_name}"')
  265. admin.close()
  266. @pytest.mark.integration
  267. def test_concurrent_alembic_upgrades_are_serialized():
  268. admin_url = os.environ.get("TEST_POSTGRES_ADMIN_URL")
  269. if not admin_url:
  270. pytest.skip("TEST_POSTGRES_ADMIN_URL is not configured")
  271. parsed = make_url(admin_url)
  272. database_name = f"dataops_concurrent_{uuid.uuid4().hex[:12]}"
  273. target_url = parsed.set(database=database_name).render_as_string(
  274. hide_password=False
  275. )
  276. connection = psycopg2.connect(admin_url)
  277. connection.autocommit = True
  278. try:
  279. with connection.cursor() as cursor:
  280. cursor.execute(f'CREATE DATABASE "{database_name}"')
  281. env = os.environ.copy()
  282. env["SQLALCHEMY_DATABASE_URI"] = target_url
  283. command = [
  284. str(ROOT / ".venv" / "bin" / "alembic"),
  285. "-c",
  286. "alembic.ini",
  287. "upgrade",
  288. "head",
  289. ]
  290. processes = [
  291. subprocess.Popen(
  292. command,
  293. cwd=ROOT,
  294. env=env,
  295. stdout=subprocess.PIPE,
  296. stderr=subprocess.STDOUT,
  297. text=True,
  298. )
  299. for _ in range(2)
  300. ]
  301. results = [process.communicate(timeout=120) for process in processes]
  302. failures = [
  303. output
  304. for process, (output, _) in zip(processes, results)
  305. if process.returncode != 0
  306. ]
  307. assert failures == []
  308. engine = create_engine(target_url)
  309. try:
  310. assert set(inspect(engine).get_table_names(schema="public")) >= (
  311. EXPECTED_UPGRADED_TABLES | {"alembic_version"}
  312. )
  313. finally:
  314. engine.dispose()
  315. finally:
  316. with connection.cursor() as cursor:
  317. cursor.execute(
  318. "SELECT pg_terminate_backend(pid) FROM pg_stat_activity "
  319. "WHERE datname = %s AND pid <> pg_backend_pid()",
  320. (database_name,),
  321. )
  322. cursor.execute(f'DROP DATABASE IF EXISTS "{database_name}"')
  323. connection.close()