ask_views.py 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159
  1. # -*- coding: utf-8 -*-
  2. r"""问答管线(P12):从详情层 `serve.py` 搬来的 `_ask_worker`/`_ask_step_text` 与相关常量。
  3. ★搬运改动(仅三处,其余逐字保留):
  4. 1. `_ask_worker` 加参数 `st`(详情层是"线程 + 全局 ASK",这里改"同步 + 传入状态字典");
  5. 2. 体内 `ASK['x']` → `st['x']`(`_on_step` 在 worker 内定义,闭包自然拿到同一个 st);
  6. 3. 对外只暴露 `ask_run(q, model, cat, xrev, lang)`:建 st、按 `ask_start` 的 hints 逻辑拼词、跑管线、收尾 state。
  7. """
  8. from __future__ import annotations
  9. import pathlib
  10. import sys
  11. from typing import Any
  12. from app_common.app_common_guanlan.api import paths as _P
  13. ASK_CAT_HINTS = {
  14. '设备类': '(设备状态类: 从本体的机组档案、对象查询与机制检索三条路取证)',
  15. '检修排程决策': '(检修排程类: 用 ont_scenario/ont_get 取排期沙盘与验收树; 建议级动作注明待人裁)',
  16. '损失分析': '(损失分析类: 用 ont_loss/ont_scenario 取证; 电量kWh为主, 金额必带电价区间0.391~0.85口径; 限电类损失≠设备责任)',
  17. '安全风险预测': '(风险预测类: 用 ont_risk/ont_mechanism_search 取证; 区分升级中/稳定高位/盲区, 未列≠正常, 不可判≠阴性)',
  18. }
  19. ASK_CAT_HINTS_EN = {
  20. '设备类': '(Equipment state: use unit profiles, object queries and mechanism search for evidence)',
  21. '检修排程决策': '(Maintenance scheduling: use scenarios and acceptance trees; mark advisory actions as pending human decision)',
  22. '损失分析': '(Loss analysis: use loss accounts and scenarios; energy in kWh/MWh, monetary value must carry the tariff-range caveat)',
  23. '安全风险预测': '(Risk prediction: distinguish escalating, stable-high and blind-spot cases; no listing is not proof of normality)',
  24. }
  25. MODEL_SPEC = {
  26. 'qwen3:8b': dict(nm='轻量 8B', loc='本地', t='约 30 到 60 秒', use='默认档: 快速简答问答 (用户裁 2026-09-07)'),
  27. 'qwen3.8:27b': dict(nm='标准 27B', loc='本地', t='约 2 到 4 分钟', use='复杂问题档: 多台横向比对与机制推断 (用户裁 2026-09-07)'),
  28. 'qwen3:30b-a3b': dict(nm='快 30B (a3b)', loc='本地', t='约 20 到 60 秒', use='8B 校闸未过时的自动升档目标; 也可直接选'),
  29. 'deepseek-r1:14b': dict(nm='本机 DeepSeek 14B', loc='本地', t='约 1 到 3 分钟', use='交叉审核审核档 (作答请优先 qwen 27B; 用户裁 2026-09-07)'),
  30. 'deepseek-r1:32b': dict(nm='本机 DeepSeek 32B', loc='本地', t='约 3 到 8 分钟', use='备用: 更强更慢, 引用格式不稳 (约三成概率不出文)'),
  31. 'deepseek-chat': dict(nm='DeepSeek', loc='云端', t='约 3 到 10 秒', use='默认档, 状态查询与机制推断',
  32. warn='内容会外发到模型厂商, 场名与台号已去标识化'),
  33. 'kimi-fast': dict(nm='Kimi', loc='云端', t='约 5 到 20 秒', use='长文归纳',
  34. warn='内容会外发到模型厂商, 场名与台号已去标识化'),
  35. 'qwen-plus': dict(nm='通义千问', loc='云端', t='约 5 到 20 秒', use='备用云端档',
  36. warn='内容会外发到模型厂商, 场名与台号已去标识化'),
  37. 'deepseek-v4-pro': dict(nm='云端 DS', loc='云端', t='约 2 到 6 分钟', use='复杂推理与交叉审核',
  38. warn='内容会外发到模型厂商, 场名与台号已去标识化'),
  39. }
  40. def _ask_step_text(k, d='', lang='zh'):
  41. k = '' if k is None else str(k)
  42. d = '' if d is None else str(d)
  43. if lang != 'en':
  44. return f'{k}· {d}' if d else k
  45. key = _ASK_STEP_KEY_EN.get(k, k)
  46. detail = d.strip()
  47. m = re.search(r'第\s*(\d+\s*/\s*\d+)\s*轮', detail)
  48. if m:
  49. detail = 'round ' + m.group(1).replace(' ', '')
  50. else:
  51. m = re.search(r'(\d{1,2})\s*一站式证据包', detail)
  52. if m:
  53. alias = 'focus unit'
  54. try:
  55. from src.windscada.deid import unitize_answer
  56. alias = unitize_answer(f"WTG{int(m.group(1)):02d}")
  57. except Exception:
  58. pass
  59. detail = f'{alias} evidence pack'
  60. elif '闭环' in detail and '证据' in detail:
  61. n = ''.join(re.findall(r'\d+', detail))
  62. detail = f'closure and trend evidence: {n} records' if n else 'closure and trend evidence'
  63. elif '全场' in detail and '判级' in detail:
  64. n = ''.join(re.findall(r'\d+', detail))
  65. detail = f'fleet ratings: {n} units' if n else 'fleet ratings'
  66. elif '换模型' in detail or '复核' in detail:
  67. detail = 'reviewing with another model'
  68. elif re.search(r'[\u3400-\u9fff\u3000-\u303f\uff00-\uffef]', detail):
  69. detail = ''
  70. if re.search(r'[\u3400-\u9fff\u3000-\u303f\uff00-\uffef]', key):
  71. key = 'Working'
  72. if re.search(r'[\u3400-\u9fff\u3000-\u303f\uff00-\uffef]', detail):
  73. detail = ''
  74. return f'{key} · {detail}' if detail else key
  75. def _ask_worker(q, model, lang='zh', st=None):
  76. import subprocess, threading, time as _t
  77. out = _P.ont() / 'ask_last.out'
  78. try:
  79. if MODEL_SPEC.get(model, {}).get('loc') == '本地': # 本地档 (qwen3 / deepseek-r1) 走直连快路径; 原按 'qwen' 前缀硬判, 加 r1 档后改按档表 loc
  80. # 本地模型走直连快路径 (fast_agent: ollama chat+tools 同进程, 砍 dsh/node/MCP 三层;
  81. # 实测 8b 116s→19s; 同一套工具与闸, 服务端强制双闸). deepseek 云端仍走 dsh 完整编排.
  82. threading.Thread(target=_warm, args=(model,), daemon=True).start()
  83. sys.path.insert(0, str(pathlib.Path('.').resolve()))
  84. sys.path.insert(0, str(pathlib.Path('scripts').resolve()))
  85. from src.ontology.fast_agent import answer_escalate as _fast # 2026-09-07: 8B 拦了自动升 a3b 重答 (用户令)
  86. def _on_step(k, d=''):
  87. st['step'] = _ask_step_text(k, d, lang)
  88. st.setdefault('steps', []).append(st['step'])
  89. def _on_partial(txt):
  90. # 先快后准 (用户裁 2026-09-07): 30 s 内先出初答 (契约直出), 状态 partial, 前端继续轮询; 27B 出来后整段替换
  91. ASK.update(answer=txt, state='partial', partial=True)
  92. r = _fast(q, model, on_step=_on_step, lang=lang, on_partial=_on_partial)
  93. _esc = f" · 升档 {r['escalated_from']}→{r['model']}" if r.get('escalated_from') else ''
  94. ans = r['answer'] + f"\n\n({r['gate']} · {r['rounds']}轮 · 直连快路径{_esc})"
  95. out.write_text(ans, encoding='utf-8') # 答案是中文: 缺 encoding 在中文 Windows 上按 GBK 写, 可能崩
  96. from ontology_p2_verify import verify as _v0
  97. ASK.update(answer=ans, verify=_v0(ans, lang=lang), partial=False, refined_from=r.get('escalated_from'))
  98. # 交叉审核: 换一个本地模型复核同一份证据 (同模型自审 = 自己给自己判卷)
  99. if st.get('xrev'):
  100. _on_step('交叉审核', '换模型复核中')
  101. from src.ontology.fast_agent import cross_review as _xr
  102. # 审核档 = DeepSeek r1:14b (用户裁 2026-09-07); 作答本身是 DeepSeek 时换 qwen 27B 审 —— 始终跨家族 (同家族自审 = 自己判卷)
  103. _rm = 'deepseek-r1:14b' if not model.startswith('deepseek') else 'qwen3.8:27b'
  104. st['review'] = _xr(q, r['answer'], r.get('facts') or [], model=_rm)
  105. ASK.update(state='done')
  106. return
  107. if False:
  108. q = '/no_think ' + q # (快路径内已带 /no_think; dsh 路径仅 deepseek)
  109. cmd = ['deploy/dsh/dsh_run.sh', '--profile', 'headless',
  110. '--patch', 'deploy/dsh/headless.patch.ontology.yml']
  111. mp = ASK_MODELS.get(model)
  112. if mp:
  113. cmd += ['--patch', mp]
  114. if lang == 'en':
  115. cmd.append(q + ' Answer in English only. End every conclusion sentence with [object/id]. Finally run claim_check.')
  116. else:
  117. cmd.append(q + ' 每个结论句末尾用[对象id]标注来源。最后调用claim_check。')
  118. # 2026-09-16: 走 src.proc.run —— 本函数跑在 detail 服务进程里 (无可见控制台), 裸 spawn
  119. # 会让 Windows 给这条"云端档问答"新建可见控制台窗口 (本地档不走这里)。
  120. from src import proc as _proc
  121. r = _proc.run(cmd, capture_output=True, text=True, timeout=600)
  122. ans = r.stdout.strip() or ('(空输出) stderr: ' + r.stderr[-500:])
  123. out.write_text(ans, encoding='utf-8')
  124. sys.path.insert(0, str(pathlib.Path('scripts').resolve()))
  125. from ontology_p2_verify import verify as _v
  126. ASK.update(answer=ans, verify=_v(ans, lang=lang), state='done')
  127. except Exception as e:
  128. _msg = str(e)
  129. if 'Connection refused' in _msg or 'Errno 61' in _msg or '11434' in _msg:
  130. _msg = '本机模型未启动 (Ollama 127.0.0.1:11434 拒连); 离线版不转云, 请先启动 Ollama' # 2026-09-07 真停实测: 原样抛 urlopen 错误给用户
  131. ASK.update(state='error', err=_msg[:300])
  132. def ask_run(q: str, model: str = 'qwen3:8b', cat: str = '', xrev: bool = False, lang: str = 'zh') -> dict[str, Any]:
  133. """同步跑一次问答,返回末态 dict(前端契约不变:状态由 Java 持有,客户端轮询 /api/ask_status)。"""
  134. st: dict[str, Any] = dict(state='running', q=q, model=model, answer='', verify=None, err='',
  135. step='启动' if lang != 'en' else 'Starting', steps=[], review=None,
  136. xrev=bool(xrev), lang=lang)
  137. hints = ASK_CAT_HINTS_EN if lang == 'en' else ASK_CAT_HINTS
  138. q2 = (hints.get(cat, '') + ' ' + q).strip()
  139. try:
  140. _ask_worker(q2, model, lang, st)
  141. except Exception as e: # noqa: BLE001
  142. st['err'] = str(e)[:300]
  143. st['state'] = 'error' if st.get('err') else 'done'
  144. return st