## 主要修复 - 修复Creo软件运行状态检测失败问题 - 添加完整的软件停止功能支持 - 改进多进程软件的进程管理逻辑 ## 技术改进 - 更新软件配置支持多进程名称检测 - 优化进程停止逻辑,增加超时配置 - 新增 stop_software WebSocket消息类型 - 完善错误处理和日志记录 ## 配置更新 - configs/software_config.yaml: 支持进程名称列表和停止超时 - 添加Revit 2017配置支持 ## 文档更新 - README.md: 更新软件配置说明和API列表 - frontend-api-docs.md: 添加停止软件API文档 - CHECKPOINT.md: 记录修复进展和解决方案 🤖 Generated with [Claude Code](https://claude.ai/code) Co-Authored-By: Claude <noreply@anthropic.com>
509 lines
20 KiB
Python
509 lines
20 KiB
Python
"""
|
||
WebSocket API路由
|
||
提供WebSocket连接端点
|
||
"""
|
||
from fastapi import APIRouter, WebSocket, WebSocketDisconnect, Depends, Query
|
||
from app.core.websocket_manager import websocket_manager, MessageType
|
||
import json
|
||
import uuid
|
||
import logging
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
# WebSocket消息类型常量
|
||
class WSMessageType:
|
||
PING = "ping"
|
||
GET_STATUS = "get_status"
|
||
GET_SOFTWARE_LIST = "get_software_list"
|
||
START_SOFTWARE = "start_software"
|
||
STOP_SOFTWARE = "stop_software"
|
||
RESTART_SOFTWARE = "restart_software"
|
||
LOG_OPERATION = "log_operation"
|
||
# 新增日志查询相关消息类型
|
||
QUERY_LOGS = "query_logs"
|
||
GET_LOG_BY_ID = "get_log_by_id"
|
||
GET_LOG_STATS = "get_log_stats"
|
||
CLEANUP_LOGS = "cleanup_logs"
|
||
GET_OPERATION_TYPES = "get_operation_types"
|
||
|
||
router = APIRouter()
|
||
|
||
|
||
@router.websocket("/connect")
|
||
async def websocket_endpoint(
|
||
websocket: WebSocket,
|
||
client_id: str = Query(None, description="客户端ID,如果不提供将自动生成"),
|
||
user_id: str = Query(None, description="用户ID,用于用户认证")
|
||
):
|
||
"""
|
||
WebSocket连接端点
|
||
|
||
连接参数:
|
||
- client_id: 客户端唯一标识符,如果不提供将自动生成
|
||
- user_id: 用户ID,用于消息推送和权限控制
|
||
|
||
连接URL示例:
|
||
ws://localhost:8000/api/v1/ws/connect?client_id=client123&user_id=user456
|
||
"""
|
||
|
||
# 如果没有提供client_id,自动生成一个
|
||
if not client_id:
|
||
client_id = str(uuid.uuid4())
|
||
|
||
try:
|
||
# 建立WebSocket连接
|
||
await websocket_manager.connect(websocket, client_id, user_id)
|
||
|
||
# 发送欢迎消息
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.INFO,
|
||
"message": f"欢迎连接!客户端ID: {client_id}",
|
||
"client_id": client_id,
|
||
"user_id": user_id,
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
|
||
# 保持连接并处理消息
|
||
while True:
|
||
try:
|
||
# 接收客户端消息
|
||
data = await websocket.receive_text()
|
||
|
||
# 解析JSON消息
|
||
try:
|
||
message = json.loads(data)
|
||
await handle_client_message(message, client_id, user_id)
|
||
except json.JSONDecodeError:
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.ERROR,
|
||
"message": "消息格式错误,请发送有效的JSON",
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
|
||
except WebSocketDisconnect:
|
||
logger.info(f"客户端 {client_id} 主动断开连接")
|
||
break
|
||
except Exception as e:
|
||
logger.error(f"处理客户端 {client_id} 消息时发生错误: {e}")
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.ERROR,
|
||
"message": f"服务器处理消息时发生错误: {str(e)}",
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
break
|
||
|
||
except Exception as e:
|
||
logger.error(f"WebSocket连接发生错误: {e}")
|
||
finally:
|
||
# 清理连接
|
||
websocket_manager.disconnect(client_id)
|
||
|
||
|
||
async def handle_client_message(message: dict, client_id: str, user_id: str):
|
||
"""
|
||
处理客户端发送的消息
|
||
|
||
支持的消息类型:
|
||
- ping: 心跳检测
|
||
- get_status: 获取服务状态
|
||
- get_software_list: 获取软件列表
|
||
- start_software: 启动软件
|
||
- stop_software: 停止软件
|
||
- restart_software: 重启软件
|
||
- log_operation: 记录用户操作日志
|
||
- query_logs: 查询操作日志
|
||
- get_log_by_id: 根据ID获取日志
|
||
- get_log_stats: 获取日志统计信息
|
||
- cleanup_logs: 清理过期日志
|
||
- get_operation_types: 获取操作类型列表
|
||
"""
|
||
from app.core.software_manager import software_manager
|
||
from app.core.log_manager import log_manager
|
||
from app.models.operation_log import ActionType, OperationStatus
|
||
|
||
message_type = message.get("type")
|
||
|
||
if message_type == WSMessageType.PING:
|
||
# 心跳响应
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.HEARTBEAT,
|
||
"message": "pong",
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
|
||
elif message_type == WSMessageType.GET_STATUS:
|
||
# 获取服务状态
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.INFO,
|
||
"message": "服务状态正常",
|
||
"data": {
|
||
"active_connections": websocket_manager.get_active_connections_count(),
|
||
"connected_users": websocket_manager.get_connected_users()
|
||
},
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
|
||
elif message_type == WSMessageType.GET_SOFTWARE_LIST:
|
||
# 获取软件列表
|
||
try:
|
||
software_list = await software_manager.get_software_list()
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.SOFTWARE_LIST_UPDATE,
|
||
"data": {"software_list": software_list},
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
except Exception as e:
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.ERROR,
|
||
"message": f"获取软件列表失败: {str(e)}",
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
|
||
elif message_type == WSMessageType.START_SOFTWARE:
|
||
# 启动软件
|
||
software_id = message.get("software_id")
|
||
if not software_id:
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.ERROR,
|
||
"message": "缺少参数: software_id",
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
return
|
||
|
||
try:
|
||
task = await software_manager.start_software(software_id)
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.INFO,
|
||
"message": f"软件 {software_id} 启动任务已创建",
|
||
"data": {
|
||
"task_id": task.id,
|
||
"software_id": software_id,
|
||
"status": task.status.value
|
||
},
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
except Exception as e:
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.ERROR,
|
||
"message": f"启动软件失败: {str(e)}",
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
|
||
elif message_type == WSMessageType.STOP_SOFTWARE:
|
||
# 停止软件
|
||
software_id = message.get("software_id")
|
||
if not software_id:
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.ERROR,
|
||
"message": "缺少参数: software_id",
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
return
|
||
|
||
try:
|
||
task = await software_manager.stop_software(software_id)
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.INFO,
|
||
"message": f"软件 {software_id} 停止任务已创建",
|
||
"data": {
|
||
"task_id": task.id,
|
||
"software_id": software_id,
|
||
"status": task.status.value
|
||
},
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
except Exception as e:
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.ERROR,
|
||
"message": f"停止软件失败: {str(e)}",
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
|
||
elif message_type == WSMessageType.RESTART_SOFTWARE:
|
||
# 重启软件
|
||
software_id = message.get("software_id")
|
||
if not software_id:
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.ERROR,
|
||
"message": "缺少参数: software_id",
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
return
|
||
|
||
try:
|
||
task = await software_manager.restart_software(software_id)
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.INFO,
|
||
"message": f"软件 {software_id} 重启任务已创建",
|
||
"data": {
|
||
"task_id": task.id,
|
||
"software_id": software_id,
|
||
"status": task.status.value
|
||
},
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
except Exception as e:
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.ERROR,
|
||
"message": f"重启软件失败: {str(e)}",
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
|
||
elif message_type == WSMessageType.LOG_OPERATION:
|
||
# 记录用户操作日志
|
||
operation = message.get("operation")
|
||
details = message.get("details", "")
|
||
|
||
if not operation:
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.ERROR,
|
||
"message": "缺少参数: operation",
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
return
|
||
|
||
try:
|
||
# 提取日志参数
|
||
log_params = {
|
||
"action_type": ActionType(message.get("action_type")) if message.get("action_type") else None,
|
||
"target_object": message.get("target_object"),
|
||
"software_version": message.get("software_version"),
|
||
"status": OperationStatus(message.get("status", "success")),
|
||
"duration": message.get("duration"),
|
||
"operation_category": message.get("operation_category"),
|
||
"extra_data": {
|
||
"file_path": message.get("file_path"),
|
||
"file_size": message.get("file_size"),
|
||
"batch_count": message.get("batch_count"),
|
||
"export_format": message.get("export_format"),
|
||
**{k: v for k, v in message.items() if k.startswith("custom_")}
|
||
}
|
||
}
|
||
|
||
# 清理空值
|
||
log_params = {k: v for k, v in log_params.items() if v is not None}
|
||
if log_params.get("extra_data"):
|
||
log_params["extra_data"] = {k: v for k, v in log_params["extra_data"].items() if v is not None}
|
||
if not log_params["extra_data"]: # 如果extra_data为空,删除它
|
||
del log_params["extra_data"]
|
||
|
||
# 记录日志
|
||
log_id = await log_manager.log_user_operation(
|
||
operation=operation,
|
||
details=details,
|
||
user_id=user_id or "anonymous",
|
||
client_id=client_id,
|
||
**log_params
|
||
)
|
||
|
||
# 响应成功
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.LOG_RECORDED,
|
||
"message": "操作日志已记录",
|
||
"data": {
|
||
"log_id": log_id,
|
||
"operation": operation
|
||
},
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
|
||
except Exception as e:
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.ERROR,
|
||
"message": f"记录日志失败: {str(e)}",
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
|
||
elif message_type == WSMessageType.QUERY_LOGS:
|
||
# 查询操作日志
|
||
try:
|
||
from app.models.operation_log import LogFilter, LogType, LogLevel
|
||
from datetime import datetime
|
||
|
||
# 构建查询过滤器
|
||
filter_params = LogFilter(
|
||
log_type=LogType(message.get("log_type")) if message.get("log_type") else None,
|
||
operation=message.get("operation"),
|
||
user_id=message.get("user_id_filter"), # 避免与当前user_id冲突
|
||
client_id=message.get("client_id_filter"),
|
||
level=LogLevel(message.get("level")) if message.get("level") else None,
|
||
start_time=datetime.fromisoformat(message.get("start_time")) if message.get("start_time") else None,
|
||
end_time=datetime.fromisoformat(message.get("end_time")) if message.get("end_time") else None,
|
||
limit=message.get("limit", 100),
|
||
offset=message.get("offset", 0)
|
||
)
|
||
|
||
logs = await log_manager.query_logs(filter_params)
|
||
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.INFO,
|
||
"message": "日志查询成功",
|
||
"data": {
|
||
"logs": [log.dict() for log in logs],
|
||
"total": len(logs),
|
||
"limit": filter_params.limit,
|
||
"offset": filter_params.offset
|
||
},
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
|
||
except Exception as e:
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.ERROR,
|
||
"message": f"查询日志失败: {str(e)}",
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
|
||
elif message_type == WSMessageType.GET_LOG_BY_ID:
|
||
# 根据ID获取日志
|
||
log_id = message.get("log_id")
|
||
if not log_id:
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.ERROR,
|
||
"message": "缺少参数: log_id",
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
return
|
||
|
||
try:
|
||
log = await log_manager.get_log_by_id(log_id)
|
||
if not log:
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.ERROR,
|
||
"message": "日志不存在",
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
return
|
||
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.INFO,
|
||
"message": "获取日志成功",
|
||
"data": {"log": log.dict()},
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
|
||
except Exception as e:
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.ERROR,
|
||
"message": f"获取日志失败: {str(e)}",
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
|
||
elif message_type == WSMessageType.GET_LOG_STATS:
|
||
# 获取日志统计信息
|
||
try:
|
||
from datetime import datetime, timedelta
|
||
from app.models.operation_log import LogFilter, LogType, LogLevel
|
||
|
||
now = datetime.now()
|
||
start_time = now - timedelta(hours=24)
|
||
|
||
# 查询不同类型的日志数量
|
||
system_filter = LogFilter(
|
||
log_type=LogType.SYSTEM_OPERATION,
|
||
start_time=start_time,
|
||
limit=1000
|
||
)
|
||
system_logs = await log_manager.query_logs(system_filter)
|
||
|
||
user_filter = LogFilter(
|
||
log_type=LogType.USER_OPERATION,
|
||
start_time=start_time,
|
||
limit=1000
|
||
)
|
||
user_logs = await log_manager.query_logs(user_filter)
|
||
|
||
error_filter = LogFilter(
|
||
level=LogLevel.ERROR,
|
||
start_time=start_time,
|
||
limit=1000
|
||
)
|
||
error_logs = await log_manager.query_logs(error_filter)
|
||
|
||
stats = {
|
||
"period": "24小时",
|
||
"system_operations": len(system_logs),
|
||
"user_operations": len(user_logs),
|
||
"error_logs": len(error_logs),
|
||
"total_logs": len(system_logs) + len(user_logs),
|
||
"timestamp": now.isoformat()
|
||
}
|
||
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.INFO,
|
||
"message": "获取统计信息成功",
|
||
"data": {"stats": stats},
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
|
||
except Exception as e:
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.ERROR,
|
||
"message": f"获取统计信息失败: {str(e)}",
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
|
||
elif message_type == WSMessageType.CLEANUP_LOGS:
|
||
# 清理过期日志
|
||
try:
|
||
await log_manager.cleanup_old_logs()
|
||
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.INFO,
|
||
"message": "过期日志清理完成",
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
|
||
except Exception as e:
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.ERROR,
|
||
"message": f"清理日志失败: {str(e)}",
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
|
||
elif message_type == WSMessageType.GET_OPERATION_TYPES:
|
||
# 获取操作类型列表
|
||
try:
|
||
from app.models.operation_log import LogFilter
|
||
|
||
# 查询最近的日志以获取操作类型
|
||
filter_params = LogFilter(limit=1000)
|
||
recent_logs = await log_manager.query_logs(filter_params)
|
||
|
||
# 统计操作类型
|
||
operations = set()
|
||
categories = set()
|
||
|
||
for log in recent_logs:
|
||
operations.add(log.operation)
|
||
if log.operation_category:
|
||
categories.add(log.operation_category)
|
||
|
||
result = {
|
||
"operations": sorted(list(operations)),
|
||
"categories": sorted(list(categories)),
|
||
"total_operations": len(operations),
|
||
"total_categories": len(categories)
|
||
}
|
||
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.INFO,
|
||
"message": "获取操作类型成功",
|
||
"data": result,
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
|
||
except Exception as e:
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.ERROR,
|
||
"message": f"获取操作类型失败: {str(e)}",
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id)
|
||
|
||
else:
|
||
# 未知消息类型
|
||
await websocket_manager.send_personal_message({
|
||
"type": MessageType.ERROR,
|
||
"message": f"未知的消息类型: {message_type}",
|
||
"timestamp": websocket_manager._get_timestamp()
|
||
}, client_id) |