# -*- coding: utf-8 -*- """ Agent 基类 支持逐 Token 实时流式打字机输出 (Token-level UI Streaming) 捕获思维链 (reasoning_content) 与最终生成结果 (content) 逐字推送到前端 """ import json import logging import time from typing import Optional, Callable logger = logging.getLogger(__name__) class BaseAgent: """智能体基类,支持 Token 级流式 LLM 推理与降级逻辑""" def __init__(self, name: str, system_prompt: str, role_icon: str = "🤖"): self.name = name self.system_prompt = system_prompt self.role_icon = role_icon self._client = None # 推理链记录 self.reasoning_trace = [] # 逐 Token 实时回调函数: callback(token_type: "reasoning"|"content", token_text: str) self.on_token_callback: Optional[Callable[[str, str], None]] = None def _get_client(self): """延迟初始化 OpenAI 客户端""" if self._client is not None: return self._client try: import httpx from openai import OpenAI import os from config import ( DEEPSEEK_API_KEY, DEEPSEEK_BASE_URL, DEEPSEEK_MODEL, VOLCENGINE_API_KEY, VOLCENGINE_BASE_URL, VOLCENGINE_MODEL, OPENAI_API_KEY, OPENAI_BASE_URL, OPENAI_MODEL ) api_key = os.environ.get("DEEPSEEK_API_KEY", "") if api_key: base_url = DEEPSEEK_BASE_URL model = DEEPSEEK_MODEL else: api_key = VOLCENGINE_API_KEY base_url = VOLCENGINE_BASE_URL model = VOLCENGINE_MODEL if not api_key: api_key = OPENAI_API_KEY base_url = OPENAI_BASE_URL model = OPENAI_MODEL if api_key: http_client = httpx.Client(trust_env=False, timeout=60.0) self._client = OpenAI( api_key=api_key, base_url=base_url, http_client=http_client, ) self._model = model logger.info(f"[{self.name}] 已成功连接大模型服务") return self._client except ImportError: logger.warning("openai 库未安装") except Exception as e: logger.warning(f"初始化 LLM 客户端失败: {e}") return None def _trace(self, step: str, content: str): """记录推理链步骤""" entry = { "timestamp": time.strftime("%H:%M:%S"), "step": step, "content": content, "agent": self.name, "icon": self.role_icon } self.reasoning_trace.append(entry) def infer(self, prompt: str, temperature: float = 0.1, max_retries: int = 0) -> str: """ 执行 SSE 流式 LLM 推理 (stream=True) 逐 Token 实时推送到 on_token_callback 渲染打字机效果 """ self.reasoning_trace = [] self._trace("📝 构建 Context", f"准备【{self.name}】数据与 Prompt") client = self._get_client() if client is None: self._trace("⚠️ 状态通知", "大模型未就绪,切换至专家规则引擎") logger.info(f"[{self.name}] LLM 不可用,降级到规则引擎") fallback = self.fallback_inference(prompt) self._trace("🔧 专家引擎输出", fallback) return fallback self._trace("🔗 大模型连接", "已连接大模型推理服务") for attempt in range(max_retries + 1): try: self._trace("🚀 发起流式推理", "正在建立 SSE 流式传输通道...") t0 = time.time() stream_resp = client.chat.completions.create( model=self._model, messages=[ {"role": "system", "content": self.system_prompt}, {"role": "user", "content": prompt}, ], temperature=temperature, max_tokens=2048, stream=True, timeout=60, ) full_content = [] reasoning_chunks = [] for chunk in stream_resp: if not chunk.choices: continue delta = chunk.choices[0].delta # 1. 逐 Token 提取深度思考过程 (reasoning_content) reasoning_piece = getattr(delta, "reasoning_content", None) or getattr(delta, "reasoning", None) if reasoning_piece: reasoning_chunks.append(reasoning_piece) if self.on_token_callback: try: self.on_token_callback("reasoning", reasoning_piece) except Exception: pass # 2. 逐 Token 提取正式回答内容 (content) content_piece = delta.content if content_piece: full_content.append(content_piece) if self.on_token_callback: try: self.on_token_callback("content", content_piece) except Exception: pass elapsed = time.time() - t0 final_text = "".join(full_content) full_reasoning = "".join(reasoning_chunks) if full_reasoning: self._trace("🧠 完整思维链", full_reasoning) if final_text.strip(): self._trace("✅ 流式生成完毕", f"耗时 {elapsed:.1f}s | 产出 {len(final_text)} 字符") self._trace("📄 原始推理输出", final_text) return final_text else: self._trace("⚠️ 输出为空", "流式生成无有效内容,降级到专家引擎") except Exception as e: error_msg = str(e) self._trace("⚡ 流式传输异常", f"连接中断: {error_msg}") logger.warning(f"[{self.name}] 流式调用失败: {e}") self._trace("🛡️ 安全降级", "无缝切换至离线风控规则引擎") fallback = self.fallback_inference(prompt) self._trace("🔧 专家引擎输出", fallback) return fallback def infer_json(self, prompt: str, temperature: float = 0.1) -> dict: """ 执行 LLM 推理并解析为 JSON """ result = self.infer(prompt, temperature) try: if "```json" in result: json_str = result.split("```json")[1].split("```")[0].strip() parsed = json.loads(json_str) self._trace("✅ 结构解析", "从 Markdown 成功提取 JSON 数据") return parsed elif "```" in result: json_str = result.split("```")[1].split("```")[0].strip() parsed = json.loads(json_str) self._trace("✅ 结构解析", "从代码块成功提取 JSON 数据") return parsed else: parsed = json.loads(result) self._trace("✅ 结构解析", "直接解析 JSON 成功") return parsed except (json.JSONDecodeError, IndexError): self._trace("⚠️ 格式适配", "启用自动结构修正") logger.warning(f"[{self.name}] JSON 解析失败,返回原始文本") return {"raw_response": result, "parse_error": True} def fallback_inference(self, prompt: str) -> str: """规则引擎降级推理""" return json.dumps({"error": "大模型服务不可用,规则引擎未实现"}, ensure_ascii=False) def __repr__(self): return f"{self.role_icon} {self.name}"