test_semantic_governance.py 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456
  1. from __future__ import annotations
  2. from copy import deepcopy
  3. from datetime import UTC, datetime
  4. import pytest
  5. ACTOR_UID = "01900000-0000-7000-8000-000000001001"
  6. REVIEWER_UID = "01900000-0000-7000-8000-000000001002"
  7. OWNER_UID = "01900000-0000-7000-8000-000000001003"
  8. DOMAIN_UID = "01900000-0000-7000-8000-000000001004"
  9. SOURCE_ASSET_UID = "01900000-0000-7000-8000-000000001005"
  10. DATA_ELEMENT_UID = "01900000-0000-7000-8000-000000001006"
  11. STANDARD_VERSION_UID = "01900000-0000-7000-8000-000000001007"
  12. class MemorySemanticRepository:
  13. def __init__(self):
  14. self.assets = {}
  15. self.versions = {}
  16. self.reviews = []
  17. self.audits = []
  18. self.mappings = {}
  19. self.active_metadata_assets = {
  20. SOURCE_ASSET_UID: {
  21. "uid": SOURCE_ASSET_UID,
  22. "snapshot": {
  23. "fields": [
  24. {"name": "material_code", "data_type": "varchar"},
  25. {"name": "status", "data_type": "varchar"},
  26. ]
  27. },
  28. }
  29. }
  30. self.data_elements = {
  31. DATA_ELEMENT_UID: {
  32. "uid": DATA_ELEMENT_UID,
  33. "status": "published",
  34. "current_version": 3,
  35. }
  36. }
  37. self.standard_versions = {
  38. STANDARD_VERSION_UID: {
  39. "id": STANDARD_VERSION_UID,
  40. "status": "published",
  41. }
  42. }
  43. def get_by_code(self, asset_kind, code):
  44. return next(
  45. (
  46. deepcopy(item)
  47. for item in self.assets.values()
  48. if item["asset_kind"] == asset_kind and item["code"] == code
  49. ),
  50. None,
  51. )
  52. def save_new(self, asset, version, audit, links):
  53. self.assets[asset["uid"]] = deepcopy(asset)
  54. self.versions[(asset["uid"], version["version"])] = deepcopy(version)
  55. self.audits.append(deepcopy(audit))
  56. self._links = deepcopy(links)
  57. return self.get(asset["uid"])
  58. def get(self, uid):
  59. item = self.assets.get(uid)
  60. if item is None:
  61. return None
  62. version = self.versions[(uid, item["current_version"])]
  63. return {**deepcopy(item), "definition": deepcopy(version["definition"])}
  64. def list(self, *, asset_kind=None, status=None):
  65. items = [self.get(uid) for uid in self.assets]
  66. if asset_kind:
  67. items = [item for item in items if item["asset_kind"] == asset_kind]
  68. if status:
  69. items = [item for item in items if item["status"] == status]
  70. return items
  71. def list_versions(self, uid):
  72. return [
  73. deepcopy(version)
  74. for (asset_uid, _version), version in sorted(self.versions.items())
  75. if asset_uid == uid
  76. ]
  77. def list_reviews(self, uid):
  78. return [deepcopy(item) for item in self.reviews if item["asset_uid"] == uid]
  79. def list_audits(self, uid):
  80. return [deepcopy(item) for item in self.audits if item["target_uid"] == uid]
  81. def transition(self, asset, version, audit, review=None):
  82. self.assets[asset["uid"]] = deepcopy(asset)
  83. self.versions[(asset["uid"], version["version"])] = deepcopy(version)
  84. self.audits.append(deepcopy(audit))
  85. if review:
  86. self.reviews.append(deepcopy(review))
  87. return self.get(asset["uid"])
  88. def append_version(self, asset, version, audit, links):
  89. self.assets[asset["uid"]] = deepcopy(asset)
  90. self.versions[(asset["uid"], version["version"])] = deepcopy(version)
  91. self.audits.append(deepcopy(audit))
  92. self._links = deepcopy(links)
  93. return self.get(asset["uid"])
  94. def active_assets_exist(self, uids):
  95. return set(uids).issubset(self.active_metadata_assets)
  96. def published_data_elements_exist(self, uids):
  97. return all(
  98. self.data_elements.get(uid, {}).get("status") == "published"
  99. for uid in uids
  100. )
  101. def published_standard_versions_exist(self, uids):
  102. return all(
  103. self.standard_versions.get(uid, {}).get("status") == "published"
  104. for uid in uids
  105. )
  106. def published_code_exists(self, code_set_uid, code):
  107. asset = self.assets.get(code_set_uid)
  108. if not asset or asset["status"] != "published":
  109. return False
  110. definition = self.versions[
  111. (code_set_uid, asset["current_version"])
  112. ]["definition"]
  113. return code in {item["code"] for item in definition["values"]}
  114. def get_active_metadata_asset(self, uid):
  115. return deepcopy(self.active_metadata_assets.get(uid))
  116. def get_data_element(self, uid):
  117. return deepcopy(self.data_elements.get(uid))
  118. def save_field_mapping(self, mapping, audit):
  119. key = (
  120. mapping["asset_uid"],
  121. mapping["field_name"],
  122. mapping["data_element_uid"],
  123. )
  124. if key in self.mappings:
  125. return deepcopy(self.mappings[key])
  126. self.mappings[key] = deepcopy(mapping)
  127. self.audits.append(deepcopy(audit))
  128. return deepcopy(mapping)
  129. def list_field_mappings(self, *, asset_uid=None, data_element_uid=None):
  130. result = list(self.mappings.values())
  131. if asset_uid:
  132. result = [item for item in result if item["asset_uid"] == asset_uid]
  133. if data_element_uid:
  134. result = [
  135. item
  136. for item in result
  137. if item["data_element_uid"] == data_element_uid
  138. ]
  139. return deepcopy(result)
  140. def _service(repository=None, events=None):
  141. from app.core.data_research.semantic_governance import (
  142. SemanticGovernanceService,
  143. )
  144. ids = (f"01900000-0000-7000-8000-{index:012d}" for index in range(2001, 9999))
  145. event_list = events if events is not None else []
  146. return SemanticGovernanceService(
  147. repository or MemorySemanticRepository(),
  148. uid_factory=ids.__next__,
  149. now_factory=lambda: datetime(2026, 7, 31, 9, 0, tzinfo=UTC),
  150. outbox_enqueue=lambda **event: event_list.append(deepcopy(event)),
  151. )
  152. def _term_payload(**overrides):
  153. payload = {
  154. "code": "TERM_MATERIAL",
  155. "name": "物料",
  156. "owner_uid": OWNER_UID,
  157. "business_domain_uid": DOMAIN_UID,
  158. "definition": "企业采购、库存和维修过程中统一识别的物资。",
  159. "aliases": ["备件", "物资"],
  160. "related_asset_uids": [SOURCE_ASSET_UID],
  161. "related_data_element_uids": [DATA_ELEMENT_UID],
  162. "standard_version_uids": [STANDARD_VERSION_UID],
  163. }
  164. payload.update(overrides)
  165. return payload
  166. def _publish(service, asset):
  167. submitted = service.submit(
  168. asset["uid"],
  169. expected_version=asset["current_version"],
  170. actor_uid=ACTOR_UID,
  171. )
  172. approved = service.review(
  173. asset["uid"],
  174. expected_version=submitted["current_version"],
  175. decision="approve",
  176. reason="定义和关联证据完整",
  177. actor_uid=REVIEWER_UID,
  178. )
  179. return service.publish(
  180. asset["uid"],
  181. expected_version=approved["current_version"],
  182. actor_uid=REVIEWER_UID,
  183. )
  184. def test_business_term_uses_distinct_review_and_publishes_knowledge_event():
  185. repository = MemorySemanticRepository()
  186. events = []
  187. service = _service(repository, events)
  188. draft = service.create_draft(
  189. "business_term",
  190. _term_payload(),
  191. actor_uid=ACTOR_UID,
  192. )
  193. with pytest.raises(ValueError, match="self-review"):
  194. service.review(
  195. draft["uid"],
  196. expected_version=1,
  197. decision="approve",
  198. reason="不能自审",
  199. actor_uid=ACTOR_UID,
  200. )
  201. published = _publish(service, draft)
  202. assert published["status"] == "published"
  203. assert published["definition"]["aliases"] == ["备件", "物资"]
  204. assert [item["action"] for item in repository.list_audits(draft["uid"])] == [
  205. "created",
  206. "submitted",
  207. "approved",
  208. "published",
  209. ]
  210. assert events == [
  211. {
  212. "aggregate_type": "semantic_asset",
  213. "aggregate_id": draft["uid"],
  214. "event_type": "semantic_asset.version_published",
  215. "payload": {
  216. "uid": draft["uid"],
  217. "asset_kind": "business_term",
  218. "version": 1,
  219. },
  220. }
  221. ]
  222. def test_code_set_validates_parent_and_cross_asset_references():
  223. repository = MemorySemanticRepository()
  224. service = _service(repository)
  225. with pytest.raises(ValueError, match="parent_code"):
  226. service.create_draft(
  227. "code_set",
  228. {
  229. "code": "MATERIAL_STATUS",
  230. "name": "物料状态",
  231. "owner_uid": OWNER_UID,
  232. "business_domain_uid": DOMAIN_UID,
  233. "definition": "物料生命周期状态。",
  234. "values": [
  235. {"code": "ACTIVE", "name": "有效", "parent_code": "MISSING"}
  236. ],
  237. },
  238. actor_uid=ACTOR_UID,
  239. )
  240. code_set = service.create_draft(
  241. "code_set",
  242. {
  243. "code": "MATERIAL_STATUS",
  244. "name": "物料状态",
  245. "owner_uid": OWNER_UID,
  246. "business_domain_uid": DOMAIN_UID,
  247. "definition": "物料生命周期状态。",
  248. "values": [
  249. {"code": "ALL", "name": "全部"},
  250. {"code": "ACTIVE", "name": "有效", "parent_code": "ALL"},
  251. {"code": "FROZEN", "name": "冻结", "parent_code": "ALL"},
  252. ],
  253. },
  254. actor_uid=ACTOR_UID,
  255. )
  256. published_codes = _publish(service, code_set)
  257. term = service.create_draft(
  258. "business_term",
  259. _term_payload(
  260. code="TERM_ACTIVE_MATERIAL",
  261. code_references=[
  262. {"code_set_uid": published_codes["uid"], "code": "ACTIVE"}
  263. ],
  264. ),
  265. actor_uid=ACTOR_UID,
  266. )
  267. assert term["definition"]["code_references"][0]["code"] == "ACTIVE"
  268. with pytest.raises(ValueError, match="published code"):
  269. service.create_draft(
  270. "business_term",
  271. _term_payload(
  272. code="TERM_BAD_STATUS",
  273. code_references=[
  274. {"code_set_uid": published_codes["uid"], "code": "UNKNOWN"}
  275. ],
  276. ),
  277. actor_uid=ACTOR_UID,
  278. )
  279. def test_metric_definition_keeps_dimensions_and_formula_without_execution():
  280. service = _service()
  281. metric = service.create_draft(
  282. "metric",
  283. {
  284. "code": "MATERIAL_COMPLETENESS",
  285. "name": "物料主数据完整率",
  286. "owner_uid": OWNER_UID,
  287. "business_domain_uid": DOMAIN_UID,
  288. "definition": "关键字段完整的有效物料占全部有效物料的比例。",
  289. "formula": "关键字段完整的有效物料数 / 全部有效物料数 * 100%",
  290. "unit": "%",
  291. "dimensions": [
  292. {"code": "category", "name": "物料分类"},
  293. {"code": "organization", "name": "组织"},
  294. ],
  295. "related_asset_uids": [SOURCE_ASSET_UID],
  296. "related_data_element_uids": [DATA_ELEMENT_UID],
  297. },
  298. actor_uid=ACTOR_UID,
  299. )
  300. assert metric["definition"]["dimensions"][0]["code"] == "category"
  301. assert "execution_sql" not in metric["definition"]
  302. with pytest.raises(ValueError, match="analysis execution"):
  303. service.create_draft(
  304. "metric",
  305. {
  306. **_term_payload(code="BAD_METRIC"),
  307. "formula": "select count(*) from material",
  308. "execution_sql": "select count(*) from material",
  309. },
  310. actor_uid=ACTOR_UID,
  311. )
  312. def test_rollback_appends_published_version_and_preserves_history():
  313. repository = MemorySemanticRepository()
  314. events = []
  315. service = _service(repository, events)
  316. first = _publish(
  317. service,
  318. service.create_draft("business_term", _term_payload(), actor_uid=ACTOR_UID),
  319. )
  320. revised = service.revise(
  321. first["uid"],
  322. _term_payload(definition="第二版定义。"),
  323. expected_version=1,
  324. reason="更新口径",
  325. actor_uid=ACTOR_UID,
  326. )
  327. second = _publish(service, revised)
  328. rolled_back = service.rollback(
  329. second["uid"],
  330. target_version=1,
  331. expected_version=2,
  332. reason="第二版口径不适用",
  333. actor_uid=REVIEWER_UID,
  334. )
  335. assert rolled_back["current_version"] == 3
  336. assert rolled_back["status"] == "published"
  337. assert rolled_back["definition"]["definition"] == _term_payload()["definition"]
  338. versions = repository.list_versions(first["uid"])
  339. assert [item["version"] for item in versions] == [1, 2, 3]
  340. assert versions[-1]["rollback_from_version"] == 1
  341. assert repository.list_audits(first["uid"])[-1]["action"] == "rolled_back"
  342. assert events[-1]["payload"]["version"] == 3
  343. def test_physical_field_mapping_requires_published_element_and_real_field():
  344. repository = MemorySemanticRepository()
  345. service = _service(repository)
  346. mapping = service.map_physical_field(
  347. {
  348. "asset_uid": SOURCE_ASSET_UID,
  349. "field_name": "material_code",
  350. "data_element_uid": DATA_ELEMENT_UID,
  351. "owner_uid": OWNER_UID,
  352. "evidence": {"catalog_snapshot_uid": "snapshot-1"},
  353. },
  354. actor_uid=ACTOR_UID,
  355. )
  356. assert mapping["status"] == "published"
  357. assert mapping["data_element_version"] == 3
  358. assert service.list_field_mappings(asset_uid=SOURCE_ASSET_UID) == [mapping]
  359. with pytest.raises(ValueError, match="physical field"):
  360. service.map_physical_field(
  361. {
  362. "asset_uid": SOURCE_ASSET_UID,
  363. "field_name": "missing_field",
  364. "data_element_uid": DATA_ELEMENT_UID,
  365. "owner_uid": OWNER_UID,
  366. },
  367. actor_uid=ACTOR_UID,
  368. )
  369. repository.data_elements[DATA_ELEMENT_UID]["status"] = "draft"
  370. with pytest.raises(ValueError, match="published data element"):
  371. service.map_physical_field(
  372. {
  373. "asset_uid": SOURCE_ASSET_UID,
  374. "field_name": "status",
  375. "data_element_uid": DATA_ELEMENT_UID,
  376. "owner_uid": OWNER_UID,
  377. },
  378. actor_uid=ACTOR_UID,
  379. )
  380. def test_semantic_exchange_exports_only_published_assets_as_json_and_rdf():
  381. from app.core.data_research.semantic_governance import SemanticExchange
  382. repository = MemorySemanticRepository()
  383. service = _service(repository)
  384. published = _publish(
  385. service,
  386. service.create_draft("business_term", _term_payload(), actor_uid=ACTOR_UID),
  387. )
  388. service.create_draft(
  389. "business_term",
  390. _term_payload(code="TERM_DRAFT"),
  391. actor_uid=ACTOR_UID,
  392. )
  393. exchange = SemanticExchange()
  394. json_document = exchange.export_json(service.list(status="published"))
  395. rdf_document = exchange.export_rdf(service.list(status="published"))
  396. assert published["code"].encode() in json_document
  397. assert b"TERM_DRAFT" not in json_document
  398. assert b"rdf:RDF" in rdf_document
  399. assert b"BusinessTerm" in rdf_document
  400. assert b"password" not in json_document.lower()
  401. assert b"password" not in rdf_document.lower()