cases.py 19 KB

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