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 test_failed_collection_quality_sla_and_recovery_persist_incidents( 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.events.data_observability import DataObservabilityService from app.core.events.data_observability_repository import ( SqlAlchemyDataObservabilityRepository, ) 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()) failed_run_uid = str(uuid.uuid4()) recovered_run_uid = str(uuid.uuid4()) asset_uid = str(uuid.uuid4()) template_uid = str(uuid.uuid4()) version_uid = str(uuid.uuid4()) quality_run_uids = [str(uuid.uuid4()), str(uuid.uuid4())] quality_event_uids = [str(uuid.uuid4()), str(uuid.uuid4())] 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-wp05-integration-hash', 'active' ) """ ), { "uid": uid, "username": f"wp05-{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"WP05 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', CAST(:scope AS jsonb), '{}'::jsonb, CAST(:owner_uid AS uuid), TRUE, 1, CAST(:created_by AS uuid) ) """ ), { "uid": plan_uid, "source_uid": source_uid, "name": "设备主数据采集", "scope": ( '{"business_domain_uid":"device",' '"business_domain_name":"设备域",' '"data_product_uid":"device-ledger",' '"data_product_name":"设备台账",' '"user_impact_group_uid":"maintenance-ops",' '"user_impact_group_name":"设备运维人员"}' ), "owner_uid": owner_uid, "created_by": actor_uid, }, ) for uid, batch, status, finished, failure in ( ( failed_run_uid, "wp05-failed", "failed", now - timedelta(minutes=10), "SOURCE_TIMEOUT", ), ( recovered_run_uid, "wp05-recovered", "completed", now - timedelta(minutes=5), None, ), ): 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, failure_code, failure_reason, actor_uid, started_at, finished_at ) VALUES ( CAST(:uid AS uuid), CAST(:plan_uid AS uuid), :batch_key, :status, 1, '{}'::jsonb, '{}'::jsonb, :snapshot_hash, '{}'::jsonb, :failure_code, :failure_reason, CAST(:actor_uid AS uuid), :started_at, :finished_at ) """ ), { "uid": uid, "plan_uid": plan_uid, "batch_key": batch, "status": status, "snapshot_hash": "5" * 64 if status == "completed" else None, "failure_code": failure, "failure_reason": "source timeout" if failure else None, "actor_uid": actor_uid, "started_at": finished - timedelta(minutes=1), "finished_at": finished, }, ) 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, 'maintenance', 'device_ledger', '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}:maintenance.device_ledger", "content_hash": "6" * 64, "snapshot": ( '{"business_domain_uid":"device",' '"business_domain_name":"设备域",' '"user_impact_group_uid":"maintenance-ops",' '"user_impact_group_name":"设备运维人员"}' ), "last_run_uid": recovered_run_uid, }, ) db.session.execute( text( """ INSERT INTO public.quality_templates ( uid, code, name, owner_uid, status, current_version, created_by ) VALUES ( CAST(:uid AS uuid), :code, 'WP05 quality', CAST(:owner_uid AS uuid), 'published', 1, CAST(:created_by AS uuid) ) """ ), { "uid": template_uid, "code": f"WP05_{template_uid[:8].upper()}", "owner_uid": owner_uid, "created_by": actor_uid, }, ) db.session.execute( text( """ INSERT INTO public.quality_template_versions ( uid, template_uid, version, status, definition, content_hash, created_by, published_by, published_at ) VALUES ( CAST(:uid AS uuid), CAST(:template_uid AS uuid), 1, 'published', '{}'::jsonb, :content_hash, CAST(:created_by AS uuid), CAST(:created_by AS uuid), :published_at ) """ ), { "uid": version_uid, "template_uid": template_uid, "content_hash": "7" * 64, "created_by": actor_uid, "published_at": now - timedelta(minutes=30), }, ) db.session.execute( text( """ UPDATE public.quality_templates SET active_version_uid = CAST(:version_uid AS uuid) WHERE uid = CAST(:template_uid AS uuid) """ ), {"version_uid": version_uid, "template_uid": template_uid}, ) for index, (run_uid, event_uid, event_status, score) in enumerate( zip( quality_run_uids, quality_event_uids, ("violated", "recovered"), (68, 92), strict=True, ) ): created_at = now - timedelta(minutes=4 - index) db.session.execute( text( """ INSERT INTO public.quality_profile_runs ( uid, template_uid, template_version_uid, template_hash, asset_uid, source_uid, business_domain_uid, batch_key, status, row_count, score, source_observed_at, comparison, profile, field_bindings, finding_count, deterministic, created_by, created_at ) VALUES ( CAST(:uid AS uuid), CAST(:template_uid AS uuid), CAST(:version_uid AS uuid), :template_hash, CAST(:asset_uid AS uuid), CAST(:source_uid AS uuid), 'device', :batch_key, 'success', 10, :score, :observed_at, '{}'::jsonb, '{}'::jsonb, '{}'::jsonb, 0, TRUE, CAST(:created_by AS uuid), :created_at ) """ ), { "uid": run_uid, "template_uid": template_uid, "version_uid": version_uid, "template_hash": "7" * 64, "asset_uid": asset_uid, "source_uid": source_uid, "batch_key": f"quality-{index}", "score": score, "observed_at": created_at, "created_by": actor_uid, "created_at": created_at, }, ) db.session.execute( text( """ INSERT INTO public.quality_sla_events ( uid, run_uid, asset_uid, sla_type, status, severity, actual, threshold, owner_uid, escalation_level, evidence, created_at ) VALUES ( CAST(:uid AS uuid), CAST(:run_uid AS uuid), CAST(:asset_uid AS uuid), 'quality_score', :status, :severity, :actual, 80, CAST(:owner_uid AS uuid), :escalation_level, CAST(:evidence AS jsonb), :created_at ) """ ), { "uid": event_uid, "run_uid": run_uid, "asset_uid": asset_uid, "status": event_status, "severity": ( "critical" if event_status == "violated" else "info" ), "actual": score, "owner_uid": owner_uid, "escalation_level": ( 1 if event_status == "violated" else 0 ), "evidence": '{"deterministic":true}', "created_at": created_at, }, ) db.session.commit() observability = DataObservabilityService( SqlAlchemyDataObservabilityRepository(db.session), commit=db.session.commit, rollback=db.session.rollback, ) result = observability.collect(actor_uid=actor_uid) incidents = observability.list_incidents() assert result["processed"] == 4 assert result["created_alerts"] == 2 assert result["recovered_alerts"] == 2 assert len(incidents) == 2 assert all(item["status"] == "monitoring" for item in incidents) ingestion = next( item for item in incidents if "采集" in item["title"] ) detail = observability.incident_detail(ingestion["uid"]) assert detail["alerts"][0]["status"] == "recovered" assert { item["target_type"] for item in detail["impacts"] } >= { "service", "business_domain", "data_product", "user_group", } closed = observability.close_incident( ingestion["uid"], { "root_cause": "源系统连接超时", "user_impact": "设备运维人员短时看到旧台账", "corrective_actions": ["增加连接超时监控"], "closure_evidence": [ {"kind": "recovery_run", "ref": recovered_run_uid} ], }, actor_uid=actor_uid, ) assert closed["status"] == "closed" assert closed["postmortem"]["closure_evidence"] finally: with app.app_context(): db.session.rollback() db.session.execute( text( """ DELETE FROM public.data_incident_postmortems WHERE created_by = CAST(:actor_uid AS uuid); DELETE FROM public.data_incident_timeline WHERE actor_uid = CAST(:actor_uid AS uuid); DELETE FROM public.data_observability_source_events WHERE source_uid IN ( :failed_run_uid, :recovered_run_uid, :quality_event_uid_1, :quality_event_uid_2 ); DELETE FROM public.data_observability_alerts WHERE owner_uid = CAST(:owner_uid AS uuid); DELETE FROM public.data_incident_impacts WHERE incident_uid IN ( SELECT uid FROM public.data_incidents WHERE created_by = CAST(:actor_uid AS uuid) ); DELETE FROM public.data_incidents WHERE created_by = CAST(:actor_uid AS uuid); DELETE FROM public.quality_sla_events WHERE uid IN ( CAST(:quality_event_uid_1 AS uuid), CAST(:quality_event_uid_2 AS uuid) ); DELETE FROM public.quality_profile_runs WHERE uid IN ( CAST(:quality_run_uid_1 AS uuid), CAST(:quality_run_uid_2 AS uuid) ); UPDATE public.quality_templates SET active_version_uid = NULL WHERE uid = CAST(:template_uid AS uuid); DELETE FROM public.quality_template_versions WHERE uid = CAST(:version_uid AS uuid); DELETE FROM public.quality_templates WHERE uid = CAST(:template_uid AS uuid); DELETE FROM public.active_metadata_assets WHERE uid = CAST(:asset_uid AS uuid); DELETE FROM public.active_metadata_runs WHERE uid IN ( CAST(:failed_run_uid AS uuid), CAST(:recovered_run_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) ); """ ), { "actor_uid": actor_uid, "owner_uid": owner_uid, "source_uid": source_uid, "plan_uid": plan_uid, "failed_run_uid": failed_run_uid, "recovered_run_uid": recovered_run_uid, "asset_uid": asset_uid, "template_uid": template_uid, "version_uid": version_uid, "quality_run_uid_1": quality_run_uids[0], "quality_run_uid_2": quality_run_uids[1], "quality_event_uid_1": quality_event_uids[0], "quality_event_uid_2": quality_event_uids[1], }, ) db.session.commit()