// message-service.uts // 实时消息中枢:统一接管 WebSocket 连接、消息分发与会话管理。 // // 职责: // 1. 在用户已登录且已拿到后端 appConf.socketUrl 后,自动注入 token 建立连接; // 2. 监听 wsManager 的 open/message/close/error/reconnect 事件; // 3. 将服务器下发的消息按 header.type 分发为应用级事件(uni.$emit), // 各页面(聊天/通知/系统)只需监听对应事件即可,无需关心底层 WS; // 4. 维护内存中的会话列表与消息列表,提供排序、未读、置顶、免打扰等会话管理; // 5. 提供 sendChat() 便捷发送单聊/群聊消息。 // // 注意:连接鉴权采用与后端约定一致的 Authorization 头(Bearer token), // 由 wsManager.connect 透传到 uni.connectSocket 的 header 中。 import { wsManager } from './websocket.uts' import { IMessageBase, IMessageHeader, IMessagePayload, IMessagePreview, IMessageSession, IMessageServiceCallbacks, IChatMessage, IConnectOptions } from './socket-types.uts' import { getUserToken } from '@/stores/user.uts' import userState from '@/stores/user.uts' class MessageService { private initialized = false private connected = false private token = '' // 会话管理状态 private callbacks : IMessageServiceCallbacks = {} private sessions : Map = new Map() private messages : Map = new Map() // session_id -> messages[] private currentUserId = userState.uid // 注册 WS 事件监听(幂等,只执行一次) init() : void { if (this.initialized == true) { return } this.initialized = true wsManager.on({ onOpen: (res : any) => { this.connected = true console.log('[MessageService] WS 已连接') uni.$emit('onWsConnected', res) }, onMessage: (msg : IMessageBase) => { // 1) 维护会话与消息列表 this.handleWebSocketMessage(msg) // 2) 应用级事件分发 this.dispatch(msg) }, onClose: (res : any) => { this.connected = false console.log('[MessageService] WS 已关闭') uni.$emit('onWsDisconnected', res) }, onError: (err : any) => { console.error('[MessageService] WS 错误', err) uni.$emit('onWsError', err) }, onReconnect: (res : any) => { console.log('[MessageService] WS 重连中', res) uni.$emit('onWsReconnect', res) }, onReconnectFailed: (res : any) => { console.error('[MessageService] WS 重连失败', res) uni.$emit('onWsReconnectFailed', res) } }) } /** * 已登录时按需建立 WS 连接(幂等)。 * @param socketUrl 后端 appConf.socketUrl,为空则跳过 */ connectIfNeeded(socketUrl : string | null) : void { this.init() if (socketUrl == null || socketUrl.length == 0) { console.warn('[MessageService] 未配置 socketUrl,跳过 WS 连接') return } const tokenResult = getUserToken() if (tokenResult == null || tokenResult.accessToken == null || tokenResult.accessToken.length == 0) { console.warn('[MessageService] 未登录或 token 为空,跳过 WS 连接') return } if (wsManager.isConnected() == true) { return } this.token = tokenResult.accessToken! const options : IConnectOptions = { reconnectInterval: 3000, heartbeatInterval: 20000, maxReconnectCount: 5, header: { 'Authorization': `Bearer ${this.token}` } } wsManager.connect(socketUrl, options) } // 主动断开 disconnect() : void { wsManager.close(1000, 'user logout') this.connected = false } isConnected() : boolean { return this.connected } /** * 发送聊天消息(单聊/群聊)。 * 真实送达以服务器回执为准,本方法仅保证消息已提交至 WS。 */ sendChat(receiver : number, sessionId : string, contentType : string, content : any) : void { const header : IMessageHeader = { id: `c_${Date.now()}_${Math.floor(Math.random() * 1000)}`, type: 'chat', target: 'single', seqno: 1, timestamp: Date.now(), source: 'client', version: '1.0.0', token: this.token } const msg : IMessageBase = { header: header, payload: { uid: userState.uid, receiver: receiver, session_id: sessionId, type: contentType, status: 'sending', content: content } as any } wsManager.send(msg) } // 按消息类型分发为应用级事件 private dispatch(msg : IMessageBase) : void { const type = msg.header.type switch (type) { case 'chat': uni.$emit('onChatMessage', msg) break case 'notify': uni.$emit('onWsNotify', msg) break case 'status': uni.$emit('onWsStatus', msg) break case 'system': uni.$emit('onWsSystem', msg) break case 'pong': // 心跳回包,wsManager 内部已处理,无需外发 break default: uni.$emit('onWsMessage', msg) } } /** * 处理 WebSocket 消息,维护会话与消息列表 */ private handleWebSocketMessage(message : IMessageBase) : void { if (message == null) { return } const { header, payload } = message const session_id = this.getSessionId(header.target, payload) // 1. 更新会话信息 this.updateSession(session_id, message) // 2. 添加到消息列表 this.addMessageToSession(session_id, message) // 3. 触发 UI 更新 uni.$emit('session_updated', session_id) } /** * 获取会话ID */ private getSessionId(target : string, payload : IMessagePayload) : string { const sessionPayload = payload as IChatMessage const sessionId = sessionPayload.session_id ?? '' switch (target) { case 'single': return 'single' + sessionId case 'group': return 'group' + sessionId case 'room': return 'room' + sessionId case 'system': default: return 'system' + sessionId } } /** * 获取会话名称 */ private getSessionName(payload : IMessagePayload, type : string) : string { const chatPayload = payload as IChatMessage switch (type) { case 'single': return chatPayload.nickname ?? '未知用户' case 'group': return '群聊' case 'room': return chatPayload.nickname ?? '房间' case 'system': return '系统消息' default: return '会话' } } /** * 创建会话 */ private createSession(session_id : string, message : IMessageBase) : IMessageSession { const { header, payload } = message const session_type = header.target const chatPayload = payload as IChatMessage return { session_id, session_type, session_name: this.getSessionName(payload, session_type), session_avatar: chatPayload.avatar ?? 'https://cdn.example.com/default_avatar.png', last_message: this.createMessagePreview(message), unread_count: 1, is_pinned: false, is_muted: false, last_active_time: message.header.timestamp, created_time: message.header.timestamp, status: 'normal', ext: {} } } /** * 创建消息预览 */ private createMessagePreview(message : IMessageBase) : IMessagePreview { const { header, payload } = message const chatPayload = payload as IChatMessage return { id: header.id, type: chatPayload.type ?? 'text', content_preview: this.formatContentPreview(payload), timestamp: header.timestamp, msg_status: chatPayload.status ?? 'delivered', is_mentioned: this.isMentioned(payload), is_important: this.isImportant(payload), ext: chatPayload.ext ?? {} } } /** * 判断是否@我 */ private isMentioned(payload : IMessagePayload) : boolean { const mentionPayload = payload as IChatMessage if (mentionPayload.at_user == null) return false else if (mentionPayload.at_user == 'all') return true else if ((mentionPayload.at_user as number[]).includes(this.currentUserId)) return true else return false } /** * 判断是否重要消息 */ private isImportant(payload : IMessagePayload) : boolean { const chatPayload = payload as IChatMessage // 红包、转账、重要通知等 const importantTypes = ['red_packet', 'transfer', 'system_notify'] return importantTypes.includes(chatPayload.type ?? '') } /** * 格式化消息预览内容 */ private formatContentPreview(payload : IMessagePayload) : string { const chatPayload = payload as IChatMessage if (chatPayload.content == null) return '[消息]' const type = chatPayload.type ?? 'text' switch (type) { case 'text': return '[文本消息]' case 'image': return '[图片]' case 'audio': return '[语音]' case 'video': return '[视频]' case 'file': return '[文件]' case 'location': return '[位置]' case 'emoji': return '[表情]' case 'sticker': return '[贴纸]' case 'red_packet': return '[红包]' case 'recall': return '[消息已撤回]' case 'system': return '[系统消息]' default: return '[消息]' } } /** * 更新会话信息 */ private updateSession(session_id : string, message : IMessageBase) : void { const { header, payload } = message const chatPayload = payload as IChatMessage // 获取或创建会话 let session = this.sessions.get(session_id) if (session == null) { session = this.createSession(session_id, message) } // 更新最后消息 session.last_message = this.createMessagePreview(message) session.last_active_time = header.timestamp // 增加未读数(如果不是自己发送的消息) if (chatPayload.uid != this.currentUserId) { session.unread_count++ } this.sessions.set(session_id, session) } /** * 添加消息到会话 */ private addMessageToSession(session_id : string, message : IMessageBase) : void { if (!this.messages.has(session_id)) { this.messages.set(session_id, []) } const sessionMessages = this.messages.get(session_id)! sessionMessages.push(message) // 保持消息按时间排序 sessionMessages.sort((a, b) => a.header.timestamp - b.header.timestamp) // 限制每会话最多保存100条消息 if (sessionMessages.length > 100) { sessionMessages.splice(0, sessionMessages.length - 100) } } /** * 获取排序后的会话列表 */ getSortedSessions() : IMessageSession[] { const sessions : IMessageSession[] = [] this.sessions.forEach((value, key) => { sessions.push(value) }) // 排序规则: // 1. 置顶的在前 // 2. 未读消息数多的在前 // 3. 未读且被@的在前 // 4. 最后活跃时间新的在前 sessions.sort((a, b) => { // 置顶排序 if (a.is_pinned != b.is_pinned) { return a.is_pinned ? -1 : 1 } // 未读数排序 if (a.unread_count != b.unread_count) { return b.unread_count - a.unread_count } // @消息排序 const aHasMention = a.last_message?.is_mentioned ?? false const bHasMention = b.last_message?.is_mentioned ?? false if (aHasMention != bHasMention) { return aHasMention ? -1 : 1 } // 最后活跃时间排序 return b.last_active_time - a.last_active_time }) return sessions } /** * 获取会话消息列表 */ getSessionMessages(session_id : string) : any[] { return (this.messages.get(session_id) as any[] | null) ?? [] } /** * 清空会话未读数 */ clearUnread(session_id : string) : void { const session = this.sessions.get(session_id) if (session != null) { session.unread_count = 0 uni.$emit('session_updated', session_id) } } /** * 置顶/取消置顶 */ togglePin(session_id : string, is_pinned : boolean = true) : void { const session = this.sessions.get(session_id) if (session != null) { session.is_pinned = is_pinned uni.$emit('session_updated', session_id) } } /** * 免打扰/取消免打扰 */ toggleMute(session_id : string, is_muted : boolean = true) : void { const session = this.sessions.get(session_id) if (session != null) { session.is_muted = is_muted uni.$emit('session_updated', session_id) } } /** * 注册消息回调 */ on(callbacks : IMessageServiceCallbacks) : void { if (callbacks.onNewMessage != null) this.callbacks.onNewMessage = callbacks.onNewMessage if (callbacks.onMessageStatusChange != null) this.callbacks.onMessageStatusChange = callbacks.onMessageStatusChange } } export { IMessageSession } export const messageService = new MessageService() export const messageManager = messageService export default messageService