test_active_metadata_postgres.py 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320
  1. from __future__ import annotations
  2. import os
  3. import uuid
  4. import pytest
  5. from sqlalchemy import text
  6. pytestmark = pytest.mark.integration
  7. def test_active_metadata_is_idempotent_traceable_and_correctable(monkeypatch):
  8. database_url = os.environ.get("TEST_DATABASE_URL")
  9. if not database_url:
  10. pytest.skip("TEST_DATABASE_URL is required")
  11. monkeypatch.setenv("DATABASE_URL", database_url)
  12. from app import create_app, db
  13. from app.core.meta_data.active_metadata import ActiveMetadataService
  14. from app.core.meta_data.active_metadata_repository import (
  15. SqlAlchemyActiveMetadataRepository,
  16. )
  17. app = create_app()
  18. app.config.update(TESTING=True)
  19. actor_uid = str(uuid.uuid4())
  20. owner_uid = str(uuid.uuid4())
  21. source_uid = str(uuid.uuid4())
  22. plan_uid = None
  23. try:
  24. with app.app_context():
  25. for uid, name in ((actor_uid, "actor"), (owner_uid, "owner")):
  26. db.session.execute(
  27. text(
  28. """
  29. INSERT INTO public.users (
  30. id, username, display_name, password_hash, status
  31. ) VALUES (
  32. CAST(:id AS uuid), :username, :username,
  33. 'p2-wp02-integration-hash', 'active'
  34. )
  35. """
  36. ),
  37. {"id": uid, "username": f"wp02-{name}-{uid[:8]}"},
  38. )
  39. db.session.execute(
  40. text(
  41. """
  42. INSERT INTO public.ingestion_sources (
  43. uid, source_type, name, config, permission_scope,
  44. status, created_by
  45. ) VALUES (
  46. CAST(:uid AS uuid), 'database', :name,
  47. '{}'::jsonb, '{}'::jsonb, 'active', :created_by
  48. )
  49. """
  50. ),
  51. {
  52. "uid": source_uid,
  53. "name": f"WP02 source {source_uid[:8]}",
  54. "created_by": actor_uid,
  55. },
  56. )
  57. db.session.commit()
  58. repository = SqlAlchemyActiveMetadataRepository(db.session)
  59. service = ActiveMetadataService(repository)
  60. plan = service.create_plan(
  61. {
  62. "source_uid": source_uid,
  63. "name": "WP02 集成主动发现",
  64. "source_kind": "database",
  65. "schedule_type": "interval",
  66. "schedule_expression": "PT30M",
  67. "discovery_mode": "cursor",
  68. "scope": {"schemas": ["public"]},
  69. "owner_uid": owner_uid,
  70. },
  71. actor_uid=actor_uid,
  72. )
  73. db.session.commit()
  74. plan_uid = plan["uid"]
  75. first_payload = {
  76. "snapshot": {
  77. "assets": [
  78. {
  79. "name": "customers",
  80. "namespace": "public",
  81. "asset_type": "table",
  82. "fields": [
  83. {
  84. "name": "id",
  85. "data_type": "bigint",
  86. "nullable": False,
  87. "ordinal_position": 1,
  88. },
  89. {
  90. "name": "email",
  91. "data_type": "varchar",
  92. "nullable": True,
  93. "ordinal_position": 2,
  94. },
  95. ],
  96. },
  97. {
  98. "name": "orders",
  99. "namespace": "public",
  100. "asset_type": "table",
  101. "fields": [],
  102. },
  103. ]
  104. },
  105. "cursor_after": {"catalog_version": 1},
  106. "lineage_sql": [
  107. {
  108. "dialect": "postgres",
  109. "sql": (
  110. "INSERT INTO mart.customer_email (customer_id, email) "
  111. "SELECT c.id, c.email FROM public.customers c"
  112. ),
  113. },
  114. "not valid lineage SQL",
  115. ],
  116. "health_signals": [
  117. {
  118. "asset_key": f"{source_uid}:public.customers",
  119. "signal_type": "quality",
  120. "value": 0.98,
  121. "status": "healthy",
  122. },
  123. {
  124. "asset_key": f"{source_uid}:public.customers",
  125. "signal_type": "freshness",
  126. "value": 60,
  127. "status": "healthy",
  128. },
  129. {
  130. "asset_key": f"{source_uid}:public.customers",
  131. "signal_type": "usage",
  132. "value": 12,
  133. "status": "healthy",
  134. },
  135. ],
  136. }
  137. first = service.execute_source_plans(
  138. source_uid,
  139. first_payload["snapshot"],
  140. batch_key="wp02-batch-1",
  141. actor_uid=actor_uid,
  142. cursor_after=first_payload["cursor_after"],
  143. )[0]
  144. # SQL lineage and health use the same batch contract as the
  145. # automatic snapshot projection and are appended by the runner.
  146. first = service.execute(
  147. plan_uid,
  148. first_payload,
  149. batch_key="wp02-batch-1-enriched",
  150. actor_uid=actor_uid,
  151. )
  152. db.session.commit()
  153. replay = service.execute(
  154. plan_uid,
  155. first_payload,
  156. batch_key="wp02-batch-1-enriched",
  157. actor_uid=actor_uid,
  158. )
  159. db.session.commit()
  160. assert replay["uid"] == first["uid"]
  161. second_payload = {
  162. "snapshot": {
  163. "assets": [
  164. {
  165. "name": "customers",
  166. "namespace": "public",
  167. "asset_type": "table",
  168. "fields": [
  169. {
  170. "name": "id",
  171. "data_type": "bigint",
  172. "nullable": False,
  173. "ordinal_position": 1,
  174. },
  175. {
  176. "name": "email",
  177. "data_type": "text",
  178. "nullable": True,
  179. "ordinal_position": 2,
  180. },
  181. ],
  182. }
  183. ]
  184. },
  185. "cursor_after": {"catalog_version": 2},
  186. }
  187. second = service.execute(
  188. plan_uid,
  189. second_payload,
  190. batch_key="wp02-batch-2",
  191. actor_uid=actor_uid,
  192. )
  193. db.session.commit()
  194. assert second["cursor_before"] == {"catalog_version": 1}
  195. assets = service.list_assets(source_uid)
  196. customer = next(item for item in assets if item["name"] == "customers")
  197. orders = next(item for item in assets if item["name"] == "orders")
  198. assert customer["current_version"] == 2
  199. assert orders["lifecycle_status"] == "deletion_candidate"
  200. changes = service.list_changes(second["uid"])
  201. assert {(item["change_type"], item["field_name"]) for item in changes} == {
  202. ("field_changed", "email"),
  203. ("deletion_candidate", None),
  204. }
  205. lineage = service.list_lineage(first["uid"])
  206. assert sum(item["parse_status"] == "resolved" for item in lineage) == 2
  207. assert sum(item["parse_status"] == "failed" for item in lineage) == 1
  208. failed = service.record_failure(
  209. plan_uid,
  210. batch_key="wp02-batch-3",
  211. error_code="TIMEOUT",
  212. failure_reason="collector timed out",
  213. actor_uid=actor_uid,
  214. )
  215. db.session.commit()
  216. assert failed["status"] == "failed"
  217. correction = service.submit_correction(
  218. customer["uid"],
  219. {
  220. "field_name": "email",
  221. "proposed_value": {"comment": "客户邮箱"},
  222. "reason": "补充业务定义",
  223. "assignee_uid": owner_uid,
  224. },
  225. actor_uid=actor_uid,
  226. )
  227. db.session.commit()
  228. resolved = service.resolve_correction(
  229. correction["uid"],
  230. expected_version=1,
  231. decision="accept",
  232. resolution={"comment": "客户邮箱"},
  233. actor_uid=owner_uid,
  234. )
  235. db.session.commit()
  236. assert resolved["current_version"] == 2
  237. counts = db.session.execute(
  238. text(
  239. """
  240. SELECT
  241. (SELECT count(*) FROM public.active_metadata_runs
  242. WHERE plan_uid = CAST(:plan_uid AS uuid)) AS runs,
  243. (SELECT count(*) FROM public.active_metadata_asset_versions v
  244. JOIN public.active_metadata_assets a ON a.uid = v.asset_uid
  245. WHERE a.source_uid = CAST(:source_uid AS uuid)) AS versions,
  246. (SELECT count(*) FROM public.active_metadata_correction_audits ca
  247. JOIN public.active_metadata_corrections c
  248. ON c.uid = ca.correction_uid
  249. JOIN public.active_metadata_assets a ON a.uid = c.asset_uid
  250. WHERE a.source_uid = CAST(:source_uid AS uuid)) AS audits
  251. """
  252. ),
  253. {"plan_uid": plan_uid, "source_uid": source_uid},
  254. ).mappings().one()
  255. assert dict(counts) == {"runs": 4, "versions": 3, "audits": 2}
  256. finally:
  257. with app.app_context():
  258. db.session.rollback()
  259. if plan_uid:
  260. for statement in (
  261. "DELETE FROM public.active_metadata_correction_audits "
  262. "WHERE correction_uid IN (SELECT c.uid FROM public.active_metadata_corrections c "
  263. "JOIN public.active_metadata_assets a ON a.uid = c.asset_uid "
  264. "WHERE a.source_uid = CAST(:source_uid AS uuid))",
  265. "DELETE FROM public.active_metadata_corrections "
  266. "WHERE asset_uid IN (SELECT uid FROM public.active_metadata_assets "
  267. "WHERE source_uid = CAST(:source_uid AS uuid))",
  268. "DELETE FROM public.active_metadata_health_signals "
  269. "WHERE asset_uid IN (SELECT uid FROM public.active_metadata_assets "
  270. "WHERE source_uid = CAST(:source_uid AS uuid))",
  271. "DELETE FROM public.active_metadata_lineage "
  272. "WHERE run_uid IN (SELECT uid FROM public.active_metadata_runs "
  273. "WHERE plan_uid = CAST(:plan_uid AS uuid))",
  274. "DELETE FROM public.active_metadata_changes "
  275. "WHERE run_uid IN (SELECT uid FROM public.active_metadata_runs "
  276. "WHERE plan_uid = CAST(:plan_uid AS uuid))",
  277. "DELETE FROM public.active_metadata_asset_versions "
  278. "WHERE asset_uid IN (SELECT uid FROM public.active_metadata_assets "
  279. "WHERE source_uid = CAST(:source_uid AS uuid))",
  280. "DELETE FROM public.active_metadata_assets "
  281. "WHERE source_uid = CAST(:source_uid AS uuid)",
  282. "DELETE FROM public.active_metadata_runs "
  283. "WHERE plan_uid = CAST(:plan_uid AS uuid)",
  284. "DELETE FROM public.active_metadata_plans "
  285. "WHERE uid = CAST(:plan_uid AS uuid)",
  286. ):
  287. db.session.execute(
  288. text(statement),
  289. {"plan_uid": plan_uid, "source_uid": source_uid},
  290. )
  291. db.session.execute(
  292. text(
  293. "DELETE FROM public.ingestion_sources "
  294. "WHERE uid = CAST(:source_uid AS uuid)"
  295. ),
  296. {"source_uid": source_uid},
  297. )
  298. db.session.execute(
  299. text(
  300. "DELETE FROM public.users "
  301. "WHERE id IN (CAST(:actor_uid AS uuid), CAST(:owner_uid AS uuid))"
  302. ),
  303. {"actor_uid": actor_uid, "owner_uid": owner_uid},
  304. )
  305. db.session.commit()