Files

452 lines
12 KiB
Plaintext

// 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<string, IMessageSession> = new Map()
private messages : Map<string, IMessageBase[]> = 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