test_device_entity_resolution.py 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485
  1. from __future__ import annotations
  2. from dataclasses import replace
  3. from datetime import datetime
  4. import pytest
  5. from app.core.data_research.device_assets import (
  6. DeviceAssetDetail,
  7. DeviceAssetMappingRecord,
  8. DeviceAssetRecord,
  9. )
  10. SOURCE_A = "00000000-0000-0000-0000-000000000101"
  11. SOURCE_B = "00000000-0000-0000-0000-000000000102"
  12. LEFT_UID = "00000000-0000-7000-8000-000000000201"
  13. RIGHT_UID = "00000000-0000-7000-8000-000000000202"
  14. def asset(
  15. uid,
  16. source_uid,
  17. source_code,
  18. *,
  19. name="一号循环泵",
  20. asset_type="device",
  21. location="动力车间",
  22. organization="设备动力部",
  23. responsible_person="张工",
  24. model="P-100",
  25. ):
  26. now = datetime(2026, 7, 29, 9, 0)
  27. record = DeviceAssetRecord(
  28. uid=uid,
  29. asset_type=asset_type,
  30. name=name,
  31. status="active",
  32. current_version=1,
  33. content_hash="a" * 64,
  34. location=location,
  35. organization=organization,
  36. responsible_person=responsible_person,
  37. attributes={"model": model},
  38. created_by="editor-1",
  39. updated_by="editor-1",
  40. created_at=now,
  41. updated_at=now,
  42. )
  43. mapping = DeviceAssetMappingRecord(
  44. uid=f"{uid[:-3]}3{uid[-2:]}",
  45. asset_uid=uid,
  46. source_uid=source_uid,
  47. source_entity="asset.equipment",
  48. asset_type=asset_type,
  49. source_code=source_code,
  50. source_updated_at=now,
  51. first_seen_at=now,
  52. last_seen_at=now,
  53. )
  54. return DeviceAssetDetail(asset=record, mappings=(mapping,))
  55. class MemoryResolutionRepository:
  56. def __init__(self, assets=()):
  57. self.assets = {item.asset.uid: item for item in assets}
  58. self.candidates = {}
  59. self.reviews = []
  60. self.merge_records = []
  61. self.rollback_records = []
  62. def list_matchable_assets(self, asset_type, *, limit):
  63. return [
  64. item
  65. for item in self.assets.values()
  66. if item.asset.asset_type == asset_type
  67. ][:limit]
  68. def get_asset_detail(self, uid):
  69. return self.assets.get(uid)
  70. def find_open_pair(self, left_uid, right_uid):
  71. return next(
  72. (
  73. item
  74. for item in self.candidates.values()
  75. if item.left_asset_uid == left_uid
  76. and item.right_asset_uid == right_uid
  77. and item.status in {"pending", "merged"}
  78. ),
  79. None,
  80. )
  81. def create_candidate(self, record):
  82. self.candidates[record.uid] = record
  83. return record
  84. def search_candidates(self, filters, *, page, page_size):
  85. records = list(self.candidates.values())
  86. for name in ("status", "suggestion_source"):
  87. if filters.get(name):
  88. records = [
  89. item
  90. for item in records
  91. if getattr(item, name) == filters[name]
  92. ]
  93. start = (page - 1) * page_size
  94. return records[start : start + page_size], len(records)
  95. def get_candidate(self, uid, *, for_update=False):
  96. del for_update
  97. return self.candidates.get(uid)
  98. def update_candidate(self, record):
  99. self.candidates[record.uid] = record
  100. return record
  101. def append_review(self, record):
  102. self.reviews.append(record)
  103. return record
  104. def active_merge_for_member(self, asset_uid):
  105. rolled_back = {item.merge_uid for item in self.rollback_records}
  106. return next(
  107. (
  108. item
  109. for item in self.merge_records
  110. if item.member_asset_uid == asset_uid
  111. and item.uid not in rolled_back
  112. ),
  113. None,
  114. )
  115. def create_merge(self, record):
  116. self.merge_records.append(record)
  117. return record
  118. def get_merge(self, uid, *, for_update=False):
  119. del for_update
  120. return next(
  121. (item for item in self.merge_records if item.uid == uid),
  122. None,
  123. )
  124. def create_rollback(self, record):
  125. self.rollback_records.append(record)
  126. return record
  127. def list_reviews(self, candidate_uid):
  128. return [
  129. item for item in self.reviews if item.candidate_uid == candidate_uid
  130. ]
  131. def list_merges(self, candidate_uid):
  132. return [
  133. item
  134. for item in self.merge_records
  135. if item.candidate_uid == candidate_uid
  136. ]
  137. def list_rollbacks(self, merge_uid):
  138. return [
  139. item
  140. for item in self.rollback_records
  141. if item.merge_uid == merge_uid
  142. ]
  143. def service(repository, *, authorized=True, auto_merge_enabled=False):
  144. from app.core.data_research.device_entity_resolution import (
  145. DeviceEntityForbidden,
  146. DeviceEntityResolutionService,
  147. )
  148. ids = iter(
  149. f"00000000-0000-7000-8000-0000000003{index:02d}"
  150. for index in range(1, 40)
  151. )
  152. clock = iter(
  153. datetime(2026, 7, 29, 10, minute)
  154. for minute in range(1, 40)
  155. )
  156. def authorize(actor_uid):
  157. if not authorized:
  158. raise DeviceEntityForbidden(
  159. f"{actor_uid} is not the accountable asset manager"
  160. )
  161. return DeviceEntityResolutionService(
  162. repository,
  163. review_authorizer=authorize,
  164. uid_factory=ids.__next__,
  165. now_factory=clock.__next__,
  166. auto_merge_enabled=auto_merge_enabled,
  167. )
  168. def test_score_is_explainable_and_uses_hand_derived_weights():
  169. from app.core.data_research.device_entity_resolution import (
  170. score_device_pair,
  171. )
  172. left = asset(LEFT_UID, SOURCE_A, "EQ-001")
  173. right = asset(
  174. RIGHT_UID,
  175. SOURCE_B,
  176. "设备-001",
  177. name=" 一号 循环泵 ",
  178. responsible_person="李工",
  179. )
  180. score = score_device_pair(left, right)
  181. assert score.confidence == pytest.approx(0.90)
  182. assert [item["signal"] for item in score.explanation] == [
  183. "name",
  184. "location",
  185. "organization",
  186. "responsible_person",
  187. "model",
  188. "source_code",
  189. ]
  190. assert score.explanation[0]["matched"] is True
  191. assert score.explanation[3]["matched"] is False
  192. assert score.explanation[5]["matched"] is False
  193. def test_generation_is_cross_source_bounded_and_idempotent_for_open_pair():
  194. left = asset(LEFT_UID, SOURCE_A, "EQ-001")
  195. right = asset(RIGHT_UID, SOURCE_B, "EQ-001")
  196. same_source = asset(
  197. "00000000-0000-7000-8000-000000000203",
  198. SOURCE_A,
  199. "EQ-003",
  200. name="三号空压机",
  201. location="空压站",
  202. organization="公用工程部",
  203. responsible_person="王工",
  204. model="AC-300",
  205. )
  206. repository = MemoryResolutionRepository((left, right, same_source))
  207. resolution = service(repository)
  208. first = resolution.generate(
  209. {"asset_type": "device", "threshold": 0.7, "limit": 20},
  210. actor_uid="editor-1",
  211. )
  212. second = resolution.generate(
  213. {"asset_type": "device", "threshold": 0.7, "limit": 20},
  214. actor_uid="editor-1",
  215. )
  216. assert first.created_count == 1
  217. assert first.existing_count == 0
  218. assert second.created_count == 0
  219. assert second.existing_count == 1
  220. assert first.records[0].left_asset_uid == LEFT_UID
  221. assert first.records[0].right_asset_uid == RIGHT_UID
  222. assert first.records[0].evidence_uids == (
  223. LEFT_UID,
  224. left.mappings[0].uid,
  225. RIGHT_UID,
  226. right.mappings[0].uid,
  227. )
  228. def test_ai_candidate_requires_provider_model_evidence_and_never_auto_merges():
  229. from app.core.data_research.device_entity_resolution import (
  230. DeviceEntityInvalid,
  231. )
  232. repository = MemoryResolutionRepository(
  233. (
  234. asset(LEFT_UID, SOURCE_A, "EQ-001"),
  235. asset(RIGHT_UID, SOURCE_B, "DEVICE-001"),
  236. )
  237. )
  238. resolution = service(
  239. repository,
  240. authorized=True,
  241. auto_merge_enabled=True,
  242. )
  243. with pytest.raises(DeviceEntityInvalid, match="model_provider"):
  244. resolution.submit_ai_candidate(
  245. {
  246. "left_asset_uid": LEFT_UID,
  247. "right_asset_uid": RIGHT_UID,
  248. "confidence": 0.999,
  249. "model_name": "entity-match-v1",
  250. "evidence_uids": ["evidence-1"],
  251. "explanation": "同一台设备",
  252. },
  253. actor_uid="editor-1",
  254. )
  255. candidate = resolution.submit_ai_candidate(
  256. {
  257. "left_asset_uid": LEFT_UID,
  258. "right_asset_uid": RIGHT_UID,
  259. "confidence": 0.999,
  260. "model_provider": "governed-provider",
  261. "model_name": "entity-match-v1",
  262. "evidence_uids": ["evidence-1"],
  263. "explanation": "同一台设备",
  264. },
  265. actor_uid="editor-1",
  266. )
  267. assert candidate.status == "pending"
  268. assert candidate.suggestion_source == "ai"
  269. assert repository.merge_records == []
  270. def test_review_requires_accountable_manager_and_appends_merge_evidence():
  271. repository = MemoryResolutionRepository(
  272. (
  273. asset(LEFT_UID, SOURCE_A, "EQ-001"),
  274. asset(RIGHT_UID, SOURCE_B, "EQ-001"),
  275. )
  276. )
  277. blocked = service(repository, authorized=False)
  278. candidate = blocked.generate(
  279. {"asset_type": "device", "threshold": 0.7},
  280. actor_uid="editor-1",
  281. ).records[0]
  282. from app.core.data_research.device_entity_resolution import (
  283. DeviceEntityForbidden,
  284. )
  285. with pytest.raises(DeviceEntityForbidden, match="accountable"):
  286. blocked.review(
  287. candidate.uid,
  288. {
  289. "decision": "approve",
  290. "canonical_asset_uid": LEFT_UID,
  291. "expected_version": 1,
  292. "reason": "跨系统编码和型号一致",
  293. },
  294. actor_uid="admin-1",
  295. )
  296. approved, review, merge = service(repository).review(
  297. candidate.uid,
  298. {
  299. "decision": "approve",
  300. "canonical_asset_uid": LEFT_UID,
  301. "expected_version": 1,
  302. "reason": "跨系统编码和型号一致",
  303. },
  304. actor_uid="admin-1",
  305. )
  306. assert approved.status == "merged"
  307. assert approved.canonical_asset_uid == LEFT_UID
  308. assert review.decision == "approve"
  309. assert merge.canonical_asset_uid == LEFT_UID
  310. assert merge.member_asset_uid == RIGHT_UID
  311. assert merge.snapshot["member"]["uid"] == RIGHT_UID
  312. assert len(repository.reviews) == 1
  313. assert len(repository.merge_records) == 1
  314. def test_rollback_is_append_only_and_restores_candidate_state():
  315. repository = MemoryResolutionRepository(
  316. (
  317. asset(LEFT_UID, SOURCE_A, "EQ-001"),
  318. asset(RIGHT_UID, SOURCE_B, "EQ-001"),
  319. )
  320. )
  321. resolution = service(repository)
  322. candidate = resolution.generate(
  323. {"asset_type": "device", "threshold": 0.7},
  324. actor_uid="editor-1",
  325. ).records[0]
  326. _candidate, _review, merge = resolution.review(
  327. candidate.uid,
  328. {
  329. "decision": "approve",
  330. "canonical_asset_uid": LEFT_UID,
  331. "expected_version": 1,
  332. "reason": "匹配证据充分",
  333. },
  334. actor_uid="admin-1",
  335. )
  336. rolled_back, rollback = resolution.rollback(
  337. merge.uid,
  338. {
  339. "expected_version": 2,
  340. "reason": "现场确认不是同一台设备",
  341. },
  342. actor_uid="admin-1",
  343. )
  344. assert rolled_back.status == "rolled_back"
  345. assert rolled_back.current_version == 3
  346. assert rollback.merge_uid == merge.uid
  347. assert rollback.snapshot["merge"]["member_asset_uid"] == RIGHT_UID
  348. assert len(repository.merge_records) == 1
  349. assert len(repository.rollback_records) == 1
  350. assert repository.active_merge_for_member(RIGHT_UID) is None
  351. def test_rule_auto_merge_is_default_off_and_requires_strict_threshold():
  352. left = asset(LEFT_UID, SOURCE_A, "EQ-001")
  353. right = asset(RIGHT_UID, SOURCE_B, "EQ-001")
  354. default_repository = MemoryResolutionRepository((left, right))
  355. default_result = service(default_repository).generate(
  356. {"asset_type": "device", "threshold": 0.7},
  357. actor_uid="admin-1",
  358. )
  359. assert default_result.records[0].status == "pending"
  360. enabled_repository = MemoryResolutionRepository((left, right))
  361. enabled_result = service(
  362. enabled_repository,
  363. auto_merge_enabled=True,
  364. ).generate(
  365. {"asset_type": "device", "threshold": 0.7},
  366. actor_uid="admin-1",
  367. )
  368. assert enabled_result.records[0].status == "merged"
  369. assert enabled_repository.reviews[0].decision == "auto_approve"
  370. assert enabled_repository.merge_records[0].member_asset_uid == RIGHT_UID
  371. def test_review_rejects_stale_version_and_active_member_conflict():
  372. from app.core.data_research.device_entity_resolution import (
  373. DeviceEntityConflict,
  374. )
  375. repository = MemoryResolutionRepository(
  376. (
  377. asset(LEFT_UID, SOURCE_A, "EQ-001"),
  378. asset(RIGHT_UID, SOURCE_B, "EQ-001"),
  379. )
  380. )
  381. resolution = service(repository)
  382. candidate = resolution.generate(
  383. {"asset_type": "device", "threshold": 0.7},
  384. actor_uid="editor-1",
  385. ).records[0]
  386. with pytest.raises(DeviceEntityConflict, match="version"):
  387. resolution.review(
  388. candidate.uid,
  389. {
  390. "decision": "reject",
  391. "expected_version": 2,
  392. "reason": "版本已变化",
  393. },
  394. actor_uid="admin-1",
  395. )
  396. repository.candidates[candidate.uid] = replace(
  397. candidate,
  398. current_version=1,
  399. )
  400. resolution.review(
  401. candidate.uid,
  402. {
  403. "decision": "approve",
  404. "canonical_asset_uid": LEFT_UID,
  405. "expected_version": 1,
  406. "reason": "证据充分",
  407. },
  408. actor_uid="admin-1",
  409. )
  410. with pytest.raises(DeviceEntityConflict, match="active merge"):
  411. resolution.submit_ai_candidate(
  412. {
  413. "left_asset_uid": LEFT_UID,
  414. "right_asset_uid": RIGHT_UID,
  415. "confidence": 0.9,
  416. "model_provider": "governed-provider",
  417. "model_name": "entity-match-v1",
  418. "evidence_uids": ["evidence-1"],
  419. "explanation": "重复候选",
  420. },
  421. actor_uid="editor-1",
  422. )