test_governance_metrics_postgres.py 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424
  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_governance_metrics_apply_scope_and_return_safe_details(monkeypatch):
  9. platform_url = os.environ.get("TEST_DATABASE_URL")
  10. if not platform_url:
  11. pytest.skip("TEST_DATABASE_URL is required")
  12. monkeypatch.setenv("DATABASE_URL", platform_url)
  13. from app import create_app, db
  14. from app.core.data_research.governance_metric_repository import (
  15. SqlAlchemyGovernanceMetricRepository,
  16. )
  17. from app.core.data_research.governance_metrics import GovernanceMetricAccess
  18. from app.models.data_research import (
  19. DeviceAsset,
  20. DeviceAssetSourceMapping,
  21. DeviceEntityMatchCandidate,
  22. DeviceEntityMatchReview,
  23. DeviceEntityMergeEvent,
  24. DeviceEntityMergeRollback,
  25. DeviceQualityIssue,
  26. DeviceQualityProfile,
  27. DeviceQualityProfileVersion,
  28. DeviceQualityRun,
  29. IngestionSource,
  30. )
  31. app = create_app()
  32. app.config.update(TESTING=True)
  33. suffix = uuid.uuid4().hex[:10]
  34. actor_uid = str(uuid.uuid4())
  35. now = datetime.now(UTC).replace(microsecond=0)
  36. domain_a = str(uuid.uuid4())
  37. domain_b = str(uuid.uuid4())
  38. source_uids = {
  39. "a": str(uuid.uuid4()),
  40. "b": str(uuid.uuid4()),
  41. "u": str(uuid.uuid4()),
  42. }
  43. asset_uids = {
  44. key: str(uuid.uuid4())
  45. for key in ("a1", "a2", "a3", "b", "u", "retired")
  46. }
  47. mapping_uids = {}
  48. candidate_uids = []
  49. merge_uids = []
  50. issue_uids = []
  51. profile_uid = str(uuid.uuid4())
  52. profile_version_uid = str(uuid.uuid4())
  53. run_uid = str(uuid.uuid4())
  54. baseline_admin = None
  55. try:
  56. with app.app_context():
  57. baseline_admin = SqlAlchemyGovernanceMetricRepository(
  58. db.session
  59. ).summary_counts(GovernanceMetricAccess(global_access=True))
  60. db.session.execute(
  61. text(
  62. """
  63. INSERT INTO public.users (
  64. id, username, display_name, password_hash, status
  65. ) VALUES (
  66. CAST(:id AS uuid), :username, :username,
  67. 'integration-test', 'active'
  68. )
  69. """
  70. ),
  71. {"id": actor_uid, "username": f"wp11-actor-{suffix}"},
  72. )
  73. for key, scope in (
  74. ("a", {"business_domains": [domain_a]}),
  75. ("b", {"business_domains": [domain_b]}),
  76. ("u", {}),
  77. ):
  78. db.session.add(
  79. IngestionSource(
  80. uid=source_uids[key],
  81. source_type="database",
  82. name=f"WP11 source {key} {suffix}",
  83. config={"password": f"source-secret-{key}"},
  84. permission_scope=scope,
  85. status="active",
  86. created_by=actor_uid,
  87. )
  88. )
  89. db.session.flush()
  90. for key, _source_key, complete, status in (
  91. ("a1", "a", True, "active"),
  92. ("a2", "a", False, "active"),
  93. ("a3", "a", True, "active"),
  94. ("b", "b", True, "active"),
  95. ("u", "u", True, "active"),
  96. ("retired", "a", True, "retired"),
  97. ):
  98. db.session.add(
  99. DeviceAsset(
  100. uid=asset_uids[key],
  101. asset_type="device",
  102. name=f"WP11设备-{suffix}-{key}",
  103. status=status,
  104. current_version=1,
  105. content_hash=(key[0] * 64),
  106. location=f"{key}车间",
  107. organization="设备部" if complete else " ",
  108. responsible_person="张工" if complete else None,
  109. attributes={"password": f"asset-secret-{key}"},
  110. created_by=actor_uid,
  111. updated_by=actor_uid,
  112. updated_at=now,
  113. )
  114. )
  115. db.session.flush()
  116. for key, source_key, _complete, _status in (
  117. ("a1", "a", True, "active"),
  118. ("a2", "a", False, "active"),
  119. ("a3", "a", True, "active"),
  120. ("b", "b", True, "active"),
  121. ("u", "u", True, "active"),
  122. ("retired", "a", True, "retired"),
  123. ):
  124. mapping_uid = str(uuid.uuid4())
  125. mapping_uids[key] = mapping_uid
  126. db.session.add(
  127. DeviceAssetSourceMapping(
  128. uid=mapping_uid,
  129. asset_uid=asset_uids[key],
  130. source_uid=source_uids[source_key],
  131. source_entity="asset.equipment",
  132. asset_type="device",
  133. source_code=f"EQ-WP11-{suffix}-{key}",
  134. source_updated_at=now,
  135. first_seen_at=now,
  136. last_seen_at=now,
  137. )
  138. )
  139. db.session.flush()
  140. def add_merge(left_key, right_key, *, rolled_back=False):
  141. candidate_uid = str(uuid.uuid4())
  142. review_uid = str(uuid.uuid4())
  143. merge_uid = str(uuid.uuid4())
  144. candidate_uids.append(candidate_uid)
  145. merge_uids.append(merge_uid)
  146. db.session.add(
  147. DeviceEntityMatchCandidate(
  148. uid=candidate_uid,
  149. left_asset_uid=asset_uids[left_key],
  150. right_asset_uid=asset_uids[right_key],
  151. canonical_asset_uid=asset_uids[left_key],
  152. status="rolled_back" if rolled_back else "merged",
  153. suggestion_source="manual",
  154. confidence=1,
  155. explanation=[],
  156. evidence_uids=[],
  157. current_version=2 if rolled_back else 1,
  158. created_by=actor_uid,
  159. )
  160. )
  161. db.session.flush()
  162. db.session.add(
  163. DeviceEntityMatchReview(
  164. uid=review_uid,
  165. candidate_uid=candidate_uid,
  166. version=1,
  167. decision="approve",
  168. reason="WP11 integration",
  169. actor_uid=actor_uid,
  170. )
  171. )
  172. db.session.flush()
  173. db.session.add(
  174. DeviceEntityMergeEvent(
  175. uid=merge_uid,
  176. candidate_uid=candidate_uid,
  177. canonical_asset_uid=asset_uids[left_key],
  178. member_asset_uid=asset_uids[right_key],
  179. review_uid=review_uid,
  180. snapshot={"password": "merge-secret"},
  181. actor_uid=actor_uid,
  182. )
  183. )
  184. db.session.flush()
  185. if rolled_back:
  186. db.session.add(
  187. DeviceEntityMergeRollback(
  188. uid=str(uuid.uuid4()),
  189. merge_uid=merge_uid,
  190. candidate_uid=candidate_uid,
  191. reason="WP11 rollback",
  192. snapshot={"password": "rollback-secret"},
  193. actor_uid=actor_uid,
  194. )
  195. )
  196. add_merge("a1", "a2")
  197. add_merge("a3", "b")
  198. add_merge("a1", "a3", rolled_back=True)
  199. db.session.add(
  200. DeviceQualityProfile(
  201. uid=profile_uid,
  202. code=f"wp11-{suffix}",
  203. name="WP11 integration",
  204. created_by=actor_uid,
  205. )
  206. )
  207. db.session.flush()
  208. db.session.add(
  209. DeviceQualityProfileVersion(
  210. uid=profile_version_uid,
  211. profile_uid=profile_uid,
  212. version=1,
  213. status="published",
  214. rules=[],
  215. content_hash=suffix.ljust(64, "0"),
  216. created_by=actor_uid,
  217. published_by=actor_uid,
  218. published_at=now,
  219. )
  220. )
  221. db.session.flush()
  222. db.session.add(
  223. DeviceQualityRun(
  224. uid=run_uid,
  225. policy_version_uid=profile_version_uid,
  226. policy_hash=suffix.ljust(64, "0"),
  227. source_uid=source_uids["a"],
  228. status="success",
  229. total_assets=3,
  230. total_violations=3,
  231. score=50,
  232. created_by=actor_uid,
  233. )
  234. )
  235. db.session.flush()
  236. for index, (asset_key, source_key, status, occurrence) in enumerate(
  237. (
  238. ("a1", "a", "closed", 1),
  239. ("a2", "a", "open", 2),
  240. ("b", "b", "closed", 1),
  241. ),
  242. start=1,
  243. ):
  244. issue_uid = str(uuid.uuid4())
  245. issue_uids.append(issue_uid)
  246. db.session.add(
  247. DeviceQualityIssue(
  248. uid=issue_uid,
  249. issue_code=f"W11{suffix[:6]}{index:02d}",
  250. source_violation_uid=str(uuid.uuid4()),
  251. source_run_uid=run_uid,
  252. rule_code="asset_context_complete",
  253. severity="error",
  254. priority="high",
  255. asset_uid=asset_uids[asset_key],
  256. field_name="organization",
  257. source_uid=source_uids[source_key],
  258. source_mapping_uid=mapping_uids[asset_key],
  259. message=f"issue-secret-{asset_key}",
  260. evidence={"password": f"evidence-secret-{asset_key}"},
  261. recurrence_key=(f"{suffix}-{asset_key}").ljust(64, "0"),
  262. occurrence_number=occurrence,
  263. status=status,
  264. assignee_uid=actor_uid,
  265. due_at=now - timedelta(days=1),
  266. current_version=1,
  267. created_by=actor_uid,
  268. updated_by=actor_uid,
  269. created_at=now,
  270. updated_at=now,
  271. closed_at=now if status == "closed" else None,
  272. )
  273. )
  274. db.session.commit()
  275. repository = SqlAlchemyGovernanceMetricRepository(db.session)
  276. admin = GovernanceMetricAccess(global_access=True)
  277. domain_viewer = GovernanceMetricAccess(
  278. global_access=False,
  279. business_domain_uids=(domain_a,),
  280. )
  281. assert baseline_admin is not None
  282. assert repository.summary_counts(admin) == {
  283. "asset_completeness": (
  284. baseline_admin["asset_completeness"][0] + 4,
  285. baseline_admin["asset_completeness"][1] + 5,
  286. ),
  287. "responsibility_coverage": (
  288. baseline_admin["responsibility_coverage"][0] + 4,
  289. baseline_admin["responsibility_coverage"][1] + 5,
  290. ),
  291. "entity_mapping": (
  292. baseline_admin["entity_mapping"][0] + 4,
  293. baseline_admin["entity_mapping"][1] + 5,
  294. ),
  295. "issue_closure": (
  296. baseline_admin["issue_closure"][0] + 2,
  297. baseline_admin["issue_closure"][1] + 3,
  298. ),
  299. "issue_recurrence": (
  300. baseline_admin["issue_recurrence"][0] + 1,
  301. baseline_admin["issue_recurrence"][1] + 3,
  302. ),
  303. }
  304. assert repository.summary_counts(domain_viewer) == {
  305. "asset_completeness": (2, 3),
  306. "responsibility_coverage": (2, 3),
  307. "entity_mapping": (2, 3),
  308. "issue_closure": (1, 2),
  309. "issue_recurrence": (1, 2),
  310. }
  311. incomplete, incomplete_total = repository.list_details(
  312. domain_viewer,
  313. metric="asset_completeness",
  314. state="incomplete",
  315. page=1,
  316. page_size=20,
  317. )
  318. mapped, mapped_total = repository.list_details(
  319. domain_viewer,
  320. metric="entity_mapping",
  321. state="mapped",
  322. page=1,
  323. page_size=20,
  324. )
  325. recurrent, recurrent_total = repository.list_details(
  326. domain_viewer,
  327. metric="issue_recurrence",
  328. state="recurrent",
  329. page=1,
  330. page_size=20,
  331. )
  332. assert incomplete_total == 1
  333. assert incomplete[0]["asset_uid"] == asset_uids["a2"]
  334. assert incomplete[0]["missing_fields"] == [
  335. "organization",
  336. "responsible_person",
  337. ]
  338. assert mapped_total == 2
  339. assert {item["asset_uid"] for item in mapped} == {
  340. asset_uids["a1"],
  341. asset_uids["a2"],
  342. }
  343. assert recurrent_total == 1
  344. assert recurrent[0]["issue_uid"] == issue_uids[1]
  345. assert recurrent[0]["overdue"] is True
  346. serialized = repr((incomplete, mapped, recurrent))
  347. for secret in (
  348. "source-secret",
  349. "asset-secret",
  350. "merge-secret",
  351. "issue-secret",
  352. "evidence-secret",
  353. "permission_scope",
  354. ):
  355. assert secret not in serialized
  356. finally:
  357. with app.app_context():
  358. for model, column, values in (
  359. (DeviceQualityIssue, DeviceQualityIssue.uid, issue_uids),
  360. (
  361. DeviceEntityMergeRollback,
  362. DeviceEntityMergeRollback.merge_uid,
  363. merge_uids,
  364. ),
  365. (
  366. DeviceEntityMergeEvent,
  367. DeviceEntityMergeEvent.uid,
  368. merge_uids,
  369. ),
  370. (
  371. DeviceEntityMatchReview,
  372. DeviceEntityMatchReview.candidate_uid,
  373. candidate_uids,
  374. ),
  375. (
  376. DeviceEntityMatchCandidate,
  377. DeviceEntityMatchCandidate.uid,
  378. candidate_uids,
  379. ),
  380. ):
  381. if values:
  382. db.session.query(model).filter(column.in_(values)).delete(
  383. synchronize_session=False
  384. )
  385. db.session.query(DeviceQualityRun).filter_by(uid=run_uid).delete()
  386. db.session.query(DeviceQualityProfileVersion).filter_by(
  387. uid=profile_version_uid
  388. ).delete()
  389. db.session.query(DeviceQualityProfile).filter_by(
  390. uid=profile_uid
  391. ).delete()
  392. if asset_uids:
  393. db.session.query(DeviceAssetSourceMapping).filter(
  394. DeviceAssetSourceMapping.asset_uid.in_(asset_uids.values())
  395. ).delete(synchronize_session=False)
  396. db.session.query(DeviceAsset).filter(
  397. DeviceAsset.uid.in_(asset_uids.values())
  398. ).delete(synchronize_session=False)
  399. db.session.query(IngestionSource).filter(
  400. IngestionSource.uid.in_(source_uids.values())
  401. ).delete(synchronize_session=False)
  402. db.session.execute(
  403. text(
  404. "DELETE FROM public.users "
  405. "WHERE id = CAST(:id AS uuid)"
  406. ),
  407. {"id": actor_uid},
  408. )
  409. db.session.commit()