| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531 |
- from __future__ import annotations
- from dataclasses import replace
- from datetime import UTC, datetime, timedelta
- import pytest
- ACTOR_UID = "00000000-0000-7000-8000-000000000901"
- ASSIGNEE_UID = "00000000-0000-7000-8000-000000000902"
- REVIEWER_UID = "00000000-0000-7000-8000-000000000903"
- ADMIN_UID = "00000000-0000-7000-8000-000000000904"
- RUN_UID = "00000000-0000-7000-8000-000000000905"
- VERIFY_RUN_UID = "00000000-0000-7000-8000-000000000906"
- VIOLATION_UID = "00000000-0000-7000-8000-000000000907"
- ASSET_UID = "00000000-0000-7000-8000-000000000908"
- class MemoryQualityIssueRepository:
- def __init__(self, violation):
- self.violations = {violation.uid: violation}
- self.runs = {violation.run_uid, VERIFY_RUN_UID}
- self.active_users = {ACTOR_UID, ASSIGNEE_UID, REVIEWER_UID, ADMIN_UID}
- self.issues = {}
- self.remediations = {}
- self.timeline = {}
- def get_violations(self, violation_uids):
- return [
- self.violations[uid]
- for uid in violation_uids
- if uid in self.violations
- ]
- def find_open_by_recurrence_key(self, recurrence_key):
- return next(
- (
- item
- for item in self.issues.values()
- if item.recurrence_key == recurrence_key
- and item.status != "closed"
- ),
- None,
- )
- def recurrence_count(self, recurrence_key):
- return sum(
- item.recurrence_key == recurrence_key
- for item in self.issues.values()
- )
- def create_issue(self, issue, event):
- self.issues[issue.uid] = issue
- self.timeline[issue.uid] = [event]
- return issue
- def get_issue(self, issue_uid, *, for_update=False):
- del for_update
- return self.issues.get(issue_uid)
- def active_user_exists(self, user_uid):
- return user_uid in self.active_users
- def run_exists(self, run_uid):
- return run_uid in self.runs
- def transition(
- self,
- issue,
- *,
- expected_version,
- event,
- remediation=None,
- remediation_review=None,
- ):
- current = self.issues[issue.uid]
- if current.current_version != expected_version:
- from app.core.data_research.errors import QualityIssueConflict
- raise QualityIssueConflict("quality issue version is stale")
- self.issues[issue.uid] = issue
- self.timeline[issue.uid].append(event)
- if remediation is not None:
- self.remediations[issue.uid] = remediation
- if remediation_review is not None:
- current_round = self.remediations[issue.uid]
- self.remediations[issue.uid] = replace(
- current_round,
- **remediation_review,
- )
- return issue
- def latest_remediation(self, issue_uid):
- return self.remediations.get(issue_uid)
- def list_issues(self, **filters):
- records = list(reversed(tuple(self.issues.values())))
- status = filters.get("status")
- assignee_uid = filters.get("assignee_uid")
- overdue_only = filters.get("overdue_only")
- now = filters.get("now")
- if status:
- records = [item for item in records if item.status == status]
- if assignee_uid:
- records = [
- item for item in records
- if item.assignee_uid == assignee_uid
- ]
- if overdue_only:
- records = [
- item for item in records
- if item.status != "closed"
- and item.due_at is not None
- and item.due_at < now
- ]
- page = filters["page"]
- page_size = filters["page_size"]
- start = (page - 1) * page_size
- return records[start : start + page_size], len(records)
- def list_timeline(self, issue_uid):
- return tuple(self.timeline.get(issue_uid, ()))
- def statistics(self, *, now):
- records = list(self.issues.values())
- recurrence_groups = {}
- for item in records:
- recurrence_groups[item.recurrence_key] = (
- recurrence_groups.get(item.recurrence_key, 0) + 1
- )
- return {
- "total": len(records),
- "open": sum(item.status != "closed" for item in records),
- "closed": sum(item.status == "closed" for item in records),
- "overdue": sum(
- item.status != "closed"
- and item.due_at is not None
- and item.due_at < now
- for item in records
- ),
- "recurrent": sum(value > 1 for value in recurrence_groups.values()),
- "recurrent_issues": sum(
- item.occurrence_number > 1 for item in records
- ),
- "recurrence_rate": (
- round(
- sum(item.occurrence_number > 1 for item in records)
- / len(records),
- 6,
- )
- if records else 0.0
- ),
- }
- @pytest.fixture()
- def violation():
- from app.core.data_research.device_quality import (
- DeviceQualityViolationRecord,
- )
- now = datetime(2026, 7, 29, 9, tzinfo=UTC)
- return DeviceQualityViolationRecord(
- uid=VIOLATION_UID,
- run_uid=RUN_UID,
- rule_code="asset_context_complete",
- severity="error",
- asset_uid=ASSET_UID,
- field_name="location,organization",
- source_uid=None,
- source_mapping_uid=None,
- message="设备台账位置、组织或责任人不完整",
- evidence={
- "asset_uid": ASSET_UID,
- "missing_fields": ["location", "organization"],
- },
- created_at=now,
- expires_at=now + timedelta(days=30),
- )
- @pytest.fixture()
- def clock():
- return datetime(2026, 7, 29, 10, tzinfo=UTC)
- def make_service(repository, clock):
- from app.core.data_research.quality_issues import QualityIssueService
- identifiers = iter(
- f"00000000-0000-7000-8000-000000000{value}"
- for value in range(920, 999)
- )
- def authorize_review(actor_uid):
- from app.core.data_research.errors import QualityIssueForbidden
- if actor_uid != REVIEWER_UID:
- raise QualityIssueForbidden("accountable reviewer required")
- return QualityIssueService(
- repository,
- review_authorizer=authorize_review,
- is_admin=lambda uid: uid == ADMIN_UID,
- uid_factory=identifiers.__next__,
- now_factory=lambda: clock,
- )
- def create_issue(service, *, due_at=None):
- result = service.import_violations(
- violation_uids=[VIOLATION_UID],
- priority="high",
- due_at=due_at,
- actor_uid=ACTOR_UID,
- )
- return result.records[0]
- def assign_start_submit(service, issue):
- assigned = service.assign(
- issue.uid,
- assignee_uid=ASSIGNEE_UID,
- due_at=issue.due_at,
- expected_version=issue.current_version,
- actor_uid=ACTOR_UID,
- )
- started = service.start(
- issue.uid,
- expected_version=assigned.current_version,
- actor_uid=ASSIGNEE_UID,
- )
- submitted = service.submit(
- issue.uid,
- summary="已补充设备位置和所属组织并完成来源核对",
- evidence_refs=[
- {
- "type": "device_asset_version",
- "uid": ASSET_UID,
- "version": 2,
- }
- ],
- expected_version=started.current_version,
- actor_uid=ASSIGNEE_UID,
- )
- return submitted
- def test_import_snapshots_violation_and_deduplicates_unresolved_identity(
- violation,
- clock,
- ):
- repository = MemoryQualityIssueRepository(violation)
- service = make_service(repository, clock)
- first = service.import_violations(
- violation_uids=[VIOLATION_UID],
- priority="high",
- due_at=clock + timedelta(days=3),
- actor_uid=ACTOR_UID,
- )
- second = service.import_violations(
- violation_uids=[VIOLATION_UID],
- priority="critical",
- due_at=clock + timedelta(days=1),
- actor_uid=ACTOR_UID,
- )
- assert first.created_count == 1
- assert second.created_count == 0
- assert second.existing_count == 1
- issue = first.records[0]
- assert issue.source_violation_uid == VIOLATION_UID
- assert issue.status == "open"
- assert issue.priority == "high"
- assert issue.occurrence_number == 1
- assert issue.evidence["missing_fields"] == ["location", "organization"]
- assert repository.timeline[issue.uid][0].action == "created"
- def test_import_is_bounded_and_requires_every_violation(violation, clock):
- from app.core.data_research.errors import QualityIssueInvalid
- service = make_service(MemoryQualityIssueRepository(violation), clock)
- with pytest.raises(QualityIssueInvalid, match="between 1 and 100"):
- service.import_violations(
- violation_uids=[],
- priority="high",
- due_at=None,
- actor_uid=ACTOR_UID,
- )
- with pytest.raises(QualityIssueInvalid, match="not found"):
- service.import_violations(
- violation_uids=[
- "00000000-0000-7000-8000-000000000999"
- ],
- priority="high",
- due_at=None,
- actor_uid=ACTOR_UID,
- )
- def test_assignment_rejects_disabled_user_and_stale_version(violation, clock):
- from app.core.data_research.errors import (
- QualityIssueConflict,
- QualityIssueInvalid,
- )
- repository = MemoryQualityIssueRepository(violation)
- service = make_service(repository, clock)
- issue = create_issue(service)
- with pytest.raises(QualityIssueInvalid, match="active user"):
- service.assign(
- issue.uid,
- assignee_uid="00000000-0000-7000-8000-000000000999",
- due_at=None,
- expected_version=1,
- actor_uid=ACTOR_UID,
- )
- assigned = service.assign(
- issue.uid,
- assignee_uid=ASSIGNEE_UID,
- due_at=clock + timedelta(days=2),
- expected_version=1,
- actor_uid=ACTOR_UID,
- )
- assert assigned.status == "assigned"
- assert assigned.assignee_uid == ASSIGNEE_UID
- assert assigned.current_version == 2
- with pytest.raises(QualityIssueConflict, match="stale"):
- service.start(
- issue.uid,
- expected_version=1,
- actor_uid=ASSIGNEE_UID,
- )
- def test_only_assignee_or_admin_can_work_and_submission_is_bounded(
- violation,
- clock,
- ):
- from app.core.data_research.errors import (
- QualityIssueForbidden,
- QualityIssueInvalid,
- )
- repository = MemoryQualityIssueRepository(violation)
- service = make_service(repository, clock)
- issue = create_issue(service)
- service.assign(
- issue.uid,
- assignee_uid=ASSIGNEE_UID,
- due_at=None,
- expected_version=1,
- actor_uid=ACTOR_UID,
- )
- with pytest.raises(QualityIssueForbidden, match="assignee"):
- service.start(
- issue.uid,
- expected_version=2,
- actor_uid=ACTOR_UID,
- )
- started = service.start(
- issue.uid,
- expected_version=2,
- actor_uid=ADMIN_UID,
- )
- with pytest.raises(QualityIssueInvalid, match="summary"):
- service.submit(
- issue.uid,
- summary="",
- evidence_refs=[],
- expected_version=started.current_version,
- actor_uid=ASSIGNEE_UID,
- )
- submitted = service.submit(
- issue.uid,
- summary="已完成字段补录",
- evidence_refs=[{"type": "note", "uid": "evidence-1"}],
- expected_version=started.current_version,
- actor_uid=ASSIGNEE_UID,
- )
- assert submitted.status == "pending_review"
- assert repository.remediations[issue.uid].review_status == "pending"
- def test_independent_accountable_review_closes_or_rejects(violation, clock):
- from app.core.data_research.errors import QualityIssueForbidden
- repository = MemoryQualityIssueRepository(violation)
- service = make_service(repository, clock)
- issue = create_issue(service)
- submitted = assign_start_submit(service, issue)
- with pytest.raises(QualityIssueForbidden, match="own remediation"):
- service.review(
- issue.uid,
- verification_result="passed",
- note="复验通过",
- verification_run_uid=VERIFY_RUN_UID,
- expected_version=submitted.current_version,
- actor_uid=ASSIGNEE_UID,
- )
- with pytest.raises(QualityIssueForbidden, match="accountable"):
- service.review(
- issue.uid,
- verification_result="passed",
- note="复验通过",
- verification_run_uid=VERIFY_RUN_UID,
- expected_version=submitted.current_version,
- actor_uid=ACTOR_UID,
- )
- rejected = service.review(
- issue.uid,
- verification_result="failed",
- note="复验仍缺少组织字段",
- verification_run_uid=VERIFY_RUN_UID,
- expected_version=submitted.current_version,
- actor_uid=REVIEWER_UID,
- )
- assert rejected.status == "in_progress"
- resubmitted = service.submit(
- issue.uid,
- summary="已再次补充所属组织",
- evidence_refs=[],
- expected_version=rejected.current_version,
- actor_uid=ASSIGNEE_UID,
- )
- closed = service.review(
- issue.uid,
- verification_result="passed",
- note="人工核对和复验均通过",
- verification_run_uid=VERIFY_RUN_UID,
- expected_version=resubmitted.current_version,
- actor_uid=REVIEWER_UID,
- )
- assert closed.status == "closed"
- assert closed.closed_at == clock
- assert repository.remediations[issue.uid].review_status == "passed"
- def test_close_reopen_and_later_import_create_a_recurrence(violation, clock):
- repository = MemoryQualityIssueRepository(violation)
- service = make_service(repository, clock)
- issue = create_issue(service, due_at=clock - timedelta(hours=1))
- submitted = assign_start_submit(service, issue)
- closed = service.review(
- issue.uid,
- verification_result="passed",
- note="复验通过",
- verification_run_uid=None,
- expected_version=submitted.current_version,
- actor_uid=REVIEWER_UID,
- )
- reopened = service.reopen(
- issue.uid,
- reason="现场抽查发现字段再次缺失",
- due_at=clock + timedelta(days=1),
- expected_version=closed.current_version,
- actor_uid=ACTOR_UID,
- )
- assert reopened.status == "assigned"
- assert reopened.closed_at is None
- assert repository.timeline[issue.uid][-1].action == "reopened"
- submitted_again = service.submit(
- issue.uid,
- summary="重新补录",
- evidence_refs=[],
- expected_version=reopened.current_version,
- actor_uid=ASSIGNEE_UID,
- )
- closed_again = service.review(
- issue.uid,
- verification_result="passed",
- note="复验通过",
- verification_run_uid=None,
- expected_version=submitted_again.current_version,
- actor_uid=REVIEWER_UID,
- )
- assert closed_again.status == "closed"
- new_violation = replace(
- violation,
- uid="00000000-0000-7000-8000-000000000909",
- run_uid=VERIFY_RUN_UID,
- )
- repository.violations[new_violation.uid] = new_violation
- recurrence = service.import_violations(
- violation_uids=[new_violation.uid],
- priority="high",
- due_at=None,
- actor_uid=ACTOR_UID,
- ).records[0]
- assert recurrence.uid != issue.uid
- assert recurrence.occurrence_number == 2
- def test_list_statistics_and_timeline_expose_overdue_and_recurrence(
- violation,
- clock,
- ):
- repository = MemoryQualityIssueRepository(violation)
- service = make_service(repository, clock)
- issue = create_issue(service, due_at=clock - timedelta(minutes=1))
- records, total = service.issues(
- status=None,
- assignee_uid=None,
- overdue_only=True,
- page=1,
- page_size=20,
- )
- statistics = service.statistics()
- timeline = service.timeline(issue.uid)
- assert total == 1
- assert records[0].is_overdue(clock)
- assert statistics == {
- "total": 1,
- "open": 1,
- "closed": 0,
- "overdue": 1,
- "recurrent": 0,
- "recurrent_issues": 0,
- "recurrence_rate": 0.0,
- }
- assert [item.action for item in timeline] == ["created"]
|