zhangqian
2024-11-19 b272fec78e30d1a10f3ab761684a119193391296
app/api/chat.py
@@ -10,8 +10,10 @@
from app.models.agent_model import AgentModel, AgentType
from app.models.base_model import get_db
from app.models.user_model import UserModel
from app.service.dialog import update_session_history
from app.service.basic import BasicService
from app.service.ragflow import RagflowService
from app.service.token import get_bisheng_token, get_ragflow_token
from app.service.service_token import get_bisheng_token, get_ragflow_token
router = APIRouter()
@@ -44,6 +46,7 @@
        try:
            async def forward_to_ragflow():
                while True:
                    is_new = False
                    message = await websocket.receive_json()
                    print(f"Received from client {chat_id}: {message}")
                    chat_history = message.get('chatHistory', [])
@@ -51,8 +54,10 @@
                    if len(chat_history) == 0:
                        chat_history = await ragflow_service.get_session_history(token, chat_id)
                        if len(chat_history) == 0:
                            is_new = True
                            chat_history = await ragflow_service.set_session(token, agent_id,
                                                                             message, chat_id, True)
                            # print("chat_history------------------------", chat_history)
                            if len(chat_history) == 0:
                                result = {"message": "内部错误:创建会话失败", "type": "close"}
                                await websocket.send_json(result)
@@ -64,6 +69,7 @@
                                "doc_ids": message.get("doc_ids", []),
                                "role": "user"
                            })
                    complete_response = ""
                    async for rag_response in ragflow_service.chat(token, chat_id, chat_history):
                        try:
                            if rag_response[:5] == "data:":
@@ -72,8 +78,9 @@
                            else:
                                # 否则,保持原样
                                text = rag_response
                            complete_response += text
                            try:
                                json_data = json.loads(text)
                                json_data = json.loads(complete_response)
                                data = json_data.get("data")
                                if data is True:  # 完成输出
                                    result = {"message": "", "type": "close"}
@@ -82,17 +89,20 @@
                                    result = {"message": "内部错误:" + answer, "type": "message"}
                                else:  # 正常输出
                                    answer = data.get("answer", "")
                                    reference = data.get("reference", "")
                                    result = {"message": answer, "type": "message", "reference": reference }
                                    reference = data.get("reference", {})
                                    result = {"message": answer, "type": "message", "reference": reference}
                                await websocket.send_json(result)
                            except json.JSONDecodeError:
                                print(f"Error decode ragflow response: {text}")
                                pass
                                complete_response = ""
                            except json.JSONDecodeError as e:
                                print(f"Error decoding JSON: {e}")
                                # print(f"Response text: {text}")
                        except Exception as e2:
                            result = {"message": f"内部错误: {e2}", "type": "close"}
                            await websocket.send_json(result)
                            print(f"Error process message of ragflow: {e2}")
                    dialog_chat_history = await ragflow_service.get_session_history(token, chat_id, 1)
                    await update_session_history(db, dialog_chat_history, current_user.id, is_new)
            # 启动任务处理客户端消息
            tasks = [
                asyncio.create_task(forward_to_ragflow())
@@ -187,6 +197,45 @@
                            await task
                        except asyncio.CancelledError:
                            pass
    elif agent_type == AgentType.BASIC:
        try:
            while True:
                # 接收前端消息
                message = await websocket.receive_json()
                question = message.get("message")
                if not question:
                    await websocket.send_json({"message": "Invalid request", "type": "error"})
                    continue
                service = BasicService(base_url=settings.basic_base_url)
                complete_response = ""
                async for result in service.excel_talk(question, chat_id):
                    try:
                        if result[:5] == "data:":
                            # 如果是,则截取掉前5个字符,并去除首尾空白符
                            text = result[5:].strip()
                        else:
                            # 否则,保持原样
                            text = result
                        complete_response += text
                        try:
                            json_data = json.loads(complete_response)
                            output = json_data.get("output", "")
                            result = {"message": output, "type": "message"}
                            await websocket.send_json(result | json_data)
                            complete_response = ""
                        except json.JSONDecodeError as e:
                            print(f"Error decoding JSON: {e}")
                            print(f"Response text: {text}")
                    except Exception as e2:
                        result = {"message": f"内部错误: {e2}", "type": "close"}
                        await websocket.send_json(result)
                        print(f"Error process message of basic agent: {e2}")
        except Exception as e:
            await websocket.send_json({"message": str(e), "type": "error"})
        finally:
            await websocket.close()
            print(f"Client {agent_id} disconnected")
    else:
        ret = {"message": "Agent not found", "type": "close"}
        await websocket.send_json(ret)