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