| ¶Ô±ÈÐÂÎļþ |
| | |
| | | import type { |
| | | Group, |
| | | ImGroupMessageNotification, |
| | | ImMessageReadNotification, |
| | | ImMessageReceiptNotification, |
| | | ImNoConversationNotification, |
| | | ImNotificationWebSocketDTO, |
| | | ImPrivateMessageNotification, |
| | | Message, |
| | | WebSocketFrame |
| | | } from '../types' |
| | | |
| | | import type { ImChannelMessageApi } from '#/api/im/message/channel' |
| | | |
| | | import { acceptHMRUpdate, defineStore } from 'pinia' |
| | | |
| | | import { readChannelMessages as apiReadChannelMessages } from '#/api/im/message/channel' |
| | | import { readGroupMessages as apiReadGroupMessages } from '#/api/im/message/group' |
| | | import { readPrivateMessages as apiReadPrivateMessages } from '#/api/im/message/private' |
| | | import { getCurrentUserId, getRefreshToken } from '#/views/im/utils/auth' |
| | | |
| | | import { buildChannelConversationStub } from '../../utils/channel' |
| | | import { |
| | | MESSAGE_GROUP_READ_ENABLED, |
| | | MESSAGE_PRIVATE_READ_ENABLED, |
| | | WS_RECONNECT_BASE_MS, |
| | | WS_RECONNECT_JITTER_MS, |
| | | WS_RECONNECT_MAX_MS |
| | | } from '../../utils/config' |
| | | import { |
| | | ImContentType, |
| | | ImConversationType, |
| | | ImMessageReceiptStatus, |
| | | ImMessageStatus, |
| | | ImRtcCallMediaType, |
| | | ImRtcParticipantStatus, |
| | | ImWebSocketMessageType, |
| | | isFriendChatTip, |
| | | isFriendNotification, |
| | | isGroupRequestNotification, |
| | | isNormalMessage |
| | | } from '../../utils/constants' |
| | | import { |
| | | getPrivateMessagePeerId, |
| | | parseRtcCallPayload, |
| | | playAudioTip, |
| | | resolveCallEndReasonText |
| | | } from '../../utils/message' |
| | | import { getFriendDisplayName, getGroupDisplayName } from '../../utils/user' |
| | | import { useConversationStore } from './conversationStore' |
| | | import { type FriendNotificationPayload, useFriendStore } from './friendStore' |
| | | import { useGroupRequestStore } from './groupRequestStore' |
| | | import { useGroupStore } from './groupStore' |
| | | import { useMessageStore } from './messageStore' |
| | | import { |
| | | type ImRtcCallEndNotification, |
| | | type ImRtcCallNotification, |
| | | type ImRtcParticipantConnectedNotification, |
| | | type ImRtcParticipantDisconnectedNotification, |
| | | useRtcStore |
| | | } from './rtcStore' |
| | | |
| | | /** FRIEND_DELETE 帧 payload æ¯å¦å¸¦ clear=trueï¼clear è¯ä¹æ¯æ¸
ä¼è¯æ¬èº«ï¼è·³è¿æ°æ³¡æ¸²æ */ |
| | | const isFriendDeleteWithClear = (frame: ImPrivateMessageNotification): boolean => { |
| | | if (frame.type !== ImContentType.FRIEND_DELETE) { |
| | | return false |
| | | } |
| | | try { |
| | | const payload = JSON.parse(frame.content || '{}') as { clear?: boolean } |
| | | return payload.clear === true |
| | | } catch { |
| | | return false |
| | | } |
| | | } |
| | | |
| | | /** ä»ç§èæ¶æ¯å¸§è§£æå¥½åéç¥ payload */ |
| | | const parseFriendNotificationPayload = ( |
| | | frame: ImPrivateMessageNotification |
| | | ): FriendNotificationPayload => JSON.parse(frame.content || '{}') as FriendNotificationPayload |
| | | |
| | | /** ç§èæ¶æ¯å¸§æ¯å¦å¯æ¨æå¥½å对端 */ |
| | | const isPrivateMessageNotification = ( |
| | | frame: ImNoConversationNotification | ImPrivateMessageNotification |
| | | ): frame is ImPrivateMessageNotification => 'senderId' in frame && 'receiverId' in frame |
| | | |
| | | const RTC_LIVEKIT_PROTOCOLS = new Set(['http:', 'https:', 'ws:', 'wss:']) |
| | | const RTC_MEDIA_TYPES = new Set<number>(Object.values(ImRtcCallMediaType)) |
| | | |
| | | /** å¿½ç¥æ®é宿¶å¸§æä¹
å失败 */ |
| | | function ignoreRealtimePersistError(promise: Promise<void>): void { |
| | | void promise.catch(() => undefined) |
| | | } |
| | | |
| | | interface WebSocketListenerSet { |
| | | close: (event: CloseEvent) => void |
| | | error: (event: Event) => void |
| | | message: (event: MessageEvent) => void |
| | | open: (event: Event) => void |
| | | } |
| | | |
| | | const websocketListenerMap = new WeakMap<WebSocket, WebSocketListenerSet>() |
| | | |
| | | /** ç»å® WebSocket äºä»¶ */ |
| | | function bindWebSocketListeners(socket: WebSocket, listeners: WebSocketListenerSet): void { |
| | | websocketListenerMap.set(socket, listeners) |
| | | socket.addEventListener('open', listeners.open) |
| | | socket.addEventListener('message', listeners.message) |
| | | socket.addEventListener('close', listeners.close) |
| | | socket.addEventListener('error', listeners.error) |
| | | } |
| | | |
| | | /** è§£ç» WebSocket äºä»¶ */ |
| | | function unbindWebSocketListeners(socket: WebSocket): void { |
| | | const listeners = websocketListenerMap.get(socket) |
| | | if (!listeners) { |
| | | return |
| | | } |
| | | socket.removeEventListener('open', listeners.open) |
| | | socket.removeEventListener('message', listeners.message) |
| | | socket.removeEventListener('close', listeners.close) |
| | | socket.removeEventListener('error', listeners.error) |
| | | websocketListenerMap.delete(socket) |
| | | } |
| | | |
| | | /** æ ¡éª LiveKit è¿æ¥å°å */ |
| | | function isValidLiveKitUrl(url?: string): boolean { |
| | | if (!url) { |
| | | return false |
| | | } |
| | | try { |
| | | return RTC_LIVEKIT_PROTOCOLS.has(new URL(url).protocol) |
| | | } catch { |
| | | return false |
| | | } |
| | | } |
| | | |
| | | /** æ ¡éªæ¥çµä¿¡ä»¤è½½è· */ |
| | | function isValidRtcInvitePayload(payload: ImRtcCallNotification): boolean { |
| | | if (!payload.room || !payload.token || !isValidLiveKitUrl(payload.livekitUrl)) { |
| | | return false |
| | | } |
| | | if (!RTC_MEDIA_TYPES.has(payload.mediaType) || !payload.inviterUserId) { |
| | | return false |
| | | } |
| | | if (payload.conversationType === ImConversationType.PRIVATE) { |
| | | return true |
| | | } |
| | | return payload.conversationType === ImConversationType.GROUP && !!payload.groupId |
| | | } |
| | | |
| | | /** |
| | | * WebSocket ç§è DTO -> å端 Messageï¼targetId æ¯ä¼è¯ä¸»é®ï¼å¯¹ç«¯ userIdï¼ |
| | | * ä¸ååé人ååæ®µï¼æ¸²æå±èµ° utils/user 宿¶ç®ï¼å¤æ³¨ / 群æµç§°åæ´åå岿¶æ¯èªå¨å·æ°ï¼ |
| | | */ |
| | | const convertPrivateMessage = ( |
| | | websocketMessage: ImPrivateMessageNotification, |
| | | currentUserId: number |
| | | ): Message => ({ |
| | | id: websocketMessage.id, |
| | | clientMessageId: websocketMessage.clientMessageId, |
| | | type: websocketMessage.type, |
| | | content: websocketMessage.content, |
| | | status: websocketMessage.status, |
| | | receiptStatus: websocketMessage.receiptStatus, |
| | | sendTime: new Date(websocketMessage.sendTime).getTime(), |
| | | senderId: websocketMessage.senderId, |
| | | targetId: getPrivateMessagePeerId(websocketMessage, currentUserId), |
| | | selfSend: websocketMessage.senderId === currentUserId |
| | | }) |
| | | |
| | | /** |
| | | * WebSocket 群è DTO -> å端 Message |
| | | * 带 atUserIds / receiverUserIds ç» @ æ è®°åå®åæ¥æ¶ç¨ï¼ |
| | | * receiptStatus / readCount 让å¤ç«¯åæ¥æ¶å°èªå·±åçç¾¤æ¶æ¯æ¶åæ§ UI ç«å»å°±ææ°æ® |
| | | */ |
| | | const convertGroupMessage = ( |
| | | websocketMessage: ImGroupMessageNotification, |
| | | currentUserId: number |
| | | ): Message => ({ |
| | | id: websocketMessage.id, |
| | | clientMessageId: websocketMessage.clientMessageId, |
| | | type: websocketMessage.type, |
| | | content: websocketMessage.content, |
| | | status: websocketMessage.status, |
| | | sendTime: new Date(websocketMessage.sendTime).getTime(), |
| | | senderId: websocketMessage.senderId, |
| | | targetId: websocketMessage.groupId, |
| | | selfSend: websocketMessage.senderId === currentUserId, |
| | | atUserIds: websocketMessage.atUserIds || [], |
| | | receiverUserIds: websocketMessage.receiverUserIds || [], |
| | | receiptStatus: websocketMessage.receiptStatus, |
| | | readCount: websocketMessage.readCount |
| | | }) |
| | | |
| | | /** |
| | | * IM WebSocket Store |
| | | * |
| | | * èè´£ï¼ä¸åªæ¯è¿éä¿¡ï¼ä¹æ¯å端 IM äºä»¶çç»ä¸å
¥å£ â èå¨ conversationStore / friendStore / groupStoreï¼ï¼ |
| | | * |
| | | * 1. é¾è·¯ç®¡çï¼å»ºè¿ / æè¿ / å¿è·³ä¿æ´» / èªå¨éè¿ |
| | | * 2. 帧ååï¼dispatchFrame â dispatchPrivateFrame / dispatchGroupFrameï¼æä¼è¯åå
容类ååæµ |
| | | * 3. ç¼å²ï¼åå§åå è½½æï¼conversationStore.loading=trueï¼æåæ¶æ¯ï¼ç pull 宿åç± useMessagePuller è° flushBuffer åæ¾ |
| | | * 4. äºä»¶å¤çï¼æç±»åååå°å¯¹åº handle*ï¼èå¨ conversation / friend / group storeï¼ï¼ |
| | | * - æ®éæ¶æ¯ï¼TEXT / IMAGE / FILE / VOICE / VIDEOï¼ï¼å
¥åº + å½åä¼è¯èªå¨å·²è¯» / æç¤ºé³ |
| | | * - 已读 / åæ§ï¼READ / RECEIPTï¼ï¼å¤ç«¯å·²è¯»åæ¥ã对æ¹è¯»ååæ§ |
| | | * - 好ååæ´ï¼FRIEND_*ï¼ï¼åæ¥ friendStore + 级èå·æ°ç§èä¼è¯ï¼FRIEND_ADD / FRIEND_DELETE é¢å¤æå
¥ä¼è¯æ°æ³¡ |
| | | * - 群个人信å·ï¼GROUP_MEMBER_SETTING_UPDATEï¼ï¼åæ¥ groupStore + 级èå·æ°ç¾¤èä¼è¯ |
| | | * - 群æåæµç§°åæ´ï¼GROUP_MEMBER_NICKNAME_UPDATEï¼ï¼åæ¥ groupStoreï¼ä¸æå
¥æ¶æ¯å表 |
| | | * - 群广æäºä»¶ï¼GROUP_*ï¼ï¼èµ° handleGroupMessage + applyGroupNotification æè·¯ï¼å« DISSOLVE / QUIT / KICK èªå¤æ¸
ç¾¤ï¼ |
| | | */ |
| | | export const useImWebSocketStore = defineStore('imWebSocketStore', { |
| | | state: () => ({ |
| | | socket: null as null | WebSocket, |
| | | isConnected: false, |
| | | reconnectTimer: null as null | ReturnType<typeof setTimeout>, |
| | | /** è¿ç»éè¿å¤±è´¥æ¬¡æ°ï¼onopen æå / disconnect 䏻卿å¼åæ¸
é¶ï¼ç¨äºææ°éé¿ */ |
| | | reconnectAttempts: 0, |
| | | heartbeatTimer: null as null | ReturnType<typeof setInterval>, |
| | | messageBuffer: [] as Array< |
| | | | { conversationType: typeof ImConversationType.CHANNEL; payload: ImChannelMessageApi.ChannelMessageRespVO } |
| | | | { |
| | | conversationType: typeof ImConversationType.GROUP |
| | | payload: ImGroupMessageNotification |
| | | } |
| | | | { |
| | | conversationType: typeof ImConversationType.PRIVATE |
| | | payload: ImPrivateMessageNotification |
| | | } |
| | | > // åå§åå è½½æå
ï¼å
ææ®éæ¶æ¯ä¸¢è¿ç¼å²åºï¼pull 宿åå䏿¬¡æ§åæ¾ |
| | | }), |
| | | |
| | | actions: { |
| | | /** |
| | | * ååºç¼å²åºæ¶æ¯å¹¶æ¸
空ï¼ç± useMessagePuller å¨ pull 宿åè°ç¨ï¼ç»ä¸åæ¾ç» conversationStoreï¼ |
| | | * é
å messageBuffer å®ç°ï¼å¨ conversationStore.loading æé´æ¶å°ç WS æ¶æ¯å
æåï¼é¿å
å pull ç minId æ¸¸æ ææ¶ |
| | | */ |
| | | flushBuffer() { |
| | | const msgs = [...this.messageBuffer] |
| | | this.messageBuffer = [] |
| | | return msgs |
| | | }, |
| | | |
| | | /** ç´æ¥ä¸¢å¼ç¼å²å¸§ä¸åæ¾ï¼cancelPull / ç¦»å¼ IM è°ç¨ï¼é²æ¢ä¸æ¬¡è¿ IM ææ§ session å¸§åæ¾è¿æ° storeï¼ */ |
| | | discardBuffer() { |
| | | this.messageBuffer = [] |
| | | }, |
| | | |
| | | /** |
| | | * è¿æ¥ WebSocket |
| | | * å¤ç¨ yudao å
ç½® /infra/ws ééï¼å端éè¿ sendObject(type, content) ä¸å |
| | | * |
| | | * è°ç¨å¥çº¦ï¼åè´¦å· / token å·æ°åå¿
é¡»å
`disconnect()` å `connect()`ï¼ |
| | | * æ¬æ¹æ³ä¸æç¥ token ååï¼æ§ socket å¨ CONNECTING / OPEN ç¶æä¼ç´æ¥å¤ç¨æ§ tokenï¼å¯è½æ¿å°é误身份 |
| | | */ |
| | | connect() { |
| | | // é´æç¨ refreshTokenï¼çå½å¨ææ´é¿ï¼access token è¿æåæå¡ç«¯ä¼éè¿ frame éç¥éç»ï¼ |
| | | const refreshToken = getRefreshToken() |
| | | if (!refreshToken) { |
| | | console.warn('[IM WS] refreshToken 为空ï¼è·³è¿è¿æ¥') |
| | | return |
| | | } |
| | | // æ§ socket è¿å¨ CONNECTING / OPEN ç´æ¥å¤ç¨ï¼é¿å
å å å¤ä»½ onmessage çå¬å¯¼è´é夿¶æ¯ / æç¤ºé³ / å·²è¯»ä¸æ¥ |
| | | const existingSocket = this.socket |
| | | if ( |
| | | existingSocket && |
| | | (existingSocket.readyState === WebSocket.OPEN || |
| | | existingSocket.readyState === WebSocket.CONNECTING) |
| | | ) { |
| | | return |
| | | } |
| | | // æ§ socket å·² CLOSING / CLOSEDï¼è§£ç»åè° + æ¸
å¼ç¨å newï¼é¿å
è handler 仿æ store å¼ç¨é»ç¢ GC |
| | | if (existingSocket) { |
| | | unbindWebSocketListeners(existingSocket) |
| | | this.socket = null |
| | | } |
| | | const url = `${this.buildWsUrl()}/infra/ws?token=${refreshToken}` |
| | | const socket = new WebSocket(url) |
| | | this.socket = socket |
| | | |
| | | bindWebSocketListeners(socket, { |
| | | // è¿æ¥å»ºç«ï¼æ è®°ä¸çº¿ + å¯å¨å¿è·³ä¿æ´»ï¼éè¿éé¿è®¡æ°å½é¶ |
| | | open: () => { |
| | | this.isConnected = true |
| | | this.reconnectAttempts = 0 |
| | | console.log('[IM WS] connected') |
| | | this.startHeartbeat() |
| | | }, |
| | | |
| | | // æ¶å°å¸§ï¼'pong' æ¯å¿è·³åºçç´æ¥åæï¼å
¶ä½æ WebSocketFrame è§£æåäº¤ç» dispatchFrame åæµ |
| | | message: (event) => { |
| | | if (event.data === 'pong') { |
| | | return |
| | | } |
| | | try { |
| | | const frame = JSON.parse(event.data) as WebSocketFrame |
| | | this.dispatchFrame(frame) |
| | | } catch (error) { |
| | | console.error('[IM WS] message parse error:', error) |
| | | } |
| | | }, |
| | | |
| | | // æå¡ç«¯å
³é / ç½ç»æï¼æ è®°ä¸çº¿ï¼æææ°éé¿èªå¨éè¿ |
| | | close: () => { |
| | | this.isConnected = false |
| | | console.log('[IM WS] disconnected') |
| | | this.reconnect() |
| | | }, |
| | | |
| | | // å¼å¸¸æ¶ä¸ä¸»å¨ reconnectï¼ä¸»å¨ close() 让 close æä¸ºå¯ä¸éè¿å
¥å£ |
| | | error: (error) => { |
| | | console.error('[IM WS] error:', error) |
| | | this.isConnected = false |
| | | this.socket?.close() |
| | | } |
| | | }) |
| | | }, |
| | | |
| | | /** æ¼æ¥ WebSocket åºç¡å°å */ |
| | | buildWsUrl(): string { |
| | | // VITE_BASE_URL å¯è½æ¯ http:// æ https:// å¼å¤´ï¼æ¿æ¢æ ws:// æ wss://ï¼å¦ææ²¡é
ç½®ï¼å°±ç¨å½å页é¢çåè®® + host |
| | | const baseUrl = (import.meta as any).env?.VITE_BASE_URL as string | undefined |
| | | if (baseUrl && baseUrl.length > 0) { |
| | | return baseUrl.replace(/^http/, 'ws') |
| | | } |
| | | // å½å页é¢åè®® + hostï¼å¦ http://localhost:8080ï¼ï¼æ¿æ¢æ ws://localhost:8080 |
| | | const protocol = window.location.protocol === 'https:' ? 'wss:' : 'ws:' |
| | | const host = window.location.host |
| | | return `${protocol}//${host}` |
| | | }, |
| | | |
| | | /** |
| | | * æ IM éç¥å¸§åå |
| | | */ |
| | | dispatchFrame(frame: WebSocketFrame) { |
| | | if (frame.type !== ImWebSocketMessageType.NOTIFICATION) { |
| | | console.debug('[IM WS] æªè¯å«äºä»¶', frame) |
| | | return |
| | | } |
| | | |
| | | const notification = this.safeParse(frame.content) as ImNotificationWebSocketDTO | null |
| | | if (!notification?.payload || !notification.contentType) { |
| | | return |
| | | } |
| | | const payload = { |
| | | ...notification.payload, |
| | | type: notification.contentType |
| | | } |
| | | switch (notification.conversationType) { |
| | | case ImConversationType.CHANNEL: { |
| | | this.dispatchChannelFrame(payload as ImChannelMessageApi.ChannelMessageRespVO) |
| | | break |
| | | } |
| | | case ImConversationType.GROUP: { |
| | | this.dispatchGroupFrame(payload as ImGroupMessageNotification) |
| | | break |
| | | } |
| | | case ImConversationType.NONE: { |
| | | this.dispatchNoConversationFrame(payload as ImNoConversationNotification) |
| | | break |
| | | } |
| | | case ImConversationType.PRIVATE: { |
| | | this.dispatchPrivateFrame(payload as ImPrivateMessageNotification) |
| | | break |
| | | } |
| | | default: { |
| | | console.debug('[IM WS] æªè¯å«éç¥', notification) |
| | | } |
| | | } |
| | | }, |
| | | |
| | | /** |
| | | * æ ä¼è¯éç¥åå |
| | | */ |
| | | dispatchNoConversationFrame(websocketMessage: ImNoConversationNotification) { |
| | | if (isFriendNotification(websocketMessage.type)) { |
| | | this.handleFriendNotification(websocketMessage) |
| | | return |
| | | } |
| | | if (isGroupRequestNotification(websocketMessage.type)) { |
| | | this.handleGroupRequestNotification(websocketMessage) |
| | | return |
| | | } |
| | | switch (websocketMessage.type) { |
| | | case ImContentType.RTC_CALL: |
| | | case ImContentType.RTC_PARTICIPANT_CONNECTED: |
| | | case ImContentType.RTC_PARTICIPANT_DISCONNECTED: { |
| | | this.handleRtcSignaling(websocketMessage) |
| | | break |
| | | } |
| | | default: { |
| | | console.debug('[IM WS] æªè¯å«æ ä¼è¯éç¥', websocketMessage) |
| | | } |
| | | } |
| | | }, |
| | | |
| | | /** |
| | | * é¢é帧ååï¼æ payload.type åå° READï¼å¤ç«¯å·²è¯»åæ¥ï¼ææ®éç´ ææ¨é |
| | | */ |
| | | dispatchChannelFrame(websocketMessage: ImChannelMessageApi.ChannelMessageRespVO) { |
| | | if (websocketMessage.type === ImContentType.READ) { |
| | | this.handleChannelRead(websocketMessage) |
| | | return |
| | | } |
| | | ignoreRealtimePersistError(this.handleChannelMessage(websocketMessage)) |
| | | }, |
| | | |
| | | /** é¢é READï¼èªå·±å
¶å®ç»ç«¯å¨æé¢ééæ ä¸ºå·²è¯»ï¼æ¬ç«¯åæ¥æ¸
é¶è¯¥é¢éæªè¯» */ |
| | | handleChannelRead(websocketMessage: ImChannelMessageApi.ChannelMessageRespVO) { |
| | | void useConversationStore() |
| | | .applyConversationReadList([ |
| | | { |
| | | id: websocketMessage.id, |
| | | conversationType: ImConversationType.CHANNEL, |
| | | targetId: websocketMessage.channelId, |
| | | messageId: websocketMessage.id |
| | | } |
| | | ]) |
| | | .catch((error) => console.warn('[IM WS] é¢éå·²è¯»åæ¥å¤±è´¥', error)) |
| | | }, |
| | | |
| | | /** |
| | | * é¢éæ¶æ¯å®æ¶å
¥ä¼è¯ï¼é¢éæ¶æ¯åå + æ ç¶ææºï¼ç´æ¥ insertMessage å³å¯ |
| | | * pull ä¸ WS æ¿å°å䏿¡ id æ¶ï¼messageStore.insertMessage å
鍿 id å»éï¼ä¸ä¼éå¤ |
| | | */ |
| | | handleChannelMessage(websocketMessage: ImChannelMessageApi.ChannelMessageRespVO): Promise<void> { |
| | | const conversationStore = useConversationStore() |
| | | const messageStore = useMessageStore() |
| | | // 离线å è½½æé´å
ç¼å²ï¼ç pull 宿ååç»ä¸åæ¾ï¼é¿å
éå¤æé¡ºåºéä¹± |
| | | if (conversationStore.loading) { |
| | | this.messageBuffer.push({ |
| | | conversationType: ImConversationType.CHANNEL, |
| | | payload: websocketMessage |
| | | }) |
| | | return Promise.resolve() |
| | | } |
| | | const sendTimeMs = |
| | | typeof websocketMessage.sendTime === 'number' |
| | | ? websocketMessage.sendTime |
| | | : new Date(websocketMessage.sendTime).getTime() |
| | | const conversation = conversationStore.getConversation( |
| | | ImConversationType.CHANNEL, |
| | | websocketMessage.channelId |
| | | ) |
| | | const isActive = |
| | | conversationStore.activeConversation?.type === ImConversationType.CHANNEL && |
| | | conversationStore.activeConversation?.targetId === websocketMessage.channelId |
| | | // é¢éåå订é
ï¼receiptStatus è¡¨è¾¾ãææ¯å¦å·²è¯»è¿æ¡ãï¼ä¼è¯æå¼å³å·²è¯» DONEï¼å¦å PENDINGï¼ä¸ pull å£å¾ä¸è´ï¼ |
| | | const persistPromise = messageStore.insertMessage( |
| | | buildChannelConversationStub(websocketMessage.channelId), |
| | | { |
| | | id: websocketMessage.id, |
| | | clientMessageId: '', |
| | | type: websocketMessage.type, |
| | | content: websocketMessage.content, |
| | | status: ImMessageStatus.NORMAL, |
| | | receiptStatus: isActive ? ImMessageReceiptStatus.DONE : ImMessageReceiptStatus.PENDING, |
| | | sendTime: sendTimeMs, |
| | | senderId: 0, |
| | | targetId: websocketMessage.channelId, |
| | | selfSend: false, |
| | | materialId: websocketMessage.materialId |
| | | } |
| | | ) |
| | | if (isActive) { |
| | | // çªå£æå¼ = å·²è¯»ï¼æ¬ç«¯æ¸
æªè¯» + 䏿¥æå¡ç«¯è¯»ä½ç½®ï¼é¿å
读ä½ç½®æ»å |
| | | const readReported = conversationStore.isReportedReadPositionCovered( |
| | | ImConversationType.CHANNEL, |
| | | websocketMessage.channelId, |
| | | websocketMessage.id |
| | | ) |
| | | conversationStore.markConversationRead( |
| | | ImConversationType.CHANNEL, |
| | | websocketMessage.channelId, |
| | | websocketMessage.id |
| | | ) |
| | | if (!readReported) { |
| | | apiReadChannelMessages(websocketMessage.channelId, websocketMessage.id) |
| | | .then(() => |
| | | conversationStore.markConversationReadReported( |
| | | ImConversationType.CHANNEL, |
| | | websocketMessage.channelId, |
| | | websocketMessage.id |
| | | ) |
| | | ) |
| | | .catch((error) => { |
| | | console.warn( |
| | | '[IM WS] é¢éèªå¨å·²è¯»ä¸æ¥å¤±è´¥', |
| | | { |
| | | conversationType: ImConversationType.CHANNEL, |
| | | channelId: websocketMessage.channelId, |
| | | messageId: websocketMessage.id |
| | | }, |
| | | error |
| | | ) |
| | | }) |
| | | } |
| | | } else if (!conversation?.silent && isNormalMessage(websocketMessage.type)) { |
| | | // éå½åä¼è¯ä¸æªå
ææ°ï¼åä¸ä¸æç¤ºé³ |
| | | playAudioTip() |
| | | } |
| | | return persistPromise |
| | | }, |
| | | |
| | | /** content æ¢å¯è½å·²æ¯å¯¹è±¡ä¹å¯è½æ¯ JSON å符串ï¼åç«¯ç¨ Map åºååä¸åï¼ */ |
| | | safeParse(raw: unknown): null | Record<string, any> { |
| | | if (!raw) { |
| | | return null |
| | | } |
| | | if (typeof raw === 'object') { |
| | | return raw as Record<string, any> |
| | | } |
| | | try { |
| | | return JSON.parse(raw as string) |
| | | } catch (error) { |
| | | console.error('[IM WS] content è§£æå¤±è´¥', error) |
| | | return null |
| | | } |
| | | }, |
| | | |
| | | // ==================== æ®éæ¶æ¯ ==================== |
| | | |
| | | /** |
| | | * ç§èç»ä¸å¸§ååï¼æ payload.typeï¼ImContentTypeï¼åå°å·²è¯» / åæ§ / 好åéç¥ / æ®éæ¶æ¯ |
| | | * |
| | | * æ¶æ¯éç¥ã已读éç¥ãåæ§éç¥ç±å¤å± contentType ç»ä¸åå |
| | | */ |
| | | dispatchPrivateFrame(websocketMessage: ImPrivateMessageNotification) { |
| | | try { |
| | | switch (websocketMessage.type) { |
| | | case ImContentType.READ: { |
| | | this.handlePrivateRead(websocketMessage as ImMessageReadNotification) |
| | | break |
| | | } |
| | | case ImContentType.RECEIPT: { |
| | | this.handlePrivateReceipt(websocketMessage as ImMessageReceiptNotification) |
| | | break |
| | | } |
| | | case ImContentType.RTC_CALL_END: { |
| | | // å
¥åº + å
³ééè¯çª + 渲æè天 tipï¼ç§èåºæ¯ï¼ |
| | | this.handleRtcCallEnd(websocketMessage) |
| | | ignoreRealtimePersistError(this.handlePrivateMessage(websocketMessage)) |
| | | break |
| | | } |
| | | default: { |
| | | if (isFriendChatTip(websocketMessage.type)) { |
| | | this.handleFriendNotification(websocketMessage) |
| | | // FRIEND_DELETE ç clear=true è¯ä¹æ¯æ¸
ä¼è¯æ¬èº«ï¼è·³è¿æ°æ³¡é¿å
å¨å·²æ¸
ä¼è¯éåå
¥èææ¶æ¯ |
| | | if (!isFriendDeleteWithClear(websocketMessage)) { |
| | | ignoreRealtimePersistError(this.handlePrivateMessage(websocketMessage)) |
| | | } |
| | | } else { |
| | | // TEXT / IMAGE / FILE / VOICE / VIDEO çæ®éæ¶æ¯ |
| | | ignoreRealtimePersistError(this.handlePrivateMessage(websocketMessage)) |
| | | } |
| | | } |
| | | } |
| | | } catch (error) { |
| | | // 忡叧çå¤çå¼å¸¸ä¸åºé»æåç»å¸§ï¼æå°å®æ´ websocketMessage ä¾¿äºææ¥ |
| | | console.warn('[IM WS] dispatchPrivateFrame å¤ç失败', websocketMessage, error) |
| | | } |
| | | }, |
| | | |
| | | /** |
| | | * 群èç»ä¸å¸§ååï¼æ payload.typeï¼ImContentTypeï¼åå°å·²è¯» / åæ§ / ç¾¤ä¸ªäººä¿¡å· / æ®éæ¶æ¯ |
| | | * |
| | | * GROUP_MEMBER_SETTING_UPDATE / GROUP_MEMBER_NICKNAME_UPDATE æ¯æåèµæåæ¥ä¿¡å·ï¼å
¶å®ç¾¤å¹¿æäºä»¶èµ° handleGroupMessage å
¥åº + 触å applyGroupNotification æè·¯ |
| | | */ |
| | | dispatchGroupFrame(websocketMessage: ImGroupMessageNotification) { |
| | | try { |
| | | switch (websocketMessage.type) { |
| | | case ImContentType.GROUP_MEMBER_NICKNAME_UPDATE: { |
| | | this.handleGroupMemberNicknameUpdate(websocketMessage) |
| | | break |
| | | } |
| | | case ImContentType.GROUP_MEMBER_SETTING_UPDATE: { |
| | | this.handleGroupMemberSettingUpdate(websocketMessage) |
| | | break |
| | | } |
| | | case ImContentType.READ: { |
| | | this.handleGroupRead(websocketMessage as ImMessageReadNotification) |
| | | break |
| | | } |
| | | case ImContentType.RECEIPT: { |
| | | this.handleGroupReceipt(websocketMessage as ImMessageReceiptNotification) |
| | | break |
| | | } |
| | | case ImContentType.RTC_CALL_END: { |
| | | // å
¥åº + ç§»é¤è¶åæ¡ + å
³ééè¯çªï¼å¦æå½åå¨è¯¥ç¾¤éè¯å
ï¼ |
| | | this.handleRtcCallEnd(websocketMessage) |
| | | ignoreRealtimePersistError(this.handleGroupMessage(websocketMessage)) |
| | | break |
| | | } |
| | | case ImContentType.RTC_CALL_START: { |
| | | // å
¥åº + 渲æè天 tipï¼åæ¶ç¨ START payload å
çææå°è¶åæ¡ï¼åç» getActiveCall / åä¸è
äºä»¶åè¡¥é½æå |
| | | this.handleRtcCallStart(websocketMessage) |
| | | ignoreRealtimePersistError(this.handleGroupMessage(websocketMessage)) |
| | | break |
| | | } |
| | | default: { |
| | | // TEXT / IMAGE / FILE / VOICE / VIDEO + GROUP_* 群广æäºä»¶ |
| | | ignoreRealtimePersistError(this.handleGroupMessage(websocketMessage)) |
| | | } |
| | | } |
| | | } catch (error) { |
| | | // 忡叧çå¤çå¼å¸¸ä¸åºé»æåç»å¸§ï¼æå°å®æ´ websocketMessage ä¾¿äºææ¥ |
| | | console.warn('[IM WS] dispatchGroupFrame å¤ç失败', websocketMessage, error) |
| | | } |
| | | }, |
| | | |
| | | /** |
| | | * ç§èæ®éæ¶æ¯ï¼TEXT / IMAGE / FILE / VOICE / VIDEOï¼å
¥åº + èªå¨å·²è¯» |
| | | * |
| | | * æµç¨ï¼ |
| | | * 1. 离线å è½½æç¼å²ï¼é¿å¼ä¸ pull åå¡«çç«æï¼ |
| | | * 2. è®¡ç® selfSend / peerId ç»´åº¦ï¼æå¥½åä¿¡æ¯åå¡«å±ç¤ºå段 |
| | | * 3. æ¤å TIP ç´æ¥è½¬èµ° recallMessageï¼ä¸è¿æ¶æ¯å表 |
| | | * 4. æé å端 Messageï¼æå
¥å°å¯¹åºç§èä¼è¯ |
| | | * 5. å½åä¼è¯æ¿æ´»æ¶èªå¨ä¸æ¥å·²è¯»ï¼å¦åéå
ææ°åæç¤ºé³ |
| | | */ |
| | | handlePrivateMessage(websocketMessage: ImPrivateMessageNotification): Promise<void> { |
| | | const conversationStore = useConversationStore() |
| | | const friendStore = useFriendStore() |
| | | const currentUserId = getCurrentUserId() |
| | | |
| | | // 0. é²å¾¡å±ï¼senderId / receiverId åä¸å«å½åç¨æ·çç§èå¸§ç´æ¥ä¸¢å¼ï¼é¿å
åç«¯è·¯ç± / å¤ç«¯ä¸²å·æ±¡æä¼è¯ |
| | | // ï¼FRIEND_* çç³»ç»éç¥ä¹èµ°è¿æ¡ééï¼ä½ fromUserId=senderIdãtoUserId=receiverId 仿¯å½åç¨æ·è§è§ï¼ |
| | | if ( |
| | | currentUserId && |
| | | websocketMessage.senderId !== currentUserId && |
| | | websocketMessage.receiverId !== currentUserId |
| | | ) { |
| | | console.warn('[IM WS] 丢å¼ä¸å±äºå½åç¨æ·çç§è帧', websocketMessage) |
| | | return Promise.resolve() |
| | | } |
| | | |
| | | // 1. 离线å è½½æé´å
ç¼å²ï¼ç pull 宿ååç»ä¸åæ¾ï¼é¿å
éå¤æé¡ºåºéä¹± |
| | | if (conversationStore.loading) { |
| | | this.messageBuffer.push({ |
| | | conversationType: ImConversationType.PRIVATE, |
| | | payload: websocketMessage |
| | | }) |
| | | return Promise.resolve() |
| | | } |
| | | |
| | | // 2. selfSend / peerIdï¼èªå·±åçæ¶æ¯å±äºãåç» receiverId çä¼è¯ãï¼å«äººåçå±äºãåéè
çä¼è¯ã |
| | | const selfSend = websocketMessage.senderId === currentUserId |
| | | const peerId = getPrivateMessagePeerId(websocketMessage, currentUserId) |
| | | // æªç¥å¯¹ç«¯ï¼éç人å 好ååå
æ¶å°æ¶æ¯çåºæ¯ï¼ï¼å¼æ¥è¡¥æä¸æ¬¡ï¼ä¸æ¬¡å渲æå°±æ name/avatar |
| | | const friend = friendStore.getFriend(peerId) |
| | | if (!friend) { |
| | | friendStore.fetchFriendInfo(peerId).catch(() => undefined) |
| | | } |
| | | // ä¼è¯æ 颿°¸è¿è·ã对端ãèµ°ï¼ä¸ç®¡è°åçæ¶æ¯ï¼ï¼è¿éåªç®ä¸æ¬¡ç» insertMessage ç¨ |
| | | const peerDisplayName = friend ? getFriendDisplayName(friend) : '' |
| | | |
| | | // 3. å端æ¤åï¼ä¸å䏿¡ RECALL æ¶æ¯ï¼content 为 `{"messageId": xxx}`ï¼å¯¹é½ ImContentTypeEnum.RECALL â RecallMessageï¼ |
| | | // è¿éæ¦æªä¸æ¥æ¹èµ° recallMessageï¼æåæ¶æ¯æ´æ°ä¸º RECALL æï¼ï¼ä¸è®©å®ä½ä¸ºæ°æ¶æ¯è¿å表 |
| | | if (websocketMessage.type === ImContentType.RECALL) { |
| | | return useMessageStore().recallMessage( |
| | | ImConversationType.PRIVATE, |
| | | peerId, |
| | | websocketMessage.content |
| | | ) |
| | | } |
| | | |
| | | // 4. å端 DTO â å端 Messageï¼åéäººåæ¸²ææ¶å®æ¶ç®ï¼ä¸åå
¥æ¶æ¯å段 |
| | | const message = convertPrivateMessage(websocketMessage, currentUserId) |
| | | const persistPromise = useMessageStore().insertMessage( |
| | | { |
| | | type: ImConversationType.PRIVATE, |
| | | targetId: peerId, |
| | | name: peerDisplayName || String(peerId), |
| | | avatar: friend?.avatar || '', |
| | | silent: friend?.silent |
| | | }, |
| | | message |
| | | ) |
| | | |
| | | // 5. ä»
å¯¹æ¹æ¶æ¯æèµ°ãèªå¨å·²è¯» / æç¤ºé³ã忝ï¼èªå·±åçä¸ä¼è§¦å |
| | | if (!selfSend) { |
| | | const conversation = conversationStore.getConversation(ImConversationType.PRIVATE, peerId) |
| | | const isActive = |
| | | conversationStore.activeConversation?.type === ImConversationType.PRIVATE && |
| | | conversationStore.activeConversation?.targetId === peerId |
| | | if (isActive) { |
| | | // è天çªå£æå¼ = å®é
çå°äºï¼æ¬ç«¯æ¸
æªè¯»ï¼ç§è已读å¼å¯æ¶å䏿¥å端ï¼è®©å¯¹æ¹ UI ç«å»åå°"已读" |
| | | // 已读ä½ç½®ç´æ¥ç¨åå°çæ¶æ¯ idï¼è¿æ¡å°±æ¯å½åä¼è¯æå¤§ idï¼ |
| | | const readReported = conversationStore.isReportedReadPositionCovered( |
| | | ImConversationType.PRIVATE, |
| | | peerId, |
| | | websocketMessage.id |
| | | ) |
| | | conversationStore.markConversationRead( |
| | | ImConversationType.PRIVATE, |
| | | peerId, |
| | | websocketMessage.id |
| | | ) |
| | | if (MESSAGE_PRIVATE_READ_ENABLED && !readReported) { |
| | | apiReadPrivateMessages(peerId, websocketMessage.id) |
| | | .then(() => |
| | | conversationStore.markConversationReadReported( |
| | | ImConversationType.PRIVATE, |
| | | peerId, |
| | | websocketMessage.id |
| | | ) |
| | | ) |
| | | .catch((error) => { |
| | | console.warn( |
| | | '[IM WS] ç§èèªå¨å·²è¯»ä¸æ¥å¤±è´¥', |
| | | { |
| | | conversationType: ImConversationType.PRIVATE, |
| | | peerId, |
| | | messageId: websocketMessage.id |
| | | }, |
| | | error |
| | | ) |
| | | }) |
| | | } |
| | | } else if (!conversation?.silent && isNormalMessage(websocketMessage.type)) { |
| | | // éå½åä¼è¯ä¸æªå
ææ°ï¼åä¸ä¸æç¤ºé³ï¼å¸¦èæµï¼è¯¦è§ playAudioTipï¼ï¼FRIEND_* çç³»ç»äºä»¶ä¸å |
| | | playAudioTip() |
| | | } |
| | | } |
| | | return persistPromise |
| | | }, |
| | | |
| | | /** ç§è READ äºä»¶ï¼èªå·±çå
¶å®ç»ç«¯å¨å¯¹æ¹ä¼è¯éæ ä¸ºå·²è¯»ï¼æ¬ç«¯åæ¥æ¸
é¶æªè¯»ï¼ç§è已读å
³éæ¶å
åºå¿½ç¥ */ |
| | | handlePrivateRead(websocketMessage: ImMessageReadNotification) { |
| | | if (!MESSAGE_PRIVATE_READ_ENABLED) { |
| | | return |
| | | } |
| | | if (!websocketMessage.id || !websocketMessage.receiverId) { |
| | | return |
| | | } |
| | | void useConversationStore() |
| | | .applyConversationReadList([ |
| | | { |
| | | id: websocketMessage.id, |
| | | conversationType: ImConversationType.PRIVATE, |
| | | targetId: websocketMessage.receiverId, |
| | | messageId: websocketMessage.id |
| | | } |
| | | ]) |
| | | .catch((error) => console.warn('[IM WS] ç§èå·²è¯»åæ¥å¤±è´¥', error)) |
| | | }, |
| | | |
| | | /** |
| | | * ç§è RECEIPT äºä»¶ï¼å¯¹æ¹è¯»äºæçæ¶æ¯ï¼æå对æ¹ä¼è¯éèªå·±åçæ¶æ¯æ 为已读 |
| | | * åç«¯å° maxReadId ç¼ç å¨éç¥ç id åæ®µï¼ |
| | | * è¿éæ®æ¤å¡è¾¹çï¼é¿å
æ"åæ§å¨è·¯ä¸æ¶ååçæ¶æ¯"误æ 为已读ï¼ç§è已读å
³éæ¶å
åºå¿½ç¥ |
| | | */ |
| | | handlePrivateReceipt(websocketMessage: ImMessageReceiptNotification) { |
| | | if (!MESSAGE_PRIVATE_READ_ENABLED) { |
| | | return |
| | | } |
| | | if (!websocketMessage.id) { |
| | | return |
| | | } |
| | | if (!websocketMessage.senderId) { |
| | | return |
| | | } |
| | | useMessageStore().applyMessageReadReceipt({ |
| | | conversationType: ImConversationType.PRIVATE, |
| | | targetId: websocketMessage.senderId, |
| | | privateReadMaxId: websocketMessage.id |
| | | }) |
| | | }, |
| | | |
| | | /** |
| | | * ç¾¤èæ®éæ¶æ¯å
¥åº + èªå¨å·²è¯»ï¼ç»æä¸ handlePrivateMessage å¯¹ç§°ï¼ |
| | | * |
| | | * æµç¨ï¼ |
| | | * 1. 离线å è½½æç¼å² |
| | | * 2. æªç¥ç¾¤æ¶æç¾¤è¯¦æ
å
åº |
| | | * 3. æ¤å TIP ç´æ¥è½¬èµ° |
| | | * 4. æé Message + at åæ®µï¼æå
¥å°å¯¹åºç¾¤èä¼è¯ï¼åéäººåæ¸²ææ¶å®æ¶ç®ï¼ |
| | | * 5. å½åä¼è¯æ¿æ´»æ¶èªå¨ä¸æ¥å·²è¯»ï¼å¸¦ lastMessageIdï¼ï¼å¦åéå
ææ°åæç¤ºé³ |
| | | */ |
| | | handleGroupMessage(websocketMessage: ImGroupMessageNotification): Promise<void> { |
| | | const conversationStore = useConversationStore() |
| | | const groupStore = useGroupStore() |
| | | const currentUserId = getCurrentUserId() |
| | | const selfSend = websocketMessage.senderId === currentUserId |
| | | |
| | | // 0. é²å¾¡å±ï¼å®åç¾¤æ¶æ¯ receiverUserIds éç©ºä¸æªå
å«å½åç¨æ·æ¶ä¸¢å¼ |
| | | // èªå·±åçï¼selfSendï¼å§ç»éè¿ï¼å
¨åå¯è§ï¼receiverUserIds 为空 / 缺失ï¼ä¹éè¿ |
| | | const receiverUserIds = websocketMessage.receiverUserIds |
| | | if ( |
| | | currentUserId && |
| | | !selfSend && |
| | | Array.isArray(receiverUserIds) && |
| | | receiverUserIds.length > 0 && |
| | | !receiverUserIds.includes(currentUserId) |
| | | ) { |
| | | console.warn('[IM WS] 丢å¼ä¸å±äºå½åç¨æ·çå®åç¾¤æ¶æ¯', websocketMessage) |
| | | return Promise.resolve() |
| | | } |
| | | |
| | | // 1. 离线å è½½æç¼å²ï¼ä¸ç§èå¯¹ç§°ï¼ |
| | | if (conversationStore.loading) { |
| | | this.messageBuffer.push({ |
| | | conversationType: ImConversationType.GROUP, |
| | | payload: websocketMessage |
| | | }) |
| | | return Promise.resolve() |
| | | } |
| | | |
| | | // 2. æªç¥ç¾¤æ¶èªå¨æç¾¤è¯¦æ
+ æåï¼è¢«æå
¥ç¾¤ä½è¿æ²¡æ¶å° GROUP_CREATE æ¶çå
åºï¼ |
| | | const group = groupStore.getGroup(websocketMessage.groupId) |
| | | if (!group) { |
| | | groupStore.fetchGroupInfo(websocketMessage.groupId).catch(() => undefined) |
| | | } |
| | | |
| | | // 3. å端æ¤åï¼ä¸å䏿¡ RECALL æ¶æ¯ï¼content 为 `{"messageId": xxx}` |
| | | // è¿éæ¦æªä¸æ¥æ¹èµ° recallMessageï¼æåæ¶æ¯æ´æ°ä¸º RECALL æï¼ |
| | | if (websocketMessage.type === ImContentType.RECALL) { |
| | | return useMessageStore().recallMessage( |
| | | ImConversationType.GROUP, |
| | | websocketMessage.groupId, |
| | | websocketMessage.content |
| | | ) |
| | | } |
| | | |
| | | // 4. å端 DTO â å端 Messageï¼åéäººåæ¸²ææ¶å®æ¶ç®ï¼ä¸åå
¥æ¶æ¯å段 |
| | | const message = convertGroupMessage(websocketMessage, currentUserId) |
| | | const persistPromise = useMessageStore().insertMessage( |
| | | { |
| | | type: ImConversationType.GROUP, |
| | | targetId: websocketMessage.groupId, |
| | | name: group ? getGroupDisplayName(group) : String(websocketMessage.groupId), |
| | | avatar: group?.avatar || '', |
| | | silent: group?.silent |
| | | }, |
| | | message |
| | | ) |
| | | |
| | | // 5. ä»
å¯¹æ¹æ¶æ¯æèµ°ãèªå¨å·²è¯» / æç¤ºé³ãï¼ä¸ç§èå¯¹ç§°ï¼ |
| | | if (!selfSend) { |
| | | const conversation = conversationStore.getConversation( |
| | | ImConversationType.GROUP, |
| | | websocketMessage.groupId |
| | | ) |
| | | const isActive = |
| | | conversationStore.activeConversation?.type === ImConversationType.GROUP && |
| | | conversationStore.activeConversation?.targetId === websocketMessage.groupId |
| | | if (isActive) { |
| | | // ç¾¤å·²è¯»ä¸æ¥éè¦å¸¦ messageIdï¼ç¾¤æ¶æ¯ä»¥"读å°ç¬¬å æ¡"çæ¸¸æ 为åï¼åºå«äºç§èåªæ receiverIdï¼ï¼ç¾¤å·²è¯»å
³éæ¶ä»
æ¬å°æ¸
é¶ |
| | | const readReported = conversationStore.isReportedReadPositionCovered( |
| | | ImConversationType.GROUP, |
| | | websocketMessage.groupId, |
| | | websocketMessage.id |
| | | ) |
| | | conversationStore.markConversationRead( |
| | | ImConversationType.GROUP, |
| | | websocketMessage.groupId, |
| | | websocketMessage.id |
| | | ) |
| | | if (MESSAGE_GROUP_READ_ENABLED && !readReported) { |
| | | apiReadGroupMessages(websocketMessage.groupId, websocketMessage.id) |
| | | .then(() => |
| | | conversationStore.markConversationReadReported( |
| | | ImConversationType.GROUP, |
| | | websocketMessage.groupId, |
| | | websocketMessage.id |
| | | ) |
| | | ) |
| | | .catch((error) => { |
| | | console.warn( |
| | | '[IM WS] 群èèªå¨å·²è¯»ä¸æ¥å¤±è´¥', |
| | | { |
| | | conversationType: ImConversationType.GROUP, |
| | | groupId: websocketMessage.groupId, |
| | | messageId: websocketMessage.id |
| | | }, |
| | | error |
| | | ) |
| | | }) |
| | | } |
| | | } else if (!conversation?.silent && isNormalMessage(websocketMessage.type)) { |
| | | // GROUP_* 群广æäºä»¶çç³»ç»æ¶æ¯ä¸åæç¤ºé³ |
| | | playAudioTip() |
| | | } |
| | | } |
| | | return persistPromise |
| | | }, |
| | | |
| | | // ==================== 群è已读 / åæ§ ==================== |
| | | |
| | | /** 群è READï¼èªå·±å
¶å®ç»ç«¯å¨æç¾¤éæ ä¸ºå·²è¯»ï¼æ¬ç«¯åæ¥æ¸
é¶è¯¥ç¾¤æªè¯» + @ 红åï¼ç¾¤å·²è¯»å
³éæ¶å
åºå¿½ç¥ */ |
| | | handleGroupRead(websocketMessage: ImMessageReadNotification) { |
| | | if (!MESSAGE_GROUP_READ_ENABLED) { |
| | | return |
| | | } |
| | | const readMessageId = websocketMessage.readId || websocketMessage.id |
| | | if (!readMessageId || !websocketMessage.groupId) { |
| | | return |
| | | } |
| | | void useConversationStore() |
| | | .applyConversationReadList([ |
| | | { |
| | | id: readMessageId, |
| | | conversationType: ImConversationType.GROUP, |
| | | targetId: websocketMessage.groupId, |
| | | messageId: readMessageId |
| | | } |
| | | ]) |
| | | .catch((error) => console.warn('[IM WS] 群èå·²è¯»åæ¥å¤±è´¥', error)) |
| | | }, |
| | | |
| | | /** 群è RECEIPTï¼æ´æ°ææ¡ç¾¤æ¶æ¯ç readCount / receiptStatusï¼ç¾¤å·²è¯»å
³éæ¶å
åºå¿½ç¥ */ |
| | | handleGroupReceipt(websocketMessage: ImMessageReceiptNotification) { |
| | | if (!MESSAGE_GROUP_READ_ENABLED) { |
| | | return |
| | | } |
| | | if (!websocketMessage.id || !websocketMessage.groupId) { |
| | | return |
| | | } |
| | | useMessageStore().applyMessageReadReceipt({ |
| | | conversationType: ImConversationType.GROUP, |
| | | targetId: websocketMessage.groupId, |
| | | groupMessageId: websocketMessage.id, |
| | | readCount: websocketMessage.readCount, |
| | | receiptStatus: websocketMessage.receiptStatus |
| | | }) |
| | | }, |
| | | |
| | | // ==================== 好åéç¥ï¼1201-1210 段ä½ï¼ ==================== |
| | | |
| | | /** |
| | | * ç® FRIEND_ADD / FRIEND_DELETE 帧çã对端 userIdãï¼ |
| | | * becomeFriends åæ¡å
¥åºååæ¹æ¶å°åä¸ä»½ payloadï¼payload.friendUserId åºå®æ¯ toUserIdï¼æ¬ç«¯çæ£ç对端è¦ä»å¸§ sender / receiver 忍 |
| | | */ |
| | | computeFriendPeerId(frame: ImPrivateMessageNotification): number { |
| | | const currentUserId = getCurrentUserId() |
| | | return getPrivateMessagePeerId(frame, currentUserId) |
| | | }, |
| | | |
| | | /** |
| | | * 好åéç¥ç»ä¸å
¥å£ï¼æ type ååå° friendStore å
é¨ dispatcher |
| | | */ |
| | | handleFriendNotification( |
| | | websocketMessage: ImNoConversationNotification | ImPrivateMessageNotification |
| | | ) { |
| | | const payload = isPrivateMessageNotification(websocketMessage) |
| | | ? parseFriendNotificationPayload(websocketMessage) |
| | | : (websocketMessage as unknown as FriendNotificationPayload) |
| | | const friendStore = useFriendStore() |
| | | switch (websocketMessage.type) { |
| | | case ImContentType.FRIEND_ADD: { |
| | | friendStore.applyFriendAddNotification( |
| | | payload, |
| | | isPrivateMessageNotification(websocketMessage) |
| | | ? this.computeFriendPeerId(websocketMessage) |
| | | : payload.friendUserId |
| | | ) |
| | | break |
| | | } |
| | | case ImContentType.FRIEND_BLOCK: { |
| | | friendStore.applyFriendBlockNotification(payload) |
| | | break |
| | | } |
| | | case ImContentType.FRIEND_DELETE: { |
| | | friendStore.applyFriendDeleteNotification( |
| | | payload, |
| | | isPrivateMessageNotification(websocketMessage) |
| | | ? this.computeFriendPeerId(websocketMessage) |
| | | : payload.friendUserId |
| | | ) |
| | | break |
| | | } |
| | | case ImContentType.FRIEND_INFO_UPDATED: { |
| | | friendStore.applyFriendInfoUpdatedNotification(payload) |
| | | break |
| | | } |
| | | case ImContentType.FRIEND_REQUEST_APPROVED: { |
| | | friendStore.applyFriendRequestApprovedNotification(payload) |
| | | break |
| | | } |
| | | case ImContentType.FRIEND_REQUEST_RECEIVED: { |
| | | friendStore.applyFriendRequestReceivedNotification(payload) |
| | | break |
| | | } |
| | | case ImContentType.FRIEND_REQUEST_REJECTED: { |
| | | friendStore.applyFriendRequestRejectedNotification(payload) |
| | | break |
| | | } |
| | | case ImContentType.FRIEND_UNBLOCK: { |
| | | friendStore.applyFriendUnblockNotification(payload) |
| | | break |
| | | } |
| | | case ImContentType.FRIEND_UPDATE: { |
| | | friendStore.applyFriendUpdateNotification(payload) |
| | | break |
| | | } |
| | | default: { |
| | | console.debug('[IM WS] æªè¯å«å¥½åéç¥', websocketMessage) |
| | | } |
| | | } |
| | | }, |
| | | |
| | | // ==================== å 群ç³è¯·éç¥ï¼1503 / 1505 / 1506ï¼ ==================== |
| | | |
| | | /** |
| | | * å 群ç³è¯·éç¥ç»ä¸å
¥å£ï¼ååå° groupRequestStoreï¼é©±å¨æ¨ªå¹
+ Drawer 忥 |
| | | */ |
| | | handleGroupRequestNotification(websocketMessage: ImNoConversationNotification) { |
| | | const payload = websocketMessage as { requestId?: number } |
| | | if (!payload.requestId) { |
| | | return |
| | | } |
| | | const groupRequestStore = useGroupRequestStore() |
| | | switch (websocketMessage.type) { |
| | | case ImContentType.GROUP_REQUEST_APPROVED: |
| | | case ImContentType.GROUP_REQUEST_REJECTED: { |
| | | groupRequestStore.removeGroupRequestById(payload.requestId) |
| | | break |
| | | } |
| | | case ImContentType.GROUP_REQUEST_RECEIVED: { |
| | | groupRequestStore.addGroupRequestById(payload.requestId).catch(() => undefined) |
| | | break |
| | | } |
| | | default: { |
| | | break |
| | | } |
| | | } |
| | | }, |
| | | |
| | | // ==================== 群å
³ç³»äºä»¶ï¼æ¿è½½äºç¾¤èééï¼æ inner type åæµï¼ ==================== |
| | | |
| | | /** |
| | | * GROUP_MEMBER_SETTING_UPDATEï¼å¤ç«¯åæ¥æåä¸ªäººè®¾ç½®åæ´ï¼silent / groupRemarkï¼ |
| | | * |
| | | * payload æºå¸¦åæ´åæ®µï¼æé null åæ®µç´æ¥å±é¨æ´æ°ï¼ç䏿¬¡ fetchGroupMemberList æ¥å£ |
| | | */ |
| | | handleGroupMemberSettingUpdate(websocketMessage: ImGroupMessageNotification) { |
| | | // content è§£æå¤±è´¥ç±å¤å± dispatchGroupFrame ç try-catch å
åºï¼å« websocketMessage æå°ï¼ï¼ä¸éå¤ catch |
| | | const payload: { groupRemark?: string; silent?: boolean; } = JSON.parse( |
| | | websocketMessage.content || '{}' |
| | | ) |
| | | const groupStore = useGroupStore() |
| | | const group = groupStore.getGroup(websocketMessage.groupId) |
| | | if (!group) { |
| | | return |
| | | } |
| | | const fields: Partial<Group> = {} |
| | | if (payload.silent != null) { |
| | | fields.silent = payload.silent |
| | | } |
| | | if (payload.groupRemark != null) { |
| | | fields.groupRemark = payload.groupRemark |
| | | } |
| | | if (Object.keys(fields).length > 0) { |
| | | groupStore.updateGroupFields(websocketMessage.groupId, fields) |
| | | } |
| | | }, |
| | | |
| | | /** GROUP_MEMBER_NICKNAME_UPDATEï¼åæ¥æåå¨ç¾¤éçæµç§° */ |
| | | handleGroupMemberNicknameUpdate(websocketMessage: ImGroupMessageNotification) { |
| | | useGroupStore().applyGroupNotification( |
| | | websocketMessage.groupId, |
| | | websocketMessage.type, |
| | | websocketMessage.content |
| | | ) |
| | | }, |
| | | |
| | | // ==================== å¿è·³ / éè¿ ==================== |
| | | |
| | | /** å¿è·³å
ï¼çº¯ææ¬ 'ping'ï¼å¯¹åºæå¡ç«¯ 'pong'ï¼å端è¿å±ç¨çº¯å符串约å®ï¼é¿å
JSON è§£æå¼éï¼ */ |
| | | sendHeartBeat() { |
| | | if (this.socket && this.isConnected) { |
| | | this.socket.send('ping') |
| | | } |
| | | }, |
| | | |
| | | /** 䏻卿å¼ï¼åæ¢ç¨æ· / éåºç»å½æ¶ç¨ï¼ï¼å
³ socket + åå¿è·³ + åæ¶å¾
éè¿ */ |
| | | disconnect() { |
| | | if (this.socket) { |
| | | // close() 弿¥è§¦å onclose / onerrorï¼åè°é伿 æ¡ä»¶ reconnectï¼ |
| | | // 主å¨å
³éè·¯å¾å¿
é¡»å
å
¨é¨è§£ç»ï¼å¦å onclose ä¼å¼åèªå¨éè¿ï¼CONNECTING æé´ç in-flight message ä¹å¯è½è¢«è onmessage æéå° stale ä¸ä¸æ |
| | | unbindWebSocketListeners(this.socket) |
| | | this.socket.close() |
| | | this.socket = null |
| | | } |
| | | // onclose 已被解ç»ï¼ä¸ä¼å帮æä»¬è®¾ isConnected=falseï¼è¿éæå¨å¤ä½ |
| | | this.isConnected = false |
| | | this.stopHeartbeat() |
| | | if (this.reconnectTimer) { |
| | | clearTimeout(this.reconnectTimer) |
| | | this.reconnectTimer = null |
| | | } |
| | | // 䏻卿å¼ï¼åè´¦å· / éåºï¼ï¼æ¸
é¶éé¿è®¡æ°ï¼ä¸æ¬¡ connect 鿰仿çé´éèµ·ç® |
| | | this.reconnectAttempts = 0 |
| | | }, |
| | | |
| | | /** |
| | | * èªå¨éè¿ï¼ææ°éé¿ base * 2^attemptï¼å°é¡¶ maxï¼+ 0~jitter ms éæºåç§» |
| | | * |
| | | * onclose æ¯å¯ä¸å
¥å£ï¼onerror ä¸åè°æ¬æ¹æ³ï¼æµè§å¨è§è两è
å¿
åæ¶è§¦åï¼é¿å
è®¡æ° +2ï¼ |
| | | * ä¸è®¾æ¬¡æ°ä¸éï¼é¢çå°é¡¶å¨ WS_RECONNECT_MAX_MSï¼çº¦ 30sï¼æç»éè¯ï¼ç´å°é¾è·¯æ¢å¤æä¸»å¨ disconnect |
| | | */ |
| | | reconnect() { |
| | | this.stopHeartbeat() |
| | | if (this.reconnectTimer) { |
| | | clearTimeout(this.reconnectTimer) |
| | | this.reconnectTimer = null |
| | | } |
| | | const backoff = Math.min( |
| | | WS_RECONNECT_BASE_MS * 2 ** this.reconnectAttempts, |
| | | WS_RECONNECT_MAX_MS |
| | | ) |
| | | const delay = backoff + Math.floor(Math.random() * WS_RECONNECT_JITTER_MS) |
| | | this.reconnectAttempts++ |
| | | console.log(`[IM WS] reconnecting in ${delay}ms (attempt ${this.reconnectAttempts})`) |
| | | this.reconnectTimer = setTimeout(() => { |
| | | this.connect() |
| | | }, delay) |
| | | }, |
| | | |
| | | /** å¿è·³ 5 ç§ä¸æ¬¡ï¼ä¿æ´» + æ¢æ´»ï¼é¾è·¯æäº onclose ä¼è§¦åï¼ç± reconnect å
åºï¼ */ |
| | | startHeartbeat() { |
| | | if (this.heartbeatTimer) clearInterval(this.heartbeatTimer) |
| | | this.heartbeatTimer = setInterval(() => { |
| | | if (this.socket && this.isConnected) { |
| | | this.sendHeartBeat() |
| | | } |
| | | }, 5000) |
| | | }, |
| | | |
| | | /** åå¿è·³ï¼disconnect / éè¿åè°ï¼é¿å
è timer 卿° socket ä¸ç»§ç»è§¦å sendHeartBeat */ |
| | | stopHeartbeat() { |
| | | if (this.heartbeatTimer) { |
| | | clearInterval(this.heartbeatTimer) |
| | | this.heartbeatTimer = null |
| | | } |
| | | }, |
| | | |
| | | // ==================== 宿¶éè¯ä¿¡ä»¤åå ==================== |
| | | |
| | | /** |
| | | * éè¯ä¿¡ä»¤ååï¼1601 RTC_CALLï¼æ status åºå INVITING / JOINED / REJECTED / NO_ANSWER / LEFTï¼+ 1602 / 1603 åä¸è
å å
¥ / ç¦»å¼ |
| | | * <p> |
| | | * åä¸ dispatcherï¼æ type ååå° rtcStore |
| | | */ |
| | | handleRtcSignaling(websocketMessage: ImNoConversationNotification) { |
| | | const rtcStore = useRtcStore() |
| | | switch (websocketMessage.type) { |
| | | case ImContentType.RTC_CALL: { |
| | | const payload = websocketMessage as unknown as ImRtcCallNotification |
| | | switch (payload.status) { |
| | | case ImRtcParticipantStatus.INVITING: { |
| | | if (!isValidRtcInvitePayload(payload)) { |
| | | console.warn('[IM WS] RTC_CALL invite payload ä¸åæ³', payload) |
| | | return |
| | | } |
| | | // å½åå·²å¨éè¯ä¸ï¼å¿½ç¥æ°æ¥çµï¼å端å±é¢ä¹ä¼æç»ï¼è¿éæ¯å
åº |
| | | if (!rtcStore.isActive) { |
| | | rtcStore.showIncoming(payload) |
| | | } |
| | | break |
| | | } |
| | | case ImRtcParticipantStatus.JOINED: |
| | | case ImRtcParticipantStatus.LEFT: { |
| | | // ACCEPT / HUNGUP æä¸éè¦æ¬ç«¯é¢å¤ååºï¼rtcStore ç¶æç± 1602/1603 + END ç»´æ¤ |
| | | break |
| | | } |
| | | case ImRtcParticipantStatus.NO_ANSWER: { |
| | | // 群éè¯å人æ¯éè¶
æ¶ï¼ä¿¡ä»¤ç¬ç«ä¿çè¯ä¹ï¼å¤çä¸ REJECTED ä¸è´ |
| | | rtcStore.applyParticipantNoAnswer(payload) |
| | | break |
| | | } |
| | | case ImRtcParticipantStatus.REJECTED: { |
| | | rtcStore.applyParticipantRejected(payload) |
| | | break |
| | | } |
| | | default: { |
| | | console.warn('[IM WS] æªè¯å«ç RTC_CALL status', payload) |
| | | } |
| | | } |
| | | return |
| | | } |
| | | case ImContentType.RTC_PARTICIPANT_CONNECTED: { |
| | | const payload = websocketMessage as unknown as ImRtcParticipantConnectedNotification |
| | | if (payload?.room && payload.userId) { |
| | | rtcStore.applyParticipantConnected(payload) |
| | | } |
| | | return |
| | | } |
| | | case ImContentType.RTC_PARTICIPANT_DISCONNECTED: { |
| | | const payload = websocketMessage as unknown as ImRtcParticipantDisconnectedNotification |
| | | if (payload?.room && payload.userId) { |
| | | rtcStore.applyParticipantDisconnected(payload) |
| | | } |
| | | } |
| | | } |
| | | }, |
| | | |
| | | /** RTC_CALL_START éè¯å¼å§ */ |
| | | handleRtcCallStart(websocketMessage: ImGroupMessageNotification) { |
| | | const payload = parseRtcCallPayload(websocketMessage.content) |
| | | if (!payload?.room || !payload.mediaType || !payload.inviterUserId) { |
| | | console.warn('[IM WS] RTC_CALL_START payload ä¸åæ³', { |
| | | groupId: websocketMessage.groupId, |
| | | messageId: websocketMessage.id, |
| | | contentLength: websocketMessage.content?.length ?? 0 |
| | | }) |
| | | return |
| | | } |
| | | useRtcStore().setGroupCall({ |
| | | room: payload.room, |
| | | groupId: websocketMessage.groupId, |
| | | mediaType: payload.mediaType, |
| | | inviterId: payload.inviterUserId, |
| | | joinedUserIds: [payload.inviterUserId], |
| | | inviteeIds: [] |
| | | }) |
| | | }, |
| | | |
| | | /** |
| | | * RTC_CALL_END éè¯ç»æï¼ç§è + 群èé½èµ°è¿ä¸æ¡ï¼payload æºå¸¦ conversationType åºå |
| | | * <p> |
| | | * ç§èï¼å
³éå½åéè¯çª |
| | | * 群èï¼ç§»é¤è¶åæ¡ï¼å¦æ¬ç«¯å¨è¯¥ç¾¤éè¯å
åå
³ééè¯çª |
| | | */ |
| | | handleRtcCallEnd( |
| | | websocketMessage: ImGroupMessageNotification | ImPrivateMessageNotification |
| | | ) { |
| | | const payload = this.safeParse(websocketMessage.content) as ImRtcCallEndNotification | null |
| | | if (!payload?.room) { |
| | | return |
| | | } |
| | | const rtcStore = useRtcStore() |
| | | const isGroup = payload.conversationType === ImConversationType.GROUP |
| | | // 群éè¯ï¼ç§»é¤å¯¹åºæ¿é´çè¶åæ¡ |
| | | const groupId = (websocketMessage as ImGroupMessageNotification).groupId |
| | | if (isGroup && groupId) { |
| | | rtcStore.removeGroupCall(groupId, payload.room) |
| | | } |
| | | // éè¯çª / æ¥çµçªæååä¸ room æ¶å
³éï¼ |
| | | // RUNNING / INVITING é¶æ®µå¯¹æ¯ call.roomï¼INCOMING é¶æ®µå¯¹æ¯ incomingPayload.room |
| | | const matchCall = rtcStore.call?.room === payload.room |
| | | const matchIncoming = rtcStore.incomingPayload?.room === payload.room |
| | | if (rtcStore.isActive && (matchCall || matchIncoming)) { |
| | | const reasonText = resolveCallEndReasonText(payload.endReason) |
| | | console.info('[Call] end:', reasonText) |
| | | rtcStore.reset() |
| | | } |
| | | } |
| | | } |
| | | }) |
| | | |
| | | export const useImWebSocketStoreWithOut = () => { |
| | | return useImWebSocketStore() |
| | | } |
| | | |
| | | // dev: 让 Pinia ç actions / state æ¹å¨æ¯æ HMRï¼é¿å
æ¯æ¬¡æ¹ store é½å¾ç¡¬å· |
| | | // å¦å Vite ææ°æ¨¡åæ¨ä¸æ¥åï¼è store å®ä¾ç action éå
仿忧彿°ä½ |
| | | if (import.meta.hot) { |
| | | import.meta.hot.accept(acceptHMRUpdate(useImWebSocketStore, import.meta.hot)) |
| | | } |