# -*- 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'.*?', '', 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()