from __future__ import annotations from dataclasses import replace from datetime import datetime import pytest ACTOR_UID = "00000000-0000-7000-8000-000000000701" MANAGER_UID = "00000000-0000-7000-8000-000000000702" SOURCE_UID = "00000000-0000-7000-8000-000000000703" class MemoryDeviceQualityRepository: def __init__(self, assets=(), code_sets=None): self.profile_uid = "00000000-0000-7000-8000-000000000710" self.versions_by_uid = {} self.version_order = [] self.assets = list(assets) self.code_sets = code_sets or { "fault": {"F-001"}, "cause": {"C-001"}, "action": {"A-001"}, } self.run_records = {} self.run_results = {} self.run_violations = {} self.run_scores = {} def ensure_profile(self, *, name, actor_uid): del name, actor_uid return self.profile_uid def latest_version(self): if not self.version_order: return None return self.versions_by_uid[self.version_order[-1]] def find_version_by_hash(self, content_hash): return next( ( item for item in self.versions_by_uid.values() if item.content_hash == content_hash ), None, ) def create_version(self, record): self.versions_by_uid[record.uid] = record self.version_order.append(record.uid) return record def get_version(self, version_uid, *, for_update=False): del for_update return self.versions_by_uid.get(version_uid) def list_versions(self): return [ self.versions_by_uid[uid] for uid in reversed(self.version_order) ] def active_version(self): published = [ item for item in self.versions_by_uid.values() if item.status == "published" ] return max(published, key=lambda item: item.version, default=None) def publish_version(self, record, *, actor_uid, published_at): for uid, item in list(self.versions_by_uid.items()): if item.status == "published": self.versions_by_uid[uid] = replace( item, status="superseded", ) published = replace( record, status="published", published_by=actor_uid, published_at=published_at, ) self.versions_by_uid[record.uid] = published return published def load_assets(self, *, source_uid, limit): records = [ item for item in self.assets if source_uid is None or any( mapping.source_uid == source_uid for mapping in item.mappings ) ] return records[:limit], len(records) def published_code_sets(self): return {key: set(value) for key, value in self.code_sets.items()} def create_run(self, run, results, violations, asset_scores): self.run_records[run.uid] = run self.run_results[run.uid] = tuple(results) self.run_violations[run.uid] = tuple(violations) self.run_scores[run.uid] = tuple(asset_scores) return run def list_runs(self, *, page, page_size): records = list(reversed(tuple(self.run_records.values()))) start = (page - 1) * page_size return records[start : start + page_size], len(records) def get_run(self, run_uid): return self.run_records.get(run_uid) def list_rule_results(self, run_uid): return self.run_results.get(run_uid, ()) def list_violations( self, run_uid, *, rule_code, page, page_size, ): records = [ item for item in self.run_violations.get(run_uid, ()) if rule_code is None or item.rule_code == rule_code ] start = (page - 1) * page_size return records[start : start + page_size], len(records) def list_asset_scores(self, run_uid, *, page, page_size): records = list(self.run_scores.get(run_uid, ())) start = (page - 1) * page_size return records[start : start + page_size], len(records) def mapping( uid, source_code, *, source_uid=SOURCE_UID, asset_type="device", ): from app.core.data_research.device_quality import ( DeviceQualitySourceMapping, ) return DeviceQualitySourceMapping( uid=uid, source_uid=source_uid, source_entity="asset.equipment", asset_type=asset_type, source_code=source_code, ) def asset( uid, asset_type, source_code, *, location="动力车间", organization="设备动力部", responsible_person="张工", attributes=None, ): from app.core.data_research.device_quality import ( DeviceQualityAssetSnapshot, ) return DeviceQualityAssetSnapshot( uid=uid, asset_type=asset_type, name=f"{asset_type}-{source_code}", status="active", current_version=1, location=location, organization=organization, responsible_person=responsible_person, attributes=attributes or {}, mappings=( mapping( f"{uid}-mapping", source_code, asset_type=asset_type, ), ), ) def service(repository, *, manager_uid=MANAGER_UID): from app.core.data_research.device_quality import DeviceQualityService identifiers = iter( f"00000000-0000-7000-8000-000000000{value}" for value in range(720, 999) ) ticks = iter( datetime(2026, 7, 29, 9, minute) for minute in range(60) ) def authorize(actor_uid): from app.core.data_research.errors import DeviceQualityForbidden if actor_uid != manager_uid: raise DeviceQualityForbidden("accountable manager required") return DeviceQualityService( repository, publish_authorizer=authorize, uid_factory=identifiers.__next__, now_factory=ticks.__next__, ) def _rules_by_code(rules): return {item["code"]: item for item in rules} def test_default_policy_is_closed_weighted_and_idempotently_versioned(): repository = MemoryDeviceQualityRepository() quality = service(repository) first = quality.bootstrap(actor_uid=ACTOR_UID) second = quality.bootstrap(actor_uid=ACTOR_UID) assert first.uid == second.uid assert first.version == 1 assert first.status == "draft" assert {item["code"] for item in first.rules} == { "asset_identity_complete", "asset_context_complete", "source_code_unique_normalized", "component_parent_resolved", "fault_code_mapped", "fault_reason_action_complete", "maintenance_closed_loop", } assert sum(item["weight"] for item in first.rules if item["enabled"]) == 100 assert len(repository.versions_by_uid) == 1 @pytest.mark.parametrize( ("mutate", "message"), ( ( lambda rules: rules + [ { "code": "run_python", "enabled": True, "severity": "error", "weight": 1, "parameters": {}, } ], "rule code", ), ( lambda rules: [ {**rules[0], "weight": 1}, *rules[1:], ], "100", ), ( lambda rules: [ { **rules[0], "parameters": {"password": "secret"}, }, *rules[1:], ], "parameter", ), ( lambda rules: [ { **rules[0], "severity": "blocker", }, *rules[1:], ], "severity", ), ), ) def test_policy_revision_rejects_open_or_unsafe_rule_shapes(mutate, message): from app.core.data_research.errors import DeviceQualityInvalid repository = MemoryDeviceQualityRepository() quality = service(repository) current = quality.bootstrap(actor_uid=ACTOR_UID) with pytest.raises(DeviceQualityInvalid, match=message): quality.revise( rules=mutate([dict(item) for item in current.rules]), expected_version=1, actor_uid=ACTOR_UID, ) def test_revision_is_immutable_and_rejects_a_stale_expected_version(): from app.core.data_research.errors import DeviceQualityConflict repository = MemoryDeviceQualityRepository() quality = service(repository) first = quality.bootstrap(actor_uid=ACTOR_UID) changed = [dict(item) for item in first.rules] changed[0] = {**changed[0], "severity": "warning"} second = quality.revise( rules=changed, expected_version=1, actor_uid=ACTOR_UID, ) assert second.version == 2 assert second.uid != first.uid assert _rules_by_code(first.rules)["asset_identity_complete"]["severity"] == ( "critical" ) assert _rules_by_code(second.rules)["asset_identity_complete"]["severity"] == ( "warning" ) with pytest.raises(DeviceQualityConflict, match="stale"): quality.revise( rules=changed, expected_version=1, actor_uid=ACTOR_UID, ) def test_only_accountable_manager_can_publish_and_one_version_remains_active(): from app.core.data_research.errors import DeviceQualityForbidden repository = MemoryDeviceQualityRepository() quality = service(repository) first = quality.bootstrap(actor_uid=ACTOR_UID) with pytest.raises(DeviceQualityForbidden): quality.publish(first.uid, actor_uid=ACTOR_UID) published = quality.publish(first.uid, actor_uid=MANAGER_UID) changed = [dict(item) for item in first.rules] changed[0] = {**changed[0], "severity": "warning"} second = quality.revise( rules=changed, expected_version=1, actor_uid=ACTOR_UID, ) quality.publish(second.uid, actor_uid=MANAGER_UID) assert published.status == "published" assert repository.get_version(first.uid).status == "superseded" assert repository.active_version().uid == second.uid def test_quality_run_evaluates_ledger_fault_and_maintenance_rules_with_evidence(): records = [ asset( "asset-device", "device", "EQ-001", location="", ), asset( "asset-component", "component", "PART-001", attributes={"parent_source_code": "MISSING"}, ), asset( "asset-alarm", "alarm", "ALARM-001", attributes={ "fault_code": "F-UNKNOWN", "cause_code": "", "action_code": "A-UNKNOWN", }, ), asset( "asset-maintenance", "maintenance_record", "WO-001", attributes={ "device_source_code": "EQ-001", "fault_source_code": "ALARM-001", "status": "open", "completed_at": "", "action_code": "A-UNKNOWN", }, ), ] repository = MemoryDeviceQualityRepository(records) quality = service(repository) version = quality.bootstrap(actor_uid=ACTOR_UID) quality.publish(version.uid, actor_uid=MANAGER_UID) run = quality.run(actor_uid=ACTOR_UID) result_by_code = { item.rule_code: item for item in repository.run_results[run.uid] } violations = repository.run_violations[run.uid] assert run.status == "success" assert run.total_assets == 4 assert run.total_violations == 5 assert run.score == pytest.approx(25.0) assert result_by_code["asset_identity_complete"].violation_count == 0 assert result_by_code["asset_context_complete"].violation_count == 1 assert result_by_code["component_parent_resolved"].violation_count == 1 assert result_by_code["fault_code_mapped"].violation_count == 1 assert ( result_by_code["fault_reason_action_complete"].violation_count == 1 ) assert result_by_code["maintenance_closed_loop"].violation_count == 1 assert { item.asset_uid for item in violations } == { "asset-device", "asset-component", "asset-alarm", "asset-maintenance", } assert all(item.source_uid == SOURCE_UID for item in violations) assert all(item.source_mapping_uid.endswith("-mapping") for item in violations) assert all("password" not in str(item.evidence).lower() for item in violations) def test_normalized_source_code_uniqueness_detects_case_and_spacing_collisions(): left = asset("asset-left", "device", "EQ-001") right = asset("asset-right", "device", " eq-001 ") repository = MemoryDeviceQualityRepository([left, right]) quality = service(repository) version = quality.bootstrap(actor_uid=ACTOR_UID) quality.publish(version.uid, actor_uid=MANAGER_UID) run = quality.run(actor_uid=ACTOR_UID) unique = next( item for item in repository.run_results[run.uid] if item.rule_code == "source_code_unique_normalized" ) assert unique.evaluated_count == 2 assert unique.violation_count == 2 assert unique.pass_rate == 0 def test_zero_applicable_rules_receive_full_weight_and_not_applicable_status(): repository = MemoryDeviceQualityRepository( [asset("asset-device", "device", "EQ-001")] ) quality = service(repository) version = quality.bootstrap(actor_uid=ACTOR_UID) quality.publish(version.uid, actor_uid=MANAGER_UID) run = quality.run(actor_uid=ACTOR_UID) results = repository.run_results[run.uid] assert run.score == 100 assert next( item for item in results if item.rule_code == "maintenance_closed_loop" ).status == "not_applicable" assert repository.run_scores[run.uid][0].score == 100 def test_run_binds_exact_policy_and_keeps_samples_bounded_without_losing_counts(): records = [ asset( f"asset-{index}", "device", f"EQ-{index:04d}", location="", ) for index in range(101) ] repository = MemoryDeviceQualityRepository(records) quality = service(repository) version = quality.bootstrap(actor_uid=ACTOR_UID) quality.publish(version.uid, actor_uid=MANAGER_UID) run = quality.run(actor_uid=ACTOR_UID, source_uid=SOURCE_UID) context_result = next( item for item in repository.run_results[run.uid] if item.rule_code == "asset_context_complete" ) assert run.policy_version_uid == version.uid assert run.policy_hash == version.content_hash assert context_result.violation_count == 101 assert context_result.sampled_count == 100 assert len( [ item for item in repository.run_violations[run.uid] if item.rule_code == "asset_context_complete" ] ) == 100 assert repository.assets[0].location == "" def test_run_refuses_more_than_five_thousand_assets(): from app.core.data_research.errors import DeviceQualityInvalid repository = MemoryDeviceQualityRepository( [ asset(f"asset-{index}", "device", f"EQ-{index:05d}") for index in range(5001) ] ) quality = service(repository) version = quality.bootstrap(actor_uid=ACTOR_UID) quality.publish(version.uid, actor_uid=MANAGER_UID) with pytest.raises(DeviceQualityInvalid, match="5,000"): quality.run(actor_uid=ACTOR_UID) def test_run_rejects_an_invalid_source_uid_before_repository_access(): from app.core.data_research.errors import DeviceQualityInvalid repository = MemoryDeviceQualityRepository() quality = service(repository) version = quality.bootstrap(actor_uid=ACTOR_UID) quality.publish(version.uid, actor_uid=MANAGER_UID) with pytest.raises(DeviceQualityInvalid, match="source_uid"): quality.run(actor_uid=ACTOR_UID, source_uid="not-a-uuid")