from __future__ import annotations from copy import deepcopy from datetime import UTC, datetime import pytest ACTOR_UID = "01900000-0000-7000-8000-000000001001" REVIEWER_UID = "01900000-0000-7000-8000-000000001002" OWNER_UID = "01900000-0000-7000-8000-000000001003" DOMAIN_UID = "01900000-0000-7000-8000-000000001004" SOURCE_ASSET_UID = "01900000-0000-7000-8000-000000001005" DATA_ELEMENT_UID = "01900000-0000-7000-8000-000000001006" STANDARD_VERSION_UID = "01900000-0000-7000-8000-000000001007" class MemorySemanticRepository: def __init__(self): self.assets = {} self.versions = {} self.reviews = [] self.audits = [] self.mappings = {} self.active_metadata_assets = { SOURCE_ASSET_UID: { "uid": SOURCE_ASSET_UID, "snapshot": { "fields": [ {"name": "material_code", "data_type": "varchar"}, {"name": "status", "data_type": "varchar"}, ] }, } } self.data_elements = { DATA_ELEMENT_UID: { "uid": DATA_ELEMENT_UID, "status": "published", "current_version": 3, } } self.standard_versions = { STANDARD_VERSION_UID: { "id": STANDARD_VERSION_UID, "status": "published", } } def get_by_code(self, asset_kind, code): return next( ( deepcopy(item) for item in self.assets.values() if item["asset_kind"] == asset_kind and item["code"] == code ), None, ) def save_new(self, asset, version, audit, links): self.assets[asset["uid"]] = deepcopy(asset) self.versions[(asset["uid"], version["version"])] = deepcopy(version) self.audits.append(deepcopy(audit)) self._links = deepcopy(links) return self.get(asset["uid"]) def get(self, uid): item = self.assets.get(uid) if item is None: return None version = self.versions[(uid, item["current_version"])] return {**deepcopy(item), "definition": deepcopy(version["definition"])} def list(self, *, asset_kind=None, status=None): items = [self.get(uid) for uid in self.assets] if asset_kind: items = [item for item in items if item["asset_kind"] == asset_kind] if status: items = [item for item in items if item["status"] == status] return items def list_versions(self, uid): return [ deepcopy(version) for (asset_uid, _version), version in sorted(self.versions.items()) if asset_uid == uid ] def list_reviews(self, uid): return [deepcopy(item) for item in self.reviews if item["asset_uid"] == uid] def list_audits(self, uid): return [deepcopy(item) for item in self.audits if item["target_uid"] == uid] def transition(self, asset, version, audit, review=None): self.assets[asset["uid"]] = deepcopy(asset) self.versions[(asset["uid"], version["version"])] = deepcopy(version) self.audits.append(deepcopy(audit)) if review: self.reviews.append(deepcopy(review)) return self.get(asset["uid"]) def append_version(self, asset, version, audit, links): self.assets[asset["uid"]] = deepcopy(asset) self.versions[(asset["uid"], version["version"])] = deepcopy(version) self.audits.append(deepcopy(audit)) self._links = deepcopy(links) return self.get(asset["uid"]) def active_assets_exist(self, uids): return set(uids).issubset(self.active_metadata_assets) def published_data_elements_exist(self, uids): return all( self.data_elements.get(uid, {}).get("status") == "published" for uid in uids ) def published_standard_versions_exist(self, uids): return all( self.standard_versions.get(uid, {}).get("status") == "published" for uid in uids ) def published_code_exists(self, code_set_uid, code): asset = self.assets.get(code_set_uid) if not asset or asset["status"] != "published": return False definition = self.versions[ (code_set_uid, asset["current_version"]) ]["definition"] return code in {item["code"] for item in definition["values"]} def get_active_metadata_asset(self, uid): return deepcopy(self.active_metadata_assets.get(uid)) def get_data_element(self, uid): return deepcopy(self.data_elements.get(uid)) def save_field_mapping(self, mapping, audit): key = ( mapping["asset_uid"], mapping["field_name"], mapping["data_element_uid"], ) if key in self.mappings: return deepcopy(self.mappings[key]) self.mappings[key] = deepcopy(mapping) self.audits.append(deepcopy(audit)) return deepcopy(mapping) def list_field_mappings(self, *, asset_uid=None, data_element_uid=None): result = list(self.mappings.values()) if asset_uid: result = [item for item in result if item["asset_uid"] == asset_uid] if data_element_uid: result = [ item for item in result if item["data_element_uid"] == data_element_uid ] return deepcopy(result) def _service(repository=None, events=None): from app.core.data_research.semantic_governance import ( SemanticGovernanceService, ) ids = (f"01900000-0000-7000-8000-{index:012d}" for index in range(2001, 9999)) event_list = events if events is not None else [] return SemanticGovernanceService( repository or MemorySemanticRepository(), uid_factory=ids.__next__, now_factory=lambda: datetime(2026, 7, 31, 9, 0, tzinfo=UTC), outbox_enqueue=lambda **event: event_list.append(deepcopy(event)), ) def _term_payload(**overrides): payload = { "code": "TERM_MATERIAL", "name": "物料", "owner_uid": OWNER_UID, "business_domain_uid": DOMAIN_UID, "definition": "企业采购、库存和维修过程中统一识别的物资。", "aliases": ["备件", "物资"], "related_asset_uids": [SOURCE_ASSET_UID], "related_data_element_uids": [DATA_ELEMENT_UID], "standard_version_uids": [STANDARD_VERSION_UID], } payload.update(overrides) return payload def _publish(service, asset): submitted = service.submit( asset["uid"], expected_version=asset["current_version"], actor_uid=ACTOR_UID, ) approved = service.review( asset["uid"], expected_version=submitted["current_version"], decision="approve", reason="定义和关联证据完整", actor_uid=REVIEWER_UID, ) return service.publish( asset["uid"], expected_version=approved["current_version"], actor_uid=REVIEWER_UID, ) def test_business_term_uses_distinct_review_and_publishes_knowledge_event(): repository = MemorySemanticRepository() events = [] service = _service(repository, events) draft = service.create_draft( "business_term", _term_payload(), actor_uid=ACTOR_UID, ) with pytest.raises(ValueError, match="self-review"): service.review( draft["uid"], expected_version=1, decision="approve", reason="不能自审", actor_uid=ACTOR_UID, ) published = _publish(service, draft) assert published["status"] == "published" assert published["definition"]["aliases"] == ["备件", "物资"] assert [item["action"] for item in repository.list_audits(draft["uid"])] == [ "created", "submitted", "approved", "published", ] assert events == [ { "aggregate_type": "semantic_asset", "aggregate_id": draft["uid"], "event_type": "semantic_asset.version_published", "payload": { "uid": draft["uid"], "asset_kind": "business_term", "version": 1, }, } ] def test_code_set_validates_parent_and_cross_asset_references(): repository = MemorySemanticRepository() service = _service(repository) with pytest.raises(ValueError, match="parent_code"): service.create_draft( "code_set", { "code": "MATERIAL_STATUS", "name": "物料状态", "owner_uid": OWNER_UID, "business_domain_uid": DOMAIN_UID, "definition": "物料生命周期状态。", "values": [ {"code": "ACTIVE", "name": "有效", "parent_code": "MISSING"} ], }, actor_uid=ACTOR_UID, ) code_set = service.create_draft( "code_set", { "code": "MATERIAL_STATUS", "name": "物料状态", "owner_uid": OWNER_UID, "business_domain_uid": DOMAIN_UID, "definition": "物料生命周期状态。", "values": [ {"code": "ALL", "name": "全部"}, {"code": "ACTIVE", "name": "有效", "parent_code": "ALL"}, {"code": "FROZEN", "name": "冻结", "parent_code": "ALL"}, ], }, actor_uid=ACTOR_UID, ) published_codes = _publish(service, code_set) term = service.create_draft( "business_term", _term_payload( code="TERM_ACTIVE_MATERIAL", code_references=[ {"code_set_uid": published_codes["uid"], "code": "ACTIVE"} ], ), actor_uid=ACTOR_UID, ) assert term["definition"]["code_references"][0]["code"] == "ACTIVE" with pytest.raises(ValueError, match="published code"): service.create_draft( "business_term", _term_payload( code="TERM_BAD_STATUS", code_references=[ {"code_set_uid": published_codes["uid"], "code": "UNKNOWN"} ], ), actor_uid=ACTOR_UID, ) def test_metric_definition_keeps_dimensions_and_formula_without_execution(): service = _service() metric = service.create_draft( "metric", { "code": "MATERIAL_COMPLETENESS", "name": "物料主数据完整率", "owner_uid": OWNER_UID, "business_domain_uid": DOMAIN_UID, "definition": "关键字段完整的有效物料占全部有效物料的比例。", "formula": "关键字段完整的有效物料数 / 全部有效物料数 * 100%", "unit": "%", "dimensions": [ {"code": "category", "name": "物料分类"}, {"code": "organization", "name": "组织"}, ], "related_asset_uids": [SOURCE_ASSET_UID], "related_data_element_uids": [DATA_ELEMENT_UID], }, actor_uid=ACTOR_UID, ) assert metric["definition"]["dimensions"][0]["code"] == "category" assert "execution_sql" not in metric["definition"] with pytest.raises(ValueError, match="analysis execution"): service.create_draft( "metric", { **_term_payload(code="BAD_METRIC"), "formula": "select count(*) from material", "execution_sql": "select count(*) from material", }, actor_uid=ACTOR_UID, ) def test_rollback_appends_published_version_and_preserves_history(): repository = MemorySemanticRepository() events = [] service = _service(repository, events) first = _publish( service, service.create_draft("business_term", _term_payload(), actor_uid=ACTOR_UID), ) revised = service.revise( first["uid"], _term_payload(definition="第二版定义。"), expected_version=1, reason="更新口径", actor_uid=ACTOR_UID, ) second = _publish(service, revised) rolled_back = service.rollback( second["uid"], target_version=1, expected_version=2, reason="第二版口径不适用", actor_uid=REVIEWER_UID, ) assert rolled_back["current_version"] == 3 assert rolled_back["status"] == "published" assert rolled_back["definition"]["definition"] == _term_payload()["definition"] versions = repository.list_versions(first["uid"]) assert [item["version"] for item in versions] == [1, 2, 3] assert versions[-1]["rollback_from_version"] == 1 assert repository.list_audits(first["uid"])[-1]["action"] == "rolled_back" assert events[-1]["payload"]["version"] == 3 def test_physical_field_mapping_requires_published_element_and_real_field(): repository = MemorySemanticRepository() service = _service(repository) mapping = service.map_physical_field( { "asset_uid": SOURCE_ASSET_UID, "field_name": "material_code", "data_element_uid": DATA_ELEMENT_UID, "owner_uid": OWNER_UID, "evidence": {"catalog_snapshot_uid": "snapshot-1"}, }, actor_uid=ACTOR_UID, ) assert mapping["status"] == "published" assert mapping["data_element_version"] == 3 assert service.list_field_mappings(asset_uid=SOURCE_ASSET_UID) == [mapping] with pytest.raises(ValueError, match="physical field"): service.map_physical_field( { "asset_uid": SOURCE_ASSET_UID, "field_name": "missing_field", "data_element_uid": DATA_ELEMENT_UID, "owner_uid": OWNER_UID, }, actor_uid=ACTOR_UID, ) repository.data_elements[DATA_ELEMENT_UID]["status"] = "draft" with pytest.raises(ValueError, match="published data element"): service.map_physical_field( { "asset_uid": SOURCE_ASSET_UID, "field_name": "status", "data_element_uid": DATA_ELEMENT_UID, "owner_uid": OWNER_UID, }, actor_uid=ACTOR_UID, ) def test_semantic_exchange_exports_only_published_assets_as_json_and_rdf(): from app.core.data_research.semantic_governance import SemanticExchange repository = MemorySemanticRepository() service = _service(repository) published = _publish( service, service.create_draft("business_term", _term_payload(), actor_uid=ACTOR_UID), ) service.create_draft( "business_term", _term_payload(code="TERM_DRAFT"), actor_uid=ACTOR_UID, ) exchange = SemanticExchange() json_document = exchange.export_json(service.list(status="published")) rdf_document = exchange.export_rdf(service.list(status="published")) assert published["code"].encode() in json_document assert b"TERM_DRAFT" not in json_document assert b"rdf:RDF" in rdf_document assert b"BusinessTerm" in rdf_document assert b"password" not in json_document.lower() assert b"password" not in rdf_document.lower()