| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434 |
- # -*- coding: utf-8 -*-
- """问答微服务 — 只暴露 /api/ask 与 /api/ask_status。
- 为什么不直接把 windscada_serve.py 搬上云:
- 那份要 numpy/pandas + src.windscada 下 taxonomy/subsys/perf/design 一整套, 而问答只用到
- fast_agent → mcp_server(pandas) + audit.grounding(re)。云上少装一个包就少一处装不上的风险,
- 2 核机器编译依赖的代价也是实打实的。静态页由 nginx 直接服务, 与本进程无关。
- 模型档: deepseek-chat / qwen-plus(阿里云百炼免费额度) / kimi-fast —— 见 fast_agent._online,
- 密钥一律从环境变量读, 不落盘。哪个档没配密钥, 前端就看不到它。
- """
- import json, os, pathlib, re, sys, threading, time
- import urllib.parse
- from http.server import ThreadingHTTPServer, BaseHTTPRequestHandler
- sys.path.insert(0, str(pathlib.Path(__file__).resolve().parent))
- sys.path.insert(0, str(pathlib.Path(__file__).resolve().parent.parent))
- from src.ontology import fast_agent as F # noqa: E402
- MODELS = {
- 'deepseek-chat': dict(nm='DeepSeek', loc='云端', t='约 15~40 秒', env='DEEPSEEK_API_KEY'),
- 'qwen-plus': dict(nm='通义千问 Plus', loc='云端', t='约 15~40 秒', env='DASHSCOPE_API_KEY'),
- 'kimi-fast': dict(nm='Kimi', loc='云端', t='约 20~50 秒', env='MOONSHOT_API_KEY'),
- }
- # ── 汇报编排 (2026-09-01 移植上云) ──────────────────────────────────
- # ★为什么能移植: rpt_compose 唯一的重依赖是 fleet_view() —— 那是整套分析栈
- # (taxonomy/subsys/perf + 全部 parquet), 2 核 1G 跑不动, 当初正因此没上云。
- # 但 compose **只要它的输出不要它的计算**, 而输出只有 170~186 KB。
- # 所以: 把 fleet_view 的结果 dump 成 JSON 传上来, compose 读文件。
- # ★模型不再打 Ollama: 原实现写死 http://127.0.0.1:11434 (本地 Ollama), 云上没有。
- # 改走 fast_agent._chat —— 与问答助手同一个云端客户端, 不重写。
- RPT_BLOCK_KEYS = ['kpi', 'water', 'emon', 'blame', 'ram', 'parts', 'attrib',
- 'pareto', 'quad', 'frate', 'units', 'chain', 'loop', 'plan', 'caliber']
- RPT_SYS = ('你是风电场报告编排助手。根据用户要求做两件事, 严格按格式输出, 不要解释:\n'
- 'BLOCKS: 用逗号分隔的块名 (只能从给定清单里选, 按报告顺序排)\n'
- 'TEXT: 一段中文总结 (150~300 字, 可分 2~3 句一组)\n'
- '铁律: TEXT 里出现的每个数字都必须来自"可用数据"里, 不许自己算不许估计; '
- '没有依据的话就不写。不要写套话, 不要写"综上所述"。')
- RPT_SYS_EN = ('You compose wind-farm reports. Do two things and output strictly in this '
- 'format, with no explanation:\n'
- 'BLOCKS: comma-separated block names (only from the given list, in report order)\n'
- 'TEXT: one English summary of 150-300 words\n'
- 'Hard rule: every number in TEXT must come from the supplied data — do not compute '
- 'or estimate any number yourself; if there is no basis, leave it out. No filler.')
- FLEET_DIR = pathlib.Path(__file__).resolve().parent.parent / 'data' / 'fleet'
- def _rpt_facts(fl):
- E = (fl.get('control') or {}).get('energy') or {}
- R = fl.get('rel') or {}
- M = ((fl.get('m8') or {}).get('mtbf')) or {}
- F = fl.get('faults') or {}
- BD = ((fl.get('fus') or {}).get('链盘')) or ((fl.get('fus') or {}).get('chain')) or {}
- return [
- f"统计期 {fl.get('win')}; 38 台 4.0MW",
- f"上网电量 {E.get('act')} MWh; 理论可发 {E.get('theo')}; 损失 {E.get('loss')} MWh; 损失率 {E.get('loss_pct')}%",
- f"等效满发 {E.get('eflh')} h; 时间可用率 {(fl.get('control') or {}).get('avail')}%",
- f"总停机 {R.get('总停机h')} h = 设备类 {R.get('设备类停机h')} + 外部类 {R.get('外部类停机h')}",
- f"平均停机间隔 {M.get('MTBF_h')} h; 停机 {M.get('停机事件数')} 起; 单次平均 {M.get('MDT_h')} h",
- f"报警 {F.get('total')} 条 / {F.get('n_codes')} 个码",
- f"需行动机组 {len(BD.get('rows') or [])} 台; 卡点 {BD.get('stuck')}",
- # 占比直接给出: 不给模型就得自己算, 算对了也会被接地闸判为"未接地数字"
- "损失构成: " + '; '.join(
- f"{x.get('k')} {x.get('loss')} MWh (占损失 "
- f"{round((x.get('loss') or 0) / max(E.get('loss') or 1, 1) * 100, 1)}%)"
- for x in (E.get('items') or [])),
- f"端到端闭环 {len(BD.get('closed') or [])} 例",
- ]
- def rpt_compose(need, canvas, model, lang='zh'):
- import time as _t, re as _re
- f = FLEET_DIR / ('fleet_%s.json' % ('en' if lang == 'en' else 'zh'))
- if not f.is_file():
- return dict(err='fleet 数据未上传 (%s)' % f.name if lang != 'en'
- else 'Fleet data not uploaded (%s)' % f.name)
- fl = json.loads(f.read_text(encoding='utf-8'))
- facts = _rpt_facts(fl)
- prompt = ('可用数据:\n' + '\n'.join('- ' + x for x in facts) +
- '\n\n可选块 (只能从这里选): ' + ', '.join(RPT_BLOCK_KEYS) +
- '\n当前画布: ' + (', '.join(canvas) if canvas else '(空)') +
- '\n用户要求: ' + need)
- t0 = _t.time()
- try:
- from src.ontology import fast_agent as _F
- msgs = [dict(role='system', content=RPT_SYS_EN if lang == 'en' else RPT_SYS),
- dict(role='user', content=prompt)]
- txt = _F._chat(model, msgs, None, npred=900) or ''
- if isinstance(txt, dict):
- txt = txt.get('content') or txt.get('text') or ''
- except Exception as e:
- return dict(err='%s: %s' % (type(e).__name__, str(e)[:140]))
- txt = _re.sub(r'<think>.*?</think>', '', str(txt), flags=_re.S).strip()
- mb = _re.search(r'BLOCKS\s*[::]\s*(.+)', txt)
- mt = _re.search(r'TEXT\s*[::]\s*(.+)', txt, _re.S)
- blocks = [b.strip() for b in _re.split(r'[,,\s]+', mb.group(1))] if mb else []
- blocks = [b for b in blocks if b in RPT_BLOCK_KEYS]
- narrative = (mt.group(1).strip() if mt else '').strip()
- bad = []
- try: # 接地闸: 叙述里的数字必须在 facts 里出现
- from src.windscada.audit import grounding as _gr
- ok, b = _gr(narrative, '\n'.join(facts))
- bad = (b or {}).get('numbers') or []
- except Exception:
- pass
- return dict(blocks=blocks, narrative=narrative, model=model,
- secs=round(_t.time() - t0, 1), ungrounded=bad, raw=txt[:900])
- # ── 本体实时查询路由 ────────────────────────────────────────────────
- # 只收**纯读**的几个: 抢修链 / 故障树 / 预防链 / 对象取值 / 检索。
- # 不收写操作, 也不收需要重算的重接口 —— 这台是 2 核 1G, 边界要自己划清。
- _ONT = {'/api/ont_chain', '/api/ont_obj', '/api/ont_query', '/api/rpt_compose'}
- # ── 公网出口统一脱敏 ────────────────────────────────────────────────
- # 本服务在 nginx `/api/` 后面**裸奔**: 任何人 curl 即可调, 不需要任何凭证。
- # 而上云闸 (windscada_cloud_pack) 只管住**烤进 HTML 的那份** —— 页面的"实时优先"
- # 策略正好从闸旁边绕过去。2026-09-03 实逮:
- # curl http://<公网IP>/api/ont_chain?kind=fault&code=3225
- # → OEM 实名「上海电气」· 机型号 SWT-4.0-130 · 业主资料目录全路径
- # 「4.0MW各部件资料/SB540-A06 液压盘式制动器用户手册.pdf」
- # ⇒ 出口必须自己带闸, 不能指望调用方脱敏。
- # ★导入失败 = 拒绝启动。安全闸不可用时放行, 比没有闸更坏 (看着有闸, 实则没有)。
- from src.windscada import deid_public as _dp
- def _scrub_deep(o):
- """递归脱敏 value (不动 key —— key 是结构, 换了下游取不到)。"""
- if isinstance(o, str):
- return _dp.scrub(o)
- if isinstance(o, list):
- return [_scrub_deep(x) for x in o]
- if isinstance(o, dict):
- return {k: _scrub_deep(v) for k, v in o.items()}
- return o
- def _ont_call(path, q, lang='zh'):
- """按 path+参数派发到 mcp_server; 语言投影与主服务同口径。"""
- from src.ontology import mcp_server as _m
- if path == '/api/rpt_compose':
- if q.get('probe'):
- return dict(ok=True) # 前端探活: 返 200 即表示本部署有编排能力
- return rpt_compose(q.get('need', ''),
- [x for x in (q.get('canvas') or '').split(',') if x],
- q.get('model') or next(iter(_avail()), ''), lang)
- if path == '/api/ont_chain':
- kind = q.get('kind', 'fault')
- if kind == 'fault':
- r = _m.ont_chain_fault(code=q.get('code', ''))
- elif kind == 'tree':
- r = _m.ont_fault_tree(system=q.get('system', ''))
- elif kind == 'prev':
- r = _m.ont_chain_preventive(system=q.get('system', ''))
- elif kind == 'plan':
- r = _m.ont_maint_plan(top=int(q.get('top', 12)))
- else:
- return dict(err='unknown kind %r' % kind)
- elif path == '/api/ont_obj':
- r = _m.ont_get(id=q.get('id', ''))
- else:
- r = _m.ont_query(**{k: v for k, v in q.items() if k in ('type', 'q', 'limit')})
- if lang and lang != 'zh':
- try:
- from src.windscada import ui_en as _ue # 数据层英文投影
- r = _json_walk(r, _ue)
- except Exception:
- pass
- return r
- def _json_walk(o, ue):
- """英文投影。**先做 key_en → key 的字段投影, 再做字符串替换。**
- ★只做字符串替换是不够的: 本体里工单等对象**自带 name_en / action_en**, 那是
- 权威译名 (来自厂商对译表与构词表), 比任何事后翻译都准。首版漏了这一步, 于是
- 工单名一路是中文 ("叶片维修"), 而 name_en 明明就在同一个 dict 里。
- """
- if isinstance(o, dict):
- # ① 字段投影: 有 x_en 就用它顶掉 x, 并丢掉 _zh 副本 (与主服务 _proj 同口径)
- out = {}
- # ★只有 k_en **有值**时才让裸 k 让位。首版无条件让位, 于是 name_en 为空的行
- # 直接显示成 "None" —— 比留中文糟得多。空就回落原文, 再走字符串替换。
- langed = {k[:-3] for k in o if k.endswith('_en') and o[k] not in (None, '', [])}
- for k, v in o.items():
- if k.endswith('_zh'):
- continue
- if k.endswith('_en'):
- if v not in (None, '', []):
- out[k[:-3]] = _json_walk(v, ue)
- continue
- if k in langed: # 已有非空 k_en, 裸 k 让位
- continue
- out[k] = _json_walk(v, ue)
- return out
- return _json_walk_str(o, ue)
- def _json_walk_str(o, ue):
- """字符串与列表的英文替换; 与主服务 _en_deep 同语义。"""
- if isinstance(o, str):
- if not any('\u4e00' <= c <= '\u9fff' for c in o):
- return o
- # 与主服务 _data_text_en 同顺序: 显示层值表 → 数据层整串/模板 → 短语层
- v = ue.VALUES.get(o)
- if v:
- return v
- try:
- from src.windscada import data_tpl_en as _t
- v = _t.NOTES.get(o) or _t.render(o)
- if v:
- return v
- except Exception:
- pass
- return ue.apply(o)
- if isinstance(o, list):
- return [_json_walk(x, ue) for x in o]
- return o
- ASK = dict(state='idle', q='', model='', answer='', err='', t0=0, step='', steps=[], verify=None, lang='zh')
- LOCK = threading.Lock()
- def _avail():
- return {k: v for k, v in MODELS.items() if os.environ.get(v['env'], '').strip()}
- _CJK_RE = __import__('re').compile(r'[\u3400-\u9fff\u3000-\u303f\uff00-\uffef]')
- _REF_RE = __import__('re').compile(r'\[\s*(?:object id|source|id|对象id|对象|来源|出处)?\s*[::]?\s*([a-z]+/[A-Za-z0-9_/.\-]+)\s*\]', __import__('re').I)
- _STEP_KEY_EN = {
- '启动': 'Starting',
- '预取': 'Prefetch',
- '推理': 'Reasoning',
- '取证': 'Evidence retrieval',
- '成文': 'Drafting',
- '校闸': 'Grounding gate',
- '交叉审核': 'Cross-review',
- '完成': 'Done',
- '失败': 'Failed',
- }
- def _verify_answer(ans, lang):
- try:
- from ontology_p2_verify import verify as _verify
- return _verify(ans or '', lang='en' if lang == 'en' else 'zh')
- except Exception as e:
- refs = list(dict.fromkeys(_REF_RE.findall(ans or '')))
- return dict(引用id数=len(refs), 全通过=False, 引用了判级对象=False, 接地闸出现=False,
- 明细=[], err=str(e)[:120])
- def _force_english_answer(ans, model):
- """Last-mile English QA: rewrite only if the model leaked CJK into the public answer."""
- if not _CJK_RE.search(ans or ''):
- return ans
- refs0 = set(_REF_RE.findall(ans or ''))
- msgs = [
- dict(role='system', content=(
- 'Rewrite the supplied wind-turbine diagnostic answer for a public English page. '
- 'Output English only. Do not output Chinese, Japanese or Korean characters. '
- 'Preserve every square-bracket object id, unit id, code number, date and numeric value exactly. '
- 'Do not add new claims.')),
- dict(role='user', content=ans or ''),
- ]
- txt = F._chat(model, msgs, tools=[], npred=2600) or {}
- if isinstance(txt, dict):
- txt = txt.get('content') or txt.get('text') or ''
- txt = str(txt).strip()
- refs1 = set(_REF_RE.findall(txt))
- if txt and not _CJK_RE.search(txt) and refs0.issubset(refs1):
- return txt
- return ans
- def _public_step_text(k, v='', lang='zh'):
- """Render live model progress for the selected UI language.
- The public English page polls /api/ask_status while the model is running.
- Those short progress labels used to bypass answer-level translation and
- surfaced Chinese such as "推理 · 第 2/5 轮". Keep this helper deliberately
- small and conservative: it translates known progress labels and replaces
- unknown Chinese detail with a neutral English progress marker.
- """
- k = '' if k is None else str(k)
- v = '' if v is None else str(v)
- if lang != 'en':
- return f'{k} · {v}' if v else k
- key = _STEP_KEY_EN.get(k, k)
- detail = v.strip()
- m = re.search(r'第\s*(\d+\s*/\s*\d+)\s*轮', detail)
- if m:
- detail = 'round ' + m.group(1).replace(' ', '')
- else:
- m = re.search(r'(\d{1,2})\s*一站式证据包', detail)
- if m:
- alias = 'focus unit'
- try:
- from src.windscada.deid import unitize_answer
- alias = unitize_answer(f"WTG{int(m.group(1)):02d}")
- except Exception:
- pass
- detail = f'{alias} evidence pack'
- elif '闭环' in detail and '证据' in detail:
- n = ''.join(re.findall(r'\d+', detail))
- detail = f'closure and trend evidence: {n} records' if n else 'closure and trend evidence'
- elif '全场' in detail and '判级' in detail:
- n = ''.join(re.findall(r'\d+', detail))
- detail = f'fleet ratings: {n} units' if n else 'fleet ratings'
- elif '换模型' in detail or '复核' in detail:
- detail = 'reviewing with another model'
- elif _CJK_RE.search(detail):
- detail = ''
- if _CJK_RE.search(key):
- key = 'Working'
- if _CJK_RE.search(detail):
- detail = ''
- return f'{key} · {detail}' if detail else key
- def _worker(q, model, lang='zh'):
- def on_step(k, v):
- text = _public_step_text(k, v, lang)
- ASK['step'] = text
- ASK['steps'].append(text)
- try:
- if lang == 'en':
- try:
- from src.windscada.deid import ununitize_question
- q = ununitize_question(q)
- except Exception as e:
- print(f'[ask] 代号反解失败: {e}', flush=True)
- # 语言指令走 F.answer 的 lang 参数 (进 system 消息); 问题本身不加前缀 ——
- # 前缀在用户侧压不过中文 SYSTEM, 实测无效。
- r = F.answer(q, model=model, max_rounds=5, on_step=on_step, lang=lang)
- ans = r.get('answer') or r.get('text') or ''
- # 返回前换成代号: 页面显示 A56, 答案里也必须是 A56。只在**出口**换 ——
- # 发给模型的消息必须保留原始台号, 否则它调检索工具查不到东西。
- try:
- from src.windscada.deid import unitize_answer
- ans = unitize_answer(ans)
- except Exception as e:
- print(f'[ask] 代号映射失败: {e}', flush=True)
- if lang == 'en' and _CJK_RE.search(ans):
- ASK['step'] = 'English QA · rewriting CJK leakage'
- ans = _force_english_answer(ans, model)
- ver = _verify_answer(ans, lang)
- if lang == 'en' and _CJK_RE.search(ans):
- ASK.update(state='error', err='English QA failed: the model returned non-English text, so the answer was withheld.',
- answer='', verify=ver, grounded=r.get('grounded'), step='English QA failed')
- return
- ASK.update(state='done', answer=ans, verify=ver, grounded=r.get('grounded'), step='完成' if lang != 'en' else 'Done')
- except Exception as e:
- ASK.update(state='error', err=f'{type(e).__name__}: {str(e)[:200]}', step='失败')
- class H(BaseHTTPRequestHandler):
- def _send(self, obj, code=200):
- b = json.dumps(_scrub_deep(obj), ensure_ascii=False).encode() # 公网出口: 一律过脱敏
- self.send_response(code)
- self.send_header('Content-Type', 'application/json; charset=utf-8')
- self.send_header('Content-Length', str(len(b)))
- self.end_headers()
- self.wfile.write(b)
- def log_message(self, *a):
- pass # 默认日志把每条请求打到 stderr, systemd 里没意义
- def do_GET(self):
- u = urllib.parse.urlparse(self.path)
- if u.path == '/api/ask_status':
- d = dict(ASK)
- # ★单槽位设计 (2 核机器, 一次一问) 的副作用: 上一问的答案**不分语言**
- # 两个页面都看得到 —— 中文页上会显示英文测试留下的答案 (2026-09-01 中文版
- # 点检活逮)。这里按语言隔离: 语言不匹配就当没问过, 不把别人的答案端给你。
- # 不改单槽位本身 —— 那是为 2 核机器有意选的, 改成多槽会互相拖到都超时。
- _want = (urllib.parse.parse_qs(u.query).get('lang') or ['zh'])[0]
- if d.get('answer') and (d.get('lang') or 'zh') != _want:
- d.update(state='idle', q='', answer='', err='', step='', steps=[], verify=None)
- d['elapsed'] = int(time.time() - ASK['t0']) if ASK['t0'] else 0
- d['models'] = _avail()
- self._send(d)
- elif u.path == '/api/ask_models':
- self._send(dict(models=_avail()))
- elif u.path in _ONT:
- # ★本体实时查询 (2026-09-01)。此前这些只走**快照预烤**, 而预烤的报警码
- # 被 limit=60 截断 —— 全场 555 个码只烤了按码号排序的前 60 个, 页面上
- # "点一下就查"的示例码 (3225/64000/64001/4110/13121) 恰好都不在里面
- # ⇒ 抢修链点下去整条没反应。接成实时查就不再受"烤了多少"限制。
- # 纯读 objects.json + 已装的 pandas, 2 核 1G 扛得住 (实测 64001 秒级返回)。
- q = {k: v[0] for k, v in urllib.parse.parse_qs(u.query).items()}
- lang = q.pop('lang', 'zh')
- try:
- self._send(_ont_call(u.path, q, lang))
- except Exception as e:
- self._send(dict(err='%s: %s' % (type(e).__name__, str(e)[:160])), 500)
- else:
- self._send(dict(err='not found'), 404)
- def do_POST(self):
- if urllib.parse.urlparse(self.path).path != '/api/ask':
- return self._send(dict(err='not found'), 404)
- n = int(self.headers.get('Content-Length') or 0)
- body = json.loads(self.rfile.read(n) or b'{}')
- q = (body.get('q') or '').strip()
- model = body.get('model') or next(iter(_avail()), '')
- _qs = urllib.parse.parse_qs(urllib.parse.urlparse(self.path).query)
- lang = (body.get('lang') or (_qs.get('lang') or [''])[0] or 'zh').strip()
- if not q:
- return self._send(dict(err='问题为空' if lang != 'en' else 'The question is empty'))
- if model not in _avail():
- msg = (f'Model {model} is not configured; available: {list(_avail())}' if lang == 'en'
- else f'模型 {model} 未配置密钥; 可用 {list(_avail())}')
- return self._send(dict(err=msg))
- with LOCK: # 串行锁: 2 核机器, 并发推理会互相拖慢到都超时
- if ASK['state'] == 'running':
- return self._send(dict(err='已有问题在推理中 (一次一问)', state='running'))
- ASK.update(state='running', q=body.get('display_q') or q, model=model, answer='', err='',
- t0=time.time(), step='启动' if lang != 'en' else 'Starting', steps=[], lang=lang, verify=None)
- threading.Thread(target=_worker, args=(q, model, lang), daemon=True).start()
- self._send(dict(state='running', model=model))
- if __name__ == '__main__':
- port = int(sys.argv[1]) if len(sys.argv) > 1 else 8033
- av = _avail()
- if not av:
- print('[ask] 警告: 没有任何模型配了密钥, /api/ask 会全部拒绝', file=sys.stderr)
- print(f'[ask] :{port} 可用模型 {list(av)}', file=sys.stderr, flush=True)
- ThreadingHTTPServer(('127.0.0.1', port), H).serve_forever()
|