test_quality_issue_postgres.py 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296
  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_quality_issue_postgres_persists_lifecycle_and_recurrence(
  9. monkeypatch,
  10. ):
  11. platform_url = os.environ.get("TEST_DATABASE_URL")
  12. if not platform_url:
  13. pytest.skip("TEST_DATABASE_URL is required")
  14. monkeypatch.setenv("DATABASE_URL", platform_url)
  15. from app import create_app, db
  16. from app.core.data_research.quality_issue_repository import (
  17. SqlAlchemyQualityIssueRepository,
  18. )
  19. from app.core.data_research.quality_issues import QualityIssueService
  20. from app.models.data_research import (
  21. DeviceAsset,
  22. DeviceQualityIssue,
  23. DeviceQualityProfile,
  24. DeviceQualityProfileVersion,
  25. DeviceQualityRun,
  26. DeviceQualityViolationSample,
  27. IngestionSource,
  28. )
  29. app = create_app()
  30. app.config.update(TESTING=True)
  31. suffix = uuid.uuid4().hex[:10]
  32. actor_uid = str(uuid.uuid4())
  33. assignee_uid = str(uuid.uuid4())
  34. reviewer_uid = str(uuid.uuid4())
  35. source_uid = str(uuid.uuid4())
  36. asset_uid = str(uuid.uuid4())
  37. profile_uid = str(uuid.uuid4())
  38. version_uid = str(uuid.uuid4())
  39. run_uid = str(uuid.uuid4())
  40. verify_run_uid = str(uuid.uuid4())
  41. violation_uid = str(uuid.uuid4())
  42. recurrence_violation_uid = str(uuid.uuid4())
  43. issue_uids = []
  44. now = datetime.now(UTC)
  45. try:
  46. with app.app_context():
  47. for user_uid, username in (
  48. (actor_uid, f"wp08-actor-{suffix}"),
  49. (assignee_uid, f"wp08-assignee-{suffix}"),
  50. (reviewer_uid, f"wp08-reviewer-{suffix}"),
  51. ):
  52. db.session.execute(
  53. text(
  54. """
  55. INSERT INTO public.users (
  56. id, username, display_name, password_hash, status
  57. ) VALUES (
  58. CAST(:id AS uuid), :username, :username,
  59. 'integration-test', 'active'
  60. )
  61. """
  62. ),
  63. {"id": user_uid, "username": username},
  64. )
  65. db.session.add(
  66. IngestionSource(
  67. uid=source_uid,
  68. source_type="database",
  69. name=f"WP08 source {suffix}",
  70. config={"database_type": "postgresql"},
  71. permission_scope={},
  72. status="active",
  73. created_by=actor_uid,
  74. )
  75. )
  76. db.session.add(
  77. DeviceAsset(
  78. uid=asset_uid,
  79. asset_type="device",
  80. name=f"WP08 device {suffix}",
  81. status="active",
  82. current_version=1,
  83. content_hash="b" * 64,
  84. attributes={},
  85. created_by=actor_uid,
  86. updated_by=actor_uid,
  87. )
  88. )
  89. db.session.add(
  90. DeviceQualityProfile(
  91. uid=profile_uid,
  92. code=f"WP08_{suffix}",
  93. name="WP08 integration profile",
  94. created_by=actor_uid,
  95. )
  96. )
  97. db.session.flush()
  98. db.session.add(
  99. DeviceQualityProfileVersion(
  100. uid=version_uid,
  101. profile_uid=profile_uid,
  102. version=1,
  103. status="draft",
  104. rules=[],
  105. content_hash="c" * 64,
  106. created_by=actor_uid,
  107. )
  108. )
  109. db.session.flush()
  110. for current_run_uid in (run_uid, verify_run_uid):
  111. db.session.add(
  112. DeviceQualityRun(
  113. uid=current_run_uid,
  114. policy_version_uid=version_uid,
  115. policy_hash="c" * 64,
  116. source_uid=source_uid,
  117. status="success",
  118. total_assets=1,
  119. total_violations=1,
  120. score=80,
  121. created_by=actor_uid,
  122. )
  123. )
  124. db.session.flush()
  125. for current_uid, current_run_uid in (
  126. (violation_uid, run_uid),
  127. (recurrence_violation_uid, verify_run_uid),
  128. ):
  129. db.session.add(
  130. DeviceQualityViolationSample(
  131. uid=current_uid,
  132. run_uid=current_run_uid,
  133. rule_code="asset_context_complete",
  134. severity="error",
  135. asset_uid=asset_uid,
  136. field_name="location",
  137. source_uid=source_uid,
  138. message="设备位置缺失",
  139. evidence={
  140. "asset_uid": asset_uid,
  141. "missing_fields": ["location"],
  142. },
  143. created_at=now,
  144. expires_at=now + timedelta(days=30),
  145. )
  146. )
  147. db.session.commit()
  148. repository = SqlAlchemyQualityIssueRepository(db.session)
  149. service = QualityIssueService(
  150. repository,
  151. review_authorizer=lambda actor: (
  152. None
  153. if actor == reviewer_uid
  154. else pytest.fail("unexpected reviewer")
  155. ),
  156. commit=db.session.commit,
  157. rollback=db.session.rollback,
  158. )
  159. baseline_statistics = service.statistics()
  160. empty_records, empty_total = service.issues(
  161. status="closed",
  162. assignee_uid=assignee_uid,
  163. overdue_only=False,
  164. page=1,
  165. page_size=20,
  166. )
  167. assert empty_records == []
  168. assert empty_total == 0
  169. imported = service.import_violations(
  170. violation_uids=[violation_uid],
  171. priority="high",
  172. due_at=now - timedelta(minutes=1),
  173. actor_uid=actor_uid,
  174. )
  175. issue = imported.records[0]
  176. issue_uids.append(issue.uid)
  177. deduplicated = service.import_violations(
  178. violation_uids=[violation_uid],
  179. priority="critical",
  180. due_at=None,
  181. actor_uid=actor_uid,
  182. )
  183. assert deduplicated.existing_count == 1
  184. assigned = service.assign(
  185. issue.uid,
  186. assignee_uid=assignee_uid,
  187. due_at=issue.due_at,
  188. expected_version=1,
  189. actor_uid=actor_uid,
  190. )
  191. started = service.start(
  192. issue.uid,
  193. expected_version=assigned.current_version,
  194. actor_uid=assignee_uid,
  195. )
  196. submitted = service.submit(
  197. issue.uid,
  198. summary="设备位置已补录并复核来源",
  199. evidence_refs=[
  200. {
  201. "type": "device_asset",
  202. "uid": asset_uid,
  203. "version": 2,
  204. }
  205. ],
  206. expected_version=started.current_version,
  207. actor_uid=assignee_uid,
  208. )
  209. closed = service.review(
  210. issue.uid,
  211. verification_result="passed",
  212. note="复核通过",
  213. verification_run_uid=verify_run_uid,
  214. expected_version=submitted.current_version,
  215. actor_uid=reviewer_uid,
  216. )
  217. assert closed.status == "closed"
  218. recurrence = service.import_violations(
  219. violation_uids=[recurrence_violation_uid],
  220. priority="high",
  221. due_at=None,
  222. actor_uid=actor_uid,
  223. ).records[0]
  224. issue_uids.append(recurrence.uid)
  225. assert recurrence.occurrence_number == 2
  226. statistics = service.statistics()
  227. assert statistics["total"] == baseline_statistics["total"] + 2
  228. assert statistics["recurrent"] == baseline_statistics["recurrent"] + 1
  229. assert (
  230. statistics["recurrent_issues"]
  231. == baseline_statistics["recurrent_issues"] + 1
  232. )
  233. assert statistics["recurrence_rate"] == round(
  234. statistics["recurrent_issues"] / statistics["total"],
  235. 6,
  236. )
  237. assert [item.action for item in service.timeline(issue.uid)] == [
  238. "created",
  239. "assigned",
  240. "started",
  241. "submitted",
  242. "closed",
  243. ]
  244. loaded, remediation = service.get(issue.uid)
  245. assert loaded.evidence["missing_fields"] == ["location"]
  246. assert remediation.review_status == "passed"
  247. assert remediation.verification_run_uid == verify_run_uid
  248. persisted = db.session.get(DeviceQualityIssue, issue.uid)
  249. assert persisted.current_version == 5
  250. assert persisted.source_violation_uid == violation_uid
  251. finally:
  252. with app.app_context():
  253. if issue_uids:
  254. db.session.query(DeviceQualityIssue).filter(
  255. DeviceQualityIssue.uid.in_(issue_uids)
  256. ).delete(synchronize_session=False)
  257. db.session.query(DeviceQualityViolationSample).filter(
  258. DeviceQualityViolationSample.uid.in_(
  259. [violation_uid, recurrence_violation_uid]
  260. )
  261. ).delete(synchronize_session=False)
  262. db.session.query(DeviceQualityRun).filter(
  263. DeviceQualityRun.uid.in_([run_uid, verify_run_uid])
  264. ).delete(synchronize_session=False)
  265. db.session.query(DeviceQualityProfileVersion).filter_by(
  266. uid=version_uid
  267. ).delete(synchronize_session=False)
  268. db.session.query(DeviceQualityProfile).filter_by(
  269. uid=profile_uid
  270. ).delete(synchronize_session=False)
  271. db.session.query(DeviceAsset).filter_by(uid=asset_uid).delete(
  272. synchronize_session=False
  273. )
  274. db.session.query(IngestionSource).filter_by(
  275. uid=source_uid
  276. ).delete(synchronize_session=False)
  277. db.session.execute(
  278. text(
  279. "DELETE FROM public.users "
  280. "WHERE id = ANY(CAST(:ids AS uuid[]))"
  281. ),
  282. {"ids": [actor_uid, assignee_uid, reviewer_uid]},
  283. )
  284. db.session.commit()