test_data_observability_postgres.py 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431
  1. from __future__ import annotations
  2. import os
  3. import uuid
  4. from datetime import UTC, datetime, timedelta
  5. import pytest
  6. from sqlalchemy import text
  7. pytestmark = pytest.mark.integration
  8. def test_failed_collection_quality_sla_and_recovery_persist_incidents(
  9. monkeypatch,
  10. ):
  11. database_url = os.environ.get("TEST_DATABASE_URL")
  12. if not database_url:
  13. pytest.skip("TEST_DATABASE_URL is required")
  14. monkeypatch.setenv("DATABASE_URL", database_url)
  15. from app import create_app, db
  16. from app.core.events.data_observability import DataObservabilityService
  17. from app.core.events.data_observability_repository import (
  18. SqlAlchemyDataObservabilityRepository,
  19. )
  20. app = create_app()
  21. app.config.update(TESTING=True)
  22. actor_uid = str(uuid.uuid4())
  23. owner_uid = str(uuid.uuid4())
  24. source_uid = str(uuid.uuid4())
  25. plan_uid = str(uuid.uuid4())
  26. failed_run_uid = str(uuid.uuid4())
  27. recovered_run_uid = str(uuid.uuid4())
  28. asset_uid = str(uuid.uuid4())
  29. template_uid = str(uuid.uuid4())
  30. version_uid = str(uuid.uuid4())
  31. quality_run_uids = [str(uuid.uuid4()), str(uuid.uuid4())]
  32. quality_event_uids = [str(uuid.uuid4()), str(uuid.uuid4())]
  33. now = datetime.now(UTC)
  34. try:
  35. with app.app_context():
  36. for uid, label in ((actor_uid, "actor"), (owner_uid, "owner")):
  37. db.session.execute(
  38. text(
  39. """
  40. INSERT INTO public.users (
  41. id, username, display_name, password_hash, status
  42. ) VALUES (
  43. CAST(:uid AS uuid), :username, :username,
  44. 'p2-wp05-integration-hash', 'active'
  45. )
  46. """
  47. ),
  48. {
  49. "uid": uid,
  50. "username": f"wp05-{label}-{uid[:8]}",
  51. },
  52. )
  53. db.session.execute(
  54. text(
  55. """
  56. INSERT INTO public.ingestion_sources (
  57. uid, source_type, name, config, permission_scope,
  58. status, created_by
  59. ) VALUES (
  60. CAST(:uid AS uuid), 'database', :name,
  61. '{}'::jsonb, '{}'::jsonb, 'active', :created_by
  62. )
  63. """
  64. ),
  65. {
  66. "uid": source_uid,
  67. "name": f"WP05 source {source_uid[:8]}",
  68. "created_by": actor_uid,
  69. },
  70. )
  71. db.session.execute(
  72. text(
  73. """
  74. INSERT INTO public.active_metadata_plans (
  75. uid, source_uid, name, source_kind, schedule_type,
  76. discovery_mode, scope, cursor_state, owner_uid,
  77. enabled, current_version, created_by
  78. ) VALUES (
  79. CAST(:uid AS uuid), CAST(:source_uid AS uuid), :name,
  80. 'database', 'manual', 'snapshot',
  81. CAST(:scope AS jsonb), '{}'::jsonb,
  82. CAST(:owner_uid AS uuid), TRUE, 1,
  83. CAST(:created_by AS uuid)
  84. )
  85. """
  86. ),
  87. {
  88. "uid": plan_uid,
  89. "source_uid": source_uid,
  90. "name": "设备主数据采集",
  91. "scope": (
  92. '{"business_domain_uid":"device",'
  93. '"business_domain_name":"设备域",'
  94. '"data_product_uid":"device-ledger",'
  95. '"data_product_name":"设备台账",'
  96. '"user_impact_group_uid":"maintenance-ops",'
  97. '"user_impact_group_name":"设备运维人员"}'
  98. ),
  99. "owner_uid": owner_uid,
  100. "created_by": actor_uid,
  101. },
  102. )
  103. for uid, batch, status, finished, failure in (
  104. (
  105. failed_run_uid,
  106. "wp05-failed",
  107. "failed",
  108. now - timedelta(minutes=10),
  109. "SOURCE_TIMEOUT",
  110. ),
  111. (
  112. recovered_run_uid,
  113. "wp05-recovered",
  114. "completed",
  115. now - timedelta(minutes=5),
  116. None,
  117. ),
  118. ):
  119. db.session.execute(
  120. text(
  121. """
  122. INSERT INTO public.active_metadata_runs (
  123. uid, plan_uid, batch_key, status, attempt_count,
  124. cursor_before, cursor_after, snapshot_hash,
  125. statistics, failure_code, failure_reason,
  126. actor_uid, started_at, finished_at
  127. ) VALUES (
  128. CAST(:uid AS uuid), CAST(:plan_uid AS uuid),
  129. :batch_key, :status, 1, '{}'::jsonb, '{}'::jsonb,
  130. :snapshot_hash, '{}'::jsonb, :failure_code,
  131. :failure_reason, CAST(:actor_uid AS uuid),
  132. :started_at, :finished_at
  133. )
  134. """
  135. ),
  136. {
  137. "uid": uid,
  138. "plan_uid": plan_uid,
  139. "batch_key": batch,
  140. "status": status,
  141. "snapshot_hash": "5" * 64 if status == "completed" else None,
  142. "failure_code": failure,
  143. "failure_reason": "source timeout" if failure else None,
  144. "actor_uid": actor_uid,
  145. "started_at": finished - timedelta(minutes=1),
  146. "finished_at": finished,
  147. },
  148. )
  149. db.session.execute(
  150. text(
  151. """
  152. INSERT INTO public.active_metadata_assets (
  153. uid, source_uid, asset_key, namespace, name, asset_type,
  154. lifecycle_status, current_version, content_hash,
  155. snapshot, health, last_run_uid
  156. ) VALUES (
  157. CAST(:uid AS uuid), CAST(:source_uid AS uuid),
  158. :asset_key, 'maintenance', 'device_ledger', 'table',
  159. 'active', 1, :content_hash, CAST(:snapshot AS jsonb),
  160. '{}'::jsonb, CAST(:last_run_uid AS uuid)
  161. )
  162. """
  163. ),
  164. {
  165. "uid": asset_uid,
  166. "source_uid": source_uid,
  167. "asset_key": f"{source_uid}:maintenance.device_ledger",
  168. "content_hash": "6" * 64,
  169. "snapshot": (
  170. '{"business_domain_uid":"device",'
  171. '"business_domain_name":"设备域",'
  172. '"user_impact_group_uid":"maintenance-ops",'
  173. '"user_impact_group_name":"设备运维人员"}'
  174. ),
  175. "last_run_uid": recovered_run_uid,
  176. },
  177. )
  178. db.session.execute(
  179. text(
  180. """
  181. INSERT INTO public.quality_templates (
  182. uid, code, name, owner_uid, status, current_version,
  183. created_by
  184. ) VALUES (
  185. CAST(:uid AS uuid), :code, 'WP05 quality',
  186. CAST(:owner_uid AS uuid), 'published', 1,
  187. CAST(:created_by AS uuid)
  188. )
  189. """
  190. ),
  191. {
  192. "uid": template_uid,
  193. "code": f"WP05_{template_uid[:8].upper()}",
  194. "owner_uid": owner_uid,
  195. "created_by": actor_uid,
  196. },
  197. )
  198. db.session.execute(
  199. text(
  200. """
  201. INSERT INTO public.quality_template_versions (
  202. uid, template_uid, version, status, definition,
  203. content_hash, created_by, published_by, published_at
  204. ) VALUES (
  205. CAST(:uid AS uuid), CAST(:template_uid AS uuid), 1,
  206. 'published', '{}'::jsonb, :content_hash,
  207. CAST(:created_by AS uuid), CAST(:created_by AS uuid),
  208. :published_at
  209. )
  210. """
  211. ),
  212. {
  213. "uid": version_uid,
  214. "template_uid": template_uid,
  215. "content_hash": "7" * 64,
  216. "created_by": actor_uid,
  217. "published_at": now - timedelta(minutes=30),
  218. },
  219. )
  220. db.session.execute(
  221. text(
  222. """
  223. UPDATE public.quality_templates
  224. SET active_version_uid = CAST(:version_uid AS uuid)
  225. WHERE uid = CAST(:template_uid AS uuid)
  226. """
  227. ),
  228. {"version_uid": version_uid, "template_uid": template_uid},
  229. )
  230. for index, (run_uid, event_uid, event_status, score) in enumerate(
  231. zip(
  232. quality_run_uids,
  233. quality_event_uids,
  234. ("violated", "recovered"),
  235. (68, 92),
  236. strict=True,
  237. )
  238. ):
  239. created_at = now - timedelta(minutes=4 - index)
  240. db.session.execute(
  241. text(
  242. """
  243. INSERT INTO public.quality_profile_runs (
  244. uid, template_uid, template_version_uid,
  245. template_hash, asset_uid, source_uid,
  246. business_domain_uid, batch_key, status, row_count,
  247. score, source_observed_at, comparison, profile,
  248. field_bindings, finding_count, deterministic,
  249. created_by, created_at
  250. ) VALUES (
  251. CAST(:uid AS uuid), CAST(:template_uid AS uuid),
  252. CAST(:version_uid AS uuid), :template_hash,
  253. CAST(:asset_uid AS uuid),
  254. CAST(:source_uid AS uuid), 'device', :batch_key,
  255. 'success', 10, :score, :observed_at, '{}'::jsonb,
  256. '{}'::jsonb, '{}'::jsonb, 0, TRUE,
  257. CAST(:created_by AS uuid), :created_at
  258. )
  259. """
  260. ),
  261. {
  262. "uid": run_uid,
  263. "template_uid": template_uid,
  264. "version_uid": version_uid,
  265. "template_hash": "7" * 64,
  266. "asset_uid": asset_uid,
  267. "source_uid": source_uid,
  268. "batch_key": f"quality-{index}",
  269. "score": score,
  270. "observed_at": created_at,
  271. "created_by": actor_uid,
  272. "created_at": created_at,
  273. },
  274. )
  275. db.session.execute(
  276. text(
  277. """
  278. INSERT INTO public.quality_sla_events (
  279. uid, run_uid, asset_uid, sla_type, status,
  280. severity, actual, threshold, owner_uid,
  281. escalation_level, evidence, created_at
  282. ) VALUES (
  283. CAST(:uid AS uuid), CAST(:run_uid AS uuid),
  284. CAST(:asset_uid AS uuid), 'quality_score',
  285. :status, :severity, :actual, 80,
  286. CAST(:owner_uid AS uuid), :escalation_level,
  287. CAST(:evidence AS jsonb), :created_at
  288. )
  289. """
  290. ),
  291. {
  292. "uid": event_uid,
  293. "run_uid": run_uid,
  294. "asset_uid": asset_uid,
  295. "status": event_status,
  296. "severity": (
  297. "critical" if event_status == "violated" else "info"
  298. ),
  299. "actual": score,
  300. "owner_uid": owner_uid,
  301. "escalation_level": (
  302. 1 if event_status == "violated" else 0
  303. ),
  304. "evidence": '{"deterministic":true}',
  305. "created_at": created_at,
  306. },
  307. )
  308. db.session.commit()
  309. observability = DataObservabilityService(
  310. SqlAlchemyDataObservabilityRepository(db.session),
  311. commit=db.session.commit,
  312. rollback=db.session.rollback,
  313. )
  314. result = observability.collect(actor_uid=actor_uid)
  315. incidents = observability.list_incidents()
  316. assert result["processed"] == 4
  317. assert result["created_alerts"] == 2
  318. assert result["recovered_alerts"] == 2
  319. assert len(incidents) == 2
  320. assert all(item["status"] == "monitoring" for item in incidents)
  321. ingestion = next(
  322. item for item in incidents if "采集" in item["title"]
  323. )
  324. detail = observability.incident_detail(ingestion["uid"])
  325. assert detail["alerts"][0]["status"] == "recovered"
  326. assert {
  327. item["target_type"] for item in detail["impacts"]
  328. } >= {
  329. "service",
  330. "business_domain",
  331. "data_product",
  332. "user_group",
  333. }
  334. closed = observability.close_incident(
  335. ingestion["uid"],
  336. {
  337. "root_cause": "源系统连接超时",
  338. "user_impact": "设备运维人员短时看到旧台账",
  339. "corrective_actions": ["增加连接超时监控"],
  340. "closure_evidence": [
  341. {"kind": "recovery_run", "ref": recovered_run_uid}
  342. ],
  343. },
  344. actor_uid=actor_uid,
  345. )
  346. assert closed["status"] == "closed"
  347. assert closed["postmortem"]["closure_evidence"]
  348. finally:
  349. with app.app_context():
  350. db.session.rollback()
  351. db.session.execute(
  352. text(
  353. """
  354. DELETE FROM public.data_incident_postmortems
  355. WHERE created_by = CAST(:actor_uid AS uuid);
  356. DELETE FROM public.data_incident_timeline
  357. WHERE actor_uid = CAST(:actor_uid AS uuid);
  358. DELETE FROM public.data_observability_source_events
  359. WHERE source_uid IN (
  360. :failed_run_uid, :recovered_run_uid,
  361. :quality_event_uid_1, :quality_event_uid_2
  362. );
  363. DELETE FROM public.data_observability_alerts
  364. WHERE owner_uid = CAST(:owner_uid AS uuid);
  365. DELETE FROM public.data_incident_impacts
  366. WHERE incident_uid IN (
  367. SELECT uid FROM public.data_incidents
  368. WHERE created_by = CAST(:actor_uid AS uuid)
  369. );
  370. DELETE FROM public.data_incidents
  371. WHERE created_by = CAST(:actor_uid AS uuid);
  372. DELETE FROM public.quality_sla_events
  373. WHERE uid IN (
  374. CAST(:quality_event_uid_1 AS uuid),
  375. CAST(:quality_event_uid_2 AS uuid)
  376. );
  377. DELETE FROM public.quality_profile_runs
  378. WHERE uid IN (
  379. CAST(:quality_run_uid_1 AS uuid),
  380. CAST(:quality_run_uid_2 AS uuid)
  381. );
  382. UPDATE public.quality_templates
  383. SET active_version_uid = NULL
  384. WHERE uid = CAST(:template_uid AS uuid);
  385. DELETE FROM public.quality_template_versions
  386. WHERE uid = CAST(:version_uid AS uuid);
  387. DELETE FROM public.quality_templates
  388. WHERE uid = CAST(:template_uid AS uuid);
  389. DELETE FROM public.active_metadata_assets
  390. WHERE uid = CAST(:asset_uid AS uuid);
  391. DELETE FROM public.active_metadata_runs
  392. WHERE uid IN (
  393. CAST(:failed_run_uid AS uuid),
  394. CAST(:recovered_run_uid AS uuid)
  395. );
  396. DELETE FROM public.active_metadata_plans
  397. WHERE uid = CAST(:plan_uid AS uuid);
  398. DELETE FROM public.ingestion_sources
  399. WHERE uid = CAST(:source_uid AS uuid);
  400. DELETE FROM public.users
  401. WHERE id IN (
  402. CAST(:actor_uid AS uuid), CAST(:owner_uid AS uuid)
  403. );
  404. """
  405. ),
  406. {
  407. "actor_uid": actor_uid,
  408. "owner_uid": owner_uid,
  409. "source_uid": source_uid,
  410. "plan_uid": plan_uid,
  411. "failed_run_uid": failed_run_uid,
  412. "recovered_run_uid": recovered_run_uid,
  413. "asset_uid": asset_uid,
  414. "template_uid": template_uid,
  415. "version_uid": version_uid,
  416. "quality_run_uid_1": quality_run_uids[0],
  417. "quality_run_uid_2": quality_run_uids[1],
  418. "quality_event_uid_1": quality_event_uids[0],
  419. "quality_event_uid_2": quality_event_uids[1],
  420. },
  421. )
  422. db.session.commit()