# -*- coding: utf-8 -*- r"""问答管线(P12):从详情层 `serve.py` 搬来的 `_ask_worker`/`_ask_step_text` 与相关常量。 ★搬运改动(仅三处,其余逐字保留): 1. `_ask_worker` 加参数 `st`(详情层是"线程 + 全局 ASK",这里改"同步 + 传入状态字典"); 2. 体内 `ASK['x']` → `st['x']`(`_on_step` 在 worker 内定义,闭包自然拿到同一个 st); 3. 对外只暴露 `ask_run(q, model, cat, xrev, lang)`:建 st、按 `ask_start` 的 hints 逻辑拼词、跑管线、收尾 state。 """ from __future__ import annotations import pathlib import sys from typing import Any from app_common.app_common_guanlan.api import paths as _P ASK_CAT_HINTS = { '设备类': '(设备状态类: 从本体的机组档案、对象查询与机制检索三条路取证)', '检修排程决策': '(检修排程类: 用 ont_scenario/ont_get 取排期沙盘与验收树; 建议级动作注明待人裁)', '损失分析': '(损失分析类: 用 ont_loss/ont_scenario 取证; 电量kWh为主, 金额必带电价区间0.391~0.85口径; 限电类损失≠设备责任)', '安全风险预测': '(风险预测类: 用 ont_risk/ont_mechanism_search 取证; 区分升级中/稳定高位/盲区, 未列≠正常, 不可判≠阴性)', } ASK_CAT_HINTS_EN = { '设备类': '(Equipment state: use unit profiles, object queries and mechanism search for evidence)', '检修排程决策': '(Maintenance scheduling: use scenarios and acceptance trees; mark advisory actions as pending human decision)', '损失分析': '(Loss analysis: use loss accounts and scenarios; energy in kWh/MWh, monetary value must carry the tariff-range caveat)', '安全风险预测': '(Risk prediction: distinguish escalating, stable-high and blind-spot cases; no listing is not proof of normality)', } MODEL_SPEC = { 'qwen3:8b': dict(nm='轻量 8B', loc='本地', t='约 30 到 60 秒', use='默认档: 快速简答问答 (用户裁 2026-09-07)'), 'qwen3.8:27b': dict(nm='标准 27B', loc='本地', t='约 2 到 4 分钟', use='复杂问题档: 多台横向比对与机制推断 (用户裁 2026-09-07)'), 'qwen3:30b-a3b': dict(nm='快 30B (a3b)', loc='本地', t='约 20 到 60 秒', use='8B 校闸未过时的自动升档目标; 也可直接选'), 'deepseek-r1:14b': dict(nm='本机 DeepSeek 14B', loc='本地', t='约 1 到 3 分钟', use='交叉审核审核档 (作答请优先 qwen 27B; 用户裁 2026-09-07)'), 'deepseek-r1:32b': dict(nm='本机 DeepSeek 32B', loc='本地', t='约 3 到 8 分钟', use='备用: 更强更慢, 引用格式不稳 (约三成概率不出文)'), 'deepseek-chat': dict(nm='DeepSeek', loc='云端', t='约 3 到 10 秒', use='默认档, 状态查询与机制推断', warn='内容会外发到模型厂商, 场名与台号已去标识化'), 'kimi-fast': dict(nm='Kimi', loc='云端', t='约 5 到 20 秒', use='长文归纳', warn='内容会外发到模型厂商, 场名与台号已去标识化'), 'qwen-plus': dict(nm='通义千问', loc='云端', t='约 5 到 20 秒', use='备用云端档', warn='内容会外发到模型厂商, 场名与台号已去标识化'), 'deepseek-v4-pro': dict(nm='云端 DS', loc='云端', t='约 2 到 6 分钟', use='复杂推理与交叉审核', warn='内容会外发到模型厂商, 场名与台号已去标识化'), } def _ask_step_text(k, d='', lang='zh'): k = '' if k is None else str(k) d = '' if d is None else str(d) if lang != 'en': return f'{k}· {d}' if d else k key = _ASK_STEP_KEY_EN.get(k, k) detail = d.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 re.search(r'[\u3400-\u9fff\u3000-\u303f\uff00-\uffef]', detail): detail = '' if re.search(r'[\u3400-\u9fff\u3000-\u303f\uff00-\uffef]', key): key = 'Working' if re.search(r'[\u3400-\u9fff\u3000-\u303f\uff00-\uffef]', detail): detail = '' return f'{key} · {detail}' if detail else key def _ask_worker(q, model, lang='zh', st=None): import subprocess, threading, time as _t out = _P.ont() / 'ask_last.out' try: if MODEL_SPEC.get(model, {}).get('loc') == '本地': # 本地档 (qwen3 / deepseek-r1) 走直连快路径; 原按 'qwen' 前缀硬判, 加 r1 档后改按档表 loc # 本地模型走直连快路径 (fast_agent: ollama chat+tools 同进程, 砍 dsh/node/MCP 三层; # 实测 8b 116s→19s; 同一套工具与闸, 服务端强制双闸). deepseek 云端仍走 dsh 完整编排. threading.Thread(target=_warm, args=(model,), daemon=True).start() sys.path.insert(0, str(pathlib.Path('.').resolve())) sys.path.insert(0, str(pathlib.Path('scripts').resolve())) from src.ontology.fast_agent import answer_escalate as _fast # 2026-09-07: 8B 拦了自动升 a3b 重答 (用户令) def _on_step(k, d=''): st['step'] = _ask_step_text(k, d, lang) st.setdefault('steps', []).append(st['step']) def _on_partial(txt): # 先快后准 (用户裁 2026-09-07): 30 s 内先出初答 (契约直出), 状态 partial, 前端继续轮询; 27B 出来后整段替换 ASK.update(answer=txt, state='partial', partial=True) r = _fast(q, model, on_step=_on_step, lang=lang, on_partial=_on_partial) _esc = f" · 升档 {r['escalated_from']}→{r['model']}" if r.get('escalated_from') else '' ans = r['answer'] + f"\n\n({r['gate']} · {r['rounds']}轮 · 直连快路径{_esc})" out.write_text(ans, encoding='utf-8') # 答案是中文: 缺 encoding 在中文 Windows 上按 GBK 写, 可能崩 from ontology_p2_verify import verify as _v0 ASK.update(answer=ans, verify=_v0(ans, lang=lang), partial=False, refined_from=r.get('escalated_from')) # 交叉审核: 换一个本地模型复核同一份证据 (同模型自审 = 自己给自己判卷) if st.get('xrev'): _on_step('交叉审核', '换模型复核中') from src.ontology.fast_agent import cross_review as _xr # 审核档 = DeepSeek r1:14b (用户裁 2026-09-07); 作答本身是 DeepSeek 时换 qwen 27B 审 —— 始终跨家族 (同家族自审 = 自己判卷) _rm = 'deepseek-r1:14b' if not model.startswith('deepseek') else 'qwen3.8:27b' st['review'] = _xr(q, r['answer'], r.get('facts') or [], model=_rm) ASK.update(state='done') return if False: q = '/no_think ' + q # (快路径内已带 /no_think; dsh 路径仅 deepseek) cmd = ['deploy/dsh/dsh_run.sh', '--profile', 'headless', '--patch', 'deploy/dsh/headless.patch.ontology.yml'] mp = ASK_MODELS.get(model) if mp: cmd += ['--patch', mp] if lang == 'en': cmd.append(q + ' Answer in English only. End every conclusion sentence with [object/id]. Finally run claim_check.') else: cmd.append(q + ' 每个结论句末尾用[对象id]标注来源。最后调用claim_check。') # 2026-09-16: 走 src.proc.run —— 本函数跑在 detail 服务进程里 (无可见控制台), 裸 spawn # 会让 Windows 给这条"云端档问答"新建可见控制台窗口 (本地档不走这里)。 from src import proc as _proc r = _proc.run(cmd, capture_output=True, text=True, timeout=600) ans = r.stdout.strip() or ('(空输出) stderr: ' + r.stderr[-500:]) out.write_text(ans, encoding='utf-8') sys.path.insert(0, str(pathlib.Path('scripts').resolve())) from ontology_p2_verify import verify as _v ASK.update(answer=ans, verify=_v(ans, lang=lang), state='done') except Exception as e: _msg = str(e) if 'Connection refused' in _msg or 'Errno 61' in _msg or '11434' in _msg: _msg = '本机模型未启动 (Ollama 127.0.0.1:11434 拒连); 离线版不转云, 请先启动 Ollama' # 2026-09-07 真停实测: 原样抛 urlopen 错误给用户 ASK.update(state='error', err=_msg[:300]) def ask_run(q: str, model: str = 'qwen3:8b', cat: str = '', xrev: bool = False, lang: str = 'zh') -> dict[str, Any]: """同步跑一次问答,返回末态 dict(前端契约不变:状态由 Java 持有,客户端轮询 /api/ask_status)。""" st: dict[str, Any] = dict(state='running', q=q, model=model, answer='', verify=None, err='', step='启动' if lang != 'en' else 'Starting', steps=[], review=None, xrev=bool(xrev), lang=lang) hints = ASK_CAT_HINTS_EN if lang == 'en' else ASK_CAT_HINTS q2 = (hints.get(cat, '') + ' ' + q).strip() try: _ask_worker(q2, model, lang, st) except Exception as e: # noqa: BLE001 st['err'] = str(e)[:300] st['state'] = 'error' if st.get('err') else 'done' return st