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