| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273 |
- # -*- 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,
- },
- )
|