agent.py 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177
  1. # -*- coding: utf-8 -*-
  2. """推理循环 (最大化推理, 判断受闸):
  3. 模型 = 推理员 + 调度员: 多轮工具调用 (ollama /api/chat tools), 思考链可开; 可自由提假设、跑 adhoc、对比、组装视图;
  4. 但结论只能经 verdict_gate (确定性判据) 定级; 叙述过接地闸 (台号/数字必须来自工具事实). 模型不可用时退化为规则调度.
  5. 预算与防环 (2026-08-23 真机测试逮: 45 次调用循环三轮、并发起 3 次 analyze): 轮数 + 工具调用总数双预算; 相同 (工具, 参数) 重复调用不再执行;
  6. 预算用完强制用已取事实收口作答 (不再甩原始事实); 重任务 (analyze) 不进推理循环工具表."""
  7. import json, re, time, urllib.request
  8. from . import llm # 随迁件
  9. from src.windcms import plugins # 仍留 src.windcms
  10. from .orchestrator import _rule_plan, ctx as _ctx
  11. SYSTEM = ('你是风电 CMS 振动诊断的推理员。你可以自由提出假设、用工具取证、写 pandas 即用即抛代码对比、反复迭代。'
  12. '规则: ①任何"定级"只能来自 verdict_gate 工具, 你自己不给级别; 未过闸的假设写 INSUFFICIENT。'
  13. '②最终回答里的台号/数字只能来自工具返回 (会被接地闸检查), 并注明来自哪个工具/文档。'
  14. '③六枚举: 定论/准定论·预警/候选/参考/INSUFFICIENT/撤回; P0/P1 不归本系统; 不命名滚道/圈侧除非工具明示。'
  15. '④问某台"状态"时, turbine_summary 给的是结构化层; 若其冲击轴标 ⚠ 或收敛报告终态有专条, 必须以终态与 scalar_trend 为准, 不要只看融合表说"正常"。'
  16. '⑤对比/趋势/谱类问题: 用 render 组装视图 (可多面板对比) 给用户看, 答案里引用工具返回的逐窗数字; 不要自己做加减乘除得出新数字。⑥"不支持 X"≠"对 X 无证据力"; "独有/特异"必报全场基率+位次+落差。'
  17. '⑦工具调用有预算, **一般问题 2-3 次调用足够** (状态问=turbine_summary+verdict_gate; 对比问=compare 或 fleet_rank); 取够事实立即作答, 不要为了全面把插件都调一遍; 只有用户明确要探索/出图时才用 hypothesis/render。'
  18. '⑧口径纪律: 排名/状态/定级一律用官方口径插件 (fleet_rank=末窗同档中位, turbine_summary, verdict_gate, scalar_trend); adhoc 只作探索, 其结果不能作结论口径 (全窗中位会把急性阶跃稀释掉)。谱插件给的 "G2不过" 表示该处无特征线, 不得写成"达到判据"。'
  19. '⑨★回答风格 (读者是现场/业主): 第一行 = 状态等级 + 一句话结论。状态等级用业主口径 (台级摘要的"状态等级(业主口径)"字段, 逐字照抄): 优秀=未见异常, 保持常规监视; 良好=轻微偏离, 记录观察; 报警=存在需现场核实/处置的异常; 危险=异常严重, 建议尽快停机检查; 不可判=数据不足或监测中断。六枚举 (候选/准定论·预警等) 不打头, 放在"证据状态"括号里说明可信度 (候选=待现场核实; 准定论·预警=基本确认; 候选·新发=新出现待核实; 候选·记基线=已记为基线下窗观察; 参考=仅记录; INSUFFICIENT=数据不足; 撤回=已取消)。然后 2-4 条要点 (数字+含义), 最后建议一行。★状态等级与证据状态必须逐字照抄工具返回, 严禁自行升降级 (掉字"准定论"写成"定论"、无依据写"危险"都会被闸拦)。'
  20. '然后 2-4 条要点, 每条 = 一个数字 + 它的含义, 数字保留原值 (例: "振动峰值达到全场同类机组中位的 4 倍")。最后若有现场动作给一行建议。'
  21. '★正文禁用内部代号: env400/env850/包络域/×fleet/×/G2/L4/L0/Main_bearing_front 等英文通道名/工具名一律不出现在正文; "×fleet 14.72" 写成 "全场中位的 14.7 倍"; 通道用中文 (主轴承前、发电机驱动端)。来源单独放末行小字, 工具用中文名 (turbine_summary=台级摘要, verdict_gate=定级闸, scalar_trend=趋势, fleet_rank=全场排名, kb_search=知识库, compare=对比, spectrum=频谱)。')
  22. TOOL_CN = {'turbine_summary': '台级摘要', 'verdict_gate': '定级闸', 'scalar_trend': '趋势', 'fleet_rank': '全场排名',
  23. 'kb_search': '知识库', 'compare': '对比', 'spectrum': '频谱', 'gate_spectral': '谱线闸', 'gate_shared_component': '共享成分闸',
  24. 'hypothesis': '假设簿', 'render': '视图', 'adhoc': '即席计算', 'explain': '通俗转述', 'export': '导出'}
  25. def _friendly(ans):
  26. """答案定稿后的确定性替换: 工具英文名→中文 (提示词管不住的最后一公里)."""
  27. for k, v in TOOL_CN.items():
  28. ans = ans.replace(k, v)
  29. return ans
  30. HEAVY = {'analyze', 'job_status'} # 不进推理循环的工具 (工作台/MCP 仍可用, 且需 confirm)
  31. def _tools_spec():
  32. out = []
  33. for p in plugins.manifest():
  34. if p['name'] in HEAVY:
  35. continue
  36. props = {}
  37. for k, v in (p['params'] or {}).items():
  38. t = 'number' if isinstance(v, (int, float)) and not isinstance(v, bool) else 'boolean' if isinstance(v, bool) else 'string'
  39. props[k] = {'type': t, 'description': f'例: {v}'}
  40. out.append({'type': 'function', 'function': {'name': p['name'], 'description': p['desc'], 'parameters': {'type': 'object', 'properties': props}}})
  41. return out
  42. def _chat(messages, tools, model, think=False, timeout=900):
  43. m = llm.pick(model, fast=True) # 交互 = 快模型优先 (显式传 model 仍可指定 32b)
  44. if not m:
  45. return None, None
  46. # 2026-08-24 提速: 32B decode ~22 tok/s, num_predict 2048 打满 = 90s/步; 诊断答案 150-250 字足够
  47. # THINK_LOCKED 模型 (恒思考版, think:false 会把思考灌 content — 2026-08-26 a3b 实逮): 强制 think:true
  48. # 拿 ollama 分离, 用 content 丢 thinking; num_predict 放大补思考预算 (400 tok 全被 thinking 吃光实测)
  49. tf = llm.effective_think(m, think)
  50. body = {'model': m, 'messages': messages, 'stream': False, 'think': tf, 'keep_alive': '30m',
  51. 'options': {'temperature': 0.3, 'num_predict': 1536 if (tf and not think) else 512, 'num_ctx': 32768}}
  52. if tools:
  53. body['tools'] = tools
  54. try:
  55. req = urllib.request.Request(llm.URL + '/api/chat', data=json.dumps(body, ensure_ascii=False).encode('utf-8'), headers={'Content-Type': 'application/json'})
  56. with urllib.request.urlopen(req, timeout=timeout) as r:
  57. d = json.loads(r.read().decode('utf-8'))
  58. msg = d.get('message') or {}
  59. if msg.get('content'): # 兜底: 标签形态的思考块剥掉 (分离失效/旧模板双保险)
  60. msg['content'] = re.sub(r'<think>.*?</think>', '', msg['content'], flags=re.S).strip()
  61. return d, m
  62. except Exception:
  63. return None, m
  64. def _final(messages, facts, model, think, steps, c, m_used, t0, note):
  65. """预算用完/接地闸多次拦截时: 强制用已取事实作答 (不带工具).
  66. ★输出纪律 (2026-08-24): 最终回答 **不超过 250 字**, 只给 结论级别 + 关键数字(带来源工具) + 现场动作; 不复述工具原文、不列表格、不解释判据含义 (判据解释在报告里, 不在对话里)。"""
  67. messages.append({'role': 'user', 'content': f'{note} 现在不再调用工具。只用上面工具返回的事实, 按系统提示第⑨条风格作答: 一句话结论(级别配人话注解) + 2-4 条要点(数字+含义, 正文禁内部代号与英文通道名) + 建议 + 末行来源(中文工具名), 不超过 250 字。'})
  68. resp, m2 = _chat(messages, None, model, think=think)
  69. ans = (resp or {}).get('message', {}).get('content', '').strip() if resp else ''
  70. ok, bad = llm.grounding(ans, '\n'.join(facts)) if ans else (False, {})
  71. if ans and not ok:
  72. # 再给一次机会: 指名越界的台号/数字, 要求删除或改为引用工具原值
  73. messages.append({'role': 'assistant', 'content': ans})
  74. messages.append({'role': 'user', 'content': f'接地闸拦截: 这些不在工具事实里: 台号 {bad.get("turbines")} / 数字 {bad.get("numbers")}。不要自己计算或推算新数字; 删掉它们或改为引用工具原值后重写, 其余不变。'})
  75. steps.append({'step': 'final', 'grounding_blocked': bad})
  76. resp, m2 = _chat(messages, None, model, think=think)
  77. ans = (resp or {}).get('message', {}).get('content', '').strip() if resp else ''
  78. ok, bad = llm.grounding(ans, '\n'.join(facts)) if ans else (False, {})
  79. if ans and ok:
  80. return dict(answer=_friendly(ans), steps=steps, views=list(c['views']), hypotheses=list(c['hypotheses']), source=f'llm-final:{m2}', model=m2, seconds=round(time.time() - t0, 1))
  81. return dict(answer=('[LLM 收口叙述被接地闸拦截 ' + str(bad) + '; 以下为工具事实权威版]\n\n' if ans else '[模型未能收口; 以下为工具事实权威版]\n\n') + '\n\n'.join(facts),
  82. steps=steps, views=list(c['views']), hypotheses=list(c['hypotheses']), source=f'llm-blocked-final:{m2}', model=m2, seconds=round(time.time() - t0, 1))
  83. def run(cfg, question, model=None, max_steps=8, max_calls=12, think=False, trace=None, on_step=None):
  84. """★think 默认 False (2026-08-24 实测): 同一诊断问题 think=True 164s / False 56s (2.9×),
  85. 且 False 版定级与 verdict_gate 一致、工具用得更对路 (True 版卡在"scalar_trend 无有效窗口数据"未纠正)。
  86. 根因: 推理已在插件/库函数里做完 (结论受闸), 模型 thinking 把 num_predict 打满 (300 tok/步 vs 25-49) 而 decode 恒定 ~22 tok/s。
  87. 需要模型自由发散的场合 (hypothesis) 再显式传 think=True。"""
  88. """多轮推理. 返回 dict(answer, steps[], views[], hypotheses[], source, model, seconds)."""
  89. c = _ctx(cfg)
  90. c.setdefault('views', [])
  91. c.setdefault('hypotheses', [])
  92. c['views'].clear()
  93. c['hypotheses'].clear()
  94. tools = _tools_spec()
  95. messages = [{'role': 'system', 'content': SYSTEM}, {'role': 'user', 'content': question}]
  96. steps, facts, m_used = [], [], None
  97. seen, calls, blocked = set(), 0, 0
  98. t0 = time.time()
  99. for step in range(max_steps):
  100. resp, m_used = _chat(messages, tools, model, think=think)
  101. if resp is None:
  102. break
  103. msg = resp.get('message', {})
  104. tool_calls = msg.get('tool_calls') or []
  105. if msg.get('thinking') and trace is not None:
  106. trace.append({'step': step, 'thinking': msg['thinking'][:2000]})
  107. messages.append({'role': 'assistant', 'content': msg.get('content', ''), **({'tool_calls': tool_calls} if tool_calls else {})})
  108. if not tool_calls:
  109. answer = msg.get('content', '').strip()
  110. ok, bad = llm.grounding(answer, '\n'.join(facts))
  111. if ok or not facts:
  112. return dict(answer=_friendly(answer), steps=steps, views=list(c['views']), hypotheses=list(c['hypotheses']), source=f'llm:{m_used}', model=m_used, seconds=round(time.time() - t0, 1))
  113. blocked += 1
  114. steps.append({'step': step, 'grounding_blocked': bad})
  115. if blocked >= 2:
  116. return _final(messages, facts, model, think, steps, c, m_used, t0, f'接地闸再次拦截 ({bad}).')
  117. messages.append({'role': 'user', 'content': f'接地闸: 你的回答引入了工具事实里没有的 {bad}。只用工具返回的台号与数字重写, 不确定的写 INSUFFICIENT。'})
  118. continue
  119. for call in tool_calls:
  120. fn = call.get('function', {})
  121. name = fn.get('name')
  122. args = fn.get('arguments') or {}
  123. if isinstance(args, str):
  124. try:
  125. args = json.loads(args)
  126. except Exception:
  127. args = {}
  128. key = (name, json.dumps(args, ensure_ascii=False, sort_keys=True))
  129. if name in HEAVY:
  130. text = f'{name} 不在推理循环里 (重任务需用户在工作台确认)。'
  131. elif key in seen:
  132. text = f'重复调用 {name} {key[1]} — 结果已在上文, 不再执行。'
  133. elif calls >= max_calls:
  134. text = '工具调用预算已用完, 请基于已取事实作答。'
  135. else:
  136. seen.add(key)
  137. calls += 1
  138. r = plugins.run(c, name, **args)
  139. text = str(r.get('text', ''))[:6000]
  140. if r.get('data') is not None and name not in ('render',):
  141. from src.windcms.plugins import _cn_sensors
  142. text += '\n数据(JSON): ' + _cn_sensors(json.dumps(r['data'], ensure_ascii=False, default=str)[:3000]) # 回显也中文化: 模型会从这里往正文抄英文通道名
  143. facts.append(f'[{name} {key[1]}] {text}')
  144. steps.append({'step': step, 'tool': name, 'args': args, 'result': text[:800], 'svg': bool(r.get('svg'))})
  145. if on_step:
  146. try:
  147. on_step({'tool': name, 'args': args, 'seconds': round(time.time() - t0, 1), 'brief': text[:120]})
  148. except Exception:
  149. on_step = None # 客户端断开: 停止推送但继续算完
  150. messages.append({'role': 'tool', 'content': text, 'tool_name': name})
  151. if calls >= max_calls:
  152. return _final(messages, facts, model, think, steps, c, m_used, t0, f'工具调用预算 ({max_calls}) 已用完。')
  153. if not steps and m_used is None:
  154. plan = _rule_plan(question)
  155. for s in plan:
  156. r = plugins.run(c, s['plugin'], **s.get('args', {}))
  157. facts.append(f'[{s["plugin"]} {json.dumps(s.get("args", {}), ensure_ascii=False)}] {r.get("text", "")}')
  158. steps.append({'step': 0, 'tool': s['plugin'], 'args': s.get('args', {}), 'result': str(r.get('text', ''))[:800], 'svg': bool(r.get('svg'))})
  159. if r.get('svg'):
  160. c['views'].append(dict(title=f'{s["plugin"]} {json.dumps(s.get("args", {}), ensure_ascii=False)}', panels=[dict(plugin=s['plugin'], args=s.get('args', {}), svg=r['svg'], text=r.get('text', ''))]))
  161. return dict(answer='\n\n'.join(facts), steps=steps, views=list(c['views']), hypotheses=list(c['hypotheses']), source='structured(rule-plan)', model=m_used, seconds=round(time.time() - t0, 1))
  162. return _final(messages, facts, model, think, steps, c, m_used, t0, f'轮数 ({max_steps}) 已用完。')