| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456 |
- 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()
|