from __future__ import annotations import os import uuid from datetime import UTC, datetime, timedelta import pytest from sqlalchemy import text pytestmark = pytest.mark.integration def _definition(): return { "schema_version": "1.0", "field_roles": { "identity": {"required": True, "checks": ["unique", "pattern"]}, "category": {"required": True, "checks": ["distribution"]}, "measure": {"required": False, "checks": ["outlier"]}, }, "thresholds": { "completeness_min": 0.95, "uniqueness_min": 1, "pattern": "^[A-Z]{2,5}-[0-9]{3}$", "volume_change_max_ratio": 0.25, "distribution_drift_max": 0.2, "freshness_max_seconds": 3600, "quality_score_min": 80, }, "sample_limit": 5, } def test_same_template_persists_two_domains_and_root_cause_evidence(monkeypatch): database_url = os.environ.get("TEST_DATABASE_URL") if not database_url: pytest.skip("TEST_DATABASE_URL is required") monkeypatch.setenv("DATABASE_URL", database_url) from app import create_app, db from app.core.data_rules.quality_operations import QualityOperationsService from app.core.data_rules.quality_repository import ( SqlAlchemyQualityOperationsRepository, ) app = create_app() app.config.update(TESTING=True) actor_uid = str(uuid.uuid4()) owner_uid = str(uuid.uuid4()) source_uid = str(uuid.uuid4()) plan_uid = str(uuid.uuid4()) metadata_run_uid = str(uuid.uuid4()) device_asset_uid = str(uuid.uuid4()) parts_asset_uid = str(uuid.uuid4()) template_uid = None run_uids = [] now = datetime.now(UTC) try: with app.app_context(): for uid, label in ((actor_uid, "actor"), (owner_uid, "owner")): db.session.execute( text( """ INSERT INTO public.users ( id, username, display_name, password_hash, status ) VALUES ( CAST(:uid AS uuid), :username, :username, 'p2-wp04-integration-hash', 'active' ) """ ), { "uid": uid, "username": f"wp04-{label}-{uid[:8]}", }, ) db.session.execute( text( """ INSERT INTO public.ingestion_sources ( uid, source_type, name, config, permission_scope, status, created_by ) VALUES ( CAST(:uid AS uuid), 'database', :name, '{}'::jsonb, '{}'::jsonb, 'active', :created_by ) """ ), { "uid": source_uid, "name": f"WP04 source {source_uid[:8]}", "created_by": actor_uid, }, ) db.session.execute( text( """ INSERT INTO public.active_metadata_plans ( uid, source_uid, name, source_kind, schedule_type, discovery_mode, scope, cursor_state, owner_uid, enabled, current_version, created_by ) VALUES ( CAST(:uid AS uuid), CAST(:source_uid AS uuid), :name, 'database', 'manual', 'snapshot', '{}'::jsonb, '{}'::jsonb, CAST(:owner_uid AS uuid), TRUE, 1, CAST(:created_by AS uuid) ) """ ), { "uid": plan_uid, "source_uid": source_uid, "name": f"WP04 plan {plan_uid[:8]}", "owner_uid": owner_uid, "created_by": actor_uid, }, ) db.session.execute( text( """ INSERT INTO public.active_metadata_runs ( uid, plan_uid, batch_key, status, attempt_count, cursor_before, cursor_after, snapshot_hash, statistics, actor_uid, started_at, finished_at ) VALUES ( CAST(:uid AS uuid), CAST(:plan_uid AS uuid), 'wp04-metadata', 'completed', 1, '{}'::jsonb, '{}'::jsonb, :snapshot_hash, '{}'::jsonb, CAST(:actor_uid AS uuid), :started_at, :finished_at ) """ ), { "uid": metadata_run_uid, "plan_uid": plan_uid, "snapshot_hash": "a" * 64, "actor_uid": actor_uid, "started_at": now - timedelta(minutes=5), "finished_at": now, }, ) for asset_uid, namespace, name, domain in ( ( device_asset_uid, "maintenance", "device_ledger", "device", ), ( parts_asset_uid, "supply", "spare_parts", "spare_parts", ), ): db.session.execute( text( """ INSERT INTO public.active_metadata_assets ( uid, source_uid, asset_key, namespace, name, asset_type, lifecycle_status, current_version, content_hash, snapshot, health, last_run_uid ) VALUES ( CAST(:uid AS uuid), CAST(:source_uid AS uuid), :asset_key, :namespace, :name, 'table', 'active', 1, :content_hash, CAST(:snapshot AS jsonb), '{}'::jsonb, CAST(:last_run_uid AS uuid) ) """ ), { "uid": asset_uid, "source_uid": source_uid, "asset_key": f"{source_uid}:{namespace}.{name}", "namespace": namespace, "name": name, "content_hash": "b" * 64, "snapshot": ( '{"business_domain_uid": "' + domain + '"}' ), "last_run_uid": metadata_run_uid, }, ) db.session.execute( text( """ INSERT INTO public.active_metadata_changes ( uid, run_uid, asset_uid, asset_key, field_name, change_type, before_state, after_state, status ) VALUES ( CAST(:uid AS uuid), CAST(:run_uid AS uuid), CAST(:asset_uid AS uuid), :asset_key, 'part_category', 'field_changed', '{}'::jsonb, '{}'::jsonb, 'accepted' ) """ ), { "uid": str(uuid.uuid4()), "run_uid": metadata_run_uid, "asset_uid": parts_asset_uid, "asset_key": f"{source_uid}:supply.spare_parts", }, ) db.session.execute( text( """ INSERT INTO public.active_metadata_lineage ( uid, run_uid, parse_status, source_asset, source_field, target_asset, target_field, relation_type, evidence ) VALUES ( CAST(:uid AS uuid), CAST(:run_uid AS uuid), 'resolved', 'erp.parts', 'part_code', 'supply.spare_parts', 'part_code', 'derived_from', '{}'::jsonb ) """ ), {"uid": str(uuid.uuid4()), "run_uid": metadata_run_uid}, ) db.session.commit() repository = SqlAlchemyQualityOperationsRepository(db.session) quality = QualityOperationsService( repository, publish_authorizer=lambda _actor_uid: None, commit=db.session.commit, rollback=db.session.rollback, ) draft = quality.create_template( { "code": f"WP04_GENERIC_{source_uid[:8].upper()}", "name": "WP04 跨域通用质量模板", "owner_uid": owner_uid, "definition": _definition(), }, actor_uid=actor_uid, ) template_uid = draft["uid"] published = quality.publish_template( template_uid, expected_version=1, actor_uid=actor_uid, ) device = quality.execute( { "template_uid": template_uid, "asset_uid": device_asset_uid, "batch_key": "device-baseline", "field_bindings": { "identity": "device_code", "category": "device_type", "measure": "temperature", }, "source_observed_at": ( now - timedelta(minutes=10) ).isoformat(), "records": [ { "device_code": "DEV-001", "device_type": "pump", "temperature": 22, }, { "device_code": "DEV-002", "device_type": "fan", "temperature": 23, }, ], }, actor_uid=actor_uid, ) run_uids.append(device["uid"]) parts = quality.execute( { "template_uid": template_uid, "asset_uid": parts_asset_uid, "batch_key": "parts-degraded", "field_bindings": { "identity": "part_code", "category": "part_category", "measure": "stock_quantity", }, "source_observed_at": ( now - timedelta(hours=2) ).isoformat(), "records": [ { "part_code": "PT-001", "part_category": "seal", "stock_quantity": 1, }, { "part_code": "PT-001", "part_category": None, "stock_quantity": 2, }, { "part_code": None, "part_category": "seal", "stock_quantity": 3, }, { "part_code": "bad", "part_category": "seal", "stock_quantity": 5000, }, ], }, actor_uid=actor_uid, ) run_uids.append(parts["uid"]) detail = quality.get_run(parts["uid"]) assert published["status"] == "published" assert device["template_version_uid"] == parts["template_version_uid"] assert { device["business_domain_uid"], parts["business_domain_uid"], } == {"device", "spare_parts"} assert detail["metrics"] assert { item["finding_type"] for item in detail["findings"] } >= { "completeness", "uniqueness", "pattern", "duplicate", "outlier", "freshness", } assert all( item["evidence"]["root_cause"]["lineage"] for item in detail["findings"] ) assert all( item["evidence"]["root_cause"]["changes"] for item in detail["findings"] ) assert all( item["evidence"]["root_cause"]["runs"] for item in detail["findings"] ) assert all( item["evidence"]["root_cause"]["responsibility"]["owner_uid"] == owner_uid for item in detail["findings"] ) assert { item["sla_type"] for item in detail["sla_events"] } == {"freshness", "quality_score"} finally: with app.app_context(): if run_uids: for table in ( "quality_sla_events", "quality_findings", "quality_profile_metrics", ): db.session.execute( text( f"DELETE FROM public.{table} " "WHERE run_uid = ANY(CAST(:uids AS uuid[]))" ), {"uids": run_uids}, ) db.session.execute( text( """ DELETE FROM public.quality_profile_runs WHERE uid = ANY(CAST(:uids AS uuid[])) """ ), {"uids": run_uids}, ) if template_uid: db.session.execute( text( """ UPDATE public.quality_templates SET active_version_uid = NULL WHERE uid = CAST(:uid AS uuid) """ ), {"uid": template_uid}, ) db.session.execute( text( """ DELETE FROM public.quality_template_versions WHERE template_uid = CAST(:uid AS uuid) """ ), {"uid": template_uid}, ) db.session.execute( text( """ DELETE FROM public.quality_templates WHERE uid = CAST(:uid AS uuid) """ ), {"uid": template_uid}, ) db.session.execute( text( """ DELETE FROM public.active_metadata_lineage WHERE run_uid = CAST(:uid AS uuid); DELETE FROM public.active_metadata_changes WHERE run_uid = CAST(:uid AS uuid); DELETE FROM public.active_metadata_assets WHERE last_run_uid = CAST(:uid AS uuid); DELETE FROM public.active_metadata_runs WHERE uid = CAST(:uid AS uuid); DELETE FROM public.active_metadata_plans WHERE uid = CAST(:plan_uid AS uuid); DELETE FROM public.ingestion_sources WHERE uid = CAST(:source_uid AS uuid); DELETE FROM public.users WHERE id IN ( CAST(:actor_uid AS uuid), CAST(:owner_uid AS uuid) ); """ ), { "uid": metadata_run_uid, "plan_uid": plan_uid, "source_uid": source_uid, "actor_uid": actor_uid, "owner_uid": owner_uid, }, ) db.session.commit()