test_device_quality.py 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542
  1. from __future__ import annotations
  2. from dataclasses import replace
  3. from datetime import datetime
  4. import pytest
  5. ACTOR_UID = "00000000-0000-7000-8000-000000000701"
  6. MANAGER_UID = "00000000-0000-7000-8000-000000000702"
  7. SOURCE_UID = "00000000-0000-7000-8000-000000000703"
  8. class MemoryDeviceQualityRepository:
  9. def __init__(self, assets=(), code_sets=None):
  10. self.profile_uid = "00000000-0000-7000-8000-000000000710"
  11. self.versions_by_uid = {}
  12. self.version_order = []
  13. self.assets = list(assets)
  14. self.code_sets = code_sets or {
  15. "fault": {"F-001"},
  16. "cause": {"C-001"},
  17. "action": {"A-001"},
  18. }
  19. self.run_records = {}
  20. self.run_results = {}
  21. self.run_violations = {}
  22. self.run_scores = {}
  23. def ensure_profile(self, *, name, actor_uid):
  24. del name, actor_uid
  25. return self.profile_uid
  26. def latest_version(self):
  27. if not self.version_order:
  28. return None
  29. return self.versions_by_uid[self.version_order[-1]]
  30. def find_version_by_hash(self, content_hash):
  31. return next(
  32. (
  33. item
  34. for item in self.versions_by_uid.values()
  35. if item.content_hash == content_hash
  36. ),
  37. None,
  38. )
  39. def create_version(self, record):
  40. self.versions_by_uid[record.uid] = record
  41. self.version_order.append(record.uid)
  42. return record
  43. def get_version(self, version_uid, *, for_update=False):
  44. del for_update
  45. return self.versions_by_uid.get(version_uid)
  46. def list_versions(self):
  47. return [
  48. self.versions_by_uid[uid]
  49. for uid in reversed(self.version_order)
  50. ]
  51. def active_version(self):
  52. published = [
  53. item
  54. for item in self.versions_by_uid.values()
  55. if item.status == "published"
  56. ]
  57. return max(published, key=lambda item: item.version, default=None)
  58. def publish_version(self, record, *, actor_uid, published_at):
  59. for uid, item in list(self.versions_by_uid.items()):
  60. if item.status == "published":
  61. self.versions_by_uid[uid] = replace(
  62. item,
  63. status="superseded",
  64. )
  65. published = replace(
  66. record,
  67. status="published",
  68. published_by=actor_uid,
  69. published_at=published_at,
  70. )
  71. self.versions_by_uid[record.uid] = published
  72. return published
  73. def load_assets(self, *, source_uid, limit):
  74. records = [
  75. item
  76. for item in self.assets
  77. if source_uid is None
  78. or any(
  79. mapping.source_uid == source_uid
  80. for mapping in item.mappings
  81. )
  82. ]
  83. return records[:limit], len(records)
  84. def published_code_sets(self):
  85. return {key: set(value) for key, value in self.code_sets.items()}
  86. def create_run(self, run, results, violations, asset_scores):
  87. self.run_records[run.uid] = run
  88. self.run_results[run.uid] = tuple(results)
  89. self.run_violations[run.uid] = tuple(violations)
  90. self.run_scores[run.uid] = tuple(asset_scores)
  91. return run
  92. def list_runs(self, *, page, page_size):
  93. records = list(reversed(tuple(self.run_records.values())))
  94. start = (page - 1) * page_size
  95. return records[start : start + page_size], len(records)
  96. def get_run(self, run_uid):
  97. return self.run_records.get(run_uid)
  98. def list_rule_results(self, run_uid):
  99. return self.run_results.get(run_uid, ())
  100. def list_violations(
  101. self,
  102. run_uid,
  103. *,
  104. rule_code,
  105. page,
  106. page_size,
  107. ):
  108. records = [
  109. item
  110. for item in self.run_violations.get(run_uid, ())
  111. if rule_code is None or item.rule_code == rule_code
  112. ]
  113. start = (page - 1) * page_size
  114. return records[start : start + page_size], len(records)
  115. def list_asset_scores(self, run_uid, *, page, page_size):
  116. records = list(self.run_scores.get(run_uid, ()))
  117. start = (page - 1) * page_size
  118. return records[start : start + page_size], len(records)
  119. def mapping(
  120. uid,
  121. source_code,
  122. *,
  123. source_uid=SOURCE_UID,
  124. asset_type="device",
  125. ):
  126. from app.core.data_research.device_quality import (
  127. DeviceQualitySourceMapping,
  128. )
  129. return DeviceQualitySourceMapping(
  130. uid=uid,
  131. source_uid=source_uid,
  132. source_entity="asset.equipment",
  133. asset_type=asset_type,
  134. source_code=source_code,
  135. )
  136. def asset(
  137. uid,
  138. asset_type,
  139. source_code,
  140. *,
  141. location="动力车间",
  142. organization="设备动力部",
  143. responsible_person="张工",
  144. attributes=None,
  145. ):
  146. from app.core.data_research.device_quality import (
  147. DeviceQualityAssetSnapshot,
  148. )
  149. return DeviceQualityAssetSnapshot(
  150. uid=uid,
  151. asset_type=asset_type,
  152. name=f"{asset_type}-{source_code}",
  153. status="active",
  154. current_version=1,
  155. location=location,
  156. organization=organization,
  157. responsible_person=responsible_person,
  158. attributes=attributes or {},
  159. mappings=(
  160. mapping(
  161. f"{uid}-mapping",
  162. source_code,
  163. asset_type=asset_type,
  164. ),
  165. ),
  166. )
  167. def service(repository, *, manager_uid=MANAGER_UID):
  168. from app.core.data_research.device_quality import DeviceQualityService
  169. identifiers = iter(
  170. f"00000000-0000-7000-8000-000000000{value}"
  171. for value in range(720, 999)
  172. )
  173. ticks = iter(
  174. datetime(2026, 7, 29, 9, minute)
  175. for minute in range(60)
  176. )
  177. def authorize(actor_uid):
  178. from app.core.data_research.errors import DeviceQualityForbidden
  179. if actor_uid != manager_uid:
  180. raise DeviceQualityForbidden("accountable manager required")
  181. return DeviceQualityService(
  182. repository,
  183. publish_authorizer=authorize,
  184. uid_factory=identifiers.__next__,
  185. now_factory=ticks.__next__,
  186. )
  187. def _rules_by_code(rules):
  188. return {item["code"]: item for item in rules}
  189. def test_default_policy_is_closed_weighted_and_idempotently_versioned():
  190. repository = MemoryDeviceQualityRepository()
  191. quality = service(repository)
  192. first = quality.bootstrap(actor_uid=ACTOR_UID)
  193. second = quality.bootstrap(actor_uid=ACTOR_UID)
  194. assert first.uid == second.uid
  195. assert first.version == 1
  196. assert first.status == "draft"
  197. assert {item["code"] for item in first.rules} == {
  198. "asset_identity_complete",
  199. "asset_context_complete",
  200. "source_code_unique_normalized",
  201. "component_parent_resolved",
  202. "fault_code_mapped",
  203. "fault_reason_action_complete",
  204. "maintenance_closed_loop",
  205. }
  206. assert sum(item["weight"] for item in first.rules if item["enabled"]) == 100
  207. assert len(repository.versions_by_uid) == 1
  208. @pytest.mark.parametrize(
  209. ("mutate", "message"),
  210. (
  211. (
  212. lambda rules: rules
  213. + [
  214. {
  215. "code": "run_python",
  216. "enabled": True,
  217. "severity": "error",
  218. "weight": 1,
  219. "parameters": {},
  220. }
  221. ],
  222. "rule code",
  223. ),
  224. (
  225. lambda rules: [
  226. {**rules[0], "weight": 1},
  227. *rules[1:],
  228. ],
  229. "100",
  230. ),
  231. (
  232. lambda rules: [
  233. {
  234. **rules[0],
  235. "parameters": {"password": "secret"},
  236. },
  237. *rules[1:],
  238. ],
  239. "parameter",
  240. ),
  241. (
  242. lambda rules: [
  243. {
  244. **rules[0],
  245. "severity": "blocker",
  246. },
  247. *rules[1:],
  248. ],
  249. "severity",
  250. ),
  251. ),
  252. )
  253. def test_policy_revision_rejects_open_or_unsafe_rule_shapes(mutate, message):
  254. from app.core.data_research.errors import DeviceQualityInvalid
  255. repository = MemoryDeviceQualityRepository()
  256. quality = service(repository)
  257. current = quality.bootstrap(actor_uid=ACTOR_UID)
  258. with pytest.raises(DeviceQualityInvalid, match=message):
  259. quality.revise(
  260. rules=mutate([dict(item) for item in current.rules]),
  261. expected_version=1,
  262. actor_uid=ACTOR_UID,
  263. )
  264. def test_revision_is_immutable_and_rejects_a_stale_expected_version():
  265. from app.core.data_research.errors import DeviceQualityConflict
  266. repository = MemoryDeviceQualityRepository()
  267. quality = service(repository)
  268. first = quality.bootstrap(actor_uid=ACTOR_UID)
  269. changed = [dict(item) for item in first.rules]
  270. changed[0] = {**changed[0], "severity": "warning"}
  271. second = quality.revise(
  272. rules=changed,
  273. expected_version=1,
  274. actor_uid=ACTOR_UID,
  275. )
  276. assert second.version == 2
  277. assert second.uid != first.uid
  278. assert _rules_by_code(first.rules)["asset_identity_complete"]["severity"] == (
  279. "critical"
  280. )
  281. assert _rules_by_code(second.rules)["asset_identity_complete"]["severity"] == (
  282. "warning"
  283. )
  284. with pytest.raises(DeviceQualityConflict, match="stale"):
  285. quality.revise(
  286. rules=changed,
  287. expected_version=1,
  288. actor_uid=ACTOR_UID,
  289. )
  290. def test_only_accountable_manager_can_publish_and_one_version_remains_active():
  291. from app.core.data_research.errors import DeviceQualityForbidden
  292. repository = MemoryDeviceQualityRepository()
  293. quality = service(repository)
  294. first = quality.bootstrap(actor_uid=ACTOR_UID)
  295. with pytest.raises(DeviceQualityForbidden):
  296. quality.publish(first.uid, actor_uid=ACTOR_UID)
  297. published = quality.publish(first.uid, actor_uid=MANAGER_UID)
  298. changed = [dict(item) for item in first.rules]
  299. changed[0] = {**changed[0], "severity": "warning"}
  300. second = quality.revise(
  301. rules=changed,
  302. expected_version=1,
  303. actor_uid=ACTOR_UID,
  304. )
  305. quality.publish(second.uid, actor_uid=MANAGER_UID)
  306. assert published.status == "published"
  307. assert repository.get_version(first.uid).status == "superseded"
  308. assert repository.active_version().uid == second.uid
  309. def test_quality_run_evaluates_ledger_fault_and_maintenance_rules_with_evidence():
  310. records = [
  311. asset(
  312. "asset-device",
  313. "device",
  314. "EQ-001",
  315. location="",
  316. ),
  317. asset(
  318. "asset-component",
  319. "component",
  320. "PART-001",
  321. attributes={"parent_source_code": "MISSING"},
  322. ),
  323. asset(
  324. "asset-alarm",
  325. "alarm",
  326. "ALARM-001",
  327. attributes={
  328. "fault_code": "F-UNKNOWN",
  329. "cause_code": "",
  330. "action_code": "A-UNKNOWN",
  331. },
  332. ),
  333. asset(
  334. "asset-maintenance",
  335. "maintenance_record",
  336. "WO-001",
  337. attributes={
  338. "device_source_code": "EQ-001",
  339. "fault_source_code": "ALARM-001",
  340. "status": "open",
  341. "completed_at": "",
  342. "action_code": "A-UNKNOWN",
  343. },
  344. ),
  345. ]
  346. repository = MemoryDeviceQualityRepository(records)
  347. quality = service(repository)
  348. version = quality.bootstrap(actor_uid=ACTOR_UID)
  349. quality.publish(version.uid, actor_uid=MANAGER_UID)
  350. run = quality.run(actor_uid=ACTOR_UID)
  351. result_by_code = {
  352. item.rule_code: item
  353. for item in repository.run_results[run.uid]
  354. }
  355. violations = repository.run_violations[run.uid]
  356. assert run.status == "success"
  357. assert run.total_assets == 4
  358. assert run.total_violations == 5
  359. assert run.score == pytest.approx(25.0)
  360. assert result_by_code["asset_identity_complete"].violation_count == 0
  361. assert result_by_code["asset_context_complete"].violation_count == 1
  362. assert result_by_code["component_parent_resolved"].violation_count == 1
  363. assert result_by_code["fault_code_mapped"].violation_count == 1
  364. assert (
  365. result_by_code["fault_reason_action_complete"].violation_count == 1
  366. )
  367. assert result_by_code["maintenance_closed_loop"].violation_count == 1
  368. assert {
  369. item.asset_uid
  370. for item in violations
  371. } == {
  372. "asset-device",
  373. "asset-component",
  374. "asset-alarm",
  375. "asset-maintenance",
  376. }
  377. assert all(item.source_uid == SOURCE_UID for item in violations)
  378. assert all(item.source_mapping_uid.endswith("-mapping") for item in violations)
  379. assert all("password" not in str(item.evidence).lower() for item in violations)
  380. def test_normalized_source_code_uniqueness_detects_case_and_spacing_collisions():
  381. left = asset("asset-left", "device", "EQ-001")
  382. right = asset("asset-right", "device", " eq-001 ")
  383. repository = MemoryDeviceQualityRepository([left, right])
  384. quality = service(repository)
  385. version = quality.bootstrap(actor_uid=ACTOR_UID)
  386. quality.publish(version.uid, actor_uid=MANAGER_UID)
  387. run = quality.run(actor_uid=ACTOR_UID)
  388. unique = next(
  389. item
  390. for item in repository.run_results[run.uid]
  391. if item.rule_code == "source_code_unique_normalized"
  392. )
  393. assert unique.evaluated_count == 2
  394. assert unique.violation_count == 2
  395. assert unique.pass_rate == 0
  396. def test_zero_applicable_rules_receive_full_weight_and_not_applicable_status():
  397. repository = MemoryDeviceQualityRepository(
  398. [asset("asset-device", "device", "EQ-001")]
  399. )
  400. quality = service(repository)
  401. version = quality.bootstrap(actor_uid=ACTOR_UID)
  402. quality.publish(version.uid, actor_uid=MANAGER_UID)
  403. run = quality.run(actor_uid=ACTOR_UID)
  404. results = repository.run_results[run.uid]
  405. assert run.score == 100
  406. assert next(
  407. item
  408. for item in results
  409. if item.rule_code == "maintenance_closed_loop"
  410. ).status == "not_applicable"
  411. assert repository.run_scores[run.uid][0].score == 100
  412. def test_run_binds_exact_policy_and_keeps_samples_bounded_without_losing_counts():
  413. records = [
  414. asset(
  415. f"asset-{index}",
  416. "device",
  417. f"EQ-{index:04d}",
  418. location="",
  419. )
  420. for index in range(101)
  421. ]
  422. repository = MemoryDeviceQualityRepository(records)
  423. quality = service(repository)
  424. version = quality.bootstrap(actor_uid=ACTOR_UID)
  425. quality.publish(version.uid, actor_uid=MANAGER_UID)
  426. run = quality.run(actor_uid=ACTOR_UID, source_uid=SOURCE_UID)
  427. context_result = next(
  428. item
  429. for item in repository.run_results[run.uid]
  430. if item.rule_code == "asset_context_complete"
  431. )
  432. assert run.policy_version_uid == version.uid
  433. assert run.policy_hash == version.content_hash
  434. assert context_result.violation_count == 101
  435. assert context_result.sampled_count == 100
  436. assert len(
  437. [
  438. item
  439. for item in repository.run_violations[run.uid]
  440. if item.rule_code == "asset_context_complete"
  441. ]
  442. ) == 100
  443. assert repository.assets[0].location == ""
  444. def test_run_refuses_more_than_five_thousand_assets():
  445. from app.core.data_research.errors import DeviceQualityInvalid
  446. repository = MemoryDeviceQualityRepository(
  447. [
  448. asset(f"asset-{index}", "device", f"EQ-{index:05d}")
  449. for index in range(5001)
  450. ]
  451. )
  452. quality = service(repository)
  453. version = quality.bootstrap(actor_uid=ACTOR_UID)
  454. quality.publish(version.uid, actor_uid=MANAGER_UID)
  455. with pytest.raises(DeviceQualityInvalid, match="5,000"):
  456. quality.run(actor_uid=ACTOR_UID)
  457. def test_run_rejects_an_invalid_source_uid_before_repository_access():
  458. from app.core.data_research.errors import DeviceQualityInvalid
  459. repository = MemoryDeviceQualityRepository()
  460. quality = service(repository)
  461. version = quality.bootstrap(actor_uid=ACTOR_UID)
  462. quality.publish(version.uid, actor_uid=MANAGER_UID)
  463. with pytest.raises(DeviceQualityInvalid, match="source_uid"):
  464. quality.run(actor_uid=ACTOR_UID, source_uid="not-a-uuid")