| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485 |
- 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",
- )
|