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"]