| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181 |
- 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()
|