# -*- coding: utf-8 -*- """ 对话式数据解读接口 —— 统一入口,SSE 流式 根据用户意图自动路由到不同智能体: - 配风计划审查(上传PDF) → 审查管线 - 数据解读/需风量计算/其他 → dialog_agent(统一通风对话助手) """ import asyncio from fastapi import APIRouter, Depends, File, Form, UploadFile from fastapi.responses import StreamingResponse from pydantic import BaseModel, Field from api.auth import get_current_user from api.sse_core import sse_event_generator from api.intent import classify_intent from api.chat_model import get_chat_model from agents.vent_agent import get_dialog_agent from agents.review_agent import stream_review from db.chat_store import ( get_messages, create_session, update_session_title, get_session_mode, set_session_mode, get_user_preferences, ) router = APIRouter() # ── 权限模式 → LangGraph interrupt 配置 ── # 基于 DeepAgents / LangGraph 官方 API: # - interrupt_before: 在执行指定节点前暂停,等待人工审批 # - interrupt_after: 在执行指定节点后暂停,让用户审查结果 MODE_INTERRUPT = { "plan": {}, # 计划模式:由 request_plan_approval 工具触发程序化中断 "full": {}, # 完全访问:无中断 } class ResumeRequest(BaseModel): """中断恢复请求""" session_id: str = Field(..., description="会话 ID") thread_id: str | None = Field(None, description="LangGraph 线程 ID") action: str = Field("approve", description="approve 或 reject") # ── 会话标题生成(LLM)── async def generate_session_title(message: str) -> str: """使用 LLM 将用户首条消息压缩为会话标题(不超过20字)。 仅在开启新会话、且 messages 表中仅此一条用户消息时调用。 参考 classify_intent 的实现模式,使用 SUMMARY_MODEL(低延迟)。 Args: message: 用户的首条消息文本 Returns: 压缩后的标题字符串,不超过20字 """ try: model = get_chat_model(fast_mode=True) prompt = ( "你是一个会话标题生成器。根据用户的第一条消息,生成一个简短的会话标题。\n" "要求:\n" "- 标题不超过20个字\n" "- 保留核心含义,去掉语气词和冗余描述\n" "- 只回复标题本身,不要加任何说明、引号或标点\n" "\n" f"用户消息:{message[:200]}\n" "标题:" ) loop = asyncio.get_event_loop() result = await loop.run_in_executor( None, lambda: model.invoke([{"role": "user", "content": prompt}]) ) raw = result.content if hasattr(result, "content") else str(result) title = raw.strip().strip('"').strip("'").strip("《》").strip() if len(title) > 20: title = title[:20] return title if title else message[:20] except Exception as e: print(f"[标题生成] LLM 生成失败,回退到截断: {e}") return message[:20] + ("..." if len(message) > 20 else "") # ── POST /api/chat ── @router.post("/chat") async def chat( message: str = Form(...), session_id: str | None = Form(default=None), thread_id: str | None = Form(default=None), mode: str | None = Form(default=None), file: UploadFile | None = File(default=None), user_info: dict = Depends(get_current_user), ): """统一对话接口(SSE 流式)—— 支持文本对话和 PDF 配风计划审查。 根据用户意图自动路由到不同智能体: - 配风计划审查(上传PDF) → 审查管线(form-reviewer + data-checker + calc-verifier) - 数据解读/需风量计算/其他 → dialog_agent(统一通风对话助手) 请求头: - x-access-token: 必填,通风系统登录令牌 请求体: multipart/form-data - message: 用户消息(必填) - session_id: 可选,会话ID - thread_id: 可选,LangGraph 线程ID - mode: 可选,权限模式(plan | full),仅在新建会话时生效 - file: 可选,附件(PDF/文档等) 响应: SSE 流(text/event-stream) """ # ── 会话管理 ── is_new = not session_id user_name = user_info.get("username", "admin") session_id = session_id or create_session(user_name=user_name) thread_id = thread_id or session_id # 设置当前用户上下文(供偏好工具等获取调用者身份) from tools.vent_tools import _current_user _current_user.set(user_name) # 保存用户原始消息(后续注入偏好/计划前缀前),用于写入数据库 original_message = message # 新建会话时,若前端传了 mode 参数则应用(否则沿用 DB 默认值) if is_new and mode in ("plan", "full"): set_session_mode(session_id, mode) # 读取会话权限模式 # ⚠️ interrupt_before / interrupt_after 是 astream() 的独立 kwargs, # 不能放在 config dict 里,否则 LangGraph 不识别 current_mode = get_session_mode(session_id) interrupt_kwargs = MODE_INTERRUPT.get(current_mode, {}) config = {"configurable": {"thread_id": thread_id}} print(f"[模式] 会话 {session_id} 权限模式: {current_mode} → interrupt: {interrupt_kwargs}") # 提取附件文件名(如有) filename = file.filename if file else None # 自动设置会话标题(仅首轮对话触发,使用 LLM 压缩,不超过20字) messages = get_messages(session_id) if len(messages) <= 1: title = await generate_session_title(message) update_session_title(session_id, title) # ── 注入用户偏好记忆(自动触发)── prefs = get_user_preferences(user_name) if prefs: pref_lines = "\n".join( f"- (id:{p['id']}) [{p['category']}] {p['content']}" for p in prefs ) message = ( f"[用户偏好记忆]\n" f"以下是用户「{user_name}」保存的习惯和偏好,请在回复时主动参考:\n" f"{pref_lines}\n" f"\n" f"用户消息:{message}" ) print(f"[偏好] 已为用户 {user_name} 注入 {len(prefs)} 条偏好记忆") # ── Agent 推理意图(含附件文件名辅助判断)── intent = await classify_intent(message, filename) # ── 按意图路由 ── if intent == "review" and file is not None: # 配风计划审查 + 有文件 → 审查管线 print(f"[路由] 意图: 配风计划审查 → stream_review (file={filename})") return StreamingResponse( stream_review(file, message, session_id), media_type="text/event-stream", headers={ "Cache-Control": "no-cache", "Connection": "keep-alive", "X-Accel-Buffering": "no", "X-Session-Id": session_id, }, ) # ── 统一走通风对话助手(处理数据解读、需风量计算、审查无文件等所有场景)── agent = get_dialog_agent() agent_cn_name = "通风对话助手" if intent == "review" and file is None: # 审查意图但无文件 → 提示用户上传 print(f"[路由] 意图: 配风计划审查 → 提示用户上传 PDF") message = f"用户想进行配风计划审查,但未上传附件。请提示用户上传配风计划PDF文件。用户原始消息:{message}" else: print(f"[路由] 意图: {intent} → dialog_agent(统一)") # 如有附件,将文件名信息追加到消息中 if filename: message = f"用户上传了文件「{filename}」。{message}" # ── 计划模式:注入"先规划 → 审批 → 执行"指令 ── if current_mode == "plan" and agent_cn_name == "通风对话助手": plan_prefix = ( "【系统指令:当前处于计划模式】\n" "你需要分三步工作:\n" "1. 规划阶段:自由使用工具获取信息、读取技能、查询数据,制定详细执行计划。\n" "2. 提交审批:计划制定完毕后,必须调用 request_plan_approval 工具提交计划摘要,等待人工审批。\n" " 在此之前严禁执行任何最终操作(如生成报告、写入文件、修改数据等)。\n" "3. 执行阶段:审批通过后,按照计划逐步执行任务,无需再次申请审批。\n" "\n用户问题:" ) message = plan_prefix + message print(f"[计划模式] 已注入计划指令前缀 ({len(plan_prefix)} 字符)") return StreamingResponse( sse_event_generator(agent, message, thread_id, session_id, config, agent_cn_name, original_user_message=original_message, interrupt_before=interrupt_kwargs.get("interrupt_before"), interrupt_after=interrupt_kwargs.get("interrupt_after")), media_type="text/event-stream", headers={ "Cache-Control": "no-cache", "Connection": "keep-alive", "X-Accel-Buffering": "no", "X-Session-Id": session_id, }, ) # ── POST /api/chat/resume ── @router.post("/chat/resume") async def resume_chat( req: ResumeRequest, user_info: dict = Depends(get_current_user), ): """恢复被中断的 LangGraph 对话(Human-in-the-Loop 审批)。 当计划模式(plan)时,Agent 调用 request_plan_approval 工具后暂停, 前端展示审批 UI。用户点击批准/拒绝后调用此接口继续执行。 请求头: - x-access-token: 必填,通风系统登录令牌 """ from api.sse_core import resume_stream session_id = req.session_id thread_id = req.thread_id or session_id # plan 模式审批后清除所有配置式中断(只审批一次,之后执行到底) # ⚠️ 必须 copy,否则 pop 会原地修改模块级常量 MODE_INTERRUPT mode = get_session_mode(session_id) interrupt_kwargs = dict(MODE_INTERRUPT.get(mode, {})) if mode == "plan": interrupt_kwargs.pop("interrupt_before", None) interrupt_kwargs.pop("interrupt_after", None) config = {"configurable": {"thread_id": thread_id}} return StreamingResponse( resume_stream( get_dialog_agent(), session_id, thread_id, config, agent_cn_name="通风对话助手", action=req.action, interrupt_before=interrupt_kwargs.get("interrupt_before"), interrupt_after=interrupt_kwargs.get("interrupt_after"), ), media_type="text/event-stream", headers={ "Cache-Control": "no-cache", "Connection": "keep-alive", "X-Accel-Buffering": "no", "X-Session-Id": session_id, }, )