流式输出

This commit is contained in:
dspwasc 2025-04-29 10:21:13 +08:00
parent d769f814b4
commit 61eaec4d64
10 changed files with 1524 additions and 291 deletions

View File

@ -18,11 +18,16 @@ django.setup() # 添加这行来初始化 Django
from django.core.asgi import get_asgi_application from django.core.asgi import get_asgi_application
from channels.routing import ProtocolTypeRouter, URLRouter from channels.routing import ProtocolTypeRouter, URLRouter
from channels.auth import AuthMiddlewareStack from channels.auth import AuthMiddlewareStack
from channels.security.websocket import AllowedHostsOriginValidator
from user_management.routing import websocket_urlpatterns from user_management.routing import websocket_urlpatterns
from user_management.middleware import TokenAuthMiddleware
# 使用TokenAuthMiddleware代替AuthMiddlewareStack
application = ProtocolTypeRouter({ application = ProtocolTypeRouter({
"http": get_asgi_application(), "http": get_asgi_application(),
"websocket": AuthMiddlewareStack( "websocket": AllowedHostsOriginValidator(
URLRouter(websocket_urlpatterns) TokenAuthMiddleware(
URLRouter(websocket_urlpatterns)
)
), ),
}) })

View File

@ -41,6 +41,9 @@ ALLOWED_HOSTS = ['*'] # 仅在开发环境使用
# 服务器配置 # 服务器配置
DEBUG = False DEBUG = False
# 是否允许注册新用户
ALLOW_REGISTRATION = True
# ALLOWED_HOSTS = ['frptx.chiyong.fun', 'localhost', '127.0.0.1'] # ALLOWED_HOSTS = ['frptx.chiyong.fun', 'localhost', '127.0.0.1']
# Application definition # Application definition
@ -70,6 +73,7 @@ MIDDLEWARE = [
'django.contrib.messages.middleware.MessageMiddleware', 'django.contrib.messages.middleware.MessageMiddleware',
'django.middleware.clickjacking.XFrameOptionsMiddleware', 'django.middleware.clickjacking.XFrameOptionsMiddleware',
'user_management.middleware.UserActivityMiddleware', 'user_management.middleware.UserActivityMiddleware',
'user_management.middleware.CSRFExemptMiddleware', # 添加CSRF豁免中间件
] ]
ROOT_URLCONF = 'role_based_system.urls' ROOT_URLCONF = 'role_based_system.urls'
@ -168,7 +172,12 @@ ASGI_APPLICATION = "role_based_system.asgi.application"
# Channel Layers 配置 # Channel Layers 配置
CHANNEL_LAYERS = { CHANNEL_LAYERS = {
"default": { "default": {
"BACKEND": "channels.layers.InMemoryChannelLayer", "BACKEND": "channels_redis.core.RedisChannelLayer",
"CONFIG": {
"hosts": [("127.0.0.1", 6379)],
"capacity": 1500, # 默认100
"expiry": 60, # 默认60秒
},
}, },
} }
@ -289,3 +298,9 @@ REST_FRAMEWORK = {
'rest_framework.parsers.MultiPartParser' 'rest_framework.parsers.MultiPartParser'
], ],
} }
# DeepSeek API配置
DEEPSEEK_API_KEY = "sk-xqbujijjqqmlmlvkhvxeogqjtzslnhdtqxqgiyuhwpoqcjvf" # 请替换为您的实际有效的DeepSeek API密钥
SILICON_CLOUD_API_KEY = 'sk-xqbujijjqqmlmlvkhvxeogqjtzslnhdtqxqgiyuhwpoqcjvf'

View File

@ -4,6 +4,13 @@ from channels.db import database_sync_to_async
from channels.exceptions import StopConsumer from channels.exceptions import StopConsumer
import logging import logging
from rest_framework.authtoken.models import Token from rest_framework.authtoken.models import Token
from urllib.parse import parse_qs
from .models import ChatHistory, KnowledgeBase
import aiohttp
import asyncio
from django.conf import settings
import uuid
import traceback
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
@ -11,19 +18,20 @@ class NotificationConsumer(AsyncWebsocketConsumer):
async def connect(self): async def connect(self):
"""建立WebSocket连接""" """建立WebSocket连接"""
try: try:
# 获取token # 从URL参数中获取token
headers = dict(self.scope['headers']) query_string = self.scope.get('query_string', b'').decode()
auth_header = headers.get(b'authorization', b'').decode() query_params = parse_qs(query_string)
token_key = query_params.get('token', [''])[0]
if not auth_header.startswith('Token '): if not token_key:
logger.warning("WebSocket连接尝试但没有提供token")
await self.close() await self.close()
return return
token_key = auth_header.split(' ')[1]
# 验证token # 验证token
self.user = await self.get_user_from_token(token_key) self.user = await self.get_user_from_token(token_key)
if not self.user: if not self.user:
logger.warning(f"WebSocket连接尝试但token无效: {token_key}")
await self.close() await self.close()
return return
@ -34,8 +42,10 @@ class NotificationConsumer(AsyncWebsocketConsumer):
self.channel_name self.channel_name
) )
await self.accept() await self.accept()
logger.info(f"用户 {self.user.username} WebSocket连接成功")
except Exception as e: except Exception as e:
logger.error(f"WebSocket连接错误: {str(e)}")
await self.close() await self.close()
@database_sync_to_async @database_sync_to_async
@ -65,3 +75,816 @@ class NotificationConsumer(AsyncWebsocketConsumer):
logger.info(f"已发送通知给用户 {self.user.username}") logger.info(f"已发送通知给用户 {self.user.username}")
except Exception as e: except Exception as e:
logger.error(f"发送通知消息时发生错误: {str(e)}") logger.error(f"发送通知消息时发生错误: {str(e)}")
class ChatConsumer(AsyncWebsocketConsumer):
async def connect(self):
"""建立 WebSocket 连接"""
try:
# 从URL参数中获取token
query_string = self.scope.get('query_string', b'').decode()
query_params = parse_qs(query_string)
token_key = query_params.get('token', [''])[0]
if not token_key:
logger.warning("WebSocket连接尝试但没有提供token")
await self.close()
return
# 验证token
self.user = await self.get_user_from_token(token_key)
if not self.user:
logger.warning(f"WebSocket连接尝试但token无效: {token_key}")
await self.close()
return
# 将用户信息存储在scope中
self.scope["user"] = self.user
await self.accept()
logger.info(f"用户 {self.user.username} WebSocket连接成功")
except Exception as e:
logger.error(f"WebSocket连接错误: {str(e)}")
await self.close()
@database_sync_to_async
def get_user_from_token(self, token_key):
try:
token = Token.objects.select_related('user').get(key=token_key)
return token.user
except Token.DoesNotExist:
return None
async def disconnect(self, close_code):
"""关闭 WebSocket 连接"""
pass
async def receive(self, text_data):
"""接收消息并处理"""
try:
data = json.loads(text_data)
# 验证必要字段
if 'question' not in data or 'conversation_id' not in data:
await self.send_error("缺少必要字段")
return
# 创建问题记录
question_record = await self.create_question_record(data)
if not question_record:
return
# 开始流式处理
await self.stream_answer(question_record, data)
except Exception as e:
logger.error(f"处理消息时出错: {str(e)}")
await self.send_error(f"处理消息时出错: {str(e)}")
@database_sync_to_async
def _create_question_record_sync(self, data):
"""同步创建问题记录"""
try:
# 获取会话历史记录
conversation_id = data['conversation_id']
existing_records = ChatHistory.objects.filter(
conversation_id=conversation_id
).order_by('created_at')
# 获取或创建元数据
if existing_records.exists():
first_record = existing_records.first()
metadata = first_record.metadata or {}
dataset_ids = metadata.get('dataset_id_list', [])
knowledge_bases = []
# 验证知识库权限
for kb_id in dataset_ids:
try:
kb = KnowledgeBase.objects.get(id=kb_id)
if not self.check_knowledge_base_permission(kb, self.scope["user"], 'read'):
raise Exception(f'无权访问知识库: {kb.name}')
knowledge_bases.append(kb)
except KnowledgeBase.DoesNotExist:
raise Exception(f'知识库不存在: {kb_id}')
else:
# 新会话处理
dataset_ids = data.get('dataset_id_list', [])
if not dataset_ids:
raise Exception('新会话需要提供知识库ID')
knowledge_bases = []
for kb_id in dataset_ids:
kb = KnowledgeBase.objects.get(id=kb_id)
if not self.check_knowledge_base_permission(kb, self.scope["user"], 'read'):
raise Exception(f'无权访问知识库: {kb.name}')
knowledge_bases.append(kb)
metadata = {
'model_id': data.get('model_id', '7a214d0e-e65e-11ef-9f4a-0242ac120006'),
'dataset_id_list': [str(kb.id) for kb in knowledge_bases],
'dataset_external_id_list': [str(kb.external_id) for kb in knowledge_bases if kb.external_id],
'dataset_names': [kb.name for kb in knowledge_bases]
}
# 创建问题记录
return ChatHistory.objects.create(
user=self.scope["user"],
knowledge_base=knowledge_bases[0],
conversation_id=conversation_id,
title=data.get('title', 'New chat'),
role='user',
content=data['question'],
metadata=metadata
)
except Exception as e:
logger.error(f"创建问题记录失败: {str(e)}")
return None, str(e)
async def create_question_record(self, data):
"""异步创建问题记录"""
try:
result = await self._create_question_record_sync(data)
if isinstance(result, tuple):
_, error_message = result
await self.send_error(error_message)
return None
return result
except Exception as e:
await self.send_error(str(e))
return None
def check_knowledge_base_permission(self, kb, user, permission_type):
"""检查知识库权限"""
# 实现权限检查逻辑
return True # 临时返回 True需要根据实际情况实现
async def stream_answer(self, question_record, data):
"""流式处理回答"""
try:
# 创建 AI 回答记录
answer_record = await database_sync_to_async(ChatHistory.objects.create)(
user=self.scope["user"],
knowledge_base=question_record.knowledge_base,
conversation_id=str(question_record.conversation_id),
title=question_record.title,
parent_id=str(question_record.id),
role='assistant',
content="",
metadata=question_record.metadata
)
# 发送初始响应
await self.send_json({
'code': 200,
'message': '开始流式传输',
'data': {
'id': str(answer_record.id),
'conversation_id': str(question_record.conversation_id),
'content': '',
'is_end': False
}
})
# 调用外部 API 获取流式响应
async with aiohttp.ClientSession() as session:
# 创建聊天会话
chat_response = await session.post(
f"{settings.API_BASE_URL}/api/application/chat/open",
json={
"id": "d5d11efa-ea9a-11ef-9933-0242ac120006",
"model_id": question_record.metadata.get('model_id'),
"dataset_id_list": question_record.metadata.get('dataset_external_id_list', []),
"multiple_rounds_dialogue": False,
"dataset_setting": {
"top_n": 10,
"similarity": "0.3",
"max_paragraph_char_number": 10000,
"search_mode": "blend",
"no_references_setting": {
"value": "{question}",
"status": "ai_questioning"
}
},
"model_setting": {
"prompt": "**相关文档内容**{data} **回答要求**:如果相关文档内容中没有可用信息,请回答\"没有在知识库中查找到相关信息,建议咨询相关技术支持或参考官方文档进行操作\"。请根据相关文档内容回答用户问题。不要输出与用户问题无关的内容。请使用中文回答客户问题。**用户问题**{question}"
},
"problem_optimization": False
}
)
chat_data = await chat_response.json()
if chat_data.get('code') != 200:
raise Exception(f"创建聊天会话失败: {chat_data}")
chat_id = chat_data['data']
# 建立流式连接
async with session.post(
f"{settings.API_BASE_URL}/api/application/chat_message/{chat_id}",
json={"message": data['question'], "re_chat": False, "stream": True},
headers={"Content-Type": "application/json"}
) as response:
full_content = ""
buffer = ""
async for chunk in response.content.iter_any():
chunk_str = chunk.decode('utf-8')
buffer += chunk_str
while '\n\n' in buffer:
parts = buffer.split('\n\n', 1)
line = parts[0]
buffer = parts[1]
if line.startswith('data: '):
try:
json_str = line[6:]
chunk_data = json.loads(json_str)
if 'content' in chunk_data:
content_part = chunk_data['content']
full_content += content_part
await self.send_json({
'code': 200,
'message': 'partial',
'data': {
'id': str(answer_record.id),
'conversation_id': str(question_record.conversation_id),
'content': content_part,
'is_end': chunk_data.get('is_end', False)
}
})
if chunk_data.get('is_end', False):
# 保存完整内容
answer_record.content = full_content.strip()
await database_sync_to_async(answer_record.save)()
# 生成或获取标题
title = await self.get_or_generate_title(
question_record.conversation_id,
data['question'],
full_content.strip()
)
# 发送最终响应
await self.send_json({
'code': 200,
'message': '完成',
'data': {
'id': str(answer_record.id),
'conversation_id': str(question_record.conversation_id),
'title': title,
'dataset_id_list': question_record.metadata.get('dataset_id_list', []),
'dataset_names': question_record.metadata.get('dataset_names', []),
'role': 'assistant',
'content': full_content.strip(),
'created_at': answer_record.created_at.strftime('%Y-%m-%d %H:%M:%S'),
'is_end': True
}
})
return
except json.JSONDecodeError as e:
logger.error(f"JSON解析错误: {e}, 数据: {line}")
continue
except Exception as e:
logger.error(f"流式处理出错: {str(e)}")
await self.send_error(str(e))
# 保存已收集的内容
if 'full_content' in locals() and full_content:
try:
answer_record.content = full_content.strip()
await database_sync_to_async(answer_record.save)()
except Exception as save_error:
logger.error(f"保存部分内容失败: {str(save_error)}")
@database_sync_to_async
def get_or_generate_title(self, conversation_id, question, answer):
"""获取或生成对话标题"""
try:
# 先检查是否已有标题
current_title = ChatHistory.objects.filter(
conversation_id=str(conversation_id)
).exclude(
title__in=["New chat", "新对话", ""]
).values_list('title', flat=True).first()
if current_title:
return current_title
# 如果没有标题,生成新标题
# 这里需要实现标题生成的逻辑
generated_title = "新对话" # 临时使用默认标题
# 更新所有相关记录的标题
ChatHistory.objects.filter(
conversation_id=str(conversation_id)
).update(title=generated_title)
return generated_title
except Exception as e:
logger.error(f"获取或生成标题失败: {str(e)}")
return "新对话"
async def send_json(self, content):
"""发送 JSON 格式的消息"""
await self.send(text_data=json.dumps(content))
async def send_error(self, message):
"""发送错误消息"""
await self.send_json({
'code': 500,
'message': message,
'data': {'is_end': True}
})
class ChatStreamConsumer(AsyncWebsocketConsumer):
async def connect(self):
"""建立WebSocket连接"""
try:
# 从URL参数中获取token
query_string = self.scope.get('query_string', b'').decode()
query_params = parse_qs(query_string)
token_key = query_params.get('token', [''])[0]
if not token_key:
logger.warning("WebSocket连接尝试但没有提供token")
await self.close()
return
# 验证token
self.user = await self.get_user_from_token(token_key)
if not self.user:
logger.warning(f"WebSocket连接尝试但token无效: {token_key}")
await self.close()
return
# 将用户信息存储在scope中
self.scope["user"] = self.user
await self.accept()
logger.info(f"用户 {self.user.username} 流式输出WebSocket连接成功")
except Exception as e:
logger.error(f"WebSocket连接错误: {str(e)}")
await self.close()
@database_sync_to_async
def get_user_from_token(self, token_key):
try:
token = Token.objects.select_related('user').get(key=token_key)
return token.user
except Token.DoesNotExist:
return None
async def disconnect(self, close_code):
"""关闭WebSocket连接"""
logger.info(f"用户 {self.user.username if hasattr(self, 'user') else 'unknown'} WebSocket连接断开代码: {close_code}")
async def receive(self, text_data):
"""接收消息并处理"""
try:
data = json.loads(text_data)
# 检查必填字段
if 'question' not in data:
await self.send_error("缺少必填字段: question")
return
if 'conversation_id' not in data:
await self.send_error("缺少必填字段: conversation_id")
return
# 处理新会话或现有会话
await self.process_chat_request(data)
except Exception as e:
logger.error(f"处理消息时出错: {str(e)}")
logger.error(traceback.format_exc())
await self.send_error(f"处理消息时出错: {str(e)}")
async def process_chat_request(self, data):
"""处理聊天请求"""
try:
conversation_id = data['conversation_id']
question = data['question']
# 获取会话信息和知识库
session_info = await self.get_session_info(data)
if not session_info:
return
knowledge_bases, metadata, dataset_external_id_list = session_info
# 创建问题记录
question_record = await self.create_question_record(
conversation_id,
question,
knowledge_bases,
metadata
)
if not question_record:
return
# 创建AI回答记录
answer_record = await self.create_answer_record(
conversation_id,
question_record,
knowledge_bases,
metadata
)
# 发送初始响应
await self.send_json({
'code': 200,
'message': '开始流式传输',
'data': {
'id': str(answer_record.id),
'conversation_id': str(conversation_id),
'content': '',
'is_end': False
}
})
# 调用外部API获取流式响应
await self.stream_from_external_api(
conversation_id,
question,
dataset_external_id_list,
answer_record,
metadata,
knowledge_bases
)
except Exception as e:
logger.error(f"处理聊天请求时出错: {str(e)}")
logger.error(traceback.format_exc())
await self.send_error(f"处理聊天请求时出错: {str(e)}")
@database_sync_to_async
def get_session_info(self, data):
"""获取会话信息和知识库"""
try:
conversation_id = data['conversation_id']
# 查找该会话ID下的历史记录
existing_records = ChatHistory.objects.filter(
conversation_id=conversation_id
).order_by('created_at')
# 如果有历史记录使用第一条记录的metadata
if existing_records.exists():
first_record = existing_records.first()
metadata = first_record.metadata or {}
# 获取知识库信息
dataset_ids = metadata.get('dataset_id_list', [])
external_id_list = metadata.get('dataset_external_id_list', [])
if not dataset_ids:
logger.error('找不到会话关联的知识库信息')
return None
# 验证知识库是否存在且用户有权限
knowledge_bases = []
for kb_id in dataset_ids:
try:
kb = KnowledgeBase.objects.get(id=kb_id)
if not self.check_knowledge_base_permission(kb, self.scope["user"], 'read'):
logger.error(f'无权访问知识库: {kb.name}')
return None
knowledge_bases.append(kb)
except KnowledgeBase.DoesNotExist:
logger.error(f'知识库不存在: {kb_id}')
return None
if not external_id_list or not knowledge_bases:
logger.error('会话关联的知识库信息不完整')
return None
return knowledge_bases, metadata, external_id_list
else:
# 如果是新会话的第一条记录需要提供知识库ID
dataset_ids = []
if 'dataset_id' in data:
dataset_ids.append(str(data['dataset_id']))
elif 'dataset_id_list' in data and isinstance(data['dataset_id_list'], (list, str)):
if isinstance(data['dataset_id_list'], str):
try:
dataset_list = json.loads(data['dataset_id_list'])
if isinstance(dataset_list, list):
dataset_ids = [str(id) for id in dataset_list]
else:
dataset_ids = [str(data['dataset_id_list'])]
except json.JSONDecodeError:
dataset_ids = [str(data['dataset_id_list'])]
else:
dataset_ids = [str(id) for id in data['dataset_id_list']]
if not dataset_ids:
logger.error('新会话需要提供知识库ID')
return None
# 验证所有知识库并收集external_ids
external_id_list = []
knowledge_bases = []
for kb_id in dataset_ids:
try:
knowledge_base = KnowledgeBase.objects.filter(id=kb_id).first()
if not knowledge_base:
logger.error(f'知识库不存在: {kb_id}')
return None
knowledge_bases.append(knowledge_base)
# 使用统一的权限检查方法
if not self.check_knowledge_base_permission(knowledge_base, self.scope["user"], 'read'):
logger.error(f'无权访问知识库: {knowledge_base.name}')
return None
# 添加知识库的external_id到列表
if knowledge_base.external_id:
external_id_list.append(str(knowledge_base.external_id))
else:
logger.warning(f"知识库 {knowledge_base.id} ({knowledge_base.name}) 没有external_id")
except Exception as e:
logger.error(f"处理知识库ID出错: {str(e)}")
return None
if not external_id_list:
logger.error('没有有效的知识库external_id')
return None
# 创建metadata
metadata = {
'model_id': data.get('model_id', '7a214d0e-e65e-11ef-9f4a-0242ac120006'),
'dataset_id_list': [str(id) for id in dataset_ids],
'dataset_external_id_list': [str(id) for id in external_id_list],
'dataset_names': [kb.name for kb in knowledge_bases]
}
return knowledge_bases, metadata, external_id_list
except Exception as e:
logger.error(f"获取会话信息时出错: {str(e)}")
return None
def check_knowledge_base_permission(self, kb, user, permission_type):
"""检查知识库权限"""
# 实现权限检查逻辑
return True # 临时返回 True需要根据实际情况实现
@database_sync_to_async
def create_question_record(self, conversation_id, question, knowledge_bases, metadata):
"""创建问题记录"""
try:
title = metadata.get('title', 'New chat')
# 创建用户问题记录
return ChatHistory.objects.create(
user=self.scope["user"],
knowledge_base=knowledge_bases[0], # 使用第一个知识库作为主知识库
conversation_id=str(conversation_id),
title=title,
role='user',
content=question,
metadata=metadata
)
except Exception as e:
logger.error(f"创建问题记录时出错: {str(e)}")
return None
@database_sync_to_async
def create_answer_record(self, conversation_id, question_record, knowledge_bases, metadata):
"""创建AI回答记录"""
try:
return ChatHistory.objects.create(
user=self.scope["user"],
knowledge_base=knowledge_bases[0],
conversation_id=str(conversation_id),
title=question_record.title,
parent_id=str(question_record.id),
role='assistant',
content="", # 初始内容为空
metadata=metadata
)
except Exception as e:
logger.error(f"创建回答记录时出错: {str(e)}")
return None
async def stream_from_external_api(self, conversation_id, question, dataset_external_id_list, answer_record, metadata, knowledge_bases):
"""从外部API获取流式响应"""
try:
# 确保所有ID都是字符串
dataset_external_ids = [str(id) if isinstance(id, uuid.UUID) else id for id in dataset_external_id_list]
# 获取标题
title = answer_record.title or 'New chat'
# 异步收集完整内容,用于最后保存
full_content = ""
# 使用aiohttp进行异步HTTP请求
async with aiohttp.ClientSession() as session:
# 第一步: 创建聊天会话
async with session.post(
f"{settings.API_BASE_URL}/api/application/chat/open",
json={
"id": "d5d11efa-ea9a-11ef-9933-0242ac120006",
"model_id": metadata.get('model_id', '7a214d0e-e65e-11ef-9f4a-0242ac120006'),
"dataset_id_list": dataset_external_ids,
"multiple_rounds_dialogue": False,
"dataset_setting": {
"top_n": 10, "similarity": "0.3",
"max_paragraph_char_number": 10000,
"search_mode": "blend",
"no_references_setting": {
"value": "{question}",
"status": "ai_questioning"
}
},
"model_setting": {
"prompt": "**相关文档内容**{data} **回答要求**:如果相关文档内容中没有可用信息,请回答\"没有在知识库中查找到相关信息,建议咨询相关技术支持或参考官方文档进行操作\"。请根据相关文档内容回答用户问题。不要输出与用户问题无关的内容。请使用中文回答客户问题。**用户问题**{question}"
},
"problem_optimization": False
}
) as chat_response:
if chat_response.status != 200:
error_msg = f"外部API调用失败: {await chat_response.text()}"
logger.error(error_msg)
await self.send_error(error_msg)
return
chat_data = await chat_response.json()
if chat_data.get('code') != 200 or not chat_data.get('data'):
error_msg = f"外部API返回错误: {chat_data}"
logger.error(error_msg)
await self.send_error(error_msg)
return
chat_id = chat_data['data']
logger.info(f"成功创建聊天会话, chat_id: {chat_id}")
# 第二步: 建立流式连接
message_url = f"{settings.API_BASE_URL}/api/application/chat_message/{chat_id}"
logger.info(f"开始流式请求: {message_url}")
# 创建流式请求
async with session.post(
url=message_url,
json={"message": question, "re_chat": False, "stream": True},
headers={"Content-Type": "application/json"}
) as message_request:
if message_request.status != 200:
error_msg = f"外部API聊天消息调用失败: {message_request.status}, {await message_request.text()}"
logger.error(error_msg)
await self.send_error(error_msg)
return
# 创建一个缓冲区以处理分段的数据
buffer = ""
# 读取并处理每个响应块
logger.info("开始处理流式响应")
async for chunk in message_request.content.iter_any():
chunk_str = chunk.decode('utf-8')
buffer += chunk_str
# 检查是否有完整的数据行
while '\n\n' in buffer:
parts = buffer.split('\n\n', 1)
line = parts[0]
buffer = parts[1]
if line.startswith('data: '):
try:
# 提取JSON数据
json_str = line[6:] # 去掉 "data: " 前缀
data = json.loads(json_str)
# 记录并处理部分响应
if 'content' in data:
content_part = data['content']
full_content += content_part
# 发送部分内容
await self.send_json({
'code': 200,
'message': 'partial',
'data': {
'id': str(answer_record.id),
'conversation_id': str(conversation_id),
'content': content_part,
'is_end': data.get('is_end', False)
}
})
# 处理结束标记
if data.get('is_end', False):
logger.info("收到流式响应结束标记")
# 保存完整内容
await self.update_answer_content(answer_record.id, full_content.strip())
# 处理标题
title = await self.get_or_generate_title(
conversation_id,
question,
full_content.strip()
)
# 发送最终响应
await self.send_json({
'code': 200,
'message': '完成',
'data': {
'id': str(answer_record.id),
'conversation_id': str(conversation_id),
'title': title,
'dataset_id_list': metadata.get('dataset_id_list', []),
'dataset_names': metadata.get('dataset_names', []),
'role': 'assistant',
'content': full_content.strip(),
'created_at': answer_record.created_at.strftime('%Y-%m-%d %H:%M:%S'),
'is_end': True
}
})
return
except json.JSONDecodeError as e:
logger.error(f"JSON解析错误: {e}, 数据: {line}")
continue
except Exception as e:
logger.error(f"流式处理出错: {str(e)}")
logger.error(traceback.format_exc())
await self.send_error(str(e))
# 保存已收集的内容
if 'full_content' in locals() and full_content:
try:
await self.update_answer_content(answer_record.id, full_content.strip())
except Exception as save_error:
logger.error(f"保存部分内容失败: {str(save_error)}")
@database_sync_to_async
def update_answer_content(self, answer_id, content):
"""更新回答内容"""
try:
answer_record = ChatHistory.objects.get(id=answer_id)
answer_record.content = content
answer_record.save()
return True
except Exception as e:
logger.error(f"更新回答内容失败: {str(e)}")
return False
@database_sync_to_async
def get_or_generate_title(self, conversation_id, question, answer):
"""获取或生成对话标题"""
try:
# 先检查是否已有标题
current_title = ChatHistory.objects.filter(
conversation_id=str(conversation_id)
).exclude(
title__in=["New chat", "新对话", ""]
).values_list('title', flat=True).first()
if current_title:
return current_title
# 简单的标题生成逻辑 (可替换为调用DeepSeek API生成标题)
generated_title = question[:20] + "..." if len(question) > 20 else question
# 更新所有相关记录的标题
ChatHistory.objects.filter(
conversation_id=str(conversation_id)
).update(title=generated_title)
return generated_title
except Exception as e:
logger.error(f"获取或生成标题失败: {str(e)}")
return "新对话"
async def send_json(self, content):
"""发送JSON格式的消息"""
await self.send(text_data=json.dumps(content))
async def send_error(self, message):
"""发送错误消息"""
await self.send_json({
'code': 500,
'message': message,
'data': {'is_end': True}
})

View File

@ -6,47 +6,124 @@ from rest_framework.authtoken.models import Token
User = get_user_model() User = get_user_model()
class Command(BaseCommand): class Command(BaseCommand):
help = '创建测试用户:1个管理员2个组长4个组员' help = '创建测试用户:4个管理员7个组长4个组员'
def handle(self, *args, **kwargs): def handle(self, *args, **kwargs):
# 创建管理员 - 技术部管理员 # 创建管理员 - 4个管理员
admin, created = User.objects.get_or_create( admins = [
username='admin', {
defaults={ 'username': 'admin1',
'email': 'admin@example.com', 'password': 'admin123',
'name': '张管理', 'email': 'admin1@example.com',
'name': '张技术管理',
'department': '技术部门',
'role': 'admin',
},
{
'username': 'admin2',
'password': 'admin123',
'email': 'admin2@example.com',
'name': '王产品管理',
'department': '产品部门',
'role': 'admin',
},
{
'username': 'admin3',
'password': 'admin123',
'email': 'admin3@example.com',
'name': '李商务管理',
'department': '商务部门',
'role': 'admin',
},
{
'username': 'admin4',
'password': 'admin123',
'email': 'admin4@example.com',
'name': '赵HR管理',
'department': 'HR',
'role': 'admin', 'role': 'admin',
'is_staff': True,
'is_superuser': True,
'last_login': timezone.now()
} }
) ]
if created:
admin.set_password('admin123')
admin.save()
token = Token.objects.create(user=admin)
self.stdout.write(self.style.SUCCESS(
f'成功创建管理员用户: {admin.username}{admin.name}, Token: {token.key}'
))
else:
self.stdout.write(self.style.WARNING(f'管理员用户已存在: {admin.username}'))
# 创建组长 - 研发部组长和测试部组长 for admin_data in admins:
admin, created = User.objects.get_or_create(
username=admin_data['username'],
defaults={
'email': admin_data['email'],
'name': admin_data['name'],
'role': admin_data['role'],
'department': admin_data['department'],
'is_staff': True,
'is_superuser': True,
'last_login': timezone.now()
}
)
if created:
admin.set_password(admin_data['password'])
admin.save()
token = Token.objects.create(user=admin)
self.stdout.write(self.style.SUCCESS(
f'成功创建管理员用户: {admin.username}{admin.name}, Token: {token.key}'
))
else:
self.stdout.write(self.style.WARNING(f'管理员用户已存在: {admin.username}'))
# 创建组长 - 7个部门的组长
leaders = [ leaders = [
{ {
'username': 'leader1', 'username': 'leader1',
'password': 'leader123', 'password': 'leader123',
'email': 'leader1@example.com', 'email': 'leader1@example.com',
'name': '李研发', 'name': '陈达人',
'department': '研发部', 'department': '达人部门',
'role': 'leader' 'role': 'leader'
}, },
{ {
'username': 'leader2', 'username': 'leader2',
'password': 'leader123', 'password': 'leader123',
'email': 'leader2@example.com', 'email': 'leader2@example.com',
'name': '王测试', 'name': '刘商务',
'department': '测试部', 'department': '商务部门',
'role': 'leader'
},
{
'username': 'leader3',
'password': 'leader123',
'email': 'leader3@example.com',
'name': '杨样本',
'department': '样本中心',
'role': 'leader'
},
{
'username': 'leader4',
'password': 'leader123',
'email': 'leader4@example.com',
'name': '黄产品',
'department': '产品部门',
'role': 'leader'
},
{
'username': 'leader5',
'password': 'leader123',
'email': 'leader5@example.com',
'name': '周AI',
'department': 'AI自媒体',
'role': 'leader'
},
{
'username': 'leader6',
'password': 'leader123',
'email': 'leader6@example.com',
'name': '吴HR',
'department': 'HR',
'role': 'leader'
},
{
'username': 'leader7',
'password': 'leader123',
'email': 'leader7@example.com',
'name': '郑技术',
'department': '技术部门',
'role': 'leader' 'role': 'leader'
} }
] ]
@ -73,67 +150,5 @@ class Command(BaseCommand):
else: else:
self.stdout.write(self.style.WARNING(f'组长用户已存在: {leader.username}')) self.stdout.write(self.style.WARNING(f'组长用户已存在: {leader.username}'))
# 创建组员 - 2个开发组员2个测试组员
members = [
{
'username': 'member1',
'password': 'member123',
'email': 'member1@example.com',
'name': '赵开发',
'department': '研发部',
'role': 'member',
'group': '前端组'
},
{
'username': 'member2',
'password': 'member123',
'email': 'member2@example.com',
'name': '钱开发',
'department': '研发部',
'role': 'member',
'group': '后端组'
},
{
'username': 'member3',
'password': 'member123',
'email': 'member3@example.com',
'name': '孙测试',
'department': '测试部',
'role': 'member',
'group': '功能测试组'
},
{
'username': 'member4',
'password': 'member123',
'email': 'member4@example.com',
'name': '周测试',
'department': '测试部',
'role': 'member',
'group': '自动化测试组'
}
]
for member_data in members:
member, created = User.objects.get_or_create(
username=member_data['username'],
defaults={
'email': member_data['email'],
'name': member_data['name'],
'role': member_data['role'],
'department': member_data['department'],
'group': member_data['group'],
'is_staff': False,
'last_login': timezone.now()
}
)
if created:
member.set_password(member_data['password'])
member.save()
token = Token.objects.create(user=member)
self.stdout.write(self.style.SUCCESS(
f'成功创建组员用户: {member.username}{member.name}, Token: {token.key}'
))
else:
self.stdout.write(self.style.WARNING(f'组员用户已存在: {member.username}'))
self.stdout.write(self.style.SUCCESS('所有测试用户创建完成!')) self.stdout.write(self.style.SUCCESS('所有测试用户创建完成!'))

View File

@ -4,6 +4,10 @@ from django.contrib.auth.models import AnonymousUser
from rest_framework.authtoken.models import Token from rest_framework.authtoken.models import Token
from django.contrib.auth import get_user_model from django.contrib.auth import get_user_model
import logging import logging
import re
from django.middleware.csrf import CsrfViewMiddleware
from django.conf import settings
from django.utils.deprecation import MiddlewareMixin
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
@ -44,11 +48,28 @@ class TokenAuthMiddleware(BaseMiddleware):
scope['user'] = AnonymousUser() scope['user'] = AnonymousUser()
return await super().__call__(scope, receive, send) return await super().__call__(scope, receive, send)
class UserActivityMiddleware: class UserActivityMiddleware(MiddlewareMixin):
"""用户活动中间件""" """中间件用于记录用户活动"""
def __init__(self, get_response):
self.get_response = get_response
def __call__(self, request): def process_request(self, request):
response = self.get_response(request) # 可以在这里记录用户活动日志
return response pass
class CSRFExemptMiddleware(MiddlewareMixin):
"""为特定URL路径豁免CSRF保护的中间件"""
def process_view(self, request, callback, callback_args, callback_kwargs):
# 检查是否有CSRF豁免URL配置
if not hasattr(settings, 'CSRF_EXEMPT_URLS'):
return None
# 获取当前请求的路径
path = request.path_info.lstrip('/')
# 检查是否匹配任何豁免模式
for exempt_pattern in settings.CSRF_EXEMPT_URLS:
if re.match(exempt_pattern, path):
setattr(request, '_dont_enforce_csrf_checks', True)
break
return None

View File

@ -0,0 +1,18 @@
# Generated by Django 5.1.5 on 2025-04-23 14:20
from django.db import migrations, models
class Migration(migrations.Migration):
dependencies = [
('user_management', '0004_knowledgebasedocument'),
]
operations = [
migrations.AddField(
model_name='chathistory',
name='title',
field=models.CharField(blank=True, default='New chat', help_text='对话标题', max_length=100, null=True),
),
]

View File

@ -0,0 +1,18 @@
# Generated by Django 5.1.5 on 2025-04-23 16:51
from django.db import migrations, models
class Migration(migrations.Migration):
dependencies = [
('user_management', '0005_chathistory_title'),
]
operations = [
migrations.AddField(
model_name='knowledgebasedocument',
name='uploader_name',
field=models.CharField(default='未知用户', max_length=100, verbose_name='上传者姓名'),
),
]

View File

@ -291,6 +291,8 @@ class ChatHistory(models.Model):
knowledge_base = models.ForeignKey('KnowledgeBase', on_delete=models.CASCADE) knowledge_base = models.ForeignKey('KnowledgeBase', on_delete=models.CASCADE)
# 用于标识知识库组合的对话 # 用于标识知识库组合的对话
conversation_id = models.CharField(max_length=100, db_index=True) conversation_id = models.CharField(max_length=100, db_index=True)
# 对话标题
title = models.CharField(max_length=100, null=True, blank=True, default='New chat', help_text="对话标题")
parent_id = models.CharField(max_length=100, null=True, blank=True) parent_id = models.CharField(max_length=100, null=True, blank=True)
role = models.CharField(max_length=20, choices=ROLE_CHOICES) role = models.CharField(max_length=20, choices=ROLE_CHOICES)
content = models.TextField() content = models.TextField()
@ -391,6 +393,7 @@ class ChatHistory(models.Model):
] ]
} }
class UserProfile(models.Model): class UserProfile(models.Model):
"""用户档案模型""" """用户档案模型"""
user = models.OneToOneField(User, on_delete=models.CASCADE, related_name='profile') user = models.OneToOneField(User, on_delete=models.CASCADE, related_name='profile')
@ -682,6 +685,7 @@ class KnowledgeBaseDocument(models.Model):
document_id = models.CharField(max_length=100, verbose_name='文档ID') document_id = models.CharField(max_length=100, verbose_name='文档ID')
document_name = models.CharField(max_length=255, verbose_name='文档名称') document_name = models.CharField(max_length=255, verbose_name='文档名称')
external_id = models.CharField(max_length=100, verbose_name='外部文档ID') external_id = models.CharField(max_length=100, verbose_name='外部文档ID')
uploader_name = models.CharField(max_length=100, default="未知用户", verbose_name='上传者姓名')
status = models.CharField( status = models.CharField(
max_length=20, max_length=20,
default='active', default='active',

View File

@ -3,4 +3,6 @@ from . import consumers
websocket_urlpatterns = [ websocket_urlpatterns = [
re_path(r'ws/notifications/$', consumers.NotificationConsumer.as_asgi()), re_path(r'ws/notifications/$', consumers.NotificationConsumer.as_asgi()),
re_path(r'ws/chat/$', consumers.ChatConsumer.as_asgi()),
re_path(r'ws/chat/stream/$', consumers.ChatStreamConsumer.as_asgi()),
] ]

File diff suppressed because it is too large Load Diff