test_quality_issues.py 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531
  1. from __future__ import annotations
  2. from dataclasses import replace
  3. from datetime import UTC, datetime, timedelta
  4. import pytest
  5. ACTOR_UID = "00000000-0000-7000-8000-000000000901"
  6. ASSIGNEE_UID = "00000000-0000-7000-8000-000000000902"
  7. REVIEWER_UID = "00000000-0000-7000-8000-000000000903"
  8. ADMIN_UID = "00000000-0000-7000-8000-000000000904"
  9. RUN_UID = "00000000-0000-7000-8000-000000000905"
  10. VERIFY_RUN_UID = "00000000-0000-7000-8000-000000000906"
  11. VIOLATION_UID = "00000000-0000-7000-8000-000000000907"
  12. ASSET_UID = "00000000-0000-7000-8000-000000000908"
  13. class MemoryQualityIssueRepository:
  14. def __init__(self, violation):
  15. self.violations = {violation.uid: violation}
  16. self.runs = {violation.run_uid, VERIFY_RUN_UID}
  17. self.active_users = {ACTOR_UID, ASSIGNEE_UID, REVIEWER_UID, ADMIN_UID}
  18. self.issues = {}
  19. self.remediations = {}
  20. self.timeline = {}
  21. def get_violations(self, violation_uids):
  22. return [
  23. self.violations[uid]
  24. for uid in violation_uids
  25. if uid in self.violations
  26. ]
  27. def find_open_by_recurrence_key(self, recurrence_key):
  28. return next(
  29. (
  30. item
  31. for item in self.issues.values()
  32. if item.recurrence_key == recurrence_key
  33. and item.status != "closed"
  34. ),
  35. None,
  36. )
  37. def recurrence_count(self, recurrence_key):
  38. return sum(
  39. item.recurrence_key == recurrence_key
  40. for item in self.issues.values()
  41. )
  42. def create_issue(self, issue, event):
  43. self.issues[issue.uid] = issue
  44. self.timeline[issue.uid] = [event]
  45. return issue
  46. def get_issue(self, issue_uid, *, for_update=False):
  47. del for_update
  48. return self.issues.get(issue_uid)
  49. def active_user_exists(self, user_uid):
  50. return user_uid in self.active_users
  51. def run_exists(self, run_uid):
  52. return run_uid in self.runs
  53. def transition(
  54. self,
  55. issue,
  56. *,
  57. expected_version,
  58. event,
  59. remediation=None,
  60. remediation_review=None,
  61. ):
  62. current = self.issues[issue.uid]
  63. if current.current_version != expected_version:
  64. from app.core.data_research.errors import QualityIssueConflict
  65. raise QualityIssueConflict("quality issue version is stale")
  66. self.issues[issue.uid] = issue
  67. self.timeline[issue.uid].append(event)
  68. if remediation is not None:
  69. self.remediations[issue.uid] = remediation
  70. if remediation_review is not None:
  71. current_round = self.remediations[issue.uid]
  72. self.remediations[issue.uid] = replace(
  73. current_round,
  74. **remediation_review,
  75. )
  76. return issue
  77. def latest_remediation(self, issue_uid):
  78. return self.remediations.get(issue_uid)
  79. def list_issues(self, **filters):
  80. records = list(reversed(tuple(self.issues.values())))
  81. status = filters.get("status")
  82. assignee_uid = filters.get("assignee_uid")
  83. overdue_only = filters.get("overdue_only")
  84. now = filters.get("now")
  85. if status:
  86. records = [item for item in records if item.status == status]
  87. if assignee_uid:
  88. records = [
  89. item for item in records
  90. if item.assignee_uid == assignee_uid
  91. ]
  92. if overdue_only:
  93. records = [
  94. item for item in records
  95. if item.status != "closed"
  96. and item.due_at is not None
  97. and item.due_at < now
  98. ]
  99. page = filters["page"]
  100. page_size = filters["page_size"]
  101. start = (page - 1) * page_size
  102. return records[start : start + page_size], len(records)
  103. def list_timeline(self, issue_uid):
  104. return tuple(self.timeline.get(issue_uid, ()))
  105. def statistics(self, *, now):
  106. records = list(self.issues.values())
  107. recurrence_groups = {}
  108. for item in records:
  109. recurrence_groups[item.recurrence_key] = (
  110. recurrence_groups.get(item.recurrence_key, 0) + 1
  111. )
  112. return {
  113. "total": len(records),
  114. "open": sum(item.status != "closed" for item in records),
  115. "closed": sum(item.status == "closed" for item in records),
  116. "overdue": sum(
  117. item.status != "closed"
  118. and item.due_at is not None
  119. and item.due_at < now
  120. for item in records
  121. ),
  122. "recurrent": sum(value > 1 for value in recurrence_groups.values()),
  123. "recurrent_issues": sum(
  124. item.occurrence_number > 1 for item in records
  125. ),
  126. "recurrence_rate": (
  127. round(
  128. sum(item.occurrence_number > 1 for item in records)
  129. / len(records),
  130. 6,
  131. )
  132. if records else 0.0
  133. ),
  134. }
  135. @pytest.fixture()
  136. def violation():
  137. from app.core.data_research.device_quality import (
  138. DeviceQualityViolationRecord,
  139. )
  140. now = datetime(2026, 7, 29, 9, tzinfo=UTC)
  141. return DeviceQualityViolationRecord(
  142. uid=VIOLATION_UID,
  143. run_uid=RUN_UID,
  144. rule_code="asset_context_complete",
  145. severity="error",
  146. asset_uid=ASSET_UID,
  147. field_name="location,organization",
  148. source_uid=None,
  149. source_mapping_uid=None,
  150. message="设备台账位置、组织或责任人不完整",
  151. evidence={
  152. "asset_uid": ASSET_UID,
  153. "missing_fields": ["location", "organization"],
  154. },
  155. created_at=now,
  156. expires_at=now + timedelta(days=30),
  157. )
  158. @pytest.fixture()
  159. def clock():
  160. return datetime(2026, 7, 29, 10, tzinfo=UTC)
  161. def make_service(repository, clock):
  162. from app.core.data_research.quality_issues import QualityIssueService
  163. identifiers = iter(
  164. f"00000000-0000-7000-8000-000000000{value}"
  165. for value in range(920, 999)
  166. )
  167. def authorize_review(actor_uid):
  168. from app.core.data_research.errors import QualityIssueForbidden
  169. if actor_uid != REVIEWER_UID:
  170. raise QualityIssueForbidden("accountable reviewer required")
  171. return QualityIssueService(
  172. repository,
  173. review_authorizer=authorize_review,
  174. is_admin=lambda uid: uid == ADMIN_UID,
  175. uid_factory=identifiers.__next__,
  176. now_factory=lambda: clock,
  177. )
  178. def create_issue(service, *, due_at=None):
  179. result = service.import_violations(
  180. violation_uids=[VIOLATION_UID],
  181. priority="high",
  182. due_at=due_at,
  183. actor_uid=ACTOR_UID,
  184. )
  185. return result.records[0]
  186. def assign_start_submit(service, issue):
  187. assigned = service.assign(
  188. issue.uid,
  189. assignee_uid=ASSIGNEE_UID,
  190. due_at=issue.due_at,
  191. expected_version=issue.current_version,
  192. actor_uid=ACTOR_UID,
  193. )
  194. started = service.start(
  195. issue.uid,
  196. expected_version=assigned.current_version,
  197. actor_uid=ASSIGNEE_UID,
  198. )
  199. submitted = service.submit(
  200. issue.uid,
  201. summary="已补充设备位置和所属组织并完成来源核对",
  202. evidence_refs=[
  203. {
  204. "type": "device_asset_version",
  205. "uid": ASSET_UID,
  206. "version": 2,
  207. }
  208. ],
  209. expected_version=started.current_version,
  210. actor_uid=ASSIGNEE_UID,
  211. )
  212. return submitted
  213. def test_import_snapshots_violation_and_deduplicates_unresolved_identity(
  214. violation,
  215. clock,
  216. ):
  217. repository = MemoryQualityIssueRepository(violation)
  218. service = make_service(repository, clock)
  219. first = service.import_violations(
  220. violation_uids=[VIOLATION_UID],
  221. priority="high",
  222. due_at=clock + timedelta(days=3),
  223. actor_uid=ACTOR_UID,
  224. )
  225. second = service.import_violations(
  226. violation_uids=[VIOLATION_UID],
  227. priority="critical",
  228. due_at=clock + timedelta(days=1),
  229. actor_uid=ACTOR_UID,
  230. )
  231. assert first.created_count == 1
  232. assert second.created_count == 0
  233. assert second.existing_count == 1
  234. issue = first.records[0]
  235. assert issue.source_violation_uid == VIOLATION_UID
  236. assert issue.status == "open"
  237. assert issue.priority == "high"
  238. assert issue.occurrence_number == 1
  239. assert issue.evidence["missing_fields"] == ["location", "organization"]
  240. assert repository.timeline[issue.uid][0].action == "created"
  241. def test_import_is_bounded_and_requires_every_violation(violation, clock):
  242. from app.core.data_research.errors import QualityIssueInvalid
  243. service = make_service(MemoryQualityIssueRepository(violation), clock)
  244. with pytest.raises(QualityIssueInvalid, match="between 1 and 100"):
  245. service.import_violations(
  246. violation_uids=[],
  247. priority="high",
  248. due_at=None,
  249. actor_uid=ACTOR_UID,
  250. )
  251. with pytest.raises(QualityIssueInvalid, match="not found"):
  252. service.import_violations(
  253. violation_uids=[
  254. "00000000-0000-7000-8000-000000000999"
  255. ],
  256. priority="high",
  257. due_at=None,
  258. actor_uid=ACTOR_UID,
  259. )
  260. def test_assignment_rejects_disabled_user_and_stale_version(violation, clock):
  261. from app.core.data_research.errors import (
  262. QualityIssueConflict,
  263. QualityIssueInvalid,
  264. )
  265. repository = MemoryQualityIssueRepository(violation)
  266. service = make_service(repository, clock)
  267. issue = create_issue(service)
  268. with pytest.raises(QualityIssueInvalid, match="active user"):
  269. service.assign(
  270. issue.uid,
  271. assignee_uid="00000000-0000-7000-8000-000000000999",
  272. due_at=None,
  273. expected_version=1,
  274. actor_uid=ACTOR_UID,
  275. )
  276. assigned = service.assign(
  277. issue.uid,
  278. assignee_uid=ASSIGNEE_UID,
  279. due_at=clock + timedelta(days=2),
  280. expected_version=1,
  281. actor_uid=ACTOR_UID,
  282. )
  283. assert assigned.status == "assigned"
  284. assert assigned.assignee_uid == ASSIGNEE_UID
  285. assert assigned.current_version == 2
  286. with pytest.raises(QualityIssueConflict, match="stale"):
  287. service.start(
  288. issue.uid,
  289. expected_version=1,
  290. actor_uid=ASSIGNEE_UID,
  291. )
  292. def test_only_assignee_or_admin_can_work_and_submission_is_bounded(
  293. violation,
  294. clock,
  295. ):
  296. from app.core.data_research.errors import (
  297. QualityIssueForbidden,
  298. QualityIssueInvalid,
  299. )
  300. repository = MemoryQualityIssueRepository(violation)
  301. service = make_service(repository, clock)
  302. issue = create_issue(service)
  303. service.assign(
  304. issue.uid,
  305. assignee_uid=ASSIGNEE_UID,
  306. due_at=None,
  307. expected_version=1,
  308. actor_uid=ACTOR_UID,
  309. )
  310. with pytest.raises(QualityIssueForbidden, match="assignee"):
  311. service.start(
  312. issue.uid,
  313. expected_version=2,
  314. actor_uid=ACTOR_UID,
  315. )
  316. started = service.start(
  317. issue.uid,
  318. expected_version=2,
  319. actor_uid=ADMIN_UID,
  320. )
  321. with pytest.raises(QualityIssueInvalid, match="summary"):
  322. service.submit(
  323. issue.uid,
  324. summary="",
  325. evidence_refs=[],
  326. expected_version=started.current_version,
  327. actor_uid=ASSIGNEE_UID,
  328. )
  329. submitted = service.submit(
  330. issue.uid,
  331. summary="已完成字段补录",
  332. evidence_refs=[{"type": "note", "uid": "evidence-1"}],
  333. expected_version=started.current_version,
  334. actor_uid=ASSIGNEE_UID,
  335. )
  336. assert submitted.status == "pending_review"
  337. assert repository.remediations[issue.uid].review_status == "pending"
  338. def test_independent_accountable_review_closes_or_rejects(violation, clock):
  339. from app.core.data_research.errors import QualityIssueForbidden
  340. repository = MemoryQualityIssueRepository(violation)
  341. service = make_service(repository, clock)
  342. issue = create_issue(service)
  343. submitted = assign_start_submit(service, issue)
  344. with pytest.raises(QualityIssueForbidden, match="own remediation"):
  345. service.review(
  346. issue.uid,
  347. verification_result="passed",
  348. note="复验通过",
  349. verification_run_uid=VERIFY_RUN_UID,
  350. expected_version=submitted.current_version,
  351. actor_uid=ASSIGNEE_UID,
  352. )
  353. with pytest.raises(QualityIssueForbidden, match="accountable"):
  354. service.review(
  355. issue.uid,
  356. verification_result="passed",
  357. note="复验通过",
  358. verification_run_uid=VERIFY_RUN_UID,
  359. expected_version=submitted.current_version,
  360. actor_uid=ACTOR_UID,
  361. )
  362. rejected = service.review(
  363. issue.uid,
  364. verification_result="failed",
  365. note="复验仍缺少组织字段",
  366. verification_run_uid=VERIFY_RUN_UID,
  367. expected_version=submitted.current_version,
  368. actor_uid=REVIEWER_UID,
  369. )
  370. assert rejected.status == "in_progress"
  371. resubmitted = service.submit(
  372. issue.uid,
  373. summary="已再次补充所属组织",
  374. evidence_refs=[],
  375. expected_version=rejected.current_version,
  376. actor_uid=ASSIGNEE_UID,
  377. )
  378. closed = service.review(
  379. issue.uid,
  380. verification_result="passed",
  381. note="人工核对和复验均通过",
  382. verification_run_uid=VERIFY_RUN_UID,
  383. expected_version=resubmitted.current_version,
  384. actor_uid=REVIEWER_UID,
  385. )
  386. assert closed.status == "closed"
  387. assert closed.closed_at == clock
  388. assert repository.remediations[issue.uid].review_status == "passed"
  389. def test_close_reopen_and_later_import_create_a_recurrence(violation, clock):
  390. repository = MemoryQualityIssueRepository(violation)
  391. service = make_service(repository, clock)
  392. issue = create_issue(service, due_at=clock - timedelta(hours=1))
  393. submitted = assign_start_submit(service, issue)
  394. closed = service.review(
  395. issue.uid,
  396. verification_result="passed",
  397. note="复验通过",
  398. verification_run_uid=None,
  399. expected_version=submitted.current_version,
  400. actor_uid=REVIEWER_UID,
  401. )
  402. reopened = service.reopen(
  403. issue.uid,
  404. reason="现场抽查发现字段再次缺失",
  405. due_at=clock + timedelta(days=1),
  406. expected_version=closed.current_version,
  407. actor_uid=ACTOR_UID,
  408. )
  409. assert reopened.status == "assigned"
  410. assert reopened.closed_at is None
  411. assert repository.timeline[issue.uid][-1].action == "reopened"
  412. submitted_again = service.submit(
  413. issue.uid,
  414. summary="重新补录",
  415. evidence_refs=[],
  416. expected_version=reopened.current_version,
  417. actor_uid=ASSIGNEE_UID,
  418. )
  419. closed_again = service.review(
  420. issue.uid,
  421. verification_result="passed",
  422. note="复验通过",
  423. verification_run_uid=None,
  424. expected_version=submitted_again.current_version,
  425. actor_uid=REVIEWER_UID,
  426. )
  427. assert closed_again.status == "closed"
  428. new_violation = replace(
  429. violation,
  430. uid="00000000-0000-7000-8000-000000000909",
  431. run_uid=VERIFY_RUN_UID,
  432. )
  433. repository.violations[new_violation.uid] = new_violation
  434. recurrence = service.import_violations(
  435. violation_uids=[new_violation.uid],
  436. priority="high",
  437. due_at=None,
  438. actor_uid=ACTOR_UID,
  439. ).records[0]
  440. assert recurrence.uid != issue.uid
  441. assert recurrence.occurrence_number == 2
  442. def test_list_statistics_and_timeline_expose_overdue_and_recurrence(
  443. violation,
  444. clock,
  445. ):
  446. repository = MemoryQualityIssueRepository(violation)
  447. service = make_service(repository, clock)
  448. issue = create_issue(service, due_at=clock - timedelta(minutes=1))
  449. records, total = service.issues(
  450. status=None,
  451. assignee_uid=None,
  452. overdue_only=True,
  453. page=1,
  454. page_size=20,
  455. )
  456. statistics = service.statistics()
  457. timeline = service.timeline(issue.uid)
  458. assert total == 1
  459. assert records[0].is_overdue(clock)
  460. assert statistics == {
  461. "total": 1,
  462. "open": 1,
  463. "closed": 0,
  464. "overdue": 1,
  465. "recurrent": 0,
  466. "recurrent_issues": 0,
  467. "recurrence_rate": 0.0,
  468. }
  469. assert [item.action for item in timeline] == ["created"]