| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542 |
- 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")
|