test_quality_operations_postgres.py 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433
  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 _definition():
  9. return {
  10. "schema_version": "1.0",
  11. "field_roles": {
  12. "identity": {"required": True, "checks": ["unique", "pattern"]},
  13. "category": {"required": True, "checks": ["distribution"]},
  14. "measure": {"required": False, "checks": ["outlier"]},
  15. },
  16. "thresholds": {
  17. "completeness_min": 0.95,
  18. "uniqueness_min": 1,
  19. "pattern": "^[A-Z]{2,5}-[0-9]{3}$",
  20. "volume_change_max_ratio": 0.25,
  21. "distribution_drift_max": 0.2,
  22. "freshness_max_seconds": 3600,
  23. "quality_score_min": 80,
  24. },
  25. "sample_limit": 5,
  26. }
  27. def test_same_template_persists_two_domains_and_root_cause_evidence(monkeypatch):
  28. database_url = os.environ.get("TEST_DATABASE_URL")
  29. if not database_url:
  30. pytest.skip("TEST_DATABASE_URL is required")
  31. monkeypatch.setenv("DATABASE_URL", database_url)
  32. from app import create_app, db
  33. from app.core.data_rules.quality_operations import QualityOperationsService
  34. from app.core.data_rules.quality_repository import (
  35. SqlAlchemyQualityOperationsRepository,
  36. )
  37. app = create_app()
  38. app.config.update(TESTING=True)
  39. actor_uid = str(uuid.uuid4())
  40. owner_uid = str(uuid.uuid4())
  41. source_uid = str(uuid.uuid4())
  42. plan_uid = str(uuid.uuid4())
  43. metadata_run_uid = str(uuid.uuid4())
  44. device_asset_uid = str(uuid.uuid4())
  45. parts_asset_uid = str(uuid.uuid4())
  46. template_uid = None
  47. run_uids = []
  48. now = datetime.now(UTC)
  49. try:
  50. with app.app_context():
  51. for uid, label in ((actor_uid, "actor"), (owner_uid, "owner")):
  52. db.session.execute(
  53. text(
  54. """
  55. INSERT INTO public.users (
  56. id, username, display_name, password_hash, status
  57. ) VALUES (
  58. CAST(:uid AS uuid), :username, :username,
  59. 'p2-wp04-integration-hash', 'active'
  60. )
  61. """
  62. ),
  63. {
  64. "uid": uid,
  65. "username": f"wp04-{label}-{uid[:8]}",
  66. },
  67. )
  68. db.session.execute(
  69. text(
  70. """
  71. INSERT INTO public.ingestion_sources (
  72. uid, source_type, name, config, permission_scope,
  73. status, created_by
  74. ) VALUES (
  75. CAST(:uid AS uuid), 'database', :name,
  76. '{}'::jsonb, '{}'::jsonb, 'active', :created_by
  77. )
  78. """
  79. ),
  80. {
  81. "uid": source_uid,
  82. "name": f"WP04 source {source_uid[:8]}",
  83. "created_by": actor_uid,
  84. },
  85. )
  86. db.session.execute(
  87. text(
  88. """
  89. INSERT INTO public.active_metadata_plans (
  90. uid, source_uid, name, source_kind, schedule_type,
  91. discovery_mode, scope, cursor_state, owner_uid,
  92. enabled, current_version, created_by
  93. ) VALUES (
  94. CAST(:uid AS uuid), CAST(:source_uid AS uuid), :name,
  95. 'database', 'manual', 'snapshot', '{}'::jsonb,
  96. '{}'::jsonb, CAST(:owner_uid AS uuid), TRUE, 1,
  97. CAST(:created_by AS uuid)
  98. )
  99. """
  100. ),
  101. {
  102. "uid": plan_uid,
  103. "source_uid": source_uid,
  104. "name": f"WP04 plan {plan_uid[:8]}",
  105. "owner_uid": owner_uid,
  106. "created_by": actor_uid,
  107. },
  108. )
  109. db.session.execute(
  110. text(
  111. """
  112. INSERT INTO public.active_metadata_runs (
  113. uid, plan_uid, batch_key, status, attempt_count,
  114. cursor_before, cursor_after, snapshot_hash, statistics,
  115. actor_uid, started_at, finished_at
  116. ) VALUES (
  117. CAST(:uid AS uuid), CAST(:plan_uid AS uuid),
  118. 'wp04-metadata', 'completed', 1, '{}'::jsonb,
  119. '{}'::jsonb, :snapshot_hash, '{}'::jsonb,
  120. CAST(:actor_uid AS uuid), :started_at, :finished_at
  121. )
  122. """
  123. ),
  124. {
  125. "uid": metadata_run_uid,
  126. "plan_uid": plan_uid,
  127. "snapshot_hash": "a" * 64,
  128. "actor_uid": actor_uid,
  129. "started_at": now - timedelta(minutes=5),
  130. "finished_at": now,
  131. },
  132. )
  133. for asset_uid, namespace, name, domain in (
  134. (
  135. device_asset_uid,
  136. "maintenance",
  137. "device_ledger",
  138. "device",
  139. ),
  140. (
  141. parts_asset_uid,
  142. "supply",
  143. "spare_parts",
  144. "spare_parts",
  145. ),
  146. ):
  147. db.session.execute(
  148. text(
  149. """
  150. INSERT INTO public.active_metadata_assets (
  151. uid, source_uid, asset_key, namespace, name,
  152. asset_type, lifecycle_status, current_version,
  153. content_hash, snapshot, health, last_run_uid
  154. ) VALUES (
  155. CAST(:uid AS uuid), CAST(:source_uid AS uuid),
  156. :asset_key, :namespace, :name, 'table', 'active',
  157. 1, :content_hash, CAST(:snapshot AS jsonb),
  158. '{}'::jsonb, CAST(:last_run_uid AS uuid)
  159. )
  160. """
  161. ),
  162. {
  163. "uid": asset_uid,
  164. "source_uid": source_uid,
  165. "asset_key": f"{source_uid}:{namespace}.{name}",
  166. "namespace": namespace,
  167. "name": name,
  168. "content_hash": "b" * 64,
  169. "snapshot": (
  170. '{"business_domain_uid": "' + domain + '"}'
  171. ),
  172. "last_run_uid": metadata_run_uid,
  173. },
  174. )
  175. db.session.execute(
  176. text(
  177. """
  178. INSERT INTO public.active_metadata_changes (
  179. uid, run_uid, asset_uid, asset_key, field_name,
  180. change_type, before_state, after_state, status
  181. ) VALUES (
  182. CAST(:uid AS uuid), CAST(:run_uid AS uuid),
  183. CAST(:asset_uid AS uuid), :asset_key, 'part_category',
  184. 'field_changed', '{}'::jsonb, '{}'::jsonb, 'accepted'
  185. )
  186. """
  187. ),
  188. {
  189. "uid": str(uuid.uuid4()),
  190. "run_uid": metadata_run_uid,
  191. "asset_uid": parts_asset_uid,
  192. "asset_key": f"{source_uid}:supply.spare_parts",
  193. },
  194. )
  195. db.session.execute(
  196. text(
  197. """
  198. INSERT INTO public.active_metadata_lineage (
  199. uid, run_uid, parse_status, source_asset, source_field,
  200. target_asset, target_field, relation_type, evidence
  201. ) VALUES (
  202. CAST(:uid AS uuid), CAST(:run_uid AS uuid), 'resolved',
  203. 'erp.parts', 'part_code', 'supply.spare_parts',
  204. 'part_code', 'derived_from', '{}'::jsonb
  205. )
  206. """
  207. ),
  208. {"uid": str(uuid.uuid4()), "run_uid": metadata_run_uid},
  209. )
  210. db.session.commit()
  211. repository = SqlAlchemyQualityOperationsRepository(db.session)
  212. quality = QualityOperationsService(
  213. repository,
  214. publish_authorizer=lambda _actor_uid: None,
  215. commit=db.session.commit,
  216. rollback=db.session.rollback,
  217. )
  218. draft = quality.create_template(
  219. {
  220. "code": f"WP04_GENERIC_{source_uid[:8].upper()}",
  221. "name": "WP04 跨域通用质量模板",
  222. "owner_uid": owner_uid,
  223. "definition": _definition(),
  224. },
  225. actor_uid=actor_uid,
  226. )
  227. template_uid = draft["uid"]
  228. published = quality.publish_template(
  229. template_uid,
  230. expected_version=1,
  231. actor_uid=actor_uid,
  232. )
  233. device = quality.execute(
  234. {
  235. "template_uid": template_uid,
  236. "asset_uid": device_asset_uid,
  237. "batch_key": "device-baseline",
  238. "field_bindings": {
  239. "identity": "device_code",
  240. "category": "device_type",
  241. "measure": "temperature",
  242. },
  243. "source_observed_at": (
  244. now - timedelta(minutes=10)
  245. ).isoformat(),
  246. "records": [
  247. {
  248. "device_code": "DEV-001",
  249. "device_type": "pump",
  250. "temperature": 22,
  251. },
  252. {
  253. "device_code": "DEV-002",
  254. "device_type": "fan",
  255. "temperature": 23,
  256. },
  257. ],
  258. },
  259. actor_uid=actor_uid,
  260. )
  261. run_uids.append(device["uid"])
  262. parts = quality.execute(
  263. {
  264. "template_uid": template_uid,
  265. "asset_uid": parts_asset_uid,
  266. "batch_key": "parts-degraded",
  267. "field_bindings": {
  268. "identity": "part_code",
  269. "category": "part_category",
  270. "measure": "stock_quantity",
  271. },
  272. "source_observed_at": (
  273. now - timedelta(hours=2)
  274. ).isoformat(),
  275. "records": [
  276. {
  277. "part_code": "PT-001",
  278. "part_category": "seal",
  279. "stock_quantity": 1,
  280. },
  281. {
  282. "part_code": "PT-001",
  283. "part_category": None,
  284. "stock_quantity": 2,
  285. },
  286. {
  287. "part_code": None,
  288. "part_category": "seal",
  289. "stock_quantity": 3,
  290. },
  291. {
  292. "part_code": "bad",
  293. "part_category": "seal",
  294. "stock_quantity": 5000,
  295. },
  296. ],
  297. },
  298. actor_uid=actor_uid,
  299. )
  300. run_uids.append(parts["uid"])
  301. detail = quality.get_run(parts["uid"])
  302. assert published["status"] == "published"
  303. assert device["template_version_uid"] == parts["template_version_uid"]
  304. assert {
  305. device["business_domain_uid"],
  306. parts["business_domain_uid"],
  307. } == {"device", "spare_parts"}
  308. assert detail["metrics"]
  309. assert {
  310. item["finding_type"] for item in detail["findings"]
  311. } >= {
  312. "completeness",
  313. "uniqueness",
  314. "pattern",
  315. "duplicate",
  316. "outlier",
  317. "freshness",
  318. }
  319. assert all(
  320. item["evidence"]["root_cause"]["lineage"]
  321. for item in detail["findings"]
  322. )
  323. assert all(
  324. item["evidence"]["root_cause"]["changes"]
  325. for item in detail["findings"]
  326. )
  327. assert all(
  328. item["evidence"]["root_cause"]["runs"]
  329. for item in detail["findings"]
  330. )
  331. assert all(
  332. item["evidence"]["root_cause"]["responsibility"]["owner_uid"]
  333. == owner_uid
  334. for item in detail["findings"]
  335. )
  336. assert {
  337. item["sla_type"] for item in detail["sla_events"]
  338. } == {"freshness", "quality_score"}
  339. finally:
  340. with app.app_context():
  341. if run_uids:
  342. for table in (
  343. "quality_sla_events",
  344. "quality_findings",
  345. "quality_profile_metrics",
  346. ):
  347. db.session.execute(
  348. text(
  349. f"DELETE FROM public.{table} "
  350. "WHERE run_uid = ANY(CAST(:uids AS uuid[]))"
  351. ),
  352. {"uids": run_uids},
  353. )
  354. db.session.execute(
  355. text(
  356. """
  357. DELETE FROM public.quality_profile_runs
  358. WHERE uid = ANY(CAST(:uids AS uuid[]))
  359. """
  360. ),
  361. {"uids": run_uids},
  362. )
  363. if template_uid:
  364. db.session.execute(
  365. text(
  366. """
  367. UPDATE public.quality_templates
  368. SET active_version_uid = NULL
  369. WHERE uid = CAST(:uid AS uuid)
  370. """
  371. ),
  372. {"uid": template_uid},
  373. )
  374. db.session.execute(
  375. text(
  376. """
  377. DELETE FROM public.quality_template_versions
  378. WHERE template_uid = CAST(:uid AS uuid)
  379. """
  380. ),
  381. {"uid": template_uid},
  382. )
  383. db.session.execute(
  384. text(
  385. """
  386. DELETE FROM public.quality_templates
  387. WHERE uid = CAST(:uid AS uuid)
  388. """
  389. ),
  390. {"uid": template_uid},
  391. )
  392. db.session.execute(
  393. text(
  394. """
  395. DELETE FROM public.active_metadata_lineage
  396. WHERE run_uid = CAST(:uid AS uuid);
  397. DELETE FROM public.active_metadata_changes
  398. WHERE run_uid = CAST(:uid AS uuid);
  399. DELETE FROM public.active_metadata_assets
  400. WHERE last_run_uid = CAST(:uid AS uuid);
  401. DELETE FROM public.active_metadata_runs
  402. WHERE uid = CAST(:uid AS uuid);
  403. DELETE FROM public.active_metadata_plans
  404. WHERE uid = CAST(:plan_uid AS uuid);
  405. DELETE FROM public.ingestion_sources
  406. WHERE uid = CAST(:source_uid AS uuid);
  407. DELETE FROM public.users
  408. WHERE id IN (
  409. CAST(:actor_uid AS uuid), CAST(:owner_uid AS uuid)
  410. );
  411. """
  412. ),
  413. {
  414. "uid": metadata_run_uid,
  415. "plan_uid": plan_uid,
  416. "source_uid": source_uid,
  417. "actor_uid": actor_uid,
  418. "owner_uid": owner_uid,
  419. },
  420. )
  421. db.session.commit()