from __future__ import annotations import os import uuid from datetime import UTC, datetime import pytest from sqlalchemy import text pytestmark = pytest.mark.integration def test_device_knowledge_prefilters_source_scope_before_matching(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.knowledge.retrieval.device import ( SqlDeviceKnowledgeRepository, ) from app.models.data_research import ( DeviceAsset, DeviceAssetSourceMapping, DeviceOperationalEvent, IngestionSource, ) app = create_app() app.config.update(TESTING=True) suffix = uuid.uuid4().hex[:10] actor_uid = str(uuid.uuid4()) domain_a = str(uuid.uuid4()) domain_b = str(uuid.uuid4()) source_a = str(uuid.uuid4()) source_b = str(uuid.uuid4()) source_unscoped = str(uuid.uuid4()) asset_a = str(uuid.uuid4()) asset_b = str(uuid.uuid4()) asset_unscoped = str(uuid.uuid4()) source_uids = (source_a, source_b, source_unscoped) asset_uids = (asset_a, asset_b, asset_unscoped) now = datetime.now(UTC).replace(microsecond=0) try: with app.app_context(): db.session.execute( text( """ INSERT INTO public.users ( id, username, display_name, password_hash, status ) VALUES ( CAST(:id AS uuid), :username, :username, 'integration-test', 'active' ) """ ), { "id": actor_uid, "username": f"wp10-actor-{suffix}", }, ) for uid, label, scope in ( (source_a, "A", {"business_domains": [domain_a]}), (source_b, "B", {"business_domains": [domain_b]}), (source_unscoped, "U", {}), ): db.session.add( IngestionSource( uid=uid, source_type="database", name=f"WP10 source {label} {suffix}", config={"password": f"secret-{label}"}, permission_scope=scope, status="active", created_by=actor_uid, ) ) for uid, label, owner in ( (asset_a, "A", "张工"), (asset_b, "B", "李工"), (asset_unscoped, "U", "王工"), ): db.session.add( DeviceAsset( uid=uid, asset_type="device", name=f"WP10循环泵-{suffix}-{label}", status="active", current_version=1, content_hash=label.lower() * 64, location=f"{label}动力车间", organization="设备动力部", responsible_person=owner, attributes={"password": f"asset-secret-{label}"}, created_by=actor_uid, updated_by=actor_uid, updated_at=now, ) ) db.session.flush() for uid, asset_uid, source_uid, label in ( (str(uuid.uuid4()), asset_a, source_a, "A"), (str(uuid.uuid4()), asset_b, source_b, "B"), ( str(uuid.uuid4()), asset_unscoped, source_unscoped, "U", ), ): db.session.add( DeviceAssetSourceMapping( uid=uid, asset_uid=asset_uid, source_uid=source_uid, source_entity="asset.equipment", asset_type="device", source_code=f"EQ-WP10-{suffix}-{label}", source_updated_at=now, first_seen_at=now, last_seen_at=now, ) ) db.session.add( DeviceOperationalEvent( uid=str(uuid.uuid4()), source_uid=source_uid, source_entity="fault.events", source_code=f"FT-WP10-{suffix}-{label}", event_type="fault", asset_uid=asset_uid, title=f"{label}域轴承故障-{suffix}", severity="error", status="observed", occurred_at=now, evidence_refs={ "password": f"event-secret-{label}" }, content_hash=(label.lower() + "f") * 32, created_by=actor_uid, ) ) db.session.commit() repository = SqlDeviceKnowledgeRepository(db.session) admin_rows = repository.search( query=f"WP10循环泵-{suffix}", global_access=True, business_domain_uids=(), limit=20, ) domain_a_rows = repository.search( query=f"WP10循环泵-{suffix}", global_access=False, business_domain_uids=(domain_a,), limit=20, ) forbidden_source = repository.search( query=f"EQ-WP10-{suffix}-B", global_access=False, business_domain_uids=(domain_a,), limit=20, ) unscoped_source = repository.search( query=f"EQ-WP10-{suffix}-U", global_access=False, business_domain_uids=(domain_a,), limit=20, ) fault_rows = repository.search( query=f"A域轴承故障-{suffix}", global_access=False, business_domain_uids=(domain_a,), limit=20, ) natural_question_rows = repository.search( query=f"WP10循环泵-{suffix}-A最近有哪些故障?", global_access=False, business_domain_uids=(domain_a,), limit=20, ) detail = repository.get_detail( asset_a, global_access=False, business_domain_uids=(domain_a,), ) forbidden_detail = repository.get_detail( asset_b, global_access=False, business_domain_uids=(domain_a,), ) assert {row.asset_uid for row in admin_rows} == set(asset_uids) assert [row.asset_uid for row in domain_a_rows] == [asset_a] assert domain_a_rows[0].business_domain_uid == domain_a assert forbidden_source == () assert unscoped_source == () assert [row.asset_uid for row in fault_rows] == [asset_a] assert [row.asset_uid for row in natural_question_rows] == [ asset_a ] assert fault_rows[0].related_events == ( ("fault", f"A域轴承故障-{suffix}", f"FT-WP10-{suffix}-A"), ) assert detail is not None assert detail.source_codes == (f"EQ-WP10-{suffix}-A",) assert forbidden_detail is None serialized = repr((domain_a_rows, fault_rows, detail)) assert "secret-" not in serialized assert "permission_scope" not in serialized finally: with app.app_context(): db.session.execute( text( "DELETE FROM public.device_operational_events " "WHERE created_by = CAST(:actor_uid AS uuid)" ), {"actor_uid": actor_uid}, ) db.session.execute( text( "DELETE FROM public.device_assets " "WHERE uid = ANY(CAST(:uids AS uuid[]))" ), {"uids": list(asset_uids)}, ) db.session.execute( text( "DELETE FROM public.ingestion_sources " "WHERE uid = ANY(CAST(:uids AS uuid[]))" ), {"uids": list(source_uids)}, ) db.session.execute( text( "DELETE FROM public.users " "WHERE id = CAST(:id AS uuid)" ), {"id": actor_uid}, ) db.session.commit()