from __future__ import annotations from dataclasses import replace from datetime import datetime import pytest from app.core.data_research.device_assets import ( DeviceAssetDetail, DeviceAssetMappingRecord, DeviceAssetRecord, ) SOURCE_A = "00000000-0000-0000-0000-000000000101" SOURCE_B = "00000000-0000-0000-0000-000000000102" LEFT_UID = "00000000-0000-7000-8000-000000000201" RIGHT_UID = "00000000-0000-7000-8000-000000000202" def asset( uid, source_uid, source_code, *, name="一号循环泵", asset_type="device", location="动力车间", organization="设备动力部", responsible_person="张工", model="P-100", ): now = datetime(2026, 7, 29, 9, 0) record = DeviceAssetRecord( uid=uid, asset_type=asset_type, name=name, status="active", current_version=1, content_hash="a" * 64, location=location, organization=organization, responsible_person=responsible_person, attributes={"model": model}, created_by="editor-1", updated_by="editor-1", created_at=now, updated_at=now, ) mapping = DeviceAssetMappingRecord( uid=f"{uid[:-3]}3{uid[-2:]}", asset_uid=uid, source_uid=source_uid, source_entity="asset.equipment", asset_type=asset_type, source_code=source_code, source_updated_at=now, first_seen_at=now, last_seen_at=now, ) return DeviceAssetDetail(asset=record, mappings=(mapping,)) class MemoryResolutionRepository: def __init__(self, assets=()): self.assets = {item.asset.uid: item for item in assets} self.candidates = {} self.reviews = [] self.merge_records = [] self.rollback_records = [] def list_matchable_assets(self, asset_type, *, limit): return [ item for item in self.assets.values() if item.asset.asset_type == asset_type ][:limit] def get_asset_detail(self, uid): return self.assets.get(uid) def find_open_pair(self, left_uid, right_uid): return next( ( item for item in self.candidates.values() if item.left_asset_uid == left_uid and item.right_asset_uid == right_uid and item.status in {"pending", "merged"} ), None, ) def create_candidate(self, record): self.candidates[record.uid] = record return record def search_candidates(self, filters, *, page, page_size): records = list(self.candidates.values()) for name in ("status", "suggestion_source"): if filters.get(name): records = [ item for item in records if getattr(item, name) == filters[name] ] start = (page - 1) * page_size return records[start : start + page_size], len(records) def get_candidate(self, uid, *, for_update=False): del for_update return self.candidates.get(uid) def update_candidate(self, record): self.candidates[record.uid] = record return record def append_review(self, record): self.reviews.append(record) return record def active_merge_for_member(self, asset_uid): rolled_back = {item.merge_uid for item in self.rollback_records} return next( ( item for item in self.merge_records if item.member_asset_uid == asset_uid and item.uid not in rolled_back ), None, ) def create_merge(self, record): self.merge_records.append(record) return record def get_merge(self, uid, *, for_update=False): del for_update return next( (item for item in self.merge_records if item.uid == uid), None, ) def create_rollback(self, record): self.rollback_records.append(record) return record def list_reviews(self, candidate_uid): return [ item for item in self.reviews if item.candidate_uid == candidate_uid ] def list_merges(self, candidate_uid): return [ item for item in self.merge_records if item.candidate_uid == candidate_uid ] def list_rollbacks(self, merge_uid): return [ item for item in self.rollback_records if item.merge_uid == merge_uid ] def service(repository, *, authorized=True, auto_merge_enabled=False): from app.core.data_research.device_entity_resolution import ( DeviceEntityForbidden, DeviceEntityResolutionService, ) ids = iter( f"00000000-0000-7000-8000-0000000003{index:02d}" for index in range(1, 40) ) clock = iter( datetime(2026, 7, 29, 10, minute) for minute in range(1, 40) ) def authorize(actor_uid): if not authorized: raise DeviceEntityForbidden( f"{actor_uid} is not the accountable asset manager" ) return DeviceEntityResolutionService( repository, review_authorizer=authorize, uid_factory=ids.__next__, now_factory=clock.__next__, auto_merge_enabled=auto_merge_enabled, ) def test_score_is_explainable_and_uses_hand_derived_weights(): from app.core.data_research.device_entity_resolution import ( score_device_pair, ) left = asset(LEFT_UID, SOURCE_A, "EQ-001") right = asset( RIGHT_UID, SOURCE_B, "设备-001", name=" 一号 循环泵 ", responsible_person="李工", ) score = score_device_pair(left, right) assert score.confidence == pytest.approx(0.90) assert [item["signal"] for item in score.explanation] == [ "name", "location", "organization", "responsible_person", "model", "source_code", ] assert score.explanation[0]["matched"] is True assert score.explanation[3]["matched"] is False assert score.explanation[5]["matched"] is False def test_generation_is_cross_source_bounded_and_idempotent_for_open_pair(): left = asset(LEFT_UID, SOURCE_A, "EQ-001") right = asset(RIGHT_UID, SOURCE_B, "EQ-001") same_source = asset( "00000000-0000-7000-8000-000000000203", SOURCE_A, "EQ-003", name="三号空压机", location="空压站", organization="公用工程部", responsible_person="王工", model="AC-300", ) repository = MemoryResolutionRepository((left, right, same_source)) resolution = service(repository) first = resolution.generate( {"asset_type": "device", "threshold": 0.7, "limit": 20}, actor_uid="editor-1", ) second = resolution.generate( {"asset_type": "device", "threshold": 0.7, "limit": 20}, actor_uid="editor-1", ) assert first.created_count == 1 assert first.existing_count == 0 assert second.created_count == 0 assert second.existing_count == 1 assert first.records[0].left_asset_uid == LEFT_UID assert first.records[0].right_asset_uid == RIGHT_UID assert first.records[0].evidence_uids == ( LEFT_UID, left.mappings[0].uid, RIGHT_UID, right.mappings[0].uid, ) def test_ai_candidate_requires_provider_model_evidence_and_never_auto_merges(): from app.core.data_research.device_entity_resolution import ( DeviceEntityInvalid, ) repository = MemoryResolutionRepository( ( asset(LEFT_UID, SOURCE_A, "EQ-001"), asset(RIGHT_UID, SOURCE_B, "DEVICE-001"), ) ) resolution = service( repository, authorized=True, auto_merge_enabled=True, ) with pytest.raises(DeviceEntityInvalid, match="model_provider"): resolution.submit_ai_candidate( { "left_asset_uid": LEFT_UID, "right_asset_uid": RIGHT_UID, "confidence": 0.999, "model_name": "entity-match-v1", "evidence_uids": ["evidence-1"], "explanation": "同一台设备", }, actor_uid="editor-1", ) candidate = resolution.submit_ai_candidate( { "left_asset_uid": LEFT_UID, "right_asset_uid": RIGHT_UID, "confidence": 0.999, "model_provider": "governed-provider", "model_name": "entity-match-v1", "evidence_uids": ["evidence-1"], "explanation": "同一台设备", }, actor_uid="editor-1", ) assert candidate.status == "pending" assert candidate.suggestion_source == "ai" assert repository.merge_records == [] def test_review_requires_accountable_manager_and_appends_merge_evidence(): repository = MemoryResolutionRepository( ( asset(LEFT_UID, SOURCE_A, "EQ-001"), asset(RIGHT_UID, SOURCE_B, "EQ-001"), ) ) blocked = service(repository, authorized=False) candidate = blocked.generate( {"asset_type": "device", "threshold": 0.7}, actor_uid="editor-1", ).records[0] from app.core.data_research.device_entity_resolution import ( DeviceEntityForbidden, ) with pytest.raises(DeviceEntityForbidden, match="accountable"): blocked.review( candidate.uid, { "decision": "approve", "canonical_asset_uid": LEFT_UID, "expected_version": 1, "reason": "跨系统编码和型号一致", }, actor_uid="admin-1", ) approved, review, merge = service(repository).review( candidate.uid, { "decision": "approve", "canonical_asset_uid": LEFT_UID, "expected_version": 1, "reason": "跨系统编码和型号一致", }, actor_uid="admin-1", ) assert approved.status == "merged" assert approved.canonical_asset_uid == LEFT_UID assert review.decision == "approve" assert merge.canonical_asset_uid == LEFT_UID assert merge.member_asset_uid == RIGHT_UID assert merge.snapshot["member"]["uid"] == RIGHT_UID assert len(repository.reviews) == 1 assert len(repository.merge_records) == 1 def test_rollback_is_append_only_and_restores_candidate_state(): repository = MemoryResolutionRepository( ( asset(LEFT_UID, SOURCE_A, "EQ-001"), asset(RIGHT_UID, SOURCE_B, "EQ-001"), ) ) resolution = service(repository) candidate = resolution.generate( {"asset_type": "device", "threshold": 0.7}, actor_uid="editor-1", ).records[0] _candidate, _review, merge = resolution.review( candidate.uid, { "decision": "approve", "canonical_asset_uid": LEFT_UID, "expected_version": 1, "reason": "匹配证据充分", }, actor_uid="admin-1", ) rolled_back, rollback = resolution.rollback( merge.uid, { "expected_version": 2, "reason": "现场确认不是同一台设备", }, actor_uid="admin-1", ) assert rolled_back.status == "rolled_back" assert rolled_back.current_version == 3 assert rollback.merge_uid == merge.uid assert rollback.snapshot["merge"]["member_asset_uid"] == RIGHT_UID assert len(repository.merge_records) == 1 assert len(repository.rollback_records) == 1 assert repository.active_merge_for_member(RIGHT_UID) is None def test_rule_auto_merge_is_default_off_and_requires_strict_threshold(): left = asset(LEFT_UID, SOURCE_A, "EQ-001") right = asset(RIGHT_UID, SOURCE_B, "EQ-001") default_repository = MemoryResolutionRepository((left, right)) default_result = service(default_repository).generate( {"asset_type": "device", "threshold": 0.7}, actor_uid="admin-1", ) assert default_result.records[0].status == "pending" enabled_repository = MemoryResolutionRepository((left, right)) enabled_result = service( enabled_repository, auto_merge_enabled=True, ).generate( {"asset_type": "device", "threshold": 0.7}, actor_uid="admin-1", ) assert enabled_result.records[0].status == "merged" assert enabled_repository.reviews[0].decision == "auto_approve" assert enabled_repository.merge_records[0].member_asset_uid == RIGHT_UID def test_review_rejects_stale_version_and_active_member_conflict(): from app.core.data_research.device_entity_resolution import ( DeviceEntityConflict, ) repository = MemoryResolutionRepository( ( asset(LEFT_UID, SOURCE_A, "EQ-001"), asset(RIGHT_UID, SOURCE_B, "EQ-001"), ) ) resolution = service(repository) candidate = resolution.generate( {"asset_type": "device", "threshold": 0.7}, actor_uid="editor-1", ).records[0] with pytest.raises(DeviceEntityConflict, match="version"): resolution.review( candidate.uid, { "decision": "reject", "expected_version": 2, "reason": "版本已变化", }, actor_uid="admin-1", ) repository.candidates[candidate.uid] = replace( candidate, current_version=1, ) resolution.review( candidate.uid, { "decision": "approve", "canonical_asset_uid": LEFT_UID, "expected_version": 1, "reason": "证据充分", }, actor_uid="admin-1", ) with pytest.raises(DeviceEntityConflict, match="active merge"): resolution.submit_ai_candidate( { "left_asset_uid": LEFT_UID, "right_asset_uid": RIGHT_UID, "confidence": 0.9, "model_provider": "governed-provider", "model_name": "entity-match-v1", "evidence_uids": ["evidence-1"], "explanation": "重复候选", }, actor_uid="editor-1", )