| ¶Ô±ÈÐÂÎļþ |
| | |
| | | import type { Conversation, Message, MessageDO } from '../types' |
| | | |
| | | import { acceptHMRUpdate, defineStore } from 'pinia' |
| | | |
| | | import { getCurrentUserId } from '#/views/im/utils/auth' |
| | | |
| | | import { |
| | | IM_AT_ALL_USER_ID, |
| | | ImContentType, |
| | | ImConversationType, |
| | | ImMessageReceiptStatus, |
| | | ImMessageStatus, |
| | | isGroupNotification, |
| | | isNormalMessage |
| | | } from '../../utils/constants' |
| | | import { resolveConversationLastContent } from '../../utils/conversation' |
| | | import { |
| | | type DbTransaction, |
| | | getClientConversationId, |
| | | getClientMessageKey, |
| | | getDb, |
| | | getServerMessageKey, |
| | | parseClientConversationId, |
| | | setMessageMaxId, |
| | | StorageKeys |
| | | } from '../../utils/db' |
| | | import { |
| | | generateClientMessageId, |
| | | parseRecallMessageId, |
| | | revokeBlobUrlsInContent |
| | | } from '../../utils/message' |
| | | import { isGroupQuit, tryGetSenderDisplayName } from '../../utils/user' |
| | | import { useConversationStore } from './conversationStore' |
| | | import { useGroupStore } from './groupStore' |
| | | |
| | | const MESSAGE_CACHE_RECENT_CONVERSATION_LIMIT = 5 |
| | | const MESSAGE_CACHE_RETAIN_CONVERSATION_LIMIT = MESSAGE_CACHE_RECENT_CONVERSATION_LIMIT + 1 |
| | | const ackMergingPromises = new Map<string, Promise<void>>() |
| | | |
| | | interface MessageConversationInfo { |
| | | type: number |
| | | targetId: number |
| | | name: string |
| | | avatar: string |
| | | silent?: boolean |
| | | } |
| | | |
| | | interface PersistMessageRecordOptions { |
| | | mergeClientRecord?: boolean |
| | | } |
| | | |
| | | /** æåæ¶æ¯æ¹éå¤ç项 */ |
| | | export type PulledMessage = |
| | | | { |
| | | conversationInfo: MessageConversationInfo |
| | | kind: 'insert' |
| | | message: Message |
| | | } |
| | | | { |
| | | conversationType: number |
| | | kind: 'recall' |
| | | recallSignalContent: string |
| | | targetId: number |
| | | } |
| | | |
| | | /** è·åä¼è¯çæ¶æ¯ç¼å key */ |
| | | function getMessageCacheKey(type: number, targetId: number): string { |
| | | return getClientConversationId(type, targetId) |
| | | } |
| | | |
| | | /** çææ¶æ¯æ¬å°ä¸»é® */ |
| | | function getMessageKey( |
| | | message: Pick<Message, 'clientMessageId' | 'id'>, |
| | | conversationType: number |
| | | ): string { |
| | | return message.id |
| | | ? getServerMessageKey(conversationType, message.id) |
| | | : getClientMessageKey(message.clientMessageId) |
| | | } |
| | | |
| | | /** è¡¥é½å®¢æ·ç«¯æ¶æ¯ç¼å· */ |
| | | function ensureClientMessageId(message: Message): Message { |
| | | if (!message.clientMessageId) { |
| | | message.clientMessageId = generateClientMessageId() |
| | | } |
| | | if (!message.id) { |
| | | message.id = undefined |
| | | } |
| | | return message |
| | | } |
| | | |
| | | /** 转æ¢ä¸º IndexedDB æ¶æ¯è®°å½ */ |
| | | function buildMessageDO(message: Message, conversationType: number): MessageDO { |
| | | return { |
| | | id: message.id, |
| | | clientMessageId: message.clientMessageId, |
| | | type: message.type, |
| | | content: message.content, |
| | | status: message.status, |
| | | sendTime: message.sendTime, |
| | | senderId: message.senderId, |
| | | atUserIds: message.atUserIds ? [...message.atUserIds] : undefined, |
| | | receiverUserIds: message.receiverUserIds ? [...message.receiverUserIds] : undefined, |
| | | receiptStatus: message.receiptStatus, |
| | | readCount: message.readCount, |
| | | materialId: message.materialId, |
| | | targetId: message.targetId, |
| | | selfSend: message.selfSend, |
| | | messageKey: getMessageKey(message, conversationType), |
| | | conversationType, |
| | | clientConversationId: getClientConversationId(conversationType, message.targetId) |
| | | } |
| | | } |
| | | |
| | | /** IndexedDB æ¶æ¯è®°å½è½¬åç«¯æ¶æ¯ */ |
| | | function buildMessageFromDO(message: MessageDO): Message { |
| | | const { |
| | | messageKey: _messageKey, |
| | | conversationType: _conversationType, |
| | | clientConversationId: _clientConversationId, |
| | | ...rest |
| | | } = message |
| | | return rest |
| | | } |
| | | |
| | | /** ç®åºæ«æ¡æ¶æ¯çåé人快ç
§ */ |
| | | function deriveLastSenderDisplayName( |
| | | conversation: Conversation, |
| | | senderId: number |
| | | ): string | undefined { |
| | | // 1. ä¼å
使ç¨å½åå
åä¸ç好å / 群æåä¿¡æ¯ |
| | | const liveSenderName = tryGetSenderDisplayName(senderId, conversation.type, conversation.targetId) |
| | | if (liveSenderName) { |
| | | return liveSenderName |
| | | } |
| | | // 2. 群æåç¼å缺失æ¶å¼æ¥è¡¥é½ |
| | | if (conversation.type === ImConversationType.GROUP) { |
| | | const groupStore = useGroupStore() |
| | | const group = groupStore.getGroup(conversation.targetId) |
| | | if (!group || isGroupQuit(group)) { |
| | | return conversation.lastSenderId === senderId ? conversation.lastSenderDisplayName : undefined |
| | | } |
| | | const fetchPromise = |
| | | group?.membersLoaded && !group.membersExpired |
| | | ? groupStore.fetchGroupMember(conversation.targetId, senderId) |
| | | : groupStore.fetchGroupMemberList(conversation.targetId) |
| | | fetchPromise.catch((error) => |
| | | console.warn( |
| | | '[IM messageStore] å
åºæç¾¤æå失败', |
| | | { groupId: conversation.targetId, senderId, fullFetch: !group?.membersLoaded }, |
| | | error |
| | | ) |
| | | ) |
| | | } |
| | | return conversation.lastSenderId === senderId ? conversation.lastSenderDisplayName : undefined |
| | | } |
| | | |
| | | /** ææ¶æ¯æ´æ°ä¼è¯æè¦ */ |
| | | function applyConversationSummary(conversation: Conversation, message: Message): void { |
| | | const senderDisplayName = deriveLastSenderDisplayName(conversation, message.senderId) |
| | | conversation.lastContent = resolveConversationLastContent( |
| | | message, |
| | | conversation.type, |
| | | conversation.targetId, |
| | | senderDisplayName |
| | | ) |
| | | conversation.lastSendTime = message.sendTime || Date.now() |
| | | conversation.lastSenderId = message.senderId |
| | | conversation.lastMessageType = message.type |
| | | conversation.lastMessageId = message.id |
| | | conversation.lastClientMessageId = message.clientMessageId |
| | | conversation.lastMessageStatus = message.status |
| | | conversation.lastReceiptStatus = message.receiptStatus |
| | | conversation.lastSelfSend = message.selfSend |
| | | conversation.lastSenderDisplayName = senderDisplayName |
| | | } |
| | | |
| | | /** ææ«æ¡æ¶æ¯éç®ä¼è¯æè¦ */ |
| | | function recomputeConversationLast(conversation: Conversation, messages: Message[]): void { |
| | | const last = messages[messages.length - 1] |
| | | if (last) { |
| | | applyConversationSummary(conversation, last) |
| | | return |
| | | } |
| | | conversation.lastContent = '' |
| | | conversation.lastSendTime = 0 |
| | | conversation.lastSenderId = undefined |
| | | conversation.lastMessageType = undefined |
| | | conversation.lastMessageId = undefined |
| | | conversation.lastClientMessageId = undefined |
| | | conversation.lastMessageStatus = undefined |
| | | conversation.lastReceiptStatus = undefined |
| | | conversation.lastSelfSend = undefined |
| | | conversation.lastSenderDisplayName = undefined |
| | | } |
| | | |
| | | /** åæ¥ç¾¤ @ ç¶æ */ |
| | | function syncConversationAtFlags(conversation: Conversation, message: Message): void { |
| | | if ( |
| | | message.selfSend || |
| | | conversation.type !== ImConversationType.GROUP || |
| | | !message.atUserIds || |
| | | message.atUserIds.length === 0 |
| | | ) { |
| | | return |
| | | } |
| | | const currentUserId = getCurrentUserId() |
| | | if (currentUserId && message.atUserIds.includes(currentUserId)) { |
| | | conversation.atMe = true |
| | | } |
| | | if (message.atUserIds.includes(IM_AT_ALL_USER_ID)) { |
| | | conversation.atAll = true |
| | | } |
| | | } |
| | | |
| | | /** åºç¨æå¡ç«¯æ¶æ¯æ´æ° */ |
| | | function applyServerMessageUpdate(message: Message, updates: Partial<Message>): void { |
| | | if (updates.content && updates.content !== message.content) { |
| | | revokeBlobUrlsInContent(message.content) |
| | | } |
| | | Object.assign(message, updates) |
| | | if (updates.id === 0) { |
| | | message.id = undefined |
| | | } |
| | | if (updates.status !== undefined && updates.status !== ImMessageStatus.SENDING) { |
| | | message.uploadProgress = undefined |
| | | if (updates.status !== ImMessageStatus.FAILED) { |
| | | message._localFile = undefined |
| | | } |
| | | } |
| | | } |
| | | |
| | | /** 夿æ¯å¦ä¸ºå䏿¡æ¶æ¯ */ |
| | | function isSameMessage(left: Message, right: Message): boolean { |
| | | if (left.id && right.id && left.id === right.id) { |
| | | return true |
| | | } |
| | | return !!left.clientMessageId && left.clientMessageId === right.clientMessageId |
| | | } |
| | | |
| | | export const useMessageStore = defineStore('imMessageStore', { |
| | | state: () => ({ |
| | | messagesByConversation: {} as Record<string, Message[]>, |
| | | loadedConversationKeys: [] as string[], |
| | | privateReadMaxIds: {} as Partial<Record<number, number>>, |
| | | privateMessageMaxId: 0, |
| | | groupMessageMaxId: 0, |
| | | channelMessageMaxId: 0 |
| | | }), |
| | | |
| | | getters: { |
| | | /** è·åä¼è¯å·²å è½½æ¶æ¯ */ |
| | | getMessages: |
| | | (state) => |
| | | (clientConversationId: string): Message[] => |
| | | state.messagesByConversation[clientConversationId] || [] |
| | | }, |
| | | |
| | | actions: { |
| | | /** æ¸
ç©ºæ¶æ¯å
å */ |
| | | clear() { |
| | | Object.values(this.messagesByConversation).forEach((messages) => { |
| | | messages.forEach((message) => { |
| | | revokeBlobUrlsInContent(message.content) |
| | | message._localFile = undefined |
| | | }) |
| | | }) |
| | | this.messagesByConversation = {} |
| | | this.loadedConversationKeys = [] |
| | | this.privateReadMaxIds = {} |
| | | this.privateMessageMaxId = 0 |
| | | this.groupMessageMaxId = 0 |
| | | this.channelMessageMaxId = 0 |
| | | ackMergingPromises.clear() |
| | | }, |
| | | |
| | | /** ä» settings å è½½æ¶æ¯æ¸¸æ */ |
| | | async loadMessageCursorList() { |
| | | const db = getDb() |
| | | const [privateMaxId, groupMaxId, channelMaxId] = await Promise.all([ |
| | | db.getSetting<number>(StorageKeys.settings.privateMessageMaxId), |
| | | db.getSetting<number>(StorageKeys.settings.groupMessageMaxId), |
| | | db.getSetting<number>(StorageKeys.settings.channelMessageMaxId) |
| | | ]) |
| | | this.privateMessageMaxId = privateMaxId || 0 |
| | | this.groupMessageMaxId = groupMaxId || 0 |
| | | this.channelMessageMaxId = channelMaxId || 0 |
| | | }, |
| | | |
| | | /** æ´æ°å
忏¸æ */ |
| | | updateMessageCursor(conversationType: number, messageId?: number) { |
| | | if (!messageId) { |
| | | return |
| | | } |
| | | if (conversationType === ImConversationType.PRIVATE && messageId > this.privateMessageMaxId) { |
| | | this.privateMessageMaxId = messageId |
| | | } else if ( |
| | | conversationType === ImConversationType.GROUP && |
| | | messageId > this.groupMessageMaxId |
| | | ) { |
| | | this.groupMessageMaxId = messageId |
| | | } else if ( |
| | | conversationType === ImConversationType.CHANNEL && |
| | | messageId > this.channelMessageMaxId |
| | | ) { |
| | | this.channelMessageMaxId = messageId |
| | | } |
| | | }, |
| | | |
| | | /** è·åç§è对æ¹å·²è¯»ä½ç½®ç¼å */ |
| | | getPrivateReadMaxId(peerId: number): number | undefined { |
| | | return this.privateReadMaxIds[peerId] |
| | | }, |
| | | |
| | | /** æ´æ°ç§è对æ¹å·²è¯»ä½ç½®ç¼å */ |
| | | updatePrivateReadMaxId(peerId: number, maxReadId: null | number = 0): number { |
| | | if (!peerId) { |
| | | return 0 |
| | | } |
| | | const nextMaxReadId = maxReadId || 0 |
| | | const current = this.getPrivateReadMaxId(peerId) |
| | | if (current !== undefined && nextMaxReadId <= current) { |
| | | return current |
| | | } |
| | | this.privateReadMaxIds = { ...this.privateReadMaxIds, [peerId]: nextMaxReadId } |
| | | return nextMaxReadId |
| | | }, |
| | | |
| | | /** æ¸
空ç§è对æ¹å·²è¯»ä½ç½®ç¼å */ |
| | | clearPrivateReadMaxIdCache(): void { |
| | | this.privateReadMaxIds = {} |
| | | }, |
| | | |
| | | /** æ è®°ä¼è¯è¿æä½¿ç¨ */ |
| | | touchConversationMessageCache(clientConversationId: string) { |
| | | this.loadedConversationKeys = [ |
| | | clientConversationId, |
| | | ...this.loadedConversationKeys.filter((key) => key !== clientConversationId) |
| | | ] |
| | | // ä¿çå½åæ´»è·ä¼è¯ + æè¿æå¼è¿çä¼è¯ |
| | | const retained = this.loadedConversationKeys.slice(0, MESSAGE_CACHE_RETAIN_CONVERSATION_LIMIT) |
| | | const removed = this.loadedConversationKeys.slice(MESSAGE_CACHE_RETAIN_CONVERSATION_LIMIT) |
| | | this.loadedConversationKeys = retained |
| | | removed.forEach((key) => { |
| | | Reflect.deleteProperty(this.messagesByConversation, key) |
| | | }) |
| | | }, |
| | | |
| | | /** å è½½å½åä¼è¯æè¿æ¶æ¯ */ |
| | | async loadMoreMessageList( |
| | | clientConversationId: string, |
| | | beforeSendTime?: number, |
| | | limit = 50 |
| | | ): Promise<Message[]> { |
| | | // 1. ä» IndexedDB ååºè¯»åä¸é¡µï¼è¿ååå·²ææ¶é´ååºæå |
| | | const list = await getDb().getMessageListByConversation(clientConversationId, { |
| | | beforeSendTime, |
| | | limit |
| | | }) |
| | | // 2. åå¹¶å°å
åç¼åï¼è¿æ»¤å·²åå¨çæ¶æ¯ |
| | | const parsed = parseClientConversationId(clientConversationId) |
| | | if (!parsed) { |
| | | return [] |
| | | } |
| | | const messages = list.map((message) => buildMessageFromDO(message)) |
| | | const existing = this.messagesByConversation[clientConversationId] || [] |
| | | const existingKeys = new Set(existing.map((message) => getMessageKey(message, parsed.type))) |
| | | const fresh = messages.filter( |
| | | (message) => !existingKeys.has(getMessageKey(message, parsed.type)) |
| | | ) |
| | | this.messagesByConversation[clientConversationId] = [...fresh, ...existing].toSorted( |
| | | (messageA, messageB) => (messageA.sendTime || 0) - (messageB.sendTime || 0) |
| | | ) |
| | | this.touchConversationMessageCache(clientConversationId) |
| | | return fresh |
| | | }, |
| | | |
| | | /** ç¡®ä¿ä¼è¯æ¶æ¯å·²å è½½ */ |
| | | async ensureConversationMessageListLoaded(conversation: Conversation) { |
| | | const key = getMessageCacheKey(conversation.type, conversation.targetId) |
| | | if (this.messagesByConversation[key]) { |
| | | this.touchConversationMessageCache(key) |
| | | return |
| | | } |
| | | await this.loadMoreMessageList(key) |
| | | }, |
| | | |
| | | /** è·åå
åæ¶æ¯æ°ç» */ |
| | | getMessageList(conversationType: number, targetId: number): Message[] { |
| | | const key = getMessageCacheKey(conversationType, targetId) |
| | | if (!this.messagesByConversation[key]) { |
| | | this.messagesByConversation[key] = [] |
| | | } |
| | | this.touchConversationMessageCache(key) |
| | | return this.messagesByConversation[key] |
| | | }, |
| | | |
| | | /** æä¹
åæ¶æ¯è®°å½ */ |
| | | async saveMessageRecord( |
| | | message: Message, |
| | | conversationType: number, |
| | | tx?: DbTransaction, |
| | | options?: PersistMessageRecordOptions |
| | | ) { |
| | | const db = getDb() |
| | | const next = buildMessageDO(message, conversationType) |
| | | // æå¡ç«¯ key æ¿æ¢ client key |
| | | if (options?.mergeClientRecord && message.id && message.clientMessageId) { |
| | | const existing = await db.getByIndex<MessageDO>( |
| | | 'messages', |
| | | 'clientMessageId', |
| | | message.clientMessageId, |
| | | tx |
| | | ) |
| | | if (existing && existing.messageKey !== next.messageKey) { |
| | | await db.delete('messages', existing.messageKey, tx) |
| | | } |
| | | } |
| | | await db.put('messages', next, tx) |
| | | }, |
| | | |
| | | /** ä¿åæ¶æ¯æ¸¸æ */ |
| | | async saveMessageCursor(conversationType: number, messageId?: number, tx?: DbTransaction) { |
| | | await setMessageMaxId(conversationType, messageId, tx) |
| | | this.updateMessageCursor(conversationType, messageId) |
| | | }, |
| | | |
| | | /** åºç¨æ¤åå°å
å */ |
| | | applyRecallMessageInMemory( |
| | | conversationType: number, |
| | | targetId: number, |
| | | recallSignalContent: string |
| | | ) { |
| | | // 1. å®ä½è¢«æ¤åçåæ¶æ¯ |
| | | const messageId = parseRecallMessageId(recallSignalContent) |
| | | if (!messageId) { |
| | | return null |
| | | } |
| | | const conversationStore = useConversationStore() |
| | | const conversation = conversationStore.getConversation(conversationType, targetId) |
| | | if (!conversation) { |
| | | return null |
| | | } |
| | | const messages = this.getMessageList(conversationType, targetId) |
| | | const message = messages.find((item) => item.id === messageId) |
| | | if (!message) { |
| | | return null |
| | | } |
| | | // 2. æ´æ°æ¶æ¯åä¼è¯æè¦ |
| | | message.type = ImContentType.RECALL |
| | | message.status = ImMessageStatus.RECALL |
| | | message.content = '' |
| | | if (messages[messages.length - 1]?.id === messageId) { |
| | | recomputeConversationLast(conversation, messages) |
| | | } |
| | | return { conversation, message } |
| | | }, |
| | | |
| | | /** æ¹éåå
¥æåæ¶æ¯ */ |
| | | async applyPulledMessageList( |
| | | pulledMessages: PulledMessage[], |
| | | conversationType: number, |
| | | maxMessageId?: number |
| | | ) { |
| | | if (pulledMessages.length === 0) { |
| | | // 1. ç©ºæ¹æ¬¡åªæ¨è¿æ¸¸æ |
| | | await this.saveMessageCursor(conversationType, maxMessageId) |
| | | return |
| | | } |
| | | const conversationStore = useConversationStore() |
| | | const persistedMessages = new Map< |
| | | string, |
| | | { conversationType: number; mergeClientRecord?: boolean; message: Message; } |
| | | >() |
| | | const changedConversations = new Map<string, Conversation>() |
| | | |
| | | const addChanged = ( |
| | | conversation: Conversation, |
| | | message: Message, |
| | | options?: PersistMessageRecordOptions |
| | | ) => { |
| | | const clientConversationId = getClientConversationId( |
| | | conversation.type, |
| | | conversation.targetId |
| | | ) |
| | | changedConversations.set(clientConversationId, conversation) |
| | | persistedMessages.set(getMessageKey(message, conversation.type), { |
| | | message, |
| | | conversationType: conversation.type, |
| | | mergeClientRecord: options?.mergeClientRecord |
| | | }) |
| | | } |
| | | |
| | | // 1. å
æ´æ°å
åï¼æ¶ééè¦æä¹
åçæ¶æ¯åä¼è¯ |
| | | for (const pulledMessage of pulledMessages) { |
| | | if (pulledMessage.kind === 'recall') { |
| | | // 1.1 æ¤åä¿¡å·æ´æ°åæ¶æ¯ |
| | | const changed = this.applyRecallMessageInMemory( |
| | | pulledMessage.conversationType, |
| | | pulledMessage.targetId, |
| | | pulledMessage.recallSignalContent |
| | | ) |
| | | if (changed) { |
| | | addChanged(changed.conversation, changed.message) |
| | | } |
| | | continue |
| | | } |
| | | |
| | | const { conversationInfo } = pulledMessage |
| | | const hasServerClientMessageId = !!pulledMessage.message.clientMessageId |
| | | const message = ensureClientMessageId(pulledMessage.message) |
| | | // 1.2 ç¡®ä¿ä¼è¯åæ¶æ¯ç¼ååå¨ |
| | | const conversation = conversationStore.ensureConversation(conversationInfo) |
| | | const messages = this.getMessageList(conversationInfo.type, conversationInfo.targetId) |
| | | const existingIndex = messages.findIndex((existing) => isSameMessage(existing, message)) |
| | | if (existingIndex !== -1) { |
| | | const existing = messages[existingIndex] |
| | | if (!existing) { |
| | | continue |
| | | } |
| | | // 1.3 å·²å卿¶æ¯åå¹¶æå¡ç«¯ç¶æ |
| | | applyServerMessageUpdate(existing, message) |
| | | if (existingIndex === messages.length - 1) { |
| | | recomputeConversationLast(conversation, messages) |
| | | syncConversationAtFlags(conversation, message) |
| | | } |
| | | addChanged(conversation, existing, { |
| | | mergeClientRecord: hasServerClientMessageId |
| | | }) |
| | | continue |
| | | } |
| | | |
| | | // 1.4 æ°æ¶æ¯æ´æ°ä¼è¯æè¦åæªè¯»ç¶æ |
| | | applyConversationSummary(conversation, message) |
| | | syncConversationAtFlags(conversation, message) |
| | | const isActive = |
| | | conversationStore.activeConversation?.type === conversationInfo.type && |
| | | conversationStore.activeConversation?.targetId === conversationInfo.targetId |
| | | if ( |
| | | !message.selfSend && |
| | | !isActive && |
| | | !conversationStore.isMessageCoveredByReadPosition(conversation, message) && |
| | | isNormalMessage(message.type) && |
| | | message.status !== ImMessageStatus.RECALL |
| | | ) { |
| | | conversation.unreadCount++ |
| | | } |
| | | |
| | | // 1.5 æ°æ¶æ¯ææå¡ç«¯ id æå
¥å
åå表 |
| | | let insertIndex = messages.length |
| | | if (message.id) { |
| | | for (const [index, existing] of messages.entries()) { |
| | | if (existing.id && message.id < existing.id) { |
| | | insertIndex = index |
| | | break |
| | | } |
| | | } |
| | | } |
| | | messages.splice(insertIndex, 0, message) |
| | | addChanged(conversation, message, { |
| | | mergeClientRecord: hasServerClientMessageId && !!message.id |
| | | }) |
| | | } |
| | | |
| | | // 2. åäºå¡åå
¥æ¶æ¯ãä¼è¯æè¦å游æ |
| | | await getDb().transaction( |
| | | ['messages', 'conversations', 'settings'], |
| | | 'readwrite', |
| | | async (tx) => { |
| | | // 2.1 åå
¥æ¬æ¹åæ´æ¶æ¯ |
| | | for (const item of persistedMessages.values()) { |
| | | await this.saveMessageRecord(item.message, item.conversationType, tx, { |
| | | mergeClientRecord: item.mergeClientRecord |
| | | }) |
| | | } |
| | | // 2.2 åå
¥æ¬æ¹åæ´ä¼è¯ |
| | | await conversationStore.saveConversationRecord([...changedConversations.values()], tx) |
| | | // 2.3 åå
¥æ¬æ¹æ¸¸æ |
| | | await setMessageMaxId(conversationType, maxMessageId, tx) |
| | | } |
| | | ) |
| | | // 3. æä¹
åæå忍è¿å
忏¸æ |
| | | this.updateMessageCursor(conversationType, maxMessageId) |
| | | for (const item of persistedMessages.values()) { |
| | | this.updateMessageCursor(item.conversationType, item.message.id) |
| | | } |
| | | }, |
| | | |
| | | /** æå
¥æ¶æ¯ */ |
| | | insertMessage( |
| | | conversationInfo: MessageConversationInfo, |
| | | messageInfo: Message, |
| | | options?: { saveMaxId?: boolean } |
| | | ): Promise<void> { |
| | | const conversationStore = useConversationStore() |
| | | const hasIncomingClientMessageId = !!messageInfo.clientMessageId |
| | | const message = ensureClientMessageId(messageInfo) |
| | | // 1. å
å¤çæ¶æ¯å¸¦æ¥çç¾¤èµæåæ´ |
| | | if (conversationInfo.type === ImConversationType.GROUP && isGroupNotification(message.type)) { |
| | | useGroupStore().applyGroupNotification( |
| | | conversationInfo.targetId, |
| | | message.type, |
| | | message.content |
| | | ) |
| | | } |
| | | |
| | | // 2. ç¡®ä¿ä¼è¯åæ¶æ¯ç¼ååå¨ |
| | | const conversation = conversationStore.ensureConversation(conversationInfo) |
| | | const messages = this.getMessageList(conversationInfo.type, conversationInfo.targetId) |
| | | const existingIndex = messages.findIndex((item) => isSameMessage(item, message)) |
| | | // 3. å·²å卿¶æ¯èµ°è¦çæ´æ° |
| | | if (existingIndex !== -1) { |
| | | const existing = messages[existingIndex] |
| | | if (!existing) { |
| | | return Promise.resolve() |
| | | } |
| | | applyServerMessageUpdate(existing, message) |
| | | if (existingIndex === messages.length - 1) { |
| | | recomputeConversationLast(conversation, messages) |
| | | syncConversationAtFlags(conversation, message) |
| | | } |
| | | return getDb() |
| | | .transaction(['messages', 'conversations', 'settings'], 'readwrite', async (tx) => { |
| | | await this.saveMessageRecord(existing, conversationInfo.type, tx, { |
| | | mergeClientRecord: hasIncomingClientMessageId |
| | | }) |
| | | await conversationStore.saveConversationRecord(conversation, tx) |
| | | if (options?.saveMaxId !== false) { |
| | | await setMessageMaxId(conversationInfo.type, message.id, tx) |
| | | } |
| | | }) |
| | | .catch((error) => { |
| | | console.error('[IM messageStore] æ¶æ¯åå
¥å¤±è´¥', error) |
| | | throw error |
| | | }) |
| | | .then(() => { |
| | | this.updateMessageCursor(conversationInfo.type, message.id) |
| | | }) |
| | | } |
| | | |
| | | // 4. æ°æ¶æ¯æ´æ°ä¼è¯æè¦åæªè¯»ç¶æ |
| | | applyConversationSummary(conversation, message) |
| | | syncConversationAtFlags(conversation, message) |
| | | |
| | | const isActive = |
| | | conversationStore.activeConversation?.type === conversationInfo.type && |
| | | conversationStore.activeConversation?.targetId === conversationInfo.targetId |
| | | if ( |
| | | !message.selfSend && |
| | | !isActive && |
| | | !conversationStore.isMessageCoveredByReadPosition(conversation, message) && |
| | | isNormalMessage(message.type) && |
| | | message.status !== ImMessageStatus.RECALL |
| | | ) { |
| | | conversation.unreadCount++ |
| | | } |
| | | |
| | | // 5. æ°æ¶æ¯æ id æå
¥å°å
åæ°ç» |
| | | let insertIndex = messages.length |
| | | if (message.id) { |
| | | for (const [index, existing] of messages.entries()) { |
| | | if (existing.id && message.id < existing.id) { |
| | | insertIndex = index |
| | | break |
| | | } |
| | | } |
| | | } |
| | | messages.splice(insertIndex, 0, message) |
| | | // 6. åäºå¡åå
¥æ¶æ¯ãä¼è¯æè¦å游æ |
| | | return getDb() |
| | | .transaction(['messages', 'conversations', 'settings'], 'readwrite', async (tx) => { |
| | | await this.saveMessageRecord(message, conversationInfo.type, tx, { |
| | | mergeClientRecord: hasIncomingClientMessageId && !!message.id |
| | | }) |
| | | await conversationStore.saveConversationRecord(conversation, tx) |
| | | if (options?.saveMaxId !== false) { |
| | | await setMessageMaxId(conversationInfo.type, message.id, tx) |
| | | } |
| | | }) |
| | | .catch((error) => { |
| | | console.error('[IM messageStore] æ¶æ¯åå
¥å¤±è´¥', error) |
| | | throw error |
| | | }) |
| | | .then(() => { |
| | | this.updateMessageCursor(conversationInfo.type, message.id) |
| | | }) |
| | | }, |
| | | |
| | | /** ack åå¹¶ */ |
| | | ackMessage( |
| | | conversationType: number, |
| | | targetId: number, |
| | | clientMessageId: string, |
| | | updates: Partial<Message> |
| | | ) { |
| | | const mergeKey = `${conversationType}:${targetId}:${clientMessageId}` |
| | | const existingPromise = ackMergingPromises.get(mergeKey) |
| | | if (existingPromise) { |
| | | return existingPromise |
| | | } |
| | | const promise = this.doAckMessage( |
| | | conversationType, |
| | | targetId, |
| | | clientMessageId, |
| | | updates |
| | | ).finally(() => { |
| | | ackMergingPromises.delete(mergeKey) |
| | | }) |
| | | ackMergingPromises.set(mergeKey, promise) |
| | | return promise |
| | | }, |
| | | |
| | | /** æ§è¡ ack åå¹¶ */ |
| | | async doAckMessage( |
| | | conversationType: number, |
| | | targetId: number, |
| | | clientMessageId: string, |
| | | updates: Partial<Message> |
| | | ) { |
| | | // 1. å®ä½å¾
åå¹¶æ¶æ¯ |
| | | const conversationStore = useConversationStore() |
| | | const conversation = conversationStore.getConversation(conversationType, targetId) |
| | | if (!conversation) { |
| | | return |
| | | } |
| | | const messages = this.getMessageList(conversationType, targetId) |
| | | const message = messages.find((item) => item.clientMessageId === clientMessageId) |
| | | if (!message) { |
| | | return |
| | | } |
| | | message._ackMerging = true |
| | | try { |
| | | // 2. åå¹¶æå¡ç«¯ ack å°å
å |
| | | applyServerMessageUpdate(message, updates) |
| | | if (messages[messages.length - 1] === message) { |
| | | recomputeConversationLast(conversation, messages) |
| | | } |
| | | // 3. åäºå¡åå
¥æ¶æ¯ãä¼è¯æè¦å游æ |
| | | await getDb() |
| | | .transaction(['messages', 'conversations', 'settings'], 'readwrite', async (tx) => { |
| | | await this.saveMessageRecord(message, conversationType, tx, { |
| | | mergeClientRecord: true |
| | | }) |
| | | await conversationStore.saveConversationRecord(conversation, tx) |
| | | await setMessageMaxId(conversationType, message.id, tx) |
| | | }) |
| | | .catch((error) => { |
| | | console.error('[IM messageStore] ack åå
¥å¤±è´¥', error) |
| | | throw error |
| | | }) |
| | | this.updateMessageCursor(conversationType, message.id) |
| | | } finally { |
| | | // 4. æ¸
çåå¹¶æ è®° |
| | | message._ackMerging = false |
| | | } |
| | | }, |
| | | |
| | | /** å±é¨æ´æ°æ¶æ¯ */ |
| | | patchMessage( |
| | | conversationType: number, |
| | | targetId: number, |
| | | clientMessageId: string, |
| | | patch: Partial<Message> |
| | | ) { |
| | | const message = this.getMessageList(conversationType, targetId).find( |
| | | (item) => item.clientMessageId === clientMessageId |
| | | ) |
| | | if (!message) { |
| | | return |
| | | } |
| | | let changed = false |
| | | for (const key in patch) { |
| | | if ( |
| | | Object.prototype.hasOwnProperty.call(patch, key) && |
| | | (patch as Record<string, unknown>)[key] !== |
| | | (message as unknown as Record<string, unknown>)[key] |
| | | ) { |
| | | changed = true |
| | | break |
| | | } |
| | | } |
| | | if (changed) { |
| | | applyServerMessageUpdate(message, patch) |
| | | } |
| | | }, |
| | | |
| | | /** æ¤åæ¶æ¯ */ |
| | | async recallMessage( |
| | | conversationType: number, |
| | | targetId: number, |
| | | recallSignalContent: string |
| | | ): Promise<void> { |
| | | const conversationStore = useConversationStore() |
| | | const changed = this.applyRecallMessageInMemory( |
| | | conversationType, |
| | | targetId, |
| | | recallSignalContent |
| | | ) |
| | | if (!changed) { |
| | | return |
| | | } |
| | | await getDb() |
| | | .transaction(['messages', 'conversations'], 'readwrite', async (tx) => { |
| | | await this.saveMessageRecord(changed.message, conversationType, tx) |
| | | await conversationStore.saveConversationRecord(changed.conversation, tx) |
| | | }) |
| | | .catch((error) => { |
| | | console.error('[IM messageStore] æ¤åæ¶æ¯åå
¥å¤±è´¥', error) |
| | | throw error |
| | | }) |
| | | }, |
| | | |
| | | /** åºç¨å·²è¯»åæ§ */ |
| | | applyMessageReadReceipt(options: { |
| | | conversationType: number |
| | | groupMessageId?: number |
| | | privateReadMaxId?: number |
| | | readCount?: number |
| | | receiptStatus?: number |
| | | targetId: number |
| | | }) { |
| | | const messages = this.getMessageList(options.conversationType, options.targetId) |
| | | const changed: Message[] = [] |
| | | // 1. ç§èåæ§æ¹éæ´æ°èªå·±åéçæ¶æ¯ |
| | | if (options.conversationType === ImConversationType.PRIVATE && options.privateReadMaxId) { |
| | | this.updatePrivateReadMaxId(options.targetId, options.privateReadMaxId) |
| | | const privateReadMaxId = options.privateReadMaxId |
| | | messages.forEach((message) => { |
| | | if ( |
| | | message.selfSend && |
| | | message.id && |
| | | message.id <= privateReadMaxId && |
| | | message.receiptStatus === ImMessageReceiptStatus.PENDING |
| | | ) { |
| | | message.receiptStatus = ImMessageReceiptStatus.DONE |
| | | changed.push(message) |
| | | } |
| | | }) |
| | | } else if (options.conversationType === ImConversationType.GROUP && options.groupMessageId) { |
| | | // 2. 群èåæ§æ´æ°åæ¡æ¶æ¯ |
| | | const message = messages.find((item) => item.id === options.groupMessageId) |
| | | if (message) { |
| | | if (options.readCount !== undefined) { |
| | | message.readCount = options.readCount |
| | | } |
| | | if (options.receiptStatus !== undefined) { |
| | | message.receiptStatus = options.receiptStatus |
| | | } |
| | | changed.push(message) |
| | | } |
| | | } |
| | | if (changed.length === 0) { |
| | | return |
| | | } |
| | | // 3. åäºå¡åå
¥åæ´æ¶æ¯ |
| | | void getDb() |
| | | .transaction(['messages'], 'readwrite', async (tx) => { |
| | | for (const message of changed) { |
| | | await this.saveMessageRecord(message, options.conversationType, tx) |
| | | } |
| | | }) |
| | | .catch((error) => console.warn('[IM messageStore] åæ§åå
¥å¤±è´¥', error)) |
| | | }, |
| | | |
| | | /** åç½®å岿¶æ¯ */ |
| | | prependMessageList(conversationType: number, targetId: number, earlierMessages: Message[]) { |
| | | if (earlierMessages.length === 0) { |
| | | return |
| | | } |
| | | const messages = this.getMessageList(conversationType, targetId) |
| | | const existingIds = new Set(messages.map((message) => message.id).filter(Boolean)) |
| | | const fresh = earlierMessages |
| | | .map((message) => ensureClientMessageId(message)) |
| | | .filter((message) => message.id && !existingIds.has(message.id)) |
| | | .toSorted((messageA, messageB) => (messageA.id || 0) - (messageB.id || 0)) |
| | | if (fresh.length === 0) { |
| | | return |
| | | } |
| | | const key = getMessageCacheKey(conversationType, targetId) |
| | | this.messagesByConversation[key] = [...fresh, ...messages] |
| | | void getDb() |
| | | .transaction(['messages'], 'readwrite', async (tx) => { |
| | | for (const message of fresh) { |
| | | await this.saveMessageRecord(message, conversationType, tx) |
| | | } |
| | | }) |
| | | .catch((error) => console.warn('[IM messageStore] å岿¶æ¯åå
¥å¤±è´¥', error)) |
| | | }, |
| | | |
| | | /** å é¤åæ¡æ¶æ¯ */ |
| | | removeMessage( |
| | | conversationType: number, |
| | | targetId: number, |
| | | key: { clientMessageId?: string; id?: number; } |
| | | ) { |
| | | // 1. å®ä½ä¼è¯åæ¶æ¯ |
| | | const conversationStore = useConversationStore() |
| | | const conversation = conversationStore.getConversation(conversationType, targetId) |
| | | if (!conversation) { |
| | | return |
| | | } |
| | | const messages = this.getMessageList(conversationType, targetId) |
| | | const index = messages.findIndex((message) => { |
| | | if (key.id && message.id && message.id === key.id) { |
| | | return true |
| | | } |
| | | return !!key.clientMessageId && message.clientMessageId === key.clientMessageId |
| | | }) |
| | | if (index === -1) { |
| | | return |
| | | } |
| | | // 2. ä»å
åç§»é¤æ¶æ¯ |
| | | const [removed] = messages.splice(index, 1) |
| | | if (!removed) { |
| | | return |
| | | } |
| | | revokeBlobUrlsInContent(removed.content) |
| | | if (index === messages.length) { |
| | | recomputeConversationLast(conversation, messages) |
| | | } |
| | | // 3. å 餿¬å°è®°å½å¹¶ä¿åä¼è¯æè¦ |
| | | getDb() |
| | | .delete('messages', getMessageKey(removed, conversationType)) |
| | | .catch((error) => console.warn('[IM messageStore] æ¶æ¯å é¤å¤±è´¥', error)) |
| | | conversationStore.saveConversation(conversation) |
| | | }, |
| | | |
| | | /** å é¤ä¼è¯å
¨é¨æ¶æ¯ */ |
| | | deleteConversationMessageList(conversationType: number, targetId: number) { |
| | | // 1. æ¸
çå
åæ¶æ¯ååªä½èµæº |
| | | const clientConversationId = getClientConversationId(conversationType, targetId) |
| | | const messages = this.messagesByConversation[clientConversationId] || [] |
| | | messages.forEach((message) => { |
| | | revokeBlobUrlsInContent(message.content) |
| | | message._localFile = undefined |
| | | }) |
| | | Reflect.deleteProperty(this.messagesByConversation, clientConversationId) |
| | | this.loadedConversationKeys = this.loadedConversationKeys.filter( |
| | | (key) => key !== clientConversationId |
| | | ) |
| | | // 2. å é¤ IndexedDB æ¶æ¯ |
| | | getDb() |
| | | .deleteByIndex('messages', 'clientConversationId', clientConversationId) |
| | | .catch((error) => console.warn('[IM messageStore] ä¼è¯æ¶æ¯å é¤å¤±è´¥', error)) |
| | | } |
| | | } |
| | | }) |
| | | |
| | | export const useMessageStoreWithOut = () => useMessageStore() |
| | | |
| | | if (import.meta.hot) { |
| | | import.meta.hot.accept(acceptHMRUpdate(useMessageStore, import.meta.hot)) |
| | | } |