|
|
@@ -4,6 +4,8 @@ from fastapi import WebSocket, WebSocketDisconnect, APIRouter, Depends, HTTPExce
|
|
|
import asyncio
|
|
|
import websockets
|
|
|
from sqlalchemy.orm import Session
|
|
|
+
|
|
|
+from Log import logger
|
|
|
from app.api import get_current_user_websocket, ResponseList, get_current_user, format_file_url, process_files
|
|
|
from app.config.config import settings
|
|
|
from app.models.agent_model import AgentModel, AgentType
|
|
|
@@ -57,7 +59,7 @@ async def report_chat(websocket: WebSocket,
|
|
|
|
|
|
# 监听毕昇发来的消息并转发给客户端
|
|
|
async def forward_to_client():
|
|
|
- last_message = "step"
|
|
|
+ is_answer = False
|
|
|
while True:
|
|
|
message = await service_websocket.recv()
|
|
|
print(f"Received from bisheng: {message}")
|
|
|
@@ -67,19 +69,54 @@ async def report_chat(websocket: WebSocket,
|
|
|
msg = data.get("message", "")
|
|
|
category = data.get("category", "")
|
|
|
|
|
|
- if len(files) != 0 or (msg and category != "answer") or data["type"] == "close":
|
|
|
- if data["type"] == "close":
|
|
|
- t = "close"
|
|
|
+ if category == "question" and steps:
|
|
|
+ is_answer = False
|
|
|
+ if not steps:
|
|
|
+ steps = "\n"
|
|
|
else:
|
|
|
- t = "stream"
|
|
|
+ steps =steps + "\n"
|
|
|
+
|
|
|
+ result = {"message": steps, "type": "stream", "files": files}
|
|
|
+ await websocket.send_json(result)
|
|
|
+ if category == "answer" and not is_answer:
|
|
|
process_files(files, agent_id)
|
|
|
- result = {"message": msg, "type": t, "files": files}
|
|
|
+ if not steps.endswith("\n"):
|
|
|
+ steps+= "\n\n"
|
|
|
+ result = {"message": steps, "type": "stream", "files": files}
|
|
|
await websocket.send_json(result)
|
|
|
- elif steps and last_message == "step":
|
|
|
- result = {"step_message": steps, "type": "stream", "files": files}
|
|
|
+ if category == "answer" and is_answer:
|
|
|
+ process_files(files, agent_id)
|
|
|
+ result = {"message": "\n", "type": "stream", "files": files}
|
|
|
await websocket.send_json(result)
|
|
|
-
|
|
|
- last_message = "message" if msg else "step"
|
|
|
+ elif category == "processing":
|
|
|
+ process_files(files, agent_id)
|
|
|
+ is_answer = True
|
|
|
+ result = {"message": msg, "type": "stream", "files": files}
|
|
|
+ await websocket.send_json(result)
|
|
|
+ elif files:
|
|
|
+ process_files(files, agent_id)
|
|
|
+ result = {"message": "", "type": "stream", "files": files}
|
|
|
+ await websocket.send_json(result)
|
|
|
+ elif data["type"] == "close":
|
|
|
+ process_files(files, agent_id)
|
|
|
+ result = {"message": "", "type": "close", "files": files}
|
|
|
+ await websocket.send_json(result)
|
|
|
+ else:
|
|
|
+ logger.error("-------------------11111111111111--------------------------")
|
|
|
+ logger.error(data)
|
|
|
+ # if len(files) != 0 or (msg and category != "answer") or data["type"] == "close":
|
|
|
+ # if data["type"] == "close":
|
|
|
+ # t = "close"
|
|
|
+ # else:
|
|
|
+ # t = "stream"
|
|
|
+ # process_files(files, agent_id)
|
|
|
+ # result = {"message": msg, "type": t, "files": files}
|
|
|
+ # await websocket.send_json(result)
|
|
|
+ # elif steps and last_message == "step":
|
|
|
+ # result = {"step_message": steps, "type": "stream", "files": files}
|
|
|
+ # await websocket.send_json(result)
|
|
|
+
|
|
|
+ # last_message = "message" if msg else "step"
|
|
|
|
|
|
# 启动两个任务,分别处理客户端和服务端的消息
|
|
|
tasks = [
|