| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526 |
- # -*- coding: utf-8 -*-
- """插件注册表 (万物可插件, 即用即抛). 每个插件 = name / desc / params / run(ctx, **kw) -> dict(kind, data, svg?, text).
- 固定界面 = 固化的插件组合; 工作台 = 手动调插件; 交互 = 模型做调度选插件。判断类插件一律调 src/sop/discriminators 库函数, 禁在此重写。"""
- from .. import paths as P
- import io, json, re, contextlib, traceback
- import numpy as np
- import pandas as pd
- from .config import SENSORS, SENSOR_CN, HIGH_BIN
- from .config import worst_level
- from . import data, tcm, svg
- REGISTRY = {}
- def plugin(name, desc, params=None):
- def deco(fn):
- REGISTRY[name] = dict(name=name, desc=desc, params=params or {}, run=fn)
- return fn
- return deco
- def _tid(t):
- t = str(t).strip().upper().replace('#', '').replace('WTG', '').replace('号', '')
- return f'WTG{int(t):02d}'
- def _sensor(s):
- s = str(s)
- for k, v in SENSOR_CN.items():
- if s == k or s == v or s in v or s.lower() in k.lower():
- return k
- return s
- @plugin('scalar_trend', '某台某测点某标量在某功率档的六窗趋势 (带 TCM mask 阈值或 fleet 参考线)',
- {'turbine': 'WTG23', 'sensor': '主轴承前', 'meas': 'Peak', 'bin': HIGH_BIN})
- def scalar_trend(ctx, turbine, sensor, meas='Peak', bin=HIGH_BIN, log=None):
- t, sen = _tid(turbine), _sensor(sensor)
- q, th = tcm.trend(ctx['scalars'], ctx['masks'], 'CGN Rudong', t, sen, meas, bin)
- ref = tcm.fleet_ref(ctx['scalars'], sen, meas, bin)
- hl = ([(th.get('yellow'), '黄', 'var(--th-yellow)'), (th.get('red'), '红', 'var(--th-red)')] if th
- else [(ref['fleet_median'], 'fleet中位', 'var(--ref-line)'), (ref['fleet_p90'], 'fleet p90', 'var(--ref-line2)')])
- lg = log if log is not None else meas in ('Peak', 'Kurtosis')
- s = svg.line_chart([{'name': meas, 'x': list(q.trigger_time), 'y': list(q.scalar_value), 'marker': True}],
- hlines=hl, title=f'{t} {SENSOR_CN.get(sen, sen)} {meas} @{tcm.bin_label(bin)}', log=lg)
- per_w = q.groupby('window').scalar_value.agg(['median', 'max', 'size']).round(4)
- return dict(kind='chart', svg=s, data=per_w.reset_index().to_dict('records'),
- text=(f'{t} {SENSOR_CN.get(sen, sen)} {meas} @{tcm.bin_label(bin)}: 逐窗中位 {per_w["median"].to_dict()}; '
- f'逐窗最大 {per_w["max"].to_dict()}; fleet 同档中位 {ref["fleet_median"]:.4g}, p90 {ref["fleet_p90"]:.4g}; '
- f'mask 阈值 {th or "无 (该标量无 TCM mask)"}'))
- @plugin('spectrum', '某台某测点最新谱 (包络 env 或原始 fft) + OEM 特征线附近峰/邻域比',
- {'turbine': 'WTG16', 'sensor': '发电机DE', 'kind': 'env', 'xmax': 300})
- def spectrum(ctx, turbine, sensor, kind='env', xmax=300, meas_name=None):
- from .report import pick_spec, oem_lines
- t, sen = _tid(turbine), _sensor(sensor)
- mn = meas_name or pick_spec(ctx['cfg'], t, sen, kind)
- sp = data.spectrum(ctx['cfg'], t, sen, mn) if mn else None
- if sp is None:
- return dict(kind='text', text=f'{t} {sen} 无 {kind} 谱')
- x, y, meta = sp
- marks = oem_lines(ctx['cfg'], sen)
- s = svg.spectrum_chart(x, y, marks=marks, title=f'{t} {SENSOR_CN.get(sen, sen)} {mn} {str(meta["trigger_time"])[:16]}',
- xmax=float(xmax), ylabel=str(meta.get('y_unit', '')))
- peaks = []
- G2_MIN = 3.0 # 内部判据 (公开模式只露 过/不过)
- for hz, lab, _ in marks:
- m = (x >= hz - 1.0) & (x <= hz + 1.0)
- if m.any():
- i = np.argmax(y[m])
- nb = (np.abs(x - hz) <= 6) & ~m
- ratio = float(y[m][i] / np.median(y[nb])) if nb.any() and np.median(y[nb]) > 0 else None
- 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 '不过')))
- 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'])
- n_pass = sum(1 for p in peaks if p['G2'] == '过')
- return dict(kind='chart', svg=s, data=peaks,
- text=f'{t} {SENSOR_CN.get(sen, sen)} {mn} 最新 {meta["trigger_time"]} (rpm {float(meta["rpm"]):.0f}): OEM 线附近 峰/邻域 = {desc} | G2 过 {n_pass}/{len(peaks)} 条 ')
- @plugin('fleet_rank', '某测点某标量在某功率档的全场排名 (末窗中位 ×fleet)',
- {'sensor': '主轴承前', 'meas': 'Peak', 'bin': HIGH_BIN, 'top': 10})
- def fleet_rank(ctx, sensor, meas='Peak', bin=HIGH_BIN, top=10):
- sen = _sensor(sensor)
- st = ctx['status']
- q = st[(st.sensor_name == sen) & (st.meas_name == meas)].sort_values('x_fleet', ascending=False).head(int(top))
- rows = q[['turbine', 'median', 'x_fleet', 'mask_status', 'size']].round(4).to_dict('records')
- return dict(kind='table', data=rows,
- 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))
- def turbine_state(ctx, t):
- """业主口径状态等级 (优秀/良好/报警/危险/不可判) — 与模版报告同一套确定性映射 (report_std), 单台现算 + ctx 缓存."""
- cache = ctx.setdefault('_tstate', {})
- if t in cache:
- return cache[t]
- from .report_std import COMP, line_state, impact_state, _worst_state, evidence_state, action_level
- l6 = ctx['l6']; summ = ctx['summary']
- n = int(t[3:])
- lines = l6[l6['台'] == f'{n}#'] if len(l6) else pd.DataFrame()
- l0 = summ.get('L0', {})
- comp_state, all_lv, all_imp = {}, [], []
- for comp, sens in COMP.items():
- states = []
- for sen in sens:
- l0v = l0.get(f'{t}|{sen}', '')
- if l0v and not l0v.startswith('可用'):
- states.append('不可判')
- continue
- lv = lines[lines['测点'] == sen]['定级'].tolist() if len(lines) else []
- all_lv += lv
- states += [line_state(x) for x in lv]
- imp = impact_typology(ctx, t, sen)
- from src.sop.vib_impact import VALID_SENSORS as _IMPACT_CAL
- cal = sen in _IMPACT_CAL
- all_imp.append(imp.get('type') if cal else None)
- states.append(impact_state(imp, calibrated_sensor=cal))
- if cal and imp.get('level') in ('候选·新发', '候选'):
- all_lv.append(imp['level'])
- from .report_std import iso_state, _iso_max, ISO_C
- iso_v = _iso_max(ctx, t, sen)
- ist = iso_state(iso_v)
- if ist:
- states.append(ist)
- if iso_v >= ISO_C:
- all_lv.append('候选')
- comp_state[comp] = _worst_state(states) if states else '优秀'
- state = _worst_state(list(comp_state.values())) # 整机级按可评部件给; 单测点缺失在分部件里可见 (交付版口径, 2026-08-24 二次校准)
- evid = evidence_state(all_lv, all_imp)
- acute = '持续恶化' in all_imp
- r = dict(state=state, comp=comp_state,
- evidence=(evid if state in ('报警', '危险') else ('INSUFFICIENT' if state == '不可判' else '—')),
- action=action_level(state, evid, acute))
- cache[t] = r
- return r
- @plugin('turbine_summary', '某台结构化结论: 融合级 / L4 过闸线 / L0 可用性 / findings 键', {'turbine': 'WTG16'})
- def turbine_summary(ctx, turbine):
- t = _tid(turbine)
- n = int(t[3:])
- fu = ctx['fusion'][ctx['fusion']['台'] == t]
- l6 = ctx['l6'][ctx['l6']['台'] == f'{n}#']
- l0 = {k.split('|')[1]: v for k, v in ctx['summary'].get('L0', {}).items() if k.startswith(t + '|')}
- fk = data.findings_for_turbine(ctx['findings'], t)
- lines = l6[['测点', '线', 'hz', '域', '绝对量', '单位', 'xfleet', '选择性', '定级']].to_dict('records') if len(l6) else []
- txt = (f'{t}: 融合 {fu.iloc[0]["融合"]} (CMS {fu.iloc[0]["CMS"]} / 模型 {fu.iloc[0]["模型"]}, 机制 {fu.iloc[0]["机制"]}); 模型依据: {fu.iloc[0]["模型依据"]}. '
- + (f'CMS 告警证据力: {fu.iloc[0]["告警证据力"]}. ' if '告警证据力' in fu.columns and str(fu.iloc[0]["告警证据力"]) not in ('', 'nan') else '')
- if len(fu) else f'{t}: 无融合记录. ')
- txt += (f'L4 过闸线 {len(lines)} 条: ' +
- '; '.join(f'{SENSOR_CN.get(r["测点"], r["测点"])} {r["线"]} {r["hz"]}Hz {r["域"]} {r["绝对量"]}{r["单位"]} ×fleet {r["xfleet"]} → {r["定级"]}' for r in lines) +
- '. L0: ' + ', '.join(f'{SENSOR_CN[s]}={l0.get(s, "?")}' for s in SENSORS) + f'. findings 键: {", ".join(fk[:12])}')
- # 冲击轴 (标量面, 融合表以谱线层为主所以这里必须带出来): 末窗同档 ×fleet 最高三项
- st = ctx['status'][(ctx['status'].turbine == t) & ctx['status'].meas_name.isin(['Peak', 'Rms_HP', 'Kurtosis'])].sort_values('x_fleet', ascending=False).head(3)
- 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()]
- flag = any(i['x_fleet'] >= 2.0 for i in imp)
- 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 '】')
- # 收敛报告 v5 逐台终态: 直接解析 §一 表格行 (按台号取整行, 不靠检索)
- v5 = _v5_rows(ctx['cfg'], n)
- if v5:
- txt += ' 【收敛报告 v5 终态 (整行, 不截断——尾部常是限定语): ' + v5[0][:4000] + '】'
- else:
- txt += ' 【收敛报告 v5 无该台专条 (未列入有动作/有争议台)】'
- # 确定性"综合判断"一行放最前: 谱线层 × 冲击轴 × v5 终态 — 模型只转述, 不自己合成
- fused = fu.iloc[0]['融合'] if len(fu) else '无记录'
- if v5:
- verdict = f'以收敛报告 v5 终态为准 (该台有专条); 谱线层融合级 {fused}; 冲击轴{"显著偏离" if flag else "未见显著偏离"}'
- elif flag:
- verdict = f'谱线层融合级 {fused}, 但冲击轴显著偏离 ⇒ 不能判为正常, 状态 = 标量面候选·待 scalar_trend/时域六窗层核 (非谱线层可定)'
- else:
- verdict = f'谱线层融合级 {fused}; 冲击轴未见显著偏离'
- ts = turbine_state(ctx, t)
- comp_txt = ', '.join(f'{k}={v}' for k, v in ts['comp'].items())
- head = f"{t} 状态等级(业主口径): **{ts['state']}** (分部件: {comp_txt}); 证据状态: {ts['evidence']}; 行动等级建议: {ts['action']}。 "
- txt = head + f'综合判断 (结构化层): {verdict}。' + txt.replace('. L0: ', '. L0 各测点有无数据: ')
- 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)
- def _v5_rows(cfg, n):
- """收敛报告 v5 §一 表格里该台的行 (首列含 n#); 返回去掉表格符号的文本列表."""
- import re
- from pathlib import Path
- out = []
- for p in [Path(p) for p in cfg.get('knowledge_docs', []) if '收敛报告' in str(p)]:
- if not p.exists():
- continue
- in_sec = False
- for line in p.read_text(encoding='utf-8').split('\n'):
- if line.startswith('## 一'):
- in_sec = True
- elif line.startswith('## ') and in_sec:
- break
- if in_sec and line.startswith('|'):
- cells = [c.strip() for c in line.strip('|').split('|')]
- if cells and re.search(rf'(^|[^0-9]){n}#', cells[0].replace('*', '')):
- out.append(' | '.join(c.replace('**', '') for c in cells))
- return out
- @plugin('gate_spectral', '对一条候选谱线看最新谱的 G2 有线闸 (峰/邻域本底≥3); 全六闸走生产扫描产物',
- {'turbine': 'WTG16', 'sensor': '发电机DE', 'hz': 135.0, 'kind': 'env'})
- def gate_spectral(ctx, turbine, sensor, hz, kind='env'):
- from .report import pick_spec
- t, sen = _tid(turbine), _sensor(sensor)
- mn = pick_spec(ctx['cfg'], t, sen, kind)
- sp = data.spectrum(ctx['cfg'], t, sen, mn) if mn else None
- if sp is None:
- return dict(kind='text', text='无谱')
- x, y, meta = sp
- hz = float(hz)
- dx = float(x[1] - x[0])
- m = (x >= hz - 2 * dx) & (x <= hz + 2 * dx)
- nb = (np.abs(x - hz) <= 20 * dx) & ~m
- if not m.any():
- return dict(kind='text', text=f'{hz} Hz 超出谱范围 {x.min():.1f}-{x.max():.1f}')
- peak = float(y[m].max())
- base = float(np.median(y[nb]))
- ratio = peak / base if base > 0 else None
- return dict(kind='text', data=dict(peak=peak, baseline=base, peak_over_baseline=ratio, bin_hz=dx, meas=mn),
- text=(f'{t} {SENSOR_CN.get(sen, sen)} {mn} @{hz} Hz: 峰 {peak:.4g}, 邻域本底 {base:.4g}, 峰/本底 {ratio:.2f} '
- f'→ G2 {"过" if ratio and ratio >= 3 else "不过"}.'))
- @plugin('gate_shared_component', '同步两通道共享成分污染门 (库函数 shared_component_gate)',
- {'coh': 0.99, 'group_delay_us': -1.2, 'phase_R': 0.9999})
- def gate_shared(ctx, coh, group_delay_us, phase_R):
- from src.sop.discriminators import shared_component_gate
- r = shared_component_gate(float(coh), float(group_delay_us), float(phase_R))
- return dict(kind='text', data=r, text=f'{r["status"]}: {r["reason"]}')
- @plugin('kb_search', '知识库检索 (收敛报告/裁决/交接/现场单/六层文档/skill/findings)', {'query': '23号 主轴承 时延', 'k': 5})
- def kb_search(ctx, query, k=5):
- from . import knowledge
- hits = knowledge.search(ctx['cfg'], query, k=int(k))
- return dict(kind='table', data=[dict(score=round(h['score'], 2), source=h['source'], text=h['text'][:300]) for h in hits],
- text='\n'.join(f'[{h["source"]}] {h["text"][:240]}' for h in hits))
- @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",
- {'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)"})
- def adhoc(ctx, code):
- ns = dict(scalars=ctx['scalars'], status=ctx['status'], l6=ctx['l6'], fusion=ctx['fusion'], np=np, pd=pd,
- HIGH_BIN=HIGH_BIN, SENSORS=SENSORS, SENSOR_CN=SENSOR_CN, sensor=_sensor, result=None)
- buf = io.StringIO()
- import builtins
- safe = {k: getattr(builtins, k) for k in ['len', 'range', 'min', 'max', 'sum', 'sorted', 'round', 'abs', 'list', 'dict',
- 'float', 'int', 'str', 'print', 'enumerate', 'zip', 'set', 'tuple']}
- try:
- with contextlib.redirect_stdout(buf):
- exec(compile(str(code), '<adhoc>', 'exec'), {'__builtins__': safe}, ns)
- except Exception:
- return dict(kind='text', text='adhoc 失败: ' + traceback.format_exc()[-600:])
- r = ns.get('result')
- if isinstance(r, (pd.DataFrame, pd.Series)):
- rr = r.reset_index() if isinstance(r, pd.Series) else r
- return dict(kind='table', data=rr.head(50).round(4).to_dict('records'), text=(buf.getvalue() + rr.head(20).round(4).to_string())[:2000])
- return dict(kind='text', data=r if isinstance(r, (int, float, str, list, dict)) else str(r), text=(buf.getvalue() + str(r))[:2000])
- @plugin('compare', '多台/多测点并排对比 (同一标量同一功率档, 一图多线)', {'turbines': 'WTG23,WTG16,WTG38', 'sensor': '主轴承前', 'meas': 'Peak', 'bin': HIGH_BIN})
- def compare(ctx, turbines, sensor, meas='Peak', bin=HIGH_BIN, log=None):
- sen = _sensor(sensor)
- tids = [_tid(t) for t in str(turbines).replace(',', ',').split(',') if t.strip()]
- cols = svg.SERIES_VARS # categorical 六槽 (skill 参考实例, 浅暗两套 validator 过检; 色随台固定不随过滤重涂)
- series, summary = [], {}
- for i, t in enumerate(tids):
- q, _ = tcm.trend(ctx['scalars'], ctx['masks'], 'CGN Rudong', t, sen, meas, bin)
- if q.empty:
- continue
- series.append({'name': t, 'x': list(q.trigger_time), 'y': list(q.scalar_value), 'marker': True, 'color': cols[i % len(cols)]})
- summary[t] = q.groupby('window').scalar_value.median().round(4).to_dict()
- ref = tcm.fleet_ref(ctx['scalars'], sen, meas, bin)
- lg = log if log is not None else meas in ('Peak', 'Kurtosis')
- 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 ''
- 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}')
- @plugin('hypothesis', '假设簿: 记录一条假设 (陈述/支持证据/反证检验/当前状态); 状态只能是 候选/INSUFFICIENT/撤回 — 定级须经 verdict_gate',
- {'statement': '23# 前主轴承冲击为转子侧周期载荷调制', 'evidence': 'scalar_trend: …; kb_search: …', 'falsifier': '若对中/螺栓/不平衡现场全正常而内窥见滚道剥落则假', 'status': 'INSUFFICIENT'})
- def hypothesis(ctx, statement, evidence='', falsifier='', status='INSUFFICIENT'):
- if status not in ('候选', 'INSUFFICIENT', '撤回'):
- status = 'INSUFFICIENT'
- h = dict(id=len(ctx.setdefault('hypotheses', [])) + 1, statement=str(statement), evidence=str(evidence), falsifier=str(falsifier), status=status)
- ctx['hypotheses'].append(h)
- return dict(kind='text', data=h, text=f'假设 #{h["id"]} 已记录 [{status}]: {statement} | 证据: {evidence[:300]} | 证伪: {falsifier[:200]}')
- @plugin('verdict_gate', '定级闸 (确定性): 给某台某测点某线按生产扫描产物 (六闸+绝对锚+证据族) 返回六枚举级别; 模型不得自定级',
- {'turbine': 'WTG16', 'sensor': '发电机DE', 'hz': 135.0})
- def verdict_gate(ctx, turbine, sensor, hz=None):
- t, sen = _tid(turbine), _sensor(sensor)
- n = int(t[3:])
- l6 = ctx['l6']
- q = l6[(l6['台'] == f'{n}#') & (l6['测点'] == sen)]
- if hz is not None and len(q):
- q = q[(pd.to_numeric(q['hz'], errors='coerce') - float(hz)).abs() <= 2.0]
- if len(q):
- rows = q[['线', 'hz', '域', '绝对量', '单位', 'xfleet', '选择性', '证据族', '解封判据', '定级']].to_dict('records')
- lv = [r['定级'] for r in rows]
- best = worst_level(lv)
- 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))
- st = ctx['status'][(ctx['status'].turbine == t) & (ctx['status'].sensor_name == sen)]
- top = st.sort_values('x_fleet', ascending=False).head(3)
- imp = impact_typology(ctx, t, sen)
- return dict(kind='text', data=dict(lines=[], scalars=top[['meas_name', 'median', 'x_fleet', 'mask_status']].round(4).to_dict('records'), impact=imp),
- 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())
- + f'. 冲击轴 (时域六窗层, 库函数 vib_scalar_typology, Peak 逐日 max/中位): 分型 {imp["type"]} → 定级 {imp["level"]}; {imp["detail"]}')
- def impact_typology(ctx, t, sen, meas='Peak'):
- """冲击轴分型 → 六枚举: 薄包装, 实现在库层 src/sop/vib_impact.impact_axis_typology (fusion 脚本同调, 单一实现)."""
- from src.sop.vib_impact import impact_axis_typology
- return impact_axis_typology(ctx['scalars'], t, sen, meas, event_windows=ctx['cfg'].get('event_windows'))
- @plugin('render', '组装视图 (模型驱动前端): panels = 插件调用列表 JSON [{"plugin":..,"args":{..}}, ...], 每个面板出图/表/文; 前端画布按顺序渲染',
- {'title': '23# vs 16# 主轴承前', 'panels': '[{"plugin":"compare","args":{"turbines":"WTG23,WTG16","sensor":"主轴承前","meas":"Peak"}},{"plugin":"spectrum","args":{"turbine":"WTG23","sensor":"主轴承前","kind":"env"}}]'})
- def render(ctx, title, panels):
- try:
- spec = json.loads(panels) if isinstance(panels, str) else panels
- except Exception:
- return dict(kind='text', text='panels 不是合法 JSON')
- out = []
- for p in spec[:12]:
- name, args = p.get('plugin'), p.get('args') or {}
- if name in ('render',):
- continue
- r = run(ctx, name, **args)
- 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)))
- view = dict(id=len(ctx.setdefault('views', [])) + 1, title=str(title), panels=out)
- ctx['views'].append(view)
- try:
- vp = ctx['cfg']['out'] / 'views.json'
- # 读写都要显式 utf-8: 面板文本是中文, 缺 encoding 会按系统 locale(cp936) 读/写 → 乱码或崩
- old = json.load(open(vp, encoding='utf-8')) if vp.exists() else []
- with open(vp, 'w', encoding='utf-8') as f:
- 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)
- except Exception:
- pass
- return dict(kind='view', data=view, text=f'视图 "{title}" 已组装 {len(out)} 个面板: ' + '; '.join(f'{x["plugin"]}: {x["text"][:120]}' for x in out))
- @plugin('analyze', '原始数据 → 报告 全链 (后台任务, 约 8 分钟, 重写全部模型产物): 必须 confirm=true 才会启动; 同一时刻只允许一个; 可带 input (TCM 导出包 / 解码目录 / 波形目录) 与 window 名; 不带则在现有窗上重跑; 用 job_status 看进度',
- {'input': '', 'window': '', 'steps': '', 'confirm': False})
- def analyze(ctx, input='', window='', steps='', confirm=False):
- import subprocess, sys, os
- from .config import ROOT
- if not (confirm is True or str(confirm).lower() in ('true', '1', 'yes')):
- return dict(kind='text', text='analyze 未启动: 这是约 8 分钟、重写全部模型产物的后台任务, 需要 confirm=true (由用户明确要求时才传). 当前产物可直接用其他插件查询.')
- lock = ctx['cfg']['out'] / 'analyze.lock'
- if lock.exists():
- try:
- pid = int(lock.read_text(encoding='utf-8').strip())
- os.kill(pid, 0)
- return dict(kind='text', text=f'analyze 已在运行 (pid {pid}), 不重复启动; 用 job_status 看进度.')
- except Exception:
- lock.unlink(missing_ok=True)
- if input and not os.path.exists(os.path.expanduser(str(input))):
- return dict(kind='text', text=f'导入失败: 路径不存在 {input!r} — 请确认盘已挂载、路径完整 (可整段从访达拖入); 支持 TCM 导出包 / 解码目录 / 波形目录.')
- # 冻结版 sys.executable = 打包二进制, spawn 分析链必挂 → 优先 WINDCMS_SCRIPT_PY (launch.sh 提供, 默认 venv)
- py = os.environ.get('WINDCMS_SCRIPT_PY') or sys.executable
- cmd = [py, str(ROOT / 'scripts/windcms.py'), 'analyze', '--farm', 'rudong']
- if input:
- cmd += ['--input', os.path.expanduser(str(input)), '--window', str(window or 'w_new')]
- if steps:
- cmd += ['--steps'] + [s for s in str(steps).replace(',', ' ').split() if s]
- logp = ctx['cfg']['out'] / 'analyze_stdout.log'
- ctx['cfg']['out'].mkdir(parents=True, exist_ok=True)
- # stdout= 只用到这个文件的 fd (内容由子进程自己写), 所以这里不 with-close:
- # 句柄留着, 免得父进程提前关掉; 显式 encoding 是为了不留"缺 encoding"的门禁告警。
- # ★2026-09-16: 本函数跑在 **CMS 服务进程**里 (由 guanlan.py 以无窗口方式拉起, 自己没有可见控制台),
- # 裸 Popen 会让 Windows 给全链分析**新建一个可见控制台窗口** (从页面点"分析"就会看到黑窗)。
- # 统一走 src.proc.spawn (CREATE_NO_WINDOW)。
- from src import proc as _proc
- _logf = open(logp, 'w', encoding='utf-8')
- p = _proc.spawn(cmd, cwd=ROOT, stdout=_logf, stderr=subprocess.STDOUT)
- lock.write_text(str(p.pid), encoding='utf-8')
- return dict(kind='text', data=dict(pid=p.pid, log=str(logp)), text=f'全链已在后台启动 (pid {p.pid}); 进度看 job_status; 日志 {logp}. 完成后产物/报告/知识库自动刷新 (界面需刷新页面).')
- @plugin('farm_overview', '全场38台业主四级总览 (报警/不可判在前): 台/整机/四部件/证据/行动/焦点 — 全场类问题先调本工具保证完整性', {})
- def farm_overview(ctx):
- from .report_std import registry, STATE_ORDER
- reg, detail, _ = registry(ctx['cfg'])
- reg = reg.sort_values('状态等级', key=lambda s_: s_.map({x: i for i, x in enumerate(STATE_ORDER)}))
- lines = [f"{r['机组']} {r['状态等级']} | 前{r['主轴承前']} 后{r['主轴承后']} 齿{r['齿轮箱']} 发{r['发电机']} | 证据{r['证据状态']} 行动{r['行动等级建议']}" + (f" | 焦点 {r['状态焦点']}" if r['状态焦点'] != '—' else '')
- for _, r in reg.iterrows()]
- dist = reg['状态等级'].value_counts().to_dict()
- return dict(kind='text', data=reg.to_dict('records'), text=f'分布 {dist}\n' + '\n'.join(lines))
- @plugin('job_status', '查看后台 analyze 任务进度与上次结果 (analyze_log.json)', {})
- def job_status(ctx):
- out = ctx['cfg']['out']
- # ★ 2026-09-17 用户令 2: analyze 的运行日志与结果都改落 logs/build/ (原先落在产物仓 out/ 下);
- # 读侧同步改, 并保留旧路径兜底 —— 老机器上可能还留着那时的件, 不该一改就显示"无运行日志"。
- from src import logfile as _lf
- _b = _lf.audit_log('windcms_analyze.json').parent # logs/audit/
- logp = (_b / 'analyze_stdout.log') if (_b / 'analyze_stdout.log').exists() else (out / 'analyze_stdout.log')
- lp = (_b / 'analyze_log.json') if (_b / 'analyze_log.json').exists() else (out / 'analyze_log.json')
- # 日志与结果都是中文文本: 不写 encoding 会按系统 locale(cp936) 读 → 乱码/解码失败
- tail = open(logp, encoding='utf-8').read()[-1200:] if logp.exists() else '(无运行日志)'
- last = json.load(open(lp, encoding='utf-8')) if lp.exists() else None
- 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 '尚无完成记录'
- return dict(kind='text', data=last, text=f'{summ}\n--- 实时日志尾 ---\n{tail}')
- def context(cfg):
- s = data.load_scalars(cfg)
- masks = data.load_masks(cfg)
- l6, fusion, summ = data.load_model(cfg)
- return dict(cfg=cfg, scalars=s, masks=masks, l6=l6, fusion=fusion, summary=summ,
- findings=data.load_findings(cfg), status=tcm.status_table(s, masks, 'CGN Rudong'))
- def _cn_sensors(text):
- """工具输出正文里的英文通道名/谱种标签 → 中文 (模型会原样往答案里抄, 源头治理; data(JSON) 保留原值供机器用)."""
- for k, v in SENSOR_CN.items():
- text = text.replace(k, v)
- for a, b in [('Indicator_1P', '1倍频指标'), ('Indicator_2P', '2倍频指标'), ('Indicator_HSW', '高速轴指标'), ('Indicator', '厂商指标')]:
- text = text.replace(a, b)
- text = re.sub(r'\benv(\d+)\b', r'包络·\1', text)
- text = re.sub(r'\benv\b', '包络谱', text)
- text = re.sub(r'\braw\b', '原始谱', text)
- return text.replace('(包络域)', '').replace('(原始域)', '')
- def run(ctx, name, **kw):
- from .redact import redact_obj, redact
- if name not in REGISTRY:
- return dict(kind='text', text=f'未知插件 {name}; 可用: {list(REGISTRY)}')
- try:
- r = REGISTRY[name]['run'](ctx, **kw)
- except Exception:
- return dict(kind='text', text=f'插件 {name} 失败: ' + traceback.format_exc()[-800:])
- # 公开模式: 文本/数据里的判据阈值脱敏 (svg 内的数值标签由 svg 层按 PUBLIC 处理)
- if 'text' in r:
- r['text'] = redact(_cn_sensors(str(r['text'])))
- 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'):
- r['data'] = redact_obj(r['data'])
- return r
- @plugin('explain', '机器语言→交互语言: 把任意插件的结构化输出翻译成给业主/现场看的自然汉语 (经接地闸: 台号与数字必须来自原文, 不过闸则退回原文); source=插件名, args=其参数 JSON',
- {'source': 'turbine_summary', 'args': '{"turbine":"WTG38"}', 'audience': '业主'})
- def explain(ctx, source, args='{}', audience='业主'):
- from . import llm
- try:
- a = json.loads(args) if isinstance(args, str) else (args or {})
- except Exception:
- return dict(kind='text', text='args 不是合法 JSON')
- r = run(ctx, source, **a) # 经 run ⇒ 公开模式已脱敏
- raw = str(r.get('text', ''))
- if raw.startswith(('未知插件', '插件')) or not raw:
- return r
- prompt = (f'把下面这段风机振动诊断系统的结构化输出, 改写成给{audience}看的自然中文 (2-5 句):\n'
- '要求: 先结论后数字; 每个数字保留并注明含义; 不新增任何数字/台号/结论; 不解释判据算法; '
- '缩写展开 (如 "候选" 说成 "候选级异常, 需现场核实确认"); 语气平实, 不用套话。\n\n=== 原文 ===\n' + raw[:3500])
- txt, m = llm.generate(prompt, temperature=0.2, num_predict=400, timeout=120)
- if txt:
- ok, bad = llm.grounding(txt, raw)
- if ok:
- return dict(kind='text', data=dict(source=source, model=m, grounded=True),
- text=txt.strip() + '\n\n——以上为通俗转述 (经接地闸核对); 结构化原文按需查 ' + source + ' 插件。')
- return dict(kind='text', data=dict(source=source, grounded=False),
- text=raw + '\n\n(通俗转述未过接地闸或模型不可用, 保留结构化原文)')
- @plugin('export', '导出: 把任意插件的结果 (或整表) 存成 电子表格(xlsx) / Word(docx) / CSV, 落 out/exports/; source = 插件名, args = 该插件参数 JSON, fmt = xlsx|docx|csv',
- {'source': 'fleet_rank', 'args': '{"sensor":"主轴承后","meas":"Peak"}', 'fmt': 'xlsx', 'title': '主轴承后 Peak 全场排名'})
- def export(ctx, source, args='{}', fmt='xlsx', title=''):
- """★脱敏在出口生效: 走 run() 拿结果 ⇒ 公开模式下导出件同样不含阈值/特征频率 (client-deliverable-hide-core-parameters)."""
- import datetime
- try:
- a = json.loads(args) if isinstance(args, str) else (args or {})
- except Exception:
- return dict(kind='text', text='args 不是合法 JSON')
- r = run(ctx, source, **a) # 经 run ⇒ 已脱敏
- if str(r.get('text', '')).startswith(('未知插件', '插件')):
- return r
- data = r.get('data')
- rows = data if isinstance(data, list) else ([data] if isinstance(data, dict) else [])
- df = pd.DataFrame(rows) if rows else pd.DataFrame()
- ttl = title or f'{source} {json.dumps(a, ensure_ascii=False)}'
- out = ctx['cfg']['out'] / 'exports'
- out.mkdir(parents=True, exist_ok=True)
- stamp = datetime.datetime.now().strftime('%Y%m%d_%H%M%S')
- base = re.sub(r'[^\w一-龥-]+', '_', f'{source}_{stamp}')[:80]
- fmt = (fmt or 'xlsx').lower()
- if fmt == 'csv':
- p = out / f'{base}.csv'
- df.to_csv(p, index=False, encoding='utf-8-sig')
- elif fmt == 'docx':
- from docx import Document
- from docx.shared import Pt
- doc = Document()
- doc.styles['Normal'].font.name = 'PingFang SC'
- doc.styles['Normal'].font.size = Pt(10.5)
- doc.add_heading(ttl, level=1)
- doc.add_paragraph(str(r.get('text', ''))[:4000])
- if len(df):
- tb = doc.add_table(rows=1, cols=len(df.columns)); tb.style = 'Table Grid'
- for i, c in enumerate(df.columns):
- tb.rows[0].cells[i].text = str(c)
- for _, rr in df.iterrows():
- cells = tb.add_row().cells
- for i, c in enumerate(df.columns):
- cells[i].text = '' if pd.isna(rr[c]) else str(rr[c])
- doc.add_paragraph(f'\n生成 {stamp} · 来源插件 {source} · 参数 {json.dumps(a, ensure_ascii=False)}'
- + ('\n(公开模式:判据门槛与特征频率已隐去)' if __import__('src.windcms.redact', fromlist=['PUBLIC']).PUBLIC else ''))
- p = out / f'{base}.docx'
- doc.save(p)
- else:
- p = out / f'{base}.xlsx'
- with pd.ExcelWriter(p, engine='openpyxl') as w:
- (df if len(df) else pd.DataFrame({'说明': [str(r.get('text', ''))[:3000]]})).to_excel(w, index=False, sheet_name='数据')
- pd.DataFrame({'项': ['标题', '来源插件', '参数', '生成时间', '结论文本'],
- '值': [ttl, source, json.dumps(a, ensure_ascii=False), stamp, str(r.get('text', ''))[:3000]]}).to_excel(w, index=False, sheet_name='说明')
- return dict(kind='text', data=dict(path=str(p), rows=len(df), fmt=fmt),
- text=f'已导出 {fmt}: {p.name}({len(df)} 行)→ {p.parent}')
- def manifest():
- from .redact import redact
- return [dict(name=v['name'], desc=redact(v['desc']), params=v['params']) for v in REGISTRY.values()]
|