test_semantic_governance_postgres.py 21 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553
  1. from __future__ import annotations
  2. import json
  3. import os
  4. import pytest
  5. from sqlalchemy import text
  6. pytestmark = pytest.mark.integration
  7. def test_semantic_governance_publishes_maps_and_rolls_back(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.common.identifiers import new_governance_uid
  14. from app.core.data_research.semantic_governance import (
  15. SemanticGovernanceService,
  16. )
  17. from app.core.data_research.semantic_repository import (
  18. SqlAlchemySemanticGovernanceRepository,
  19. )
  20. from app.core.data_rules.repository import DataRuleRepository
  21. from app.core.events.outbox import claim_outbox, enqueue_outbox
  22. app = create_app()
  23. app.config.update(TESTING=True)
  24. actor_uid = new_governance_uid()
  25. owner_uid = new_governance_uid()
  26. reviewer_uid = new_governance_uid()
  27. source_uid = new_governance_uid()
  28. plan_uid = new_governance_uid()
  29. run_uid = new_governance_uid()
  30. physical_asset_uid = new_governance_uid()
  31. element_uid = new_governance_uid()
  32. element_version_uid = new_governance_uid()
  33. standard_uid = new_governance_uid()
  34. standard_id = new_governance_uid()
  35. standard_version_uid = new_governance_uid()
  36. semantic_uids: list[str] = []
  37. try:
  38. with app.app_context():
  39. for user_uid, label in (
  40. (actor_uid, "editor"),
  41. (owner_uid, "owner"),
  42. (reviewer_uid, "reviewer"),
  43. ):
  44. db.session.execute(
  45. text(
  46. """
  47. INSERT INTO public.users (
  48. id, username, display_name, password_hash, status
  49. ) VALUES (
  50. CAST(:uid AS uuid), :username, :username,
  51. 'p2-wp03-integration-hash', 'active'
  52. )
  53. """
  54. ),
  55. {
  56. "uid": user_uid,
  57. "username": f"wp03-{label}-{user_uid[:8]}",
  58. },
  59. )
  60. db.session.execute(
  61. text(
  62. """
  63. INSERT INTO public.ingestion_sources (
  64. uid, source_type, name, config, permission_scope,
  65. status, created_by
  66. ) VALUES (
  67. CAST(:uid AS uuid), 'database', :name,
  68. '{}'::jsonb, '{}'::jsonb, 'active', :created_by
  69. )
  70. """
  71. ),
  72. {
  73. "uid": source_uid,
  74. "name": f"WP03 source {source_uid[:8]}",
  75. "created_by": actor_uid,
  76. },
  77. )
  78. db.session.execute(
  79. text(
  80. """
  81. INSERT INTO public.active_metadata_plans (
  82. uid, source_uid, name, source_kind, schedule_type,
  83. schedule_expression, discovery_mode, scope, owner_uid,
  84. enabled, cursor_state, created_by
  85. ) VALUES (
  86. CAST(:uid AS uuid), CAST(:source_uid AS uuid), :name,
  87. 'database', 'manual', NULL, 'snapshot', '{}'::jsonb,
  88. CAST(:owner_uid AS uuid), TRUE, '{}'::jsonb,
  89. CAST(:actor_uid AS uuid)
  90. )
  91. """
  92. ),
  93. {
  94. "uid": plan_uid,
  95. "source_uid": source_uid,
  96. "name": "WP03 semantic mapping source",
  97. "owner_uid": owner_uid,
  98. "actor_uid": actor_uid,
  99. },
  100. )
  101. db.session.execute(
  102. text(
  103. """
  104. INSERT INTO public.active_metadata_runs (
  105. uid, plan_uid, batch_key, status, attempt_count,
  106. cursor_before, cursor_after, statistics, actor_uid,
  107. started_at, finished_at
  108. ) VALUES (
  109. CAST(:uid AS uuid), CAST(:plan_uid AS uuid), :batch_key,
  110. 'completed', 1, '{}'::jsonb, '{}'::jsonb,
  111. '{}'::jsonb, CAST(:actor_uid AS uuid),
  112. CURRENT_TIMESTAMP, CURRENT_TIMESTAMP
  113. )
  114. """
  115. ),
  116. {
  117. "uid": run_uid,
  118. "plan_uid": plan_uid,
  119. "batch_key": f"wp03-{run_uid}",
  120. "actor_uid": actor_uid,
  121. },
  122. )
  123. snapshot = {
  124. "name": "customers",
  125. "namespace": "public",
  126. "asset_type": "table",
  127. "fields": [{"name": "customer_level", "data_type": "varchar"}],
  128. }
  129. db.session.execute(
  130. text(
  131. """
  132. INSERT INTO public.active_metadata_assets (
  133. uid, source_uid, asset_key, namespace, name, asset_type,
  134. lifecycle_status, current_version, content_hash,
  135. snapshot, health, last_run_uid
  136. ) VALUES (
  137. CAST(:uid AS uuid), CAST(:source_uid AS uuid), :asset_key,
  138. 'public', 'customers', 'table', 'active', 1,
  139. :content_hash, CAST(:snapshot AS jsonb), '{}'::jsonb,
  140. CAST(:run_uid AS uuid)
  141. )
  142. """
  143. ),
  144. {
  145. "uid": physical_asset_uid,
  146. "source_uid": source_uid,
  147. "asset_key": f"{source_uid}:public.customers",
  148. "content_hash": "a" * 64,
  149. "snapshot": json.dumps(snapshot),
  150. "run_uid": run_uid,
  151. },
  152. )
  153. db.session.execute(
  154. text(
  155. """
  156. INSERT INTO public.data_elements (
  157. uid, code, current_version, status,
  158. business_domain_uids, created_by
  159. ) VALUES (
  160. CAST(:uid AS uuid), :code, 1, 'published',
  161. CAST(:domains AS jsonb), :created_by
  162. );
  163. INSERT INTO public.data_element_versions (
  164. uid, data_element_uid, version, status, snapshot,
  165. evidence_uids, created_by
  166. ) VALUES (
  167. CAST(:version_uid AS uuid), CAST(:uid AS uuid), 1,
  168. 'published', CAST(:snapshot AS jsonb), '[]'::jsonb,
  169. :created_by
  170. )
  171. """
  172. ),
  173. {
  174. "uid": element_uid,
  175. "version_uid": element_version_uid,
  176. "code": f"CUSTOMER_LEVEL_{element_uid[:8]}",
  177. "domains": json.dumps([source_uid]),
  178. "created_by": actor_uid,
  179. "snapshot": json.dumps({"name": "客户等级"}),
  180. },
  181. )
  182. db.session.execute(
  183. text(
  184. """
  185. INSERT INTO public.data_standards (
  186. id, standard_uid, name, owner_uid, status
  187. ) VALUES (
  188. CAST(:id AS uuid), CAST(:standard_uid AS uuid),
  189. '客户主数据标准', CAST(:owner_uid AS uuid), 'active'
  190. );
  191. INSERT INTO public.data_standard_versions (
  192. id, standard_uid, version_no, source_text,
  193. standard_spec, spec_hash, scope, status, created_by,
  194. published_at
  195. ) VALUES (
  196. CAST(:version_uid AS uuid),
  197. CAST(:standard_uid AS uuid), 1, '客户等级必须使用统一代码',
  198. '{}'::jsonb, :spec_hash, '{}'::jsonb, 'validated',
  199. CAST(:actor_uid AS uuid), NULL
  200. )
  201. """
  202. ),
  203. {
  204. "id": standard_id,
  205. "standard_uid": standard_uid,
  206. "version_uid": standard_version_uid,
  207. "owner_uid": owner_uid,
  208. "actor_uid": actor_uid,
  209. "spec_hash": "b" * 64,
  210. },
  211. )
  212. db.session.commit()
  213. standard = DataRuleRepository(db.session).publish_standard_version(
  214. version_id=standard_version_uid,
  215. published_by=actor_uid,
  216. )
  217. db.session.commit()
  218. assert standard["status"] == "published"
  219. repository = SqlAlchemySemanticGovernanceRepository(db.session)
  220. service = SemanticGovernanceService(
  221. repository,
  222. outbox_enqueue=lambda **event: enqueue_outbox(
  223. db.session, **event
  224. ),
  225. )
  226. code_set = service.create_draft(
  227. "code_set",
  228. {
  229. "code": f"CUSTOMER_LEVEL_{source_uid[:8]}",
  230. "name": "客户等级代码集",
  231. "owner_uid": owner_uid,
  232. "business_domain_uid": source_uid,
  233. "definition": "客户等级统一参考数据",
  234. "values": [
  235. {"code": "VIP", "name": "重要客户"},
  236. {
  237. "code": "VIP_A",
  238. "name": "A级重要客户",
  239. "parent_code": "VIP",
  240. },
  241. ],
  242. "standard_version_uids": [standard_version_uid],
  243. },
  244. actor_uid=actor_uid,
  245. )
  246. semantic_uids.append(code_set["uid"])
  247. service.submit(code_set["uid"], expected_version=1, actor_uid=actor_uid)
  248. service.review(
  249. code_set["uid"],
  250. expected_version=1,
  251. decision="approve",
  252. reason="参考数据责任人确认",
  253. actor_uid=reviewer_uid,
  254. )
  255. service.publish(
  256. code_set["uid"], expected_version=1, actor_uid=reviewer_uid
  257. )
  258. term = service.create_draft(
  259. "business_term",
  260. {
  261. "code": f"TERM_CUSTOMER_LEVEL_{source_uid[:8]}",
  262. "name": "客户等级",
  263. "owner_uid": owner_uid,
  264. "business_domain_uid": source_uid,
  265. "definition": "客户运营分层等级",
  266. "aliases": ["客群等级", "客户分层"],
  267. "related_asset_uids": [physical_asset_uid],
  268. "related_data_element_uids": [element_uid],
  269. "standard_version_uids": [standard_version_uid],
  270. "code_references": [
  271. {"code_set_uid": code_set["uid"], "code": "VIP"}
  272. ],
  273. },
  274. actor_uid=actor_uid,
  275. )
  276. semantic_uids.append(term["uid"])
  277. service.submit(term["uid"], expected_version=1, actor_uid=actor_uid)
  278. service.review(
  279. term["uid"],
  280. expected_version=1,
  281. decision="approve",
  282. reason="术语定义审核通过",
  283. actor_uid=reviewer_uid,
  284. )
  285. service.publish(term["uid"], expected_version=1, actor_uid=reviewer_uid)
  286. metric_payload = {
  287. "code": f"METRIC_VIP_RATE_{source_uid[:8]}",
  288. "name": "重要客户占比",
  289. "owner_uid": owner_uid,
  290. "business_domain_uid": source_uid,
  291. "definition": "重要客户数占全部客户数的比例",
  292. "formula": "VIP客户数 / 全部客户数",
  293. "unit": "%",
  294. "dimensions": [{"code": "region", "name": "区域"}],
  295. "related_asset_uids": [physical_asset_uid],
  296. "standard_version_uids": [standard_version_uid],
  297. }
  298. metric = service.create_draft(
  299. "metric", metric_payload, actor_uid=actor_uid
  300. )
  301. semantic_uids.append(metric["uid"])
  302. for semantic_uid in (metric["uid"],):
  303. service.submit(
  304. semantic_uid, expected_version=1, actor_uid=actor_uid
  305. )
  306. service.review(
  307. semantic_uid,
  308. expected_version=1,
  309. decision="approve",
  310. reason="指标口径审核通过",
  311. actor_uid=reviewer_uid,
  312. )
  313. service.publish(
  314. semantic_uid,
  315. expected_version=1,
  316. actor_uid=reviewer_uid,
  317. )
  318. mapping = service.map_physical_field(
  319. {
  320. "asset_uid": physical_asset_uid,
  321. "field_name": "customer_level",
  322. "data_element_uid": element_uid,
  323. "owner_uid": owner_uid,
  324. "evidence": {"source": "catalog"},
  325. },
  326. actor_uid=actor_uid,
  327. )
  328. revised = service.revise(
  329. metric["uid"],
  330. {**metric_payload, "formula": "VIP客户数 / 有效客户数"},
  331. expected_version=1,
  332. reason="排除无效客户",
  333. actor_uid=actor_uid,
  334. )
  335. service.submit(
  336. metric["uid"],
  337. expected_version=2,
  338. actor_uid=actor_uid,
  339. )
  340. service.review(
  341. metric["uid"],
  342. expected_version=2,
  343. decision="approve",
  344. reason="口径修订审核通过",
  345. actor_uid=reviewer_uid,
  346. )
  347. service.publish(
  348. metric["uid"],
  349. expected_version=2,
  350. actor_uid=reviewer_uid,
  351. )
  352. rolled_back = service.rollback(
  353. metric["uid"],
  354. target_version=1,
  355. expected_version=2,
  356. reason="恢复已验收口径",
  357. actor_uid=reviewer_uid,
  358. )
  359. db.session.commit()
  360. assert mapping["data_element_version"] == 1
  361. assert revised["current_version"] == 2
  362. assert rolled_back["current_version"] == 3
  363. assert rolled_back["definition"]["formula"] == (
  364. "VIP客户数 / 全部客户数"
  365. )
  366. assert [
  367. item["status"]
  368. for item in service.list_versions(metric["uid"])
  369. ] == ["superseded", "superseded", "published"]
  370. assert len(service.list_audits(metric["uid"])) == 9
  371. event_count = db.session.execute(
  372. text(
  373. """
  374. SELECT count(*) FROM public.outbox_events
  375. WHERE aggregate_type = 'semantic_asset'
  376. AND aggregate_id = :aggregate_id
  377. AND event_type = 'semantic_asset.version_published'
  378. """
  379. ),
  380. {"aggregate_id": metric["uid"]},
  381. ).scalar_one()
  382. assert int(event_count) == 3
  383. standard_event_count = db.session.execute(
  384. text(
  385. """
  386. SELECT count(*) FROM public.outbox_events
  387. WHERE aggregate_type = 'data_standard'
  388. AND aggregate_id = :aggregate_id
  389. AND event_type = 'data_standard.version_published'
  390. """
  391. ),
  392. {"aggregate_id": standard_uid},
  393. ).scalar_one()
  394. assert int(standard_event_count) == 1
  395. claimed = claim_outbox(
  396. db.session,
  397. limit=10,
  398. event_types=["data_standard.version_published"],
  399. )
  400. assert len(claimed) == 1
  401. assert claimed[0].event_type == "data_standard.version_published"
  402. db.session.rollback()
  403. finally:
  404. with app.app_context():
  405. db.session.rollback()
  406. if semantic_uids:
  407. params = {"uids": semantic_uids}
  408. db.session.execute(
  409. text(
  410. "DELETE FROM public.outbox_events "
  411. "WHERE aggregate_type = 'semantic_asset' "
  412. "AND aggregate_id = ANY(:uids)"
  413. ),
  414. params,
  415. )
  416. for table in (
  417. "semantic_asset_reviews",
  418. "semantic_asset_links",
  419. "semantic_asset_versions",
  420. ):
  421. db.session.execute(
  422. text(
  423. f"DELETE FROM public.{table} "
  424. "WHERE asset_uid = ANY(CAST(:uids AS uuid[]))"
  425. ),
  426. params,
  427. )
  428. db.session.execute(
  429. text(
  430. "DELETE FROM public.semantic_publication_audits "
  431. "WHERE target_uid = ANY(CAST(:uids AS uuid[]))"
  432. ),
  433. params,
  434. )
  435. db.session.execute(
  436. text(
  437. "DELETE FROM public.semantic_assets "
  438. "WHERE uid = ANY(CAST(:uids AS uuid[]))"
  439. ),
  440. params,
  441. )
  442. db.session.execute(
  443. text(
  444. "DELETE FROM public.semantic_publication_audits "
  445. "WHERE target_type = 'field_mapping' "
  446. "AND target_uid IN ("
  447. "SELECT uid FROM public.data_element_field_mappings "
  448. "WHERE asset_uid = CAST(:asset_uid AS uuid))"
  449. ),
  450. {"asset_uid": physical_asset_uid},
  451. )
  452. db.session.execute(
  453. text(
  454. "DELETE FROM public.data_element_field_mappings "
  455. "WHERE asset_uid = CAST(:asset_uid AS uuid)"
  456. ),
  457. {"asset_uid": physical_asset_uid},
  458. )
  459. db.session.execute(
  460. text(
  461. "DELETE FROM public.outbox_events "
  462. "WHERE aggregate_type = 'data_standard' "
  463. "AND aggregate_id = :aggregate_id"
  464. ),
  465. {"aggregate_id": standard_uid},
  466. )
  467. db.session.execute(
  468. text(
  469. "DELETE FROM public.semantic_publication_audits "
  470. "WHERE target_type = 'data_standard' "
  471. "AND target_uid = CAST(:target_uid AS uuid)"
  472. ),
  473. {"target_uid": standard_uid},
  474. )
  475. db.session.execute(
  476. text(
  477. "DELETE FROM public.data_standard_versions "
  478. "WHERE id = CAST(:uid AS uuid)"
  479. ),
  480. {"uid": standard_version_uid},
  481. )
  482. db.session.execute(
  483. text(
  484. "DELETE FROM public.data_standards "
  485. "WHERE id = CAST(:uid AS uuid)"
  486. ),
  487. {"uid": standard_id},
  488. )
  489. db.session.execute(
  490. text(
  491. "DELETE FROM public.data_element_versions "
  492. "WHERE data_element_uid = CAST(:uid AS uuid)"
  493. ),
  494. {"uid": element_uid},
  495. )
  496. db.session.execute(
  497. text(
  498. "DELETE FROM public.data_elements "
  499. "WHERE uid = CAST(:uid AS uuid)"
  500. ),
  501. {"uid": element_uid},
  502. )
  503. db.session.execute(
  504. text(
  505. "DELETE FROM public.active_metadata_assets "
  506. "WHERE uid = CAST(:uid AS uuid)"
  507. ),
  508. {"uid": physical_asset_uid},
  509. )
  510. db.session.execute(
  511. text(
  512. "DELETE FROM public.active_metadata_runs "
  513. "WHERE uid = CAST(:uid AS uuid)"
  514. ),
  515. {"uid": run_uid},
  516. )
  517. db.session.execute(
  518. text(
  519. "DELETE FROM public.active_metadata_plans "
  520. "WHERE uid = CAST(:uid AS uuid)"
  521. ),
  522. {"uid": plan_uid},
  523. )
  524. db.session.execute(
  525. text(
  526. "DELETE FROM public.ingestion_sources "
  527. "WHERE uid = CAST(:uid AS uuid)"
  528. ),
  529. {"uid": source_uid},
  530. )
  531. db.session.execute(
  532. text(
  533. "DELETE FROM public.users "
  534. "WHERE id = ANY(CAST(:uids AS uuid[]))"
  535. ),
  536. {"uids": [actor_uid, owner_uid, reviewer_uid]},
  537. )
  538. db.session.commit()