test_device_entity_resolution_postgres.py 6.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181
  1. from __future__ import annotations
  2. import os
  3. import uuid
  4. import pytest
  5. pytestmark = pytest.mark.integration
  6. def test_entity_resolution_persists_review_merge_and_rollback(monkeypatch):
  7. platform_url = os.environ.get("TEST_DATABASE_URL")
  8. if not platform_url:
  9. pytest.skip("TEST_DATABASE_URL is required")
  10. monkeypatch.setenv("DATABASE_URL", platform_url)
  11. from app import create_app, db
  12. from app.core.data_research.device_asset_repository import (
  13. SqlAlchemyDeviceAssetRepository,
  14. )
  15. from app.core.data_research.device_assets import DeviceAssetService
  16. from app.core.data_research.device_entity_repository import (
  17. SqlAlchemyDeviceEntityResolutionRepository,
  18. )
  19. from app.core.data_research.device_entity_resolution import (
  20. DeviceEntityResolutionService,
  21. )
  22. from app.models.data_research import (
  23. DeviceAsset,
  24. DeviceAssetSourceMapping,
  25. DeviceAssetVersion,
  26. DeviceEntityMatchCandidate,
  27. DeviceEntityMatchReview,
  28. DeviceEntityMergeEvent,
  29. DeviceEntityMergeRollback,
  30. IngestionSource,
  31. )
  32. app = create_app()
  33. app.config.update(TESTING=True)
  34. source_uids = [str(uuid.uuid4()), str(uuid.uuid4())]
  35. asset_uids = []
  36. candidate_uid = None
  37. merge_uid = None
  38. try:
  39. with app.app_context():
  40. for index, source_uid in enumerate(source_uids, start=1):
  41. db.session.add(
  42. IngestionSource(
  43. uid=source_uid,
  44. source_type="database",
  45. name=f"WP06 实体匹配测试源 {index}",
  46. config={
  47. "database_type": "postgresql",
  48. "database": f"wp06_{index}",
  49. "schema": "asset",
  50. },
  51. permission_scope={},
  52. status="active",
  53. created_by="integration-test",
  54. )
  55. )
  56. db.session.commit()
  57. assets = DeviceAssetService(
  58. SqlAlchemyDeviceAssetRepository(db.session),
  59. commit=db.session.commit,
  60. rollback=db.session.rollback,
  61. )
  62. for source_uid, source_code in zip(
  63. source_uids,
  64. ("EQ-WP06-001", "EQ-WP06-001"),
  65. strict=True,
  66. ):
  67. result = assets.import_records(
  68. {
  69. "source_uid": source_uid,
  70. "source_entity": "asset.equipment",
  71. "records": [
  72. {
  73. "asset_type": "device",
  74. "source_code": source_code,
  75. "name": "WP06 循环水泵",
  76. "location": "动力车间",
  77. "organization": "设备动力部",
  78. "responsible_person": "张工",
  79. "attributes": {"model": "P-WP06"},
  80. }
  81. ],
  82. },
  83. actor_uid="integration-test",
  84. )
  85. asset_uids.append(result.items[0].asset.uid)
  86. repository = SqlAlchemyDeviceEntityResolutionRepository(
  87. db.session
  88. )
  89. service = DeviceEntityResolutionService(
  90. repository,
  91. review_authorizer=lambda _actor_uid: None,
  92. commit=db.session.commit,
  93. rollback=db.session.rollback,
  94. )
  95. generated = service.generate(
  96. {"asset_type": "device", "threshold": 0.98, "limit": 100},
  97. actor_uid="editor-test",
  98. )
  99. owned = [
  100. item
  101. for item in generated.records
  102. if {item.left_asset_uid, item.right_asset_uid}
  103. == set(asset_uids)
  104. ]
  105. assert len(owned) == 1
  106. candidate = owned[0]
  107. candidate_uid = candidate.uid
  108. merged, review, merge = service.review(
  109. candidate.uid,
  110. {
  111. "decision": "approve",
  112. "canonical_asset_uid": asset_uids[0],
  113. "expected_version": 1,
  114. "reason": "集成测试证据一致",
  115. },
  116. actor_uid="admin-test",
  117. )
  118. merge_uid = merge.uid
  119. assert merged.status == "merged"
  120. assert review.decision == "approve"
  121. assert repository.active_merge_for_member(asset_uids[1]).uid == (
  122. merge.uid
  123. )
  124. assert merge.snapshot["canonical"]["uid"] == asset_uids[0]
  125. assert len(repository.list_reviews(candidate.uid)) == 1
  126. rolled_back, rollback = service.rollback(
  127. merge.uid,
  128. {
  129. "expected_version": 2,
  130. "reason": "集成测试回滚",
  131. },
  132. actor_uid="admin-test",
  133. )
  134. assert rolled_back.status == "rolled_back"
  135. assert rollback.snapshot["merge"]["member_asset_uid"] == (
  136. asset_uids[1]
  137. )
  138. assert repository.active_merge_for_member(asset_uids[1]) is None
  139. assert len(repository.list_rollbacks(merge.uid)) == 1
  140. finally:
  141. with app.app_context():
  142. if merge_uid:
  143. db.session.query(DeviceEntityMergeRollback).filter_by(
  144. merge_uid=merge_uid
  145. ).delete(synchronize_session=False)
  146. if candidate_uid:
  147. db.session.query(DeviceEntityMergeEvent).filter_by(
  148. candidate_uid=candidate_uid
  149. ).delete(synchronize_session=False)
  150. db.session.query(DeviceEntityMatchReview).filter_by(
  151. candidate_uid=candidate_uid
  152. ).delete(synchronize_session=False)
  153. db.session.query(DeviceEntityMatchCandidate).filter_by(
  154. uid=candidate_uid
  155. ).delete(synchronize_session=False)
  156. if asset_uids:
  157. db.session.query(DeviceAssetVersion).filter(
  158. DeviceAssetVersion.asset_uid.in_(asset_uids)
  159. ).delete(synchronize_session=False)
  160. db.session.query(DeviceAssetSourceMapping).filter(
  161. DeviceAssetSourceMapping.asset_uid.in_(asset_uids)
  162. ).delete(synchronize_session=False)
  163. db.session.query(DeviceAsset).filter(
  164. DeviceAsset.uid.in_(asset_uids)
  165. ).delete(synchronize_session=False)
  166. db.session.query(IngestionSource).filter(
  167. IngestionSource.uid.in_(source_uids)
  168. ).delete(synchronize_session=False)
  169. db.session.commit()