plugins.py 34 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526
  1. # -*- coding: utf-8 -*-
  2. """插件注册表 (万物可插件, 即用即抛). 每个插件 = name / desc / params / run(ctx, **kw) -> dict(kind, data, svg?, text).
  3. 固定界面 = 固化的插件组合; 工作台 = 手动调插件; 交互 = 模型做调度选插件。判断类插件一律调 src/sop/discriminators 库函数, 禁在此重写。"""
  4. from .. import paths as P
  5. import io, json, re, contextlib, traceback
  6. import numpy as np
  7. import pandas as pd
  8. from .config import SENSORS, SENSOR_CN, HIGH_BIN
  9. from .config import worst_level
  10. from . import data, tcm, svg
  11. REGISTRY = {}
  12. def plugin(name, desc, params=None):
  13. def deco(fn):
  14. REGISTRY[name] = dict(name=name, desc=desc, params=params or {}, run=fn)
  15. return fn
  16. return deco
  17. def _tid(t):
  18. t = str(t).strip().upper().replace('#', '').replace('WTG', '').replace('号', '')
  19. return f'WTG{int(t):02d}'
  20. def _sensor(s):
  21. s = str(s)
  22. for k, v in SENSOR_CN.items():
  23. if s == k or s == v or s in v or s.lower() in k.lower():
  24. return k
  25. return s
  26. @plugin('scalar_trend', '某台某测点某标量在某功率档的六窗趋势 (带 TCM mask 阈值或 fleet 参考线)',
  27. {'turbine': 'WTG23', 'sensor': '主轴承前', 'meas': 'Peak', 'bin': HIGH_BIN})
  28. def scalar_trend(ctx, turbine, sensor, meas='Peak', bin=HIGH_BIN, log=None):
  29. t, sen = _tid(turbine), _sensor(sensor)
  30. q, th = tcm.trend(ctx['scalars'], ctx['masks'], 'CGN Rudong', t, sen, meas, bin)
  31. ref = tcm.fleet_ref(ctx['scalars'], sen, meas, bin)
  32. hl = ([(th.get('yellow'), '黄', 'var(--th-yellow)'), (th.get('red'), '红', 'var(--th-red)')] if th
  33. else [(ref['fleet_median'], 'fleet中位', 'var(--ref-line)'), (ref['fleet_p90'], 'fleet p90', 'var(--ref-line2)')])
  34. lg = log if log is not None else meas in ('Peak', 'Kurtosis')
  35. s = svg.line_chart([{'name': meas, 'x': list(q.trigger_time), 'y': list(q.scalar_value), 'marker': True}],
  36. hlines=hl, title=f'{t} {SENSOR_CN.get(sen, sen)} {meas} @{tcm.bin_label(bin)}', log=lg)
  37. per_w = q.groupby('window').scalar_value.agg(['median', 'max', 'size']).round(4)
  38. return dict(kind='chart', svg=s, data=per_w.reset_index().to_dict('records'),
  39. text=(f'{t} {SENSOR_CN.get(sen, sen)} {meas} @{tcm.bin_label(bin)}: 逐窗中位 {per_w["median"].to_dict()}; '
  40. f'逐窗最大 {per_w["max"].to_dict()}; fleet 同档中位 {ref["fleet_median"]:.4g}, p90 {ref["fleet_p90"]:.4g}; '
  41. f'mask 阈值 {th or "无 (该标量无 TCM mask)"}'))
  42. @plugin('spectrum', '某台某测点最新谱 (包络 env 或原始 fft) + OEM 特征线附近峰/邻域比',
  43. {'turbine': 'WTG16', 'sensor': '发电机DE', 'kind': 'env', 'xmax': 300})
  44. def spectrum(ctx, turbine, sensor, kind='env', xmax=300, meas_name=None):
  45. from .report import pick_spec, oem_lines
  46. t, sen = _tid(turbine), _sensor(sensor)
  47. mn = meas_name or pick_spec(ctx['cfg'], t, sen, kind)
  48. sp = data.spectrum(ctx['cfg'], t, sen, mn) if mn else None
  49. if sp is None:
  50. return dict(kind='text', text=f'{t} {sen} 无 {kind} 谱')
  51. x, y, meta = sp
  52. marks = oem_lines(ctx['cfg'], sen)
  53. s = svg.spectrum_chart(x, y, marks=marks, title=f'{t} {SENSOR_CN.get(sen, sen)} {mn} {str(meta["trigger_time"])[:16]}',
  54. xmax=float(xmax), ylabel=str(meta.get('y_unit', '')))
  55. peaks = []
  56. G2_MIN = 3.0 # 内部判据 (公开模式只露 过/不过)
  57. for hz, lab, _ in marks:
  58. m = (x >= hz - 1.0) & (x <= hz + 1.0)
  59. if m.any():
  60. i = np.argmax(y[m])
  61. nb = (np.abs(x - hz) <= 6) & ~m
  62. ratio = float(y[m][i] / np.median(y[nb])) if nb.any() and np.median(y[nb]) > 0 else None
  63. peaks.append(dict(line=lab, hz=float(hz), peak_hz=float(x[m][i]), peak=float(y[m][i]), ratio_to_neighborhood=ratio, G2=('过' if ratio is not None and ratio >= G2_MIN else '不过')))
  64. desc = ', '.join(f'{p["line"]}@{p["peak_hz"]:.2f}Hz {p["ratio_to_neighborhood"]:.1f}× G2{p["G2"]}' for p in peaks if p['ratio_to_neighborhood'])
  65. n_pass = sum(1 for p in peaks if p['G2'] == '过')
  66. return dict(kind='chart', svg=s, data=peaks,
  67. text=f'{t} {SENSOR_CN.get(sen, sen)} {mn} 最新 {meta["trigger_time"]} (rpm {float(meta["rpm"]):.0f}): OEM 线附近 峰/邻域 = {desc} | G2 过 {n_pass}/{len(peaks)} 条 ')
  68. @plugin('fleet_rank', '某测点某标量在某功率档的全场排名 (末窗中位 ×fleet)',
  69. {'sensor': '主轴承前', 'meas': 'Peak', 'bin': HIGH_BIN, 'top': 10})
  70. def fleet_rank(ctx, sensor, meas='Peak', bin=HIGH_BIN, top=10):
  71. sen = _sensor(sensor)
  72. st = ctx['status']
  73. q = st[(st.sensor_name == sen) & (st.meas_name == meas)].sort_values('x_fleet', ascending=False).head(int(top))
  74. rows = q[['turbine', 'median', 'x_fleet', 'mask_status', 'size']].round(4).to_dict('records')
  75. return dict(kind='table', data=rows,
  76. text=f'{SENSOR_CN.get(sen, sen)} {meas} 末窗中位 ×fleet 前 {top}: ' + ', '.join(f'{r["turbine"]} {r["median"]:.4g} ({r["x_fleet"]:.2f}×)' for r in rows))
  77. def turbine_state(ctx, t):
  78. """业主口径状态等级 (优秀/良好/报警/危险/不可判) — 与模版报告同一套确定性映射 (report_std), 单台现算 + ctx 缓存."""
  79. cache = ctx.setdefault('_tstate', {})
  80. if t in cache:
  81. return cache[t]
  82. from .report_std import COMP, line_state, impact_state, _worst_state, evidence_state, action_level
  83. l6 = ctx['l6']; summ = ctx['summary']
  84. n = int(t[3:])
  85. lines = l6[l6['台'] == f'{n}#'] if len(l6) else pd.DataFrame()
  86. l0 = summ.get('L0', {})
  87. comp_state, all_lv, all_imp = {}, [], []
  88. for comp, sens in COMP.items():
  89. states = []
  90. for sen in sens:
  91. l0v = l0.get(f'{t}|{sen}', '')
  92. if l0v and not l0v.startswith('可用'):
  93. states.append('不可判')
  94. continue
  95. lv = lines[lines['测点'] == sen]['定级'].tolist() if len(lines) else []
  96. all_lv += lv
  97. states += [line_state(x) for x in lv]
  98. imp = impact_typology(ctx, t, sen)
  99. from src.sop.vib_impact import VALID_SENSORS as _IMPACT_CAL
  100. cal = sen in _IMPACT_CAL
  101. all_imp.append(imp.get('type') if cal else None)
  102. states.append(impact_state(imp, calibrated_sensor=cal))
  103. if cal and imp.get('level') in ('候选·新发', '候选'):
  104. all_lv.append(imp['level'])
  105. from .report_std import iso_state, _iso_max, ISO_C
  106. iso_v = _iso_max(ctx, t, sen)
  107. ist = iso_state(iso_v)
  108. if ist:
  109. states.append(ist)
  110. if iso_v >= ISO_C:
  111. all_lv.append('候选')
  112. comp_state[comp] = _worst_state(states) if states else '优秀'
  113. state = _worst_state(list(comp_state.values())) # 整机级按可评部件给; 单测点缺失在分部件里可见 (交付版口径, 2026-08-24 二次校准)
  114. evid = evidence_state(all_lv, all_imp)
  115. acute = '持续恶化' in all_imp
  116. r = dict(state=state, comp=comp_state,
  117. evidence=(evid if state in ('报警', '危险') else ('INSUFFICIENT' if state == '不可判' else '—')),
  118. action=action_level(state, evid, acute))
  119. cache[t] = r
  120. return r
  121. @plugin('turbine_summary', '某台结构化结论: 融合级 / L4 过闸线 / L0 可用性 / findings 键', {'turbine': 'WTG16'})
  122. def turbine_summary(ctx, turbine):
  123. t = _tid(turbine)
  124. n = int(t[3:])
  125. fu = ctx['fusion'][ctx['fusion']['台'] == t]
  126. l6 = ctx['l6'][ctx['l6']['台'] == f'{n}#']
  127. l0 = {k.split('|')[1]: v for k, v in ctx['summary'].get('L0', {}).items() if k.startswith(t + '|')}
  128. fk = data.findings_for_turbine(ctx['findings'], t)
  129. lines = l6[['测点', '线', 'hz', '域', '绝对量', '单位', 'xfleet', '选择性', '定级']].to_dict('records') if len(l6) else []
  130. txt = (f'{t}: 融合 {fu.iloc[0]["融合"]} (CMS {fu.iloc[0]["CMS"]} / 模型 {fu.iloc[0]["模型"]}, 机制 {fu.iloc[0]["机制"]}); 模型依据: {fu.iloc[0]["模型依据"]}. '
  131. + (f'CMS 告警证据力: {fu.iloc[0]["告警证据力"]}. ' if '告警证据力' in fu.columns and str(fu.iloc[0]["告警证据力"]) not in ('', 'nan') else '')
  132. if len(fu) else f'{t}: 无融合记录. ')
  133. txt += (f'L4 过闸线 {len(lines)} 条: ' +
  134. '; '.join(f'{SENSOR_CN.get(r["测点"], r["测点"])} {r["线"]} {r["hz"]}Hz {r["域"]} {r["绝对量"]}{r["单位"]} ×fleet {r["xfleet"]} → {r["定级"]}' for r in lines) +
  135. '. L0: ' + ', '.join(f'{SENSOR_CN[s]}={l0.get(s, "?")}' for s in SENSORS) + f'. findings 键: {", ".join(fk[:12])}')
  136. # 冲击轴 (标量面, 融合表以谱线层为主所以这里必须带出来): 末窗同档 ×fleet 最高三项
  137. st = ctx['status'][(ctx['status'].turbine == t) & ctx['status'].meas_name.isin(['Peak', 'Rms_HP', 'Kurtosis'])].sort_values('x_fleet', ascending=False).head(3)
  138. imp = [dict(sensor=SENSOR_CN.get(r.sensor_name, r.sensor_name), meas=r.meas_name, median=round(float(r['median']), 4), x_fleet=round(float(r.x_fleet), 2)) for _, r in st.iterrows()]
  139. flag = any(i['x_fleet'] >= 2.0 for i in imp)
  140. txt += ' 【冲击轴 (末窗 3601-4400 kW 档, 标量 ×fleet 前三): ' + '; '.join(f'{i["sensor"]} {i["meas"]} {i["median"]} ({i["x_fleet"]}×)' for i in imp) + ('; ⚠ 标量面存在显著偏离, 融合表的谱线层结论不能单独代表该台状态, 必须看 scalar_trend/收敛报告终态】' if flag else '】')
  141. # 收敛报告 v5 逐台终态: 直接解析 §一 表格行 (按台号取整行, 不靠检索)
  142. v5 = _v5_rows(ctx['cfg'], n)
  143. if v5:
  144. txt += ' 【收敛报告 v5 终态 (整行, 不截断——尾部常是限定语): ' + v5[0][:4000] + '】'
  145. else:
  146. txt += ' 【收敛报告 v5 无该台专条 (未列入有动作/有争议台)】'
  147. # 确定性"综合判断"一行放最前: 谱线层 × 冲击轴 × v5 终态 — 模型只转述, 不自己合成
  148. fused = fu.iloc[0]['融合'] if len(fu) else '无记录'
  149. if v5:
  150. verdict = f'以收敛报告 v5 终态为准 (该台有专条); 谱线层融合级 {fused}; 冲击轴{"显著偏离" if flag else "未见显著偏离"}'
  151. elif flag:
  152. verdict = f'谱线层融合级 {fused}, 但冲击轴显著偏离 ⇒ 不能判为正常, 状态 = 标量面候选·待 scalar_trend/时域六窗层核 (非谱线层可定)'
  153. else:
  154. verdict = f'谱线层融合级 {fused}; 冲击轴未见显著偏离'
  155. ts = turbine_state(ctx, t)
  156. comp_txt = ', '.join(f'{k}={v}' for k, v in ts['comp'].items())
  157. head = f"{t} 状态等级(业主口径): **{ts['state']}** (分部件: {comp_txt}); 证据状态: {ts['evidence']}; 行动等级建议: {ts['action']}。 "
  158. txt = head + f'综合判断 (结构化层): {verdict}。' + txt.replace('. L0: ', '. L0 各测点有无数据: ')
  159. return dict(kind='text', data=dict(verdict=verdict, fusion=fu.to_dict('records'), lines=lines, l0=l0, findings_keys=fk, impact_axis=imp, v5=[v[:4000] for v in v5[:1]]), text=txt)
  160. def _v5_rows(cfg, n):
  161. """收敛报告 v5 §一 表格里该台的行 (首列含 n#); 返回去掉表格符号的文本列表."""
  162. import re
  163. from pathlib import Path
  164. out = []
  165. for p in [Path(p) for p in cfg.get('knowledge_docs', []) if '收敛报告' in str(p)]:
  166. if not p.exists():
  167. continue
  168. in_sec = False
  169. for line in p.read_text(encoding='utf-8').split('\n'):
  170. if line.startswith('## 一'):
  171. in_sec = True
  172. elif line.startswith('## ') and in_sec:
  173. break
  174. if in_sec and line.startswith('|'):
  175. cells = [c.strip() for c in line.strip('|').split('|')]
  176. if cells and re.search(rf'(^|[^0-9]){n}#', cells[0].replace('*', '')):
  177. out.append(' | '.join(c.replace('**', '') for c in cells))
  178. return out
  179. @plugin('gate_spectral', '对一条候选谱线看最新谱的 G2 有线闸 (峰/邻域本底≥3); 全六闸走生产扫描产物',
  180. {'turbine': 'WTG16', 'sensor': '发电机DE', 'hz': 135.0, 'kind': 'env'})
  181. def gate_spectral(ctx, turbine, sensor, hz, kind='env'):
  182. from .report import pick_spec
  183. t, sen = _tid(turbine), _sensor(sensor)
  184. mn = pick_spec(ctx['cfg'], t, sen, kind)
  185. sp = data.spectrum(ctx['cfg'], t, sen, mn) if mn else None
  186. if sp is None:
  187. return dict(kind='text', text='无谱')
  188. x, y, meta = sp
  189. hz = float(hz)
  190. dx = float(x[1] - x[0])
  191. m = (x >= hz - 2 * dx) & (x <= hz + 2 * dx)
  192. nb = (np.abs(x - hz) <= 20 * dx) & ~m
  193. if not m.any():
  194. return dict(kind='text', text=f'{hz} Hz 超出谱范围 {x.min():.1f}-{x.max():.1f}')
  195. peak = float(y[m].max())
  196. base = float(np.median(y[nb]))
  197. ratio = peak / base if base > 0 else None
  198. return dict(kind='text', data=dict(peak=peak, baseline=base, peak_over_baseline=ratio, bin_hz=dx, meas=mn),
  199. text=(f'{t} {SENSOR_CN.get(sen, sen)} {mn} @{hz} Hz: 峰 {peak:.4g}, 邻域本底 {base:.4g}, 峰/本底 {ratio:.2f} '
  200. f'→ G2 {"过" if ratio and ratio >= 3 else "不过"}.'))
  201. @plugin('gate_shared_component', '同步两通道共享成分污染门 (库函数 shared_component_gate)',
  202. {'coh': 0.99, 'group_delay_us': -1.2, 'phase_R': 0.9999})
  203. def gate_shared(ctx, coh, group_delay_us, phase_R):
  204. from src.sop.discriminators import shared_component_gate
  205. r = shared_component_gate(float(coh), float(group_delay_us), float(phase_R))
  206. return dict(kind='text', data=r, text=f'{r["status"]}: {r["reason"]}')
  207. @plugin('kb_search', '知识库检索 (收敛报告/裁决/交接/现场单/六层文档/skill/findings)', {'query': '23号 主轴承 时延', 'k': 5})
  208. def kb_search(ctx, query, k=5):
  209. from . import knowledge
  210. hits = knowledge.search(ctx['cfg'], query, k=int(k))
  211. return dict(kind='table', data=[dict(score=round(h['score'], 2), source=h['source'], text=h['text'][:300]) for h in hits],
  212. text='\n'.join(f'[{h["source"]}] {h["text"][:240]}' for h in hits))
  213. @plugin('adhoc', "即用即抛: 在数据上下文里执行一小段 pandas 代码, 结果放 result. 变量: scalars(列 turbine/sensor_name/meas_name/trigger_time/rpm/condition_key/scalar_value/window; sensor_name 为英文键 Main_bearing_front 等, 可用 sensor('主轴承后') 转换), status(末窗 台×测点×标量 median/x_fleet/mask_status), l6(过闸线), fusion(38台融合), np, pd, HIGH_BIN, SENSORS, SENSOR_CN",
  214. {'code': "result = scalars[(scalars.sensor_name==sensor('主轴承后'))&(scalars.meas_name=='Peak')&(scalars.condition_key==HIGH_BIN)].groupby('turbine').scalar_value.median().sort_values(ascending=False).head(5)"})
  215. def adhoc(ctx, code):
  216. ns = dict(scalars=ctx['scalars'], status=ctx['status'], l6=ctx['l6'], fusion=ctx['fusion'], np=np, pd=pd,
  217. HIGH_BIN=HIGH_BIN, SENSORS=SENSORS, SENSOR_CN=SENSOR_CN, sensor=_sensor, result=None)
  218. buf = io.StringIO()
  219. import builtins
  220. safe = {k: getattr(builtins, k) for k in ['len', 'range', 'min', 'max', 'sum', 'sorted', 'round', 'abs', 'list', 'dict',
  221. 'float', 'int', 'str', 'print', 'enumerate', 'zip', 'set', 'tuple']}
  222. try:
  223. with contextlib.redirect_stdout(buf):
  224. exec(compile(str(code), '<adhoc>', 'exec'), {'__builtins__': safe}, ns)
  225. except Exception:
  226. return dict(kind='text', text='adhoc 失败: ' + traceback.format_exc()[-600:])
  227. r = ns.get('result')
  228. if isinstance(r, (pd.DataFrame, pd.Series)):
  229. rr = r.reset_index() if isinstance(r, pd.Series) else r
  230. return dict(kind='table', data=rr.head(50).round(4).to_dict('records'), text=(buf.getvalue() + rr.head(20).round(4).to_string())[:2000])
  231. return dict(kind='text', data=r if isinstance(r, (int, float, str, list, dict)) else str(r), text=(buf.getvalue() + str(r))[:2000])
  232. @plugin('compare', '多台/多测点并排对比 (同一标量同一功率档, 一图多线)', {'turbines': 'WTG23,WTG16,WTG38', 'sensor': '主轴承前', 'meas': 'Peak', 'bin': HIGH_BIN})
  233. def compare(ctx, turbines, sensor, meas='Peak', bin=HIGH_BIN, log=None):
  234. sen = _sensor(sensor)
  235. tids = [_tid(t) for t in str(turbines).replace(',', ',').split(',') if t.strip()]
  236. cols = svg.SERIES_VARS # categorical 六槽 (skill 参考实例, 浅暗两套 validator 过检; 色随台固定不随过滤重涂)
  237. series, summary = [], {}
  238. for i, t in enumerate(tids):
  239. q, _ = tcm.trend(ctx['scalars'], ctx['masks'], 'CGN Rudong', t, sen, meas, bin)
  240. if q.empty:
  241. continue
  242. series.append({'name': t, 'x': list(q.trigger_time), 'y': list(q.scalar_value), 'marker': True, 'color': cols[i % len(cols)]})
  243. summary[t] = q.groupby('window').scalar_value.median().round(4).to_dict()
  244. ref = tcm.fleet_ref(ctx['scalars'], sen, meas, bin)
  245. lg = log if log is not None else meas in ('Peak', 'Kurtosis')
  246. s = svg.line_chart(series, hlines=[(ref['fleet_median'], 'fleet中位', 'var(--ref-line)'), (ref['fleet_p90'], 'fleet p90', 'var(--ref-line2)')], title=f'对比 {SENSOR_CN.get(sen, sen)} {meas} @{tcm.bin_label(bin)}', log=lg) if series else ''
  247. return dict(kind='chart', svg=s, data=summary, text=f'对比 {SENSOR_CN.get(sen, sen)} {meas}: 逐窗中位 ' + json.dumps(summary, ensure_ascii=False) + f'; fleet 中位 {ref["fleet_median"]:.4g}, p90 {ref["fleet_p90"]:.4g}')
  248. @plugin('hypothesis', '假设簿: 记录一条假设 (陈述/支持证据/反证检验/当前状态); 状态只能是 候选/INSUFFICIENT/撤回 — 定级须经 verdict_gate',
  249. {'statement': '23# 前主轴承冲击为转子侧周期载荷调制', 'evidence': 'scalar_trend: …; kb_search: …', 'falsifier': '若对中/螺栓/不平衡现场全正常而内窥见滚道剥落则假', 'status': 'INSUFFICIENT'})
  250. def hypothesis(ctx, statement, evidence='', falsifier='', status='INSUFFICIENT'):
  251. if status not in ('候选', 'INSUFFICIENT', '撤回'):
  252. status = 'INSUFFICIENT'
  253. h = dict(id=len(ctx.setdefault('hypotheses', [])) + 1, statement=str(statement), evidence=str(evidence), falsifier=str(falsifier), status=status)
  254. ctx['hypotheses'].append(h)
  255. return dict(kind='text', data=h, text=f'假设 #{h["id"]} 已记录 [{status}]: {statement} | 证据: {evidence[:300]} | 证伪: {falsifier[:200]}')
  256. @plugin('verdict_gate', '定级闸 (确定性): 给某台某测点某线按生产扫描产物 (六闸+绝对锚+证据族) 返回六枚举级别; 模型不得自定级',
  257. {'turbine': 'WTG16', 'sensor': '发电机DE', 'hz': 135.0})
  258. def verdict_gate(ctx, turbine, sensor, hz=None):
  259. t, sen = _tid(turbine), _sensor(sensor)
  260. n = int(t[3:])
  261. l6 = ctx['l6']
  262. q = l6[(l6['台'] == f'{n}#') & (l6['测点'] == sen)]
  263. if hz is not None and len(q):
  264. q = q[(pd.to_numeric(q['hz'], errors='coerce') - float(hz)).abs() <= 2.0]
  265. if len(q):
  266. rows = q[['线', 'hz', '域', '绝对量', '单位', 'xfleet', '选择性', '证据族', '解封判据', '定级']].to_dict('records')
  267. lv = [r['定级'] for r in rows]
  268. best = worst_level(lv)
  269. return dict(kind='table', data=rows, text=f'{t} {SENSOR_CN.get(sen, sen)}' + (f' @{hz}Hz' if hz else '') + f': 过闸线 {len(rows)} 条, 最高定级 {best}; ' + '; '.join(f'{r["线"]} {r["hz"]}Hz {r["域"]} {r["绝对量"]}{r["单位"]} ×fleet {r["xfleet"]} 证据族 {r["证据族"]} 解封 {r["解封判据"]} → {r["定级"]}' for r in rows))
  270. st = ctx['status'][(ctx['status'].turbine == t) & (ctx['status'].sensor_name == sen)]
  271. top = st.sort_values('x_fleet', ascending=False).head(3)
  272. imp = impact_typology(ctx, t, sen)
  273. return dict(kind='text', data=dict(lines=[], scalars=top[['meas_name', 'median', 'x_fleet', 'mask_status']].round(4).to_dict('records'), impact=imp),
  274. text=f'{t} {SENSOR_CN.get(sen, sen)}' + (f' @{hz}Hz' if hz else '') + ': 无过闸谱线 (L4 层) ⇒ 谱线层 INSUFFICIENT; 标量面末窗 ×fleet: ' + ', '.join(f'{r.meas_name} {r.x_fleet:.2f}× ({r.mask_status})' for r in top.itertuples())
  275. + f'. 冲击轴 (时域六窗层, 库函数 vib_scalar_typology, Peak 逐日 max/中位): 分型 {imp["type"]} → 定级 {imp["level"]}; {imp["detail"]}')
  276. def impact_typology(ctx, t, sen, meas='Peak'):
  277. """冲击轴分型 → 六枚举: 薄包装, 实现在库层 src/sop/vib_impact.impact_axis_typology (fusion 脚本同调, 单一实现)."""
  278. from src.sop.vib_impact import impact_axis_typology
  279. return impact_axis_typology(ctx['scalars'], t, sen, meas, event_windows=ctx['cfg'].get('event_windows'))
  280. @plugin('render', '组装视图 (模型驱动前端): panels = 插件调用列表 JSON [{"plugin":..,"args":{..}}, ...], 每个面板出图/表/文; 前端画布按顺序渲染',
  281. {'title': '23# vs 16# 主轴承前', 'panels': '[{"plugin":"compare","args":{"turbines":"WTG23,WTG16","sensor":"主轴承前","meas":"Peak"}},{"plugin":"spectrum","args":{"turbine":"WTG23","sensor":"主轴承前","kind":"env"}}]'})
  282. def render(ctx, title, panels):
  283. try:
  284. spec = json.loads(panels) if isinstance(panels, str) else panels
  285. except Exception:
  286. return dict(kind='text', text='panels 不是合法 JSON')
  287. out = []
  288. for p in spec[:12]:
  289. name, args = p.get('plugin'), p.get('args') or {}
  290. if name in ('render',):
  291. continue
  292. r = run(ctx, name, **args)
  293. out.append(dict(plugin=name, args=args, svg=r.get('svg', ''), text=str(r.get('text', ''))[:1500], data=(r.get('data') if isinstance(r.get('data'), list) else None)))
  294. view = dict(id=len(ctx.setdefault('views', [])) + 1, title=str(title), panels=out)
  295. ctx['views'].append(view)
  296. try:
  297. vp = ctx['cfg']['out'] / 'views.json'
  298. # 读写都要显式 utf-8: 面板文本是中文, 缺 encoding 会按系统 locale(cp936) 读/写 → 乱码或崩
  299. old = json.load(open(vp, encoding='utf-8')) if vp.exists() else []
  300. with open(vp, 'w', encoding='utf-8') as f:
  301. json.dump((old + [dict(id=view['id'], title=view['title'], panels=[dict(plugin=x['plugin'], args=x['args'], svg=x['svg'], text=x['text']) for x in out])])[-20:], f, ensure_ascii=False)
  302. except Exception:
  303. pass
  304. return dict(kind='view', data=view, text=f'视图 "{title}" 已组装 {len(out)} 个面板: ' + '; '.join(f'{x["plugin"]}: {x["text"][:120]}' for x in out))
  305. @plugin('analyze', '原始数据 → 报告 全链 (后台任务, 约 8 分钟, 重写全部模型产物): 必须 confirm=true 才会启动; 同一时刻只允许一个; 可带 input (TCM 导出包 / 解码目录 / 波形目录) 与 window 名; 不带则在现有窗上重跑; 用 job_status 看进度',
  306. {'input': '', 'window': '', 'steps': '', 'confirm': False})
  307. def analyze(ctx, input='', window='', steps='', confirm=False):
  308. import subprocess, sys, os
  309. from .config import ROOT
  310. if not (confirm is True or str(confirm).lower() in ('true', '1', 'yes')):
  311. return dict(kind='text', text='analyze 未启动: 这是约 8 分钟、重写全部模型产物的后台任务, 需要 confirm=true (由用户明确要求时才传). 当前产物可直接用其他插件查询.')
  312. lock = ctx['cfg']['out'] / 'analyze.lock'
  313. if lock.exists():
  314. try:
  315. pid = int(lock.read_text(encoding='utf-8').strip())
  316. os.kill(pid, 0)
  317. return dict(kind='text', text=f'analyze 已在运行 (pid {pid}), 不重复启动; 用 job_status 看进度.')
  318. except Exception:
  319. lock.unlink(missing_ok=True)
  320. if input and not os.path.exists(os.path.expanduser(str(input))):
  321. return dict(kind='text', text=f'导入失败: 路径不存在 {input!r} — 请确认盘已挂载、路径完整 (可整段从访达拖入); 支持 TCM 导出包 / 解码目录 / 波形目录.')
  322. # 冻结版 sys.executable = 打包二进制, spawn 分析链必挂 → 优先 WINDCMS_SCRIPT_PY (launch.sh 提供, 默认 venv)
  323. py = os.environ.get('WINDCMS_SCRIPT_PY') or sys.executable
  324. cmd = [py, str(ROOT / 'scripts/windcms.py'), 'analyze', '--farm', 'rudong']
  325. if input:
  326. cmd += ['--input', os.path.expanduser(str(input)), '--window', str(window or 'w_new')]
  327. if steps:
  328. cmd += ['--steps'] + [s for s in str(steps).replace(',', ' ').split() if s]
  329. logp = ctx['cfg']['out'] / 'analyze_stdout.log'
  330. ctx['cfg']['out'].mkdir(parents=True, exist_ok=True)
  331. # stdout= 只用到这个文件的 fd (内容由子进程自己写), 所以这里不 with-close:
  332. # 句柄留着, 免得父进程提前关掉; 显式 encoding 是为了不留"缺 encoding"的门禁告警。
  333. # ★2026-09-16: 本函数跑在 **CMS 服务进程**里 (由 guanlan.py 以无窗口方式拉起, 自己没有可见控制台),
  334. # 裸 Popen 会让 Windows 给全链分析**新建一个可见控制台窗口** (从页面点"分析"就会看到黑窗)。
  335. # 统一走 src.proc.spawn (CREATE_NO_WINDOW)。
  336. from src import proc as _proc
  337. _logf = open(logp, 'w', encoding='utf-8')
  338. p = _proc.spawn(cmd, cwd=ROOT, stdout=_logf, stderr=subprocess.STDOUT)
  339. lock.write_text(str(p.pid), encoding='utf-8')
  340. return dict(kind='text', data=dict(pid=p.pid, log=str(logp)), text=f'全链已在后台启动 (pid {p.pid}); 进度看 job_status; 日志 {logp}. 完成后产物/报告/知识库自动刷新 (界面需刷新页面).')
  341. @plugin('farm_overview', '全场38台业主四级总览 (报警/不可判在前): 台/整机/四部件/证据/行动/焦点 — 全场类问题先调本工具保证完整性', {})
  342. def farm_overview(ctx):
  343. from .report_std import registry, STATE_ORDER
  344. reg, detail, _ = registry(ctx['cfg'])
  345. reg = reg.sort_values('状态等级', key=lambda s_: s_.map({x: i for i, x in enumerate(STATE_ORDER)}))
  346. lines = [f"{r['机组']} {r['状态等级']} | 前{r['主轴承前']} 后{r['主轴承后']} 齿{r['齿轮箱']} 发{r['发电机']} | 证据{r['证据状态']} 行动{r['行动等级建议']}" + (f" | 焦点 {r['状态焦点']}" if r['状态焦点'] != '—' else '')
  347. for _, r in reg.iterrows()]
  348. dist = reg['状态等级'].value_counts().to_dict()
  349. return dict(kind='text', data=reg.to_dict('records'), text=f'分布 {dist}\n' + '\n'.join(lines))
  350. @plugin('job_status', '查看后台 analyze 任务进度与上次结果 (analyze_log.json)', {})
  351. def job_status(ctx):
  352. out = ctx['cfg']['out']
  353. # ★ 2026-09-17 用户令 2: analyze 的运行日志与结果都改落 logs/build/ (原先落在产物仓 out/ 下);
  354. # 读侧同步改, 并保留旧路径兜底 —— 老机器上可能还留着那时的件, 不该一改就显示"无运行日志"。
  355. from src import logfile as _lf
  356. _b = _lf.audit_log('windcms_analyze.json').parent # logs/audit/
  357. logp = (_b / 'analyze_stdout.log') if (_b / 'analyze_stdout.log').exists() else (out / 'analyze_stdout.log')
  358. lp = (_b / 'analyze_log.json') if (_b / 'analyze_log.json').exists() else (out / 'analyze_log.json')
  359. # 日志与结果都是中文文本: 不写 encoding 会按系统 locale(cp936) 读 → 乱码/解码失败
  360. tail = open(logp, encoding='utf-8').read()[-1200:] if logp.exists() else '(无运行日志)'
  361. last = json.load(open(lp, encoding='utf-8')) if lp.exists() else None
  362. summ = (f"上次结果: {last['status']} {last['seconds']}s @ {last['finished']}; 窗 {last['windows']}; 步骤 " + ', '.join(f"{s['step']}:{s['rc']}/{s['seconds']}s" for s in last['steps'])) if last else '尚无完成记录'
  363. return dict(kind='text', data=last, text=f'{summ}\n--- 实时日志尾 ---\n{tail}')
  364. def context(cfg):
  365. s = data.load_scalars(cfg)
  366. masks = data.load_masks(cfg)
  367. l6, fusion, summ = data.load_model(cfg)
  368. return dict(cfg=cfg, scalars=s, masks=masks, l6=l6, fusion=fusion, summary=summ,
  369. findings=data.load_findings(cfg), status=tcm.status_table(s, masks, 'CGN Rudong'))
  370. def _cn_sensors(text):
  371. """工具输出正文里的英文通道名/谱种标签 → 中文 (模型会原样往答案里抄, 源头治理; data(JSON) 保留原值供机器用)."""
  372. for k, v in SENSOR_CN.items():
  373. text = text.replace(k, v)
  374. for a, b in [('Indicator_1P', '1倍频指标'), ('Indicator_2P', '2倍频指标'), ('Indicator_HSW', '高速轴指标'), ('Indicator', '厂商指标')]:
  375. text = text.replace(a, b)
  376. text = re.sub(r'\benv(\d+)\b', r'包络·\1', text)
  377. text = re.sub(r'\benv\b', '包络谱', text)
  378. text = re.sub(r'\braw\b', '原始谱', text)
  379. return text.replace('(包络域)', '').replace('(原始域)', '')
  380. def run(ctx, name, **kw):
  381. from .redact import redact_obj, redact
  382. if name not in REGISTRY:
  383. return dict(kind='text', text=f'未知插件 {name}; 可用: {list(REGISTRY)}')
  384. try:
  385. r = REGISTRY[name]['run'](ctx, **kw)
  386. except Exception:
  387. return dict(kind='text', text=f'插件 {name} 失败: ' + traceback.format_exc()[-800:])
  388. # 公开模式: 文本/数据里的判据阈值脱敏 (svg 内的数值标签由 svg 层按 PUBLIC 处理)
  389. if 'text' in r:
  390. r['text'] = redact(_cn_sensors(str(r['text'])))
  391. if r.get('data') is not None and name in ('gate_spectral', 'gate_shared_component', 'verdict_gate', 'turbine_summary', 'kb_search', 'hypothesis', 'spectrum', 'render', 'compare'):
  392. r['data'] = redact_obj(r['data'])
  393. return r
  394. @plugin('explain', '机器语言→交互语言: 把任意插件的结构化输出翻译成给业主/现场看的自然汉语 (经接地闸: 台号与数字必须来自原文, 不过闸则退回原文); source=插件名, args=其参数 JSON',
  395. {'source': 'turbine_summary', 'args': '{"turbine":"WTG38"}', 'audience': '业主'})
  396. def explain(ctx, source, args='{}', audience='业主'):
  397. from . import llm
  398. try:
  399. a = json.loads(args) if isinstance(args, str) else (args or {})
  400. except Exception:
  401. return dict(kind='text', text='args 不是合法 JSON')
  402. r = run(ctx, source, **a) # 经 run ⇒ 公开模式已脱敏
  403. raw = str(r.get('text', ''))
  404. if raw.startswith(('未知插件', '插件')) or not raw:
  405. return r
  406. prompt = (f'把下面这段风机振动诊断系统的结构化输出, 改写成给{audience}看的自然中文 (2-5 句):\n'
  407. '要求: 先结论后数字; 每个数字保留并注明含义; 不新增任何数字/台号/结论; 不解释判据算法; '
  408. '缩写展开 (如 "候选" 说成 "候选级异常, 需现场核实确认"); 语气平实, 不用套话。\n\n=== 原文 ===\n' + raw[:3500])
  409. txt, m = llm.generate(prompt, temperature=0.2, num_predict=400, timeout=120)
  410. if txt:
  411. ok, bad = llm.grounding(txt, raw)
  412. if ok:
  413. return dict(kind='text', data=dict(source=source, model=m, grounded=True),
  414. text=txt.strip() + '\n\n——以上为通俗转述 (经接地闸核对); 结构化原文按需查 ' + source + ' 插件。')
  415. return dict(kind='text', data=dict(source=source, grounded=False),
  416. text=raw + '\n\n(通俗转述未过接地闸或模型不可用, 保留结构化原文)')
  417. @plugin('export', '导出: 把任意插件的结果 (或整表) 存成 电子表格(xlsx) / Word(docx) / CSV, 落 out/exports/; source = 插件名, args = 该插件参数 JSON, fmt = xlsx|docx|csv',
  418. {'source': 'fleet_rank', 'args': '{"sensor":"主轴承后","meas":"Peak"}', 'fmt': 'xlsx', 'title': '主轴承后 Peak 全场排名'})
  419. def export(ctx, source, args='{}', fmt='xlsx', title=''):
  420. """★脱敏在出口生效: 走 run() 拿结果 ⇒ 公开模式下导出件同样不含阈值/特征频率 (client-deliverable-hide-core-parameters)."""
  421. import datetime
  422. try:
  423. a = json.loads(args) if isinstance(args, str) else (args or {})
  424. except Exception:
  425. return dict(kind='text', text='args 不是合法 JSON')
  426. r = run(ctx, source, **a) # 经 run ⇒ 已脱敏
  427. if str(r.get('text', '')).startswith(('未知插件', '插件')):
  428. return r
  429. data = r.get('data')
  430. rows = data if isinstance(data, list) else ([data] if isinstance(data, dict) else [])
  431. df = pd.DataFrame(rows) if rows else pd.DataFrame()
  432. ttl = title or f'{source} {json.dumps(a, ensure_ascii=False)}'
  433. out = ctx['cfg']['out'] / 'exports'
  434. out.mkdir(parents=True, exist_ok=True)
  435. stamp = datetime.datetime.now().strftime('%Y%m%d_%H%M%S')
  436. base = re.sub(r'[^\w一-龥-]+', '_', f'{source}_{stamp}')[:80]
  437. fmt = (fmt or 'xlsx').lower()
  438. if fmt == 'csv':
  439. p = out / f'{base}.csv'
  440. df.to_csv(p, index=False, encoding='utf-8-sig')
  441. elif fmt == 'docx':
  442. from docx import Document
  443. from docx.shared import Pt
  444. doc = Document()
  445. doc.styles['Normal'].font.name = 'PingFang SC'
  446. doc.styles['Normal'].font.size = Pt(10.5)
  447. doc.add_heading(ttl, level=1)
  448. doc.add_paragraph(str(r.get('text', ''))[:4000])
  449. if len(df):
  450. tb = doc.add_table(rows=1, cols=len(df.columns)); tb.style = 'Table Grid'
  451. for i, c in enumerate(df.columns):
  452. tb.rows[0].cells[i].text = str(c)
  453. for _, rr in df.iterrows():
  454. cells = tb.add_row().cells
  455. for i, c in enumerate(df.columns):
  456. cells[i].text = '' if pd.isna(rr[c]) else str(rr[c])
  457. doc.add_paragraph(f'\n生成 {stamp} · 来源插件 {source} · 参数 {json.dumps(a, ensure_ascii=False)}'
  458. + ('\n(公开模式:判据门槛与特征频率已隐去)' if __import__('src.windcms.redact', fromlist=['PUBLIC']).PUBLIC else ''))
  459. p = out / f'{base}.docx'
  460. doc.save(p)
  461. else:
  462. p = out / f'{base}.xlsx'
  463. with pd.ExcelWriter(p, engine='openpyxl') as w:
  464. (df if len(df) else pd.DataFrame({'说明': [str(r.get('text', ''))[:3000]]})).to_excel(w, index=False, sheet_name='数据')
  465. pd.DataFrame({'项': ['标题', '来源插件', '参数', '生成时间', '结论文本'],
  466. '值': [ttl, source, json.dumps(a, ensure_ascii=False), stamp, str(r.get('text', ''))[:3000]]}).to_excel(w, index=False, sheet_name='说明')
  467. return dict(kind='text', data=dict(path=str(p), rows=len(df), fmt=fmt),
  468. text=f'已导出 {fmt}: {p.name}({len(df)} 行)→ {p.parent}')
  469. def manifest():
  470. from .redact import redact
  471. return [dict(name=v['name'], desc=redact(v['desc']), params=v['params']) for v in REGISTRY.values()]