cases.py 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416
  1. # -*- coding: utf-8 -*-
  2. """cases.py — findings.json(每批快照) → cases.json(跨批次持续档) 生产器 (运维 agent 柱1「智能工单」).
  3. 为什么: findings.json = 每批全量快照, 无记性 → 同一前兆每月被"重新发现"(B16)。
  4. case = (farm,turbine,system) 的持续健康档: 跨批次状态延续 + 反向闭环的左半 (evidence 追加)。
  5. 本模块 = 附录A 状态机的**产线**(schemas.validate_case 是**闸**); 产出必过 validate_cases。
  6. 设计红线 (与 schemas 不变量同源):
  7. ① 诚实封顶 CANDIDATE: findings 多为 相对排名(fleet-z), 禁 fleet-z 直升 CONFIRMED(不变量①)。
  8. 升 CONFIRMED 必 own-data (finding['own_data_confirmed']=True), 生产器不伪造。
  9. ② 一 key 一档永续: 每 (farm,turbine,system) 全程只一个 case, 跨批改状态, 不重复开案(不变量⑥)。
  10. ③ TICKETED 及以后 = 人工裁决点①, 生产器**不碰** (只在 NORMAL↔WATCH↔CANDIDATE↔CONFIRMED↔REFUTED
  11. + CLOSED→RELAPSED 之间自动流转)。工单签发/销案是人的动作, 走独立入口。
  12. 纯核心 (produce_cases): now/batch_id 作参传入, 不碰文件系统/时钟 (副作用在 CLI __main__)。
  13. """
  14. try:
  15. from app_common.app_common_guanlan.api import install_root as _install_root
  16. except ImportError: # 理论不可达;包结构异常时回退到按位置上跳
  17. from pathlib import Path as _P
  18. def _install_root(_f): return _P(_f).resolve().parents[4]
  19. from src import paths as P
  20. import json
  21. import re
  22. import sys
  23. from pathlib import Path
  24. # 状态序 (advance 用): NORMAL<WATCH<CANDIDATE<CONFIRMED
  25. _STATE_ORDER = {'NORMAL': 0, 'WATCH': 1, 'CANDIDATE': 2, 'CONFIRMED': 3}
  26. # verdict → 目标态 (诚实封顶: 无 own-data 时 定论/准定论·预警 也只到 CANDIDATE)
  27. _VERDICT_STATE = {
  28. '定论': 'CANDIDATE', '准定论·预警': 'CANDIDATE', '候选': 'CANDIDATE', '参考': 'WATCH',
  29. }
  30. # 标题 → discriminator_registry 正式 id (有序命中即停; 无命中退 slug 待补)。
  31. # configs/discriminator_registry.yaml 单源; test_cases_producer_dp 断言映射目标 ∈ 注册表 (外键真)。
  32. _DISCRIMINATOR_RULES = [
  33. ('对风', 'static_yaw_assessment'), ('静态偏航', 'static_yaw_assessment'),
  34. ('三相', 'phase_imbalance'), ('相偏', 'phase_imbalance'), ('相不平衡', 'phase_imbalance'),
  35. ('1P', 'blade_imbalance_1p'),
  36. ('变桨', 'pitch_stick'), ('桨距', 'pitch_stick'),
  37. ('残差', 'nbm_residual'), ('NBM', 'nbm_residual'), ('温度', 'nbm_residual'), ('气动', 'nbm_residual'),
  38. ('偏航', 'yaw_aging_operators'),
  39. ]
  40. def discriminator_id_of(finding):
  41. """finding.title → 判据注册表正式 id (evidence.discriminator_id 外键真)。无映射→ slug(title) 待补。"""
  42. title = str(finding.get('title', ''))
  43. for kw, did in _DISCRIMINATOR_RULES:
  44. if kw in title:
  45. return did
  46. return _slug(title)
  47. def system_of(finding):
  48. """finding.title → 规范 system key (跨批次稳定, 跨 OEM 命名前缀感知)。附录A.3 system 枚举。
  49. ★前缀感知 (2026-07-07, 岳林柏第二样板逮到 farm-tuned gap): 跨 OEM 通道命名歧义 —
  50. '驱动端' 华电=发电机DE / 盂县'齿轮箱中间轴驱动端'=齿箱IMS-DE。故: **子系统前缀优先** (齿轮箱/发电机)
  51. 定归属, 再定部件; **裸命名** (无子系统前缀) 按行业约定 (低速/高速轴承=齿箱, 驱动端/自由端=发电机)。"""
  52. t = str(finding.get('title', ''))
  53. # 非温度系统 (关键词明确, 先拦)
  54. if '偏航' in t or '对风' in t:
  55. return 'yaw'
  56. if '变桨' in t or '桨距' in t:
  57. return 'pitch'
  58. if '欠发' in t:
  59. return 'underperf'
  60. if '变流' in t or '变频' in t:
  61. return 'converter'
  62. if '安全链' in t:
  63. return 'safety_chain'
  64. if '主轴承' in t:
  65. return 'main_brg'
  66. # DE/NDE 判别 (非驱动端/自由端 = NDE; 优先于 DE 判)
  67. is_nde = ('非驱动端' in t) or ('自由端' in t)
  68. is_de = ('驱动端' in t) and not is_nde
  69. # 齿轮箱子系统 (前缀优先: '齿轮箱中间轴驱动端'=齿箱 非发电机)
  70. if '齿轮箱' in t or '齿轮' in t:
  71. if '油' in t or '润滑' in t:
  72. return 'gearbox_oil'
  73. if '高速' in t:
  74. return 'gearbox_hss_brg'
  75. if '中速' in t or '中间轴' in t:
  76. return 'gearbox_ims_de' if is_de else ('gearbox_ims_nde' if is_nde else 'gearbox_ims_brg')
  77. if '低速' in t:
  78. return 'gearbox_lss_brg'
  79. return 'gearbox'
  80. # 发电机子系统 (显式前缀)
  81. if '发电机' in t:
  82. if '定子' in t:
  83. return 'gen_stator'
  84. if is_de:
  85. return 'gen_de_brg'
  86. if is_nde:
  87. return 'gen_nde_brg'
  88. return 'gen'
  89. # 裸命名 (无子系统前缀) — 行业约定
  90. if '定子' in t:
  91. return 'gen_stator'
  92. if '高速轴' in t:
  93. return 'gearbox_hss_brg'
  94. if '中速轴' in t or '中间轴' in t:
  95. return 'gearbox_ims_brg'
  96. if '低速轴' in t:
  97. return 'gearbox_lss_brg'
  98. if is_de:
  99. return 'gen_de_brg'
  100. if is_nde:
  101. return 'gen_nde_brg'
  102. if '主轴' in t:
  103. return 'main_shaft'
  104. if '传动链' in t:
  105. return 'driveline'
  106. if '联轴器' in t:
  107. return 'coupling'
  108. slug = re.sub(r'[^a-z0-9一-鿿]+', '_', t.lower()).strip('_')
  109. return slug[:32] or 'unknown'
  110. def _slug(s):
  111. return re.sub(r'[^a-z0-9一-鿿]+', '_', str(s).lower()).strip('_')[:48] or 'x'
  112. def _disc(finding):
  113. """finding['discriminator'] 归一为 dict。
  114. 该字段实测有三种写法: dict (结构化判别器, 53 场) / None (无) /
  115. **散文字符串** (rudong 13 条把它当"判别方法说明"写)。原写法 `... or {}` 只兜住
  116. None 与空串, 非空字符串是 truthy 会漏进 .get() -> AttributeError。
  117. 散文串没有 per_turbine/flagged, 语义上等同"无" -> 与 None 同样处理。
  118. """
  119. d = finding.get('discriminator')
  120. return d if isinstance(d, dict) else {}
  121. def flagged_turbines(finding):
  122. """finding 明确 flag 的台号 list (affected_machines 优先, 退 discriminator.flagged)。
  123. 空 = 该 finding 不制造 case (参考/INSUFFICIENT/健康场级 → affected null → 天然静默)。"""
  124. am = finding.get('affected_machines')
  125. if isinstance(am, str):
  126. return [am]
  127. if isinstance(am, list) and am:
  128. return [str(t) for t in am]
  129. disc = _disc(finding)
  130. fl = disc.get('flagged')
  131. if isinstance(fl, list) and fl:
  132. return [str(t) for t in fl]
  133. return []
  134. def _turbine_value(finding, tid):
  135. """从 discriminator.per_turbine 抽该台的代表值 (z 优先, 退 resid_mean)。无→None."""
  136. disc = _disc(finding)
  137. pt = disc.get('per_turbine')
  138. if isinstance(pt, list):
  139. for r in pt:
  140. if isinstance(r, dict) and str(r.get('tid')) == str(tid):
  141. return r.get('z', r.get('resid_mean'))
  142. return None
  143. def _target_state(finding):
  144. """verdict → 目标态 (own-data 才允许 CONFIRMED; 否则封顶 CANDIDATE)。"""
  145. from src.sop.schemas import norm_verdict
  146. v = norm_verdict(finding.get('verdict', ''))
  147. base = _VERDICT_STATE.get(v)
  148. if base is None:
  149. return None # INSUFFICIENT/撤回/未知 → 不开案
  150. if v in ('定论', '准定论·预警') and finding.get('own_data_confirmed'):
  151. return 'CONFIRMED'
  152. return base
  153. def _open_history(target, batch_id, now):
  154. """新档: 合成 NORMAL→…→target 的合法迁移链 (附录A.2 边)。"""
  155. chain = [('NORMAL', 'WATCH', 'signal_raised')]
  156. if _STATE_ORDER[target] >= _STATE_ORDER['CANDIDATE']:
  157. chain.append(('WATCH', 'CANDIDATE', 'evidence_added'))
  158. if target == 'CONFIRMED':
  159. chain.append(('CANDIDATE', 'CONFIRMED', 'own_data_confirmed'))
  160. return [{'ts': now, 'from_state': f, 'to_state': t, 'event': e, 'by': 'engine', 'batch_id': batch_id}
  161. for (f, t, e) in chain]
  162. def _advance_chain(cur, target):
  163. """cur→target 的合法逐步边 (只升不降; ≤cur 返回空)。CONFIRMED 不自动降。"""
  164. if _STATE_ORDER.get(target, 0) <= _STATE_ORDER.get(cur, 0):
  165. return []
  166. seq = ['NORMAL', 'WATCH', 'CANDIDATE', 'CONFIRMED']
  167. i, j = seq.index(cur), seq.index(target)
  168. events = {('NORMAL', 'WATCH'): 'signal_raised', ('WATCH', 'CANDIDATE'): 'evidence_added',
  169. ('CANDIDATE', 'CONFIRMED'): 'own_data_confirmed'}
  170. steps = []
  171. for k in range(i, j):
  172. f, t = seq[k], seq[k + 1]
  173. steps.append((f, t, events[(f, t)]))
  174. return steps
  175. def produce_cases(findings_doc, prev_cases, batch_id, now, decay_after=2):
  176. """核心 (纯函数): 一批 findings + 上批 cases → 新 cases list。产出过 validate_cases。
  177. findings_doc: findings.json 解析 dict (含 _farm/_generated/findings)。
  178. prev_cases: 上批 cases list (无则 [])。 batch_id: 本批标识 (findings._generated)。
  179. now: 本批日期 str (YYYY-MM-DD, 供 opened_at/updated_at/history.ts)。
  180. decay_after: 连续 N 批未 flag 后 WATCH/CANDIDATE 衰减回 NORMAL (CONFIRMED 不自动衰减)。
  181. """
  182. import copy
  183. farm = findings_doc.get('_farm', '?')
  184. idx = {} # (farm,turbine,system) → case
  185. for c in copy.deepcopy(prev_cases or []): # 不就地改入参 (别名安全)
  186. if isinstance(c, dict):
  187. idx[(c.get('farm'), c.get('turbine'), c.get('system'))] = c
  188. flagged_now = set()
  189. for f in findings_doc.get('findings', []):
  190. target = _target_state(f)
  191. if target is None:
  192. continue
  193. sysk = system_of(f)
  194. did = discriminator_id_of(f)
  195. for tid in flagged_turbines(f):
  196. key = (farm, tid, sysk)
  197. flagged_now.add(key)
  198. ev = {'date': now, 'discriminator_id': did, 'value': _turbine_value(f, tid),
  199. 'verdict': f.get('verdict'), 'batch_id': batch_id,
  200. 'own_data': bool(f.get('own_data_confirmed'))}
  201. prev = idx.get(key)
  202. if prev is None:
  203. idx[key] = {
  204. 'case_id': f'{farm}-{tid}-{sysk}-001',
  205. 'farm': farm, 'turbine': tid, 'system': sysk, 'state': target,
  206. 'tier': None, 'verdict': f.get('verdict'),
  207. 'opened_at': now, 'updated_at': now, 'clean_streak': 0,
  208. 'evidence': [ev], 'history': _open_history(target, batch_id, now),
  209. }
  210. else:
  211. _advance_existing(prev, target, ev, batch_id, now)
  212. # 未 flag 的活跃档 → 衰减 (记账, 达阈回 NORMAL; CONFIRMED/终态不动)
  213. for key, c in idx.items():
  214. if key in flagged_now:
  215. continue
  216. _decay(c, batch_id, now, decay_after)
  217. return list(idx.values())
  218. def _advance_existing(c, target, ev, batch_id, now):
  219. """已存档: 追加 evidence + 合法升态 (CLOSED→复发路; REFUTED 终态不重开=dedup)。"""
  220. c.setdefault('evidence', []).append(ev)
  221. c['updated_at'] = now
  222. c['clean_streak'] = 0
  223. c['verdict'] = ev.get('verdict')
  224. st = c.get('state')
  225. hist = c.setdefault('history', [])
  226. if st == 'REFUTED':
  227. return # 终态: 同信号不重登顶 (refuted_basis dedup)
  228. if st == 'CLOSED': # 复发: CLOSED→RELAPSED→(CONFIRMED 若 own-data)
  229. hist.append({'ts': now, 'from_state': 'CLOSED', 'to_state': 'RELAPSED',
  230. 'event': 'relapse_detected', 'by': 'engine', 'batch_id': batch_id})
  231. c['state'] = 'RELAPSED'
  232. c['relapse_of'] = c.get('case_id')
  233. if target == 'CONFIRMED':
  234. hist.append({'ts': now, 'from_state': 'RELAPSED', 'to_state': 'CONFIRMED',
  235. 'event': 'own_data_confirmed', 'by': 'engine', 'batch_id': batch_id})
  236. c['state'] = 'CONFIRMED'
  237. return
  238. for (f, t, e) in _advance_chain(st, target): # NORMAL/WATCH/CANDIDATE 升态
  239. hist.append({'ts': now, 'from_state': f, 'to_state': t, 'event': e,
  240. 'by': 'engine', 'batch_id': batch_id})
  241. c['state'] = t
  242. def _decay(c, batch_id, now, decay_after):
  243. """未 flag: clean_streak++; 达阈且 state∈{WATCH,CANDIDATE} → 直接回 NORMAL (附录A.2 decay 边)。"""
  244. st = c.get('state')
  245. if st not in ('WATCH', 'CANDIDATE'):
  246. return # CONFIRMED/终态/已 NORMAL 不自动衰减
  247. c['clean_streak'] = int(c.get('clean_streak', 0)) + 1
  248. c['updated_at'] = now
  249. if c['clean_streak'] >= decay_after:
  250. c.setdefault('evidence', []).append(
  251. {'date': now, 'discriminator_id': 'decay', 'value': None,
  252. 'verdict': '参考', 'batch_id': batch_id, 'own_data': False})
  253. c.setdefault('history', []).append(
  254. {'ts': now, 'from_state': st, 'to_state': 'NORMAL', 'event': 'decay_to_normal',
  255. 'by': 'engine', 'batch_id': batch_id})
  256. c['state'] = 'NORMAL'
  257. c['clean_streak'] = 0
  258. def apply_human_action(case, action, now, batch_id='human', **kw):
  259. """人工侧合法迁移入口 (生产器不碰的右半: 工单/修/回验/证伪/own-data 确诊)。
  260. 与 produce_cases(engine 侧)对称: 本函数是人的动作 → by=human 审计。返回改后 copy (不改入参)。
  261. 每步只走附录A.2 合法边; 产出仍须过 schemas.validate_case。action:
  262. confirm_own_data: CANDIDATE→CONFIRMED (kw: value/discriminator_id → 追加 own_data 证据)
  263. issue_ticket: CONFIRMED/RELAPSED→TICKETED (kw: ticket={ticket_id,issued_by,actions,...})
  264. start_repair: TICKETED→IN_REPAIR
  265. report_verify: TICKETED/IN_REPAIR→VERIFY (kw: baseline,target,window_days)
  266. close: VERIFY→CLOSED (kw: readings 非空; 自动 verify.result=pass)
  267. verify_fail: VERIFY→CONFIRMED (未根治)
  268. refute: CANDIDATE/CONFIRMED→REFUTED (kw: refuted_basis 必填)
  269. """
  270. import copy
  271. from src.sop.schemas import CASE_EDGES
  272. c = copy.deepcopy(case)
  273. st = c.get('state')
  274. hist = c.setdefault('history', [])
  275. def _step(to, event):
  276. if to not in CASE_EDGES.get(st, set()):
  277. raise ValueError(f'apply_human_action: {st}→{to} 非法边 (action={action})')
  278. hist.append({'ts': now, 'from_state': st, 'to_state': to, 'event': event,
  279. 'by': 'human', 'batch_id': batch_id})
  280. c['state'] = to
  281. c['updated_at'] = now
  282. if action == 'confirm_own_data':
  283. c.setdefault('evidence', []).append(
  284. {'date': now, 'discriminator_id': kw.get('discriminator_id', 'field_measurement'),
  285. 'value': kw.get('value'), 'verdict': '准定论·预警', 'batch_id': batch_id, 'own_data': True})
  286. _step('CONFIRMED', 'own_data_confirmed')
  287. elif action == 'issue_ticket':
  288. c['ticket'] = kw.get('ticket') or {}
  289. _step('TICKETED', 'ticket_issued')
  290. elif action == 'start_repair':
  291. _step('IN_REPAIR', 'repair_started')
  292. elif action == 'report_verify':
  293. c['verify'] = {'baseline': kw.get('baseline'), 'target': kw.get('target'),
  294. 'window_days': kw.get('window_days', 60), 'readings': kw.get('readings', []),
  295. 'result': 'pending'}
  296. _step('VERIFY', 'repair_reported')
  297. elif action == 'close':
  298. vf = c.setdefault('verify', {})
  299. vf['readings'] = kw.get('readings') or vf.get('readings') or []
  300. vf['result'] = 'pass'
  301. _step('CLOSED', 'verify_pass')
  302. elif action == 'verify_fail':
  303. c.setdefault('verify', {})['result'] = 'fail'
  304. _step('CONFIRMED', 'verify_fail')
  305. elif action == 'refute':
  306. c['refuted_basis'] = kw.get('refuted_basis') or ''
  307. _step('REFUTED', 'guard_refuted')
  308. else:
  309. raise ValueError(f'apply_human_action: 未知 action {action!r}')
  310. return c
  311. def apply_analysis(case, verdict, mechanism, now, review=None, field_needs=None,
  312. ref=None, by='deep-dive'):
  313. """规范回灌深挖结果 (Phase-2 机制归因 + 独立审) → case.analysis。返回改后 copy (不改入参)。
  314. 与 apply_human_action(状态迁移) 对称: 本函数只挂**分析层**, 不碰 case.state。
  315. ★纪律 (不变量①): analysis.verdict 是更深表征(可 > case.verdict, 如 D29 候选→深挖准定论·预警),
  316. 但升 CONFIRMED 须现场 own-data → 走 apply_human_action('confirm_own_data'), 非本函数。
  317. 产出过 schemas.validate_analysis (verdict 六枚举 + mechanism 必填)。"""
  318. import copy
  319. from src.sop.schemas import norm_verdict, VERDICTS
  320. if norm_verdict(verdict) not in VERDICTS:
  321. raise ValueError(f'apply_analysis: verdict "{verdict}" ∉ 六枚举 {sorted(VERDICTS)}')
  322. if not mechanism:
  323. raise ValueError('apply_analysis: mechanism 必填 (深挖必给机制归因)')
  324. c = copy.deepcopy(case)
  325. an = {'verdict': norm_verdict(verdict), 'mechanism': mechanism, 'review': review,
  326. 'field_needs': field_needs, 'ref': ref, 'date': now, 'by': by}
  327. c['analysis'] = {k: v for k, v in an.items() if v is not None} # 去 None 清爽
  328. c['updated_at'] = now
  329. return c
  330. # ---------------------------------------------------------------------------
  331. if __name__ == '__main__':
  332. # 控制台可能是 GBK(中文 Windows 代码页 936): 正文里的 ✔ ✗ ✅ ⚠ 这类字符编不出来会抛
  333. # UnicodeEncodeError, 脚本干成了事却以退出码 1 结束(同类坑见 src/console.py)。降级为 '?' 而不是崩;
  334. # 不用 import 是为了兼顾 python -m 与直接当脚本跑两种启动方式。
  335. import sys as _sys
  336. for _s in (_sys.stdout, _sys.stderr):
  337. try: _s.reconfigure(errors='replace')
  338. except Exception: pass
  339. # CLI: python -m src.sop.cases <farm> [--date YYYY-MM-DD]
  340. # 读 outputs/<farm>/sop/{findings.json, cases.json(若有)} → 生产 → 校验 → 写 cases.json
  341. import argparse
  342. from src.sop.schemas import validate_cases
  343. ROOT = _install_root(__file__)
  344. ap = argparse.ArgumentParser()
  345. ap.add_argument('farm')
  346. ap.add_argument('--date', default=None, help='本批日期 (默认取 findings._generated 日期)')
  347. ap.add_argument('--dry', action='store_true', help='只校验不写盘')
  348. a = ap.parse_args()
  349. fdir = P.sop(a.farm)
  350. fdoc = json.loads((fdir / 'findings.json').read_text(encoding='utf-8'))
  351. cpath = fdir / 'cases.json'
  352. prev = json.loads(cpath.read_text(encoding='utf-8')) if cpath.exists() else []
  353. batch_id = fdoc.get('_generated', 'unknown')
  354. now = a.date or str(batch_id)[:10]
  355. cases = produce_cases(fdoc, prev, batch_id, now)
  356. issues = validate_cases(cases)
  357. if issues:
  358. print(f'❌ validate_cases {len(issues)} issue:')
  359. for i in issues[:20]:
  360. print(' -', i)
  361. sys.exit(2)
  362. print(f'✅ {a.farm}: {len(cases)} cases (prev {len(prev)}), validate_cases PASS')
  363. from collections import Counter
  364. print(' 状态分布:', dict(Counter(c['state'] for c in cases)))
  365. if not a.dry:
  366. cpath.write_text(json.dumps(cases, ensure_ascii=False, indent=2), encoding='utf-8')
  367. print(f' 写盘: {cpath.relative_to(ROOT)}')