| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320 |
- from __future__ import annotations
- import os
- import uuid
- import pytest
- from sqlalchemy import text
- pytestmark = pytest.mark.integration
- def test_active_metadata_is_idempotent_traceable_and_correctable(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.meta_data.active_metadata import ActiveMetadataService
- from app.core.meta_data.active_metadata_repository import (
- SqlAlchemyActiveMetadataRepository,
- )
- 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 = None
- try:
- with app.app_context():
- for uid, name in ((actor_uid, "actor"), (owner_uid, "owner")):
- db.session.execute(
- text(
- """
- INSERT INTO public.users (
- id, username, display_name, password_hash, status
- ) VALUES (
- CAST(:id AS uuid), :username, :username,
- 'p2-wp02-integration-hash', 'active'
- )
- """
- ),
- {"id": uid, "username": f"wp02-{name}-{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"WP02 source {source_uid[:8]}",
- "created_by": actor_uid,
- },
- )
- db.session.commit()
- repository = SqlAlchemyActiveMetadataRepository(db.session)
- service = ActiveMetadataService(repository)
- plan = service.create_plan(
- {
- "source_uid": source_uid,
- "name": "WP02 集成主动发现",
- "source_kind": "database",
- "schedule_type": "interval",
- "schedule_expression": "PT30M",
- "discovery_mode": "cursor",
- "scope": {"schemas": ["public"]},
- "owner_uid": owner_uid,
- },
- actor_uid=actor_uid,
- )
- db.session.commit()
- plan_uid = plan["uid"]
- first_payload = {
- "snapshot": {
- "assets": [
- {
- "name": "customers",
- "namespace": "public",
- "asset_type": "table",
- "fields": [
- {
- "name": "id",
- "data_type": "bigint",
- "nullable": False,
- "ordinal_position": 1,
- },
- {
- "name": "email",
- "data_type": "varchar",
- "nullable": True,
- "ordinal_position": 2,
- },
- ],
- },
- {
- "name": "orders",
- "namespace": "public",
- "asset_type": "table",
- "fields": [],
- },
- ]
- },
- "cursor_after": {"catalog_version": 1},
- "lineage_sql": [
- {
- "dialect": "postgres",
- "sql": (
- "INSERT INTO mart.customer_email (customer_id, email) "
- "SELECT c.id, c.email FROM public.customers c"
- ),
- },
- "not valid lineage SQL",
- ],
- "health_signals": [
- {
- "asset_key": f"{source_uid}:public.customers",
- "signal_type": "quality",
- "value": 0.98,
- "status": "healthy",
- },
- {
- "asset_key": f"{source_uid}:public.customers",
- "signal_type": "freshness",
- "value": 60,
- "status": "healthy",
- },
- {
- "asset_key": f"{source_uid}:public.customers",
- "signal_type": "usage",
- "value": 12,
- "status": "healthy",
- },
- ],
- }
- first = service.execute_source_plans(
- source_uid,
- first_payload["snapshot"],
- batch_key="wp02-batch-1",
- actor_uid=actor_uid,
- cursor_after=first_payload["cursor_after"],
- )[0]
- # SQL lineage and health use the same batch contract as the
- # automatic snapshot projection and are appended by the runner.
- first = service.execute(
- plan_uid,
- first_payload,
- batch_key="wp02-batch-1-enriched",
- actor_uid=actor_uid,
- )
- db.session.commit()
- replay = service.execute(
- plan_uid,
- first_payload,
- batch_key="wp02-batch-1-enriched",
- actor_uid=actor_uid,
- )
- db.session.commit()
- assert replay["uid"] == first["uid"]
- second_payload = {
- "snapshot": {
- "assets": [
- {
- "name": "customers",
- "namespace": "public",
- "asset_type": "table",
- "fields": [
- {
- "name": "id",
- "data_type": "bigint",
- "nullable": False,
- "ordinal_position": 1,
- },
- {
- "name": "email",
- "data_type": "text",
- "nullable": True,
- "ordinal_position": 2,
- },
- ],
- }
- ]
- },
- "cursor_after": {"catalog_version": 2},
- }
- second = service.execute(
- plan_uid,
- second_payload,
- batch_key="wp02-batch-2",
- actor_uid=actor_uid,
- )
- db.session.commit()
- assert second["cursor_before"] == {"catalog_version": 1}
- assets = service.list_assets(source_uid)
- customer = next(item for item in assets if item["name"] == "customers")
- orders = next(item for item in assets if item["name"] == "orders")
- assert customer["current_version"] == 2
- assert orders["lifecycle_status"] == "deletion_candidate"
- changes = service.list_changes(second["uid"])
- assert {(item["change_type"], item["field_name"]) for item in changes} == {
- ("field_changed", "email"),
- ("deletion_candidate", None),
- }
- lineage = service.list_lineage(first["uid"])
- assert sum(item["parse_status"] == "resolved" for item in lineage) == 2
- assert sum(item["parse_status"] == "failed" for item in lineage) == 1
- failed = service.record_failure(
- plan_uid,
- batch_key="wp02-batch-3",
- error_code="TIMEOUT",
- failure_reason="collector timed out",
- actor_uid=actor_uid,
- )
- db.session.commit()
- assert failed["status"] == "failed"
- correction = service.submit_correction(
- customer["uid"],
- {
- "field_name": "email",
- "proposed_value": {"comment": "客户邮箱"},
- "reason": "补充业务定义",
- "assignee_uid": owner_uid,
- },
- actor_uid=actor_uid,
- )
- db.session.commit()
- resolved = service.resolve_correction(
- correction["uid"],
- expected_version=1,
- decision="accept",
- resolution={"comment": "客户邮箱"},
- actor_uid=owner_uid,
- )
- db.session.commit()
- assert resolved["current_version"] == 2
- counts = db.session.execute(
- text(
- """
- SELECT
- (SELECT count(*) FROM public.active_metadata_runs
- WHERE plan_uid = CAST(:plan_uid AS uuid)) AS runs,
- (SELECT count(*) FROM public.active_metadata_asset_versions v
- JOIN public.active_metadata_assets a ON a.uid = v.asset_uid
- WHERE a.source_uid = CAST(:source_uid AS uuid)) AS versions,
- (SELECT count(*) FROM public.active_metadata_correction_audits ca
- JOIN public.active_metadata_corrections c
- ON c.uid = ca.correction_uid
- JOIN public.active_metadata_assets a ON a.uid = c.asset_uid
- WHERE a.source_uid = CAST(:source_uid AS uuid)) AS audits
- """
- ),
- {"plan_uid": plan_uid, "source_uid": source_uid},
- ).mappings().one()
- assert dict(counts) == {"runs": 4, "versions": 3, "audits": 2}
- finally:
- with app.app_context():
- db.session.rollback()
- if plan_uid:
- for statement in (
- "DELETE FROM public.active_metadata_correction_audits "
- "WHERE correction_uid IN (SELECT c.uid FROM public.active_metadata_corrections c "
- "JOIN public.active_metadata_assets a ON a.uid = c.asset_uid "
- "WHERE a.source_uid = CAST(:source_uid AS uuid))",
- "DELETE FROM public.active_metadata_corrections "
- "WHERE asset_uid IN (SELECT uid FROM public.active_metadata_assets "
- "WHERE source_uid = CAST(:source_uid AS uuid))",
- "DELETE FROM public.active_metadata_health_signals "
- "WHERE asset_uid IN (SELECT uid FROM public.active_metadata_assets "
- "WHERE source_uid = CAST(:source_uid AS uuid))",
- "DELETE FROM public.active_metadata_lineage "
- "WHERE run_uid IN (SELECT uid FROM public.active_metadata_runs "
- "WHERE plan_uid = CAST(:plan_uid AS uuid))",
- "DELETE FROM public.active_metadata_changes "
- "WHERE run_uid IN (SELECT uid FROM public.active_metadata_runs "
- "WHERE plan_uid = CAST(:plan_uid AS uuid))",
- "DELETE FROM public.active_metadata_asset_versions "
- "WHERE asset_uid IN (SELECT uid FROM public.active_metadata_assets "
- "WHERE source_uid = CAST(:source_uid AS uuid))",
- "DELETE FROM public.active_metadata_assets "
- "WHERE source_uid = CAST(:source_uid AS uuid)",
- "DELETE FROM public.active_metadata_runs "
- "WHERE plan_uid = CAST(:plan_uid AS uuid)",
- "DELETE FROM public.active_metadata_plans "
- "WHERE uid = CAST(:plan_uid AS uuid)",
- ):
- db.session.execute(
- text(statement),
- {"plan_uid": plan_uid, "source_uid": source_uid},
- )
- db.session.execute(
- text(
- "DELETE FROM public.ingestion_sources "
- "WHERE uid = CAST(:source_uid AS uuid)"
- ),
- {"source_uid": source_uid},
- )
- db.session.execute(
- text(
- "DELETE FROM public.users "
- "WHERE id IN (CAST(:actor_uid AS uuid), CAST(:owner_uid AS uuid))"
- ),
- {"actor_uid": actor_uid, "owner_uid": owner_uid},
- )
- db.session.commit()
|