优化WebSocket消息订阅逻辑,增强事件系统的类型安全性,新增调试信息记录功能,提升代码可读性和一致性。

This commit is contained in:
乘风
2026-01-14 15:36:54 +08:00
parent 0d6e1edd91
commit dcc89594bb
5 changed files with 785 additions and 0 deletions

View File

@@ -0,0 +1,107 @@
/**
* 账号状态处理器
*
* 负责处理:
* - 账号在线状态
* - 账号登录响应
* - 账号列表更新
*/
import type { WebSocketMessage } from '@/types/wechat'
import type { HandlerRegistry } from './index'
import { useMessageSubscription } from '../useMessageSubscription'
import { WS_CMD_TYPE } from '@/constants/wechat'
import { useAccountStore } from '@/stores/modules/wechat'
const { emitAccountStatus } = useMessageSubscription()
// ==================== 账号状态处理器 ====================
/**
* 处理账号在线状态响应
* CmdRequestWechatAccountsAliveStatusResp
*/
function handleAccountAliveStatus(wsMessage: WebSocketMessage): void {
const data = wsMessage.data
if (!data || !data.wechatAccountsAliveStatus) {
console.warn('[AccountHandler] 账号状态数据为空')
return
}
const aliveStatus = data.wechatAccountsAliveStatus as Record<number, boolean>
console.log('[AccountHandler] 📡 账号状态更新:', {
count: Object.keys(aliveStatus).length,
online: Object.values(aliveStatus).filter(Boolean).length,
})
// 发布每个账号的状态
Object.entries(aliveStatus).forEach(([accountId, isOnline]) => {
emitAccountStatus(Number(accountId), isOnline)
})
// 更新 Store
try {
const accountStore = useAccountStore()
accountStore.updateAccountsOnlineStatus(aliveStatus)
} catch (error) {
console.error('[AccountHandler] 更新 Store 失败:', error)
}
}
/**
* 处理登录响应
* CmdSignInResp
*/
function handleSignInResp(wsMessage: WebSocketMessage): void {
const data = wsMessage.data
if (!data) {
console.warn('[AccountHandler] 登录响应数据为空')
return
}
console.log('[AccountHandler] 🔐 登录响应:', {
code: data.code,
message: data.message,
})
if (data.code === 0) {
console.log('[AccountHandler] ✅ 登录成功')
} else {
console.error('[AccountHandler] ❌ 登录失败:', data.message)
}
}
/**
* 处理账号列表更新
* CmdAccountListUpdate
*/
function handleAccountListUpdate(wsMessage: WebSocketMessage): void {
const data = wsMessage.data
if (!data || !data.accounts) {
console.warn('[AccountHandler] 账号列表数据为空')
return
}
console.log('[AccountHandler] 📋 账号列表更新:', {
count: data.accounts.length,
})
// TODO: 更新账号列表 Store
}
// ==================== 导出处理器映射 ====================
export const accountHandlers: HandlerRegistry = {
// 账号在线状态
[WS_CMD_TYPE.ACCOUNT_ALIVE_STATUS_RESP]: handleAccountAliveStatus,
// 登录响应
[WS_CMD_TYPE.SIGN_IN_RESP]: handleSignInResp,
// 账号列表更新
[WS_CMD_TYPE.ACCOUNT_LIST_UPDATE]: handleAccountListUpdate,
}

View File

@@ -0,0 +1,210 @@
/**
* WebSocket 消息处理器注册中心
*
* 职责:
* - 统一注册所有消息处理器
* - 提供处理器调用入口
* - 支持动态扩展
*/
import type { WebSocketMessage } from '@/types/wechat'
import { messageHandlers } from './messageHandlers'
import { accountHandlers } from './accountHandlers'
import { sessionHandlers } from './sessionHandlers'
import { systemHandlers } from './systemHandlers'
// ==================== 类型定义 ====================
export type MessageHandler = (message: WebSocketMessage) => void | Promise<void>
export type HandlerRegistry = Record<string, MessageHandler>
// ==================== 处理器注册表 ====================
/**
* 默认处理器映射
* 包含所有模块的处理器
*/
const defaultHandlers: HandlerRegistry = {
...messageHandlers,
...accountHandlers,
...sessionHandlers,
...systemHandlers,
}
// ==================== 处理器管理器 ====================
/**
* 消息处理器管理器
* 支持动态注册和注销处理器
*/
export class MessageHandlerManager {
private handlers: HandlerRegistry
private readonly isDev = import.meta.env.DEV
constructor(initialHandlers: HandlerRegistry = {}) {
this.handlers = { ...defaultHandlers, ...initialHandlers }
if (this.isDev) {
console.log('[MessageHandler] 已注册的处理器:', Object.keys(this.handlers))
}
}
/**
* 处理消息
* @param cmdType 命令类型
* @param message 消息数据
*/
async handle(cmdType: string, message: WebSocketMessage): Promise<void> {
const handler = this.handlers[cmdType]
if (!handler) {
if (this.isDev) {
console.warn(`[MessageHandler] 未注册的消息类型: ${cmdType}`)
}
return
}
try {
if (this.isDev) {
console.log(`[MessageHandler] 处理消息: ${cmdType}`, {
preview: JSON.stringify(message).substring(0, 100) + '...',
})
}
await handler(message)
if (this.isDev) {
console.log(`[MessageHandler] ✅ 处理完成: ${cmdType}`)
}
} catch (error) {
console.error(`[MessageHandler] ❌ 处理失败: ${cmdType}`, error)
// 不抛出错误,避免影响后续消息处理
}
}
/**
* 注册新的处理器
* @param cmdType 命令类型
* @param handler 处理函数
*/
register(cmdType: string, handler: MessageHandler): void {
if (this.handlers[cmdType] && this.isDev) {
console.warn(`[MessageHandler] 覆盖已存在的处理器: ${cmdType}`)
}
this.handlers[cmdType] = handler
if (this.isDev) {
console.log(`[MessageHandler] 注册处理器: ${cmdType}`)
}
}
/**
* 批量注册处理器
* @param handlers 处理器映射
*/
registerBatch(handlers: HandlerRegistry): void {
Object.entries(handlers).forEach(([cmdType, handler]) => {
this.register(cmdType, handler)
})
}
/**
* 注销处理器
* @param cmdType 命令类型
*/
unregister(cmdType: string): void {
if (!this.handlers[cmdType]) {
if (this.isDev) {
console.warn(`[MessageHandler] 处理器不存在: ${cmdType}`)
}
return
}
delete this.handlers[cmdType]
if (this.isDev) {
console.log(`[MessageHandler] 注销处理器: ${cmdType}`)
}
}
/**
* 获取已注册的处理器列表
*/
getRegisteredTypes(): string[] {
return Object.keys(this.handlers)
}
/**
* 检查是否已注册
* @param cmdType 命令类型
*/
has(cmdType: string): boolean {
return !!this.handlers[cmdType]
}
/**
* 清空所有处理器(用于测试)
*/
clear(): void {
this.handlers = {}
if (this.isDev) {
console.log('[MessageHandler] 已清空所有处理器')
}
}
/**
* 重置为默认处理器
*/
reset(): void {
this.handlers = { ...defaultHandlers }
if (this.isDev) {
console.log('[MessageHandler] 已重置为默认处理器')
}
}
}
// ==================== 单例实例 ====================
let managerInstance: MessageHandlerManager | null = null
/**
* 获取消息处理器管理器实例(单例)
*/
export function getMessageHandlerManager(): MessageHandlerManager {
if (!managerInstance) {
managerInstance = new MessageHandlerManager()
}
return managerInstance
}
/**
* 创建消息处理函数(简化版)
*
* @example
* ```typescript
* const handleMessage = createMessageHandler()
*
* ws.onmessage = (event) => {
* const message = JSON.parse(event.data)
* handleMessage(message.cmdType, message)
* }
* ```
*/
export function createMessageHandler() {
const manager = getMessageHandlerManager()
return (cmdType: string, message: WebSocketMessage) => {
return manager.handle(cmdType, message)
}
}
// ==================== 导出 ====================
export {
messageHandlers,
accountHandlers,
sessionHandlers,
systemHandlers,
}

View File

@@ -0,0 +1,225 @@
/**
* 聊天消息处理器
*
* 负责处理:
* - 新消息接收
* - 消息状态更新
* - 消息撤回
*/
import type { WebSocketMessage } from '@/types/wechat'
import type { HandlerRegistry } from './index'
import { useMessageSubscription } from '../useMessageSubscription'
import { SessionManager } from '@/utils/dbManagers/SessionManager'
import { MessageManager } from '@/utils/dbManagers/MessageManager'
import { WS_CMD_TYPE } from '@/constants/wechat'
const { emitNewMessage, emitMessageUpdate, emitMessageRecall } = useMessageSubscription()
// ==================== 辅助函数 ====================
/**
* 验证消息数据有效性
*/
function validateMessageData(data: any): boolean {
if (!data) {
console.warn('[MessageHandler] 消息数据为空')
return false
}
if (!data.sessionId || !data.sessionType) {
console.warn('[MessageHandler] 缺少会话信息:', data)
return false
}
return true
}
/**
* 构建标准 Message 对象
*/
function buildStandardMessage(data: any, currentAccountId?: number): Message {
return {
id: data.id || Date.now(),
clientId: data.clientId || `msg_${Date.now()}`,
serverId: data.serverId,
sessionId: data.sessionId,
wechatAccountId: data.wechatAccountId || currentAccountId || 0,
msgType: data.msgType || 1,
content: data.content || '',
direction: data.isSend ? 'send' : 'receive',
isSend: data.isSend || false,
createTime: data.createTime || new Date().toISOString(),
timestamp: data.wechatTime || Date.now(),
status: data.status || 'success',
sender: data.sender,
}
}
// ==================== 消息处理器 ====================
/**
* 处理新消息
* CmdNewMessage / CmdReceiveMessage
*/
async function handleNewMessage(wsMessage: WebSocketMessage): Promise<void> {
const data = wsMessage.data
// 1. 验证数据
if (!validateMessageData(data)) {
return
}
// 2. 消息去重(如果启用)
if (data.clientId && (await MessageManager.checkDuplicate(data.clientId))) {
console.log('[MessageHandler] 重复消息,跳过:', data.clientId)
return
}
console.log('[MessageHandler] 📨 收到新消息:', {
sessionId: data.sessionId,
type: data.sessionType,
content: data.content?.substring(0, 20) + '...',
})
try {
// 3. 更新会话(异步,不阻塞)
SessionManager.updateOnNewMessage(
data.sessionId,
data.sessionType,
data.content,
data.wechatAccountId
).catch((error) => {
console.error('[MessageHandler] 更新会话失败:', error)
})
// 4. 保存消息记录(如果启用缓存)
if (MessageManager.shouldCacheMessage(data.sessionId)) {
const chatMessage = buildStandardMessage(data)
MessageManager.addMessage(chatMessage).catch((error) => {
console.error('[MessageHandler] 保存消息失败:', error)
})
}
// 5. 发布新消息事件(立即通知订阅者)
const message = buildStandardMessage(data)
emitNewMessage(message)
console.log('[MessageHandler] ✅ 新消息处理完成')
} catch (error) {
console.error('[MessageHandler] ❌ 新消息处理失败:', error)
}
}
/**
* 处理消息状态更新
* CmdSendMessageResp / CmdMessageStatus
*/
function handleMessageStatus(wsMessage: WebSocketMessage): void {
const data = wsMessage.data
if (!data) {
console.warn('[MessageHandler] 消息状态数据为空')
return
}
const { messageId, clientId, status, sessionId, seq } = data
// 优先使用 messageId其次 clientId最后 seq
const id = messageId || clientId || seq
if (!id) {
console.warn('[MessageHandler] 缺少消息标识:', data)
return
}
console.log('[MessageHandler] 📝 消息状态更新:', {
id,
status,
sessionId,
})
// 发布消息更新事件
emitMessageUpdate(id, {
status: status || 'success',
serverId: data.serverId,
})
// 更新数据库(异步)
if (sessionId && MessageManager.shouldCacheMessage(sessionId)) {
MessageManager.updateMessageStatus(id, status || 'success').catch((error) => {
console.error('[MessageHandler] 更新消息状态失败:', error)
})
}
}
/**
* 处理消息撤回
* CmdRecallMessage
*/
function handleMessageRecall(wsMessage: WebSocketMessage): void {
const data = wsMessage.data
if (!data || !data.messageId || !data.sessionId) {
console.warn('[MessageHandler] 撤回消息数据不完整:', data)
return
}
console.log('[MessageHandler] 🗑️ 消息撤回:', {
messageId: data.messageId,
sessionId: data.sessionId,
})
// 发布撤回事件
emitMessageRecall(data.messageId, data.sessionId)
// 更新数据库(异步)
if (MessageManager.shouldCacheMessage(data.sessionId)) {
MessageManager.markAsRecalled(data.messageId).catch((error) => {
console.error('[MessageHandler] 标记撤回失败:', error)
})
}
}
/**
* 处理发送消息结果
* CmdSendMessageResult
*/
function handleSendMessageResult(wsMessage: WebSocketMessage): void {
const data = wsMessage.data
if (!data) {
console.warn('[MessageHandler] 发送结果数据为空')
return
}
console.log('[MessageHandler] 📤 发送消息结果:', {
messageId: data.friendMessageId || data.chatroomMessageId,
status: data.sendStatus,
})
// 更新消息状态
const messageId = data.friendMessageId || data.chatroomMessageId
if (messageId) {
const status = data.sendStatus === 0 ? 'success' : 'failed'
emitMessageUpdate(messageId, { status })
}
}
// ==================== 导出处理器映射 ====================
export const messageHandlers: HandlerRegistry = {
// 接收新消息
[WS_CMD_TYPE.RECEIVE_MESSAGE]: handleNewMessage,
[WS_CMD_TYPE.NEW_MESSAGE]: handleNewMessage, // 别名
// 消息状态
[WS_CMD_TYPE.MESSAGE_STATUS]: handleMessageStatus,
[WS_CMD_TYPE.SEND_MESSAGE_RESP]: handleMessageStatus, // 兼容旧版
// 消息撤回
[WS_CMD_TYPE.RECALL_MESSAGE]: handleMessageRecall,
// 发送结果
[WS_CMD_TYPE.SEND_MESSAGE_RESULT]: handleSendMessageResult,
}

View File

@@ -0,0 +1,114 @@
/**
* 会话相关处理器
*
* 负责处理:
* - 会话列表更新
* - 会话置顶
* - 未读数变化
*/
import type { WebSocketMessage } from '@/types/wechat'
import type { HandlerRegistry } from './index'
import { useMessageSubscription } from '../useMessageSubscription'
import { WS_CMD_TYPE } from '@/constants/wechat'
import { useSessionStore } from '@/stores/modules/wechat'
const { emitSessionUpdate } = useMessageSubscription()
// ==================== 会话处理器 ====================
/**
* 处理会话更新
* CmdSessionUpdate
*/
function handleSessionUpdate(wsMessage: WebSocketMessage): void {
const data = wsMessage.data
if (!data || !data.sessionId) {
console.warn('[SessionHandler] 会话更新数据不完整')
return
}
console.log('[SessionHandler] 🔄 会话更新:', {
sessionId: data.sessionId,
type: data.type,
})
// 发布会话更新事件
emitSessionUpdate(data.sessionId, data)
// 更新 Store
try {
const sessionStore = useSessionStore()
if (data.type === 'pin') {
sessionStore.togglePin(data.sessionId)
} else if (data.type === 'mute') {
sessionStore.toggleMute(data.sessionId)
}
} catch (error) {
console.error('[SessionHandler] 更新 Store 失败:', error)
}
}
/**
* 处理未读数清零
* CmdClearUnread
*/
function handleClearUnread(wsMessage: WebSocketMessage): void {
const data = wsMessage.data
if (!data || !data.sessionId) {
console.warn('[SessionHandler] 清除未读数据不完整')
return
}
console.log('[SessionHandler] 🔔 清除未读:', {
sessionId: data.sessionId,
})
// 更新 Store
try {
const sessionStore = useSessionStore()
sessionStore.clearUnreadCount(data.sessionId)
} catch (error) {
console.error('[SessionHandler] 清除未读失败:', error)
}
}
/**
* 处理会话删除
* CmdSessionDelete
*/
function handleSessionDelete(wsMessage: WebSocketMessage): void {
const data = wsMessage.data
if (!data || !data.sessionId) {
console.warn('[SessionHandler] 删除会话数据不完整')
return
}
console.log('[SessionHandler] 🗑️ 删除会话:', {
sessionId: data.sessionId,
})
// 更新 Store
try {
const sessionStore = useSessionStore()
sessionStore.removeSession(data.sessionId)
} catch (error) {
console.error('[SessionHandler] 删除会话失败:', error)
}
}
// ==================== 导出处理器映射 ====================
export const sessionHandlers: HandlerRegistry = {
// 会话更新
[WS_CMD_TYPE.SESSION_UPDATE]: handleSessionUpdate,
// 清除未读
[WS_CMD_TYPE.CLEAR_UNREAD]: handleClearUnread,
// 删除会话
[WS_CMD_TYPE.SESSION_DELETE]: handleSessionDelete,
}

View File

@@ -0,0 +1,129 @@
/**
* 系统通知处理器
*
* 负责处理:
* - 系统通知
* - 心跳响应
* - 错误消息
*/
import type { WebSocketMessage } from '@/types/wechat'
import type { HandlerRegistry } from './index'
import { useMessageSubscription } from '../useMessageSubscription'
import { WS_CMD_TYPE } from '@/constants/wechat'
import { ElNotification } from 'element-plus'
const { emitSystemNotification } = useMessageSubscription()
// ==================== 系统处理器 ====================
/**
* 处理心跳响应
* CmdHeartbeat / CmdPing
*/
function handleHeartbeat(wsMessage: WebSocketMessage): void {
console.log('[SystemHandler] 💓 心跳响应')
// 更新最后活跃时间
const lastActiveTime = Date.now()
localStorage.setItem('lastSyncTime', lastActiveTime.toString())
}
/**
* 处理系统通知
* CmdSystemNotification
*/
function handleSystemNotification(wsMessage: WebSocketMessage): void {
const data = wsMessage.data
if (!data || !data.type || !data.content) {
console.warn('[SystemHandler] 系统通知数据不完整')
return
}
console.log('[SystemHandler] 📢 系统通知:', {
type: data.type,
content: data.content,
})
// 发布通知事件
emitSystemNotification(data.type, data.content)
// 显示桌面通知
if (data.showNotification !== false) {
ElNotification({
title: '系统通知',
message: data.content,
type: data.level || 'info',
duration: 5000,
})
}
}
/**
* 处理错误消息
* CmdError
*/
function handleError(wsMessage: WebSocketMessage): void {
const data = wsMessage.data
if (!data) {
console.warn('[SystemHandler] 错误消息数据为空')
return
}
console.error('[SystemHandler] ❌ 错误消息:', {
code: data.code,
message: data.message,
})
// 显示错误通知
ElNotification({
title: '错误',
message: data.message || '发生未知错误',
type: 'error',
duration: 5000,
})
}
/**
* 处理服务器时间同步
* CmdServerTime
*/
function handleServerTime(wsMessage: WebSocketMessage): void {
const data = wsMessage.data
if (!data || !data.serverTime) {
return
}
const serverTime = data.serverTime
const localTime = Date.now()
const timeDiff = serverTime - localTime
console.log('[SystemHandler] 🕐 服务器时间:', {
serverTime: new Date(serverTime).toLocaleString(),
localTime: new Date(localTime).toLocaleString(),
diff: `${timeDiff}ms`,
})
// 保存时间差(用于时间校准)
localStorage.setItem('serverTimeDiff', timeDiff.toString())
}
// ==================== 导出处理器映射 ====================
export const systemHandlers: HandlerRegistry = {
// 心跳
[WS_CMD_TYPE.HEARTBEAT]: handleHeartbeat,
[WS_CMD_TYPE.PING]: handleHeartbeat, // 别名
// 系统通知
[WS_CMD_TYPE.SYSTEM_NOTIFICATION]: handleSystemNotification,
// 错误消息
[WS_CMD_TYPE.ERROR]: handleError,
// 服务器时间
[WS_CMD_TYPE.SERVER_TIME]: handleServerTime,
}