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