from __future__ import annotations import os import uuid import pytest pytestmark = pytest.mark.integration def test_entity_resolution_persists_review_merge_and_rollback(monkeypatch): platform_url = os.environ.get("TEST_DATABASE_URL") if not platform_url: pytest.skip("TEST_DATABASE_URL is required") monkeypatch.setenv("DATABASE_URL", platform_url) from app import create_app, db from app.core.data_research.device_asset_repository import ( SqlAlchemyDeviceAssetRepository, ) from app.core.data_research.device_assets import DeviceAssetService from app.core.data_research.device_entity_repository import ( SqlAlchemyDeviceEntityResolutionRepository, ) from app.core.data_research.device_entity_resolution import ( DeviceEntityResolutionService, ) from app.models.data_research import ( DeviceAsset, DeviceAssetSourceMapping, DeviceAssetVersion, DeviceEntityMatchCandidate, DeviceEntityMatchReview, DeviceEntityMergeEvent, DeviceEntityMergeRollback, IngestionSource, ) app = create_app() app.config.update(TESTING=True) source_uids = [str(uuid.uuid4()), str(uuid.uuid4())] asset_uids = [] candidate_uid = None merge_uid = None try: with app.app_context(): for index, source_uid in enumerate(source_uids, start=1): db.session.add( IngestionSource( uid=source_uid, source_type="database", name=f"WP06 实体匹配测试源 {index}", config={ "database_type": "postgresql", "database": f"wp06_{index}", "schema": "asset", }, permission_scope={}, status="active", created_by="integration-test", ) ) db.session.commit() assets = DeviceAssetService( SqlAlchemyDeviceAssetRepository(db.session), commit=db.session.commit, rollback=db.session.rollback, ) for source_uid, source_code in zip( source_uids, ("EQ-WP06-001", "EQ-WP06-001"), strict=True, ): result = assets.import_records( { "source_uid": source_uid, "source_entity": "asset.equipment", "records": [ { "asset_type": "device", "source_code": source_code, "name": "WP06 循环水泵", "location": "动力车间", "organization": "设备动力部", "responsible_person": "张工", "attributes": {"model": "P-WP06"}, } ], }, actor_uid="integration-test", ) asset_uids.append(result.items[0].asset.uid) repository = SqlAlchemyDeviceEntityResolutionRepository( db.session ) service = DeviceEntityResolutionService( repository, review_authorizer=lambda _actor_uid: None, commit=db.session.commit, rollback=db.session.rollback, ) generated = service.generate( {"asset_type": "device", "threshold": 0.98, "limit": 100}, actor_uid="editor-test", ) owned = [ item for item in generated.records if {item.left_asset_uid, item.right_asset_uid} == set(asset_uids) ] assert len(owned) == 1 candidate = owned[0] candidate_uid = candidate.uid merged, review, merge = service.review( candidate.uid, { "decision": "approve", "canonical_asset_uid": asset_uids[0], "expected_version": 1, "reason": "集成测试证据一致", }, actor_uid="admin-test", ) merge_uid = merge.uid assert merged.status == "merged" assert review.decision == "approve" assert repository.active_merge_for_member(asset_uids[1]).uid == ( merge.uid ) assert merge.snapshot["canonical"]["uid"] == asset_uids[0] assert len(repository.list_reviews(candidate.uid)) == 1 rolled_back, rollback = service.rollback( merge.uid, { "expected_version": 2, "reason": "集成测试回滚", }, actor_uid="admin-test", ) assert rolled_back.status == "rolled_back" assert rollback.snapshot["merge"]["member_asset_uid"] == ( asset_uids[1] ) assert repository.active_merge_for_member(asset_uids[1]) is None assert len(repository.list_rollbacks(merge.uid)) == 1 finally: with app.app_context(): if merge_uid: db.session.query(DeviceEntityMergeRollback).filter_by( merge_uid=merge_uid ).delete(synchronize_session=False) if candidate_uid: db.session.query(DeviceEntityMergeEvent).filter_by( candidate_uid=candidate_uid ).delete(synchronize_session=False) db.session.query(DeviceEntityMatchReview).filter_by( candidate_uid=candidate_uid ).delete(synchronize_session=False) db.session.query(DeviceEntityMatchCandidate).filter_by( uid=candidate_uid ).delete(synchronize_session=False) if asset_uids: db.session.query(DeviceAssetVersion).filter( DeviceAssetVersion.asset_uid.in_(asset_uids) ).delete(synchronize_session=False) db.session.query(DeviceAssetSourceMapping).filter( DeviceAssetSourceMapping.asset_uid.in_(asset_uids) ).delete(synchronize_session=False) db.session.query(DeviceAsset).filter( DeviceAsset.uid.in_(asset_uids) ).delete(synchronize_session=False) db.session.query(IngestionSource).filter( IngestionSource.uid.in_(source_uids) ).delete(synchronize_session=False) db.session.commit()