gaoluyang
2026-06-29 27cd042df9aca0383a49f3514bc21958dd890912
src/views/im/home/composables/useMessageSender.ts
¶Ô±ÈÐÂÎļþ
@@ -0,0 +1,337 @@
import type { Conversation, Message } from '../types'
import { readChannelMessages as apiReadChannelMessages } from '#/api/im/message/channel'
import {
  readGroupMessages as apiReadGroupMessages,
  recallGroupMessage as apiRecallGroupMessage,
  sendGroupMessage as apiSendGroupMessage
} from '#/api/im/message/group'
import {
  getPrivateMaxReadMessageId as apiGetPrivateMaxReadMessageId,
  readPrivateMessages as apiReadPrivateMessages,
  recallPrivateMessage as apiRecallPrivateMessage,
  sendPrivateMessage as apiSendPrivateMessage
} from '#/api/im/message/private'
import { getCurrentUserId } from '#/views/im/utils/auth'
import { MESSAGE_GROUP_READ_ENABLED, MESSAGE_PRIVATE_READ_ENABLED } from '../../utils/config'
import { ImContentType, ImConversationType, ImMessageStatus } from '../../utils/constants'
import { getClientConversationId } from '../../utils/db'
import {
  generateClientMessageId,
  type QuoteMessage,
  serializeMessage,
  type TextMessage,
  withQuotePayload
} from '../../utils/message'
import { useConversationStore } from '../store/conversationStore'
import { useMessageStore } from '../store/messageStore'
/** éžæ–‡æœ¬æ¶ˆæ¯çš„æ‰©å±•选项(通用) */
interface SendExtOptions {
  atUserIds?: number[] // ç¾¤èŠ @ çš„用户编号列表
  receipt?: boolean // æ˜¯å¦éœ€è¦ç¾¤å›žæ‰§ï¼ˆé»˜è®¤ false)
  targetId?: number // è¦†ç›–默认的 targetId
  /**
   * æ˜¾å¼æŒ‡å®šç›®æ ‡ä¼šè¯ï¼ˆè½¬å‘ / åç‰‡æŽ¨èåœºæ™¯ï¼‰
   *
   * ä¸ä¼ æ—¶é»˜è®¤å– conversationStore.activeConversation;传入时按本值发送 + ä¹è§‚更新到对应会话,
   * ä¸è¦æ±‚该会话当前是激活状态(适合发给「非当前会话」的多个目标)
   */
  conversation?: Conversation
  /** è¢«å¼•用消息(可选):写进 content.quote ç”¨äºŽä¹è§‚渲染,服务端按 quote.messageId åæŸ¥é‡ç®—覆盖 */
  quote?: QuoteMessage
  /**
   * å¤ç”¨å·²å­˜åœ¨çš„æœ¬åœ°å ä½æ¶ˆæ¯ clientMessageId(媒体上传场景)
   *
   * åª’体上传链路在请求服务端前已经 insertMessage äº†å ä½ï¼ˆå¸¦ blob URL + è¿›åº¦æ¡ï¼‰ï¼Œ
   * è¿™é‡Œè·³è¿‡ buildLocalMessage / insertMessage,直接拿这个 id èµ° ackMessage æ”¶å°¾ï¼Œé¿å…é‡å¤æ’入两条
   */
  existingClientMessageId?: string
}
/**
 * æ¶ˆæ¯å‘送 / æ’¤å›ž / å·²è¯» ç»„合式逻辑
 *
 * è®¾è®¡è¦ç‚¹ï¼š
 * 1. ç§èŠ / ç¾¤èŠæŽ¥å£ç­¾åå¯¹ç§°ï¼ŒæŒ‰ conversation.type åˆ†æ”¯è°ƒåº¦ï¼Œå·®å¼‚在分支内部消化
 * 2. å‘送走「乐观更新」:先 insertMessage å†™å…¥ SENDING å ä½ï¼Œè¯·æ±‚成功 ackMessage æ›´æ–°ä¸º NORMAL,失败更新为 FAILED
 * 3. æ’¤å›žä¸åšä¹è§‚更新:服务端通过 WebSocket RECALL äº‹ä»¶å›žä¼ ï¼Œç”± websocketStore ç»Ÿä¸€æ›´æ–°çŠ¶æ€ï¼Œé¿å…ç½‘ç»œå¤±è´¥åŽä¸å¯å›žé€€
 * 4. å·²è¯»ä¸ŠæŠ¥ï¼šæœ¬ç«¯ç«‹åˆ»æ¸…未读数并记录本地读位置;接口失败仅记录日志
 */
export const useMessageSender = () => {
  const conversationStore = useConversationStore()
  const messageStore = useMessageStore()
  /** æž„造本地乐观消息对象 */
  const buildLocalMessage = (opts: {
    atUserIds?: number[]
    clientMessageId: string
    content: string
    targetId: number
    type: number
  }): Message => {
    return {
      clientMessageId: opts.clientMessageId,
      type: opts.type,
      content: opts.content,
      status: ImMessageStatus.SENDING,
      sendTime: Date.now(),
      senderId: getCurrentUserId(),
      targetId: opts.targetId,
      selfSend: true,
      atUserIds: opts.atUserIds
    }
  }
  /**
   * å‘送任意类型的消息(底层实现)
   * 1. æ–‡æœ¬ã€å›¾ç‰‡ã€æ–‡ä»¶ã€è¯­éŸ³ç­‰éƒ½èµ°è¿™é‡Œ
   * 2. type / content ç”±è°ƒç”¨æ–¹æž„造
   * 3. è¿”回值:成功 true / å¤±è´¥ false(失败时本地占位已标 FAILED);参数缺失等无法发送的场景也返 false
   *    è½¬å‘ / åç‰‡æŽ¨èç­‰åœºæ™¯æŒ‰è¿”回值决定是否继续后续动作(如有留言时仅在名片成功后再发留言)
   */
  const sendRaw = async (
    type: number,
    content: string,
    options?: SendExtOptions
  ): Promise<boolean> => {
    // 1. å‚数校验:优先用显式传入的 conversation(转发场景),否则取激活会话
    const conversation = options?.conversation ?? conversationStore.activeConversation
    if (!conversation) {
      return false
    }
    const realTarget = options?.targetId || conversation.targetId
    if (!realTarget) {
      return false
    }
    // 2. å‡†å¤‡ clientMessageId:媒体上传链路在 step 1 å·²ç» insertMessage å ä½ï¼Œè¿™é‡Œç›´æŽ¥å¤ç”¨ id;其余场景走默认乐观插入
    let clientMessageId: string
    if (options?.existingClientMessageId) {
      clientMessageId = options.existingClientMessageId
      // å ä½è‹¥å·²è¢«åˆ é™¤ï¼ˆä¸Šä¼ æœŸé—´ç”¨æˆ·å³é”®åˆ é™¤ / æ’¤å›ž / removeMessage ç­‰ï¼‰åˆ™æ”¾å¼ƒå‘送,
      // å¦åˆ™ sendRaw ä»ä¼šæŠŠæ¶ˆæ¯æŽ¨åˆ°æœåŠ¡ç«¯ï¼Œå¯¼è‡´"本地无气泡 / å¯¹æ–¹å´æ”¶åˆ°ä¸€æ¡"
      const stillExists = messageStore
        .getMessageList(conversation.type, realTarget)
        .some((message) => message.clientMessageId === clientMessageId && !message._ackMerging)
      if (!stillExists) {
        return false
      }
    } else {
      clientMessageId = generateClientMessageId()
      const message = buildLocalMessage({
        clientMessageId,
        content,
        targetId: realTarget,
        type,
        atUserIds: options?.atUserIds
      })
      const conversationInfo = {
        type: conversation.type,
        targetId: realTarget,
        name: conversation.name || String(realTarget),
        avatar: conversation.avatar || ''
      }
      void messageStore.insertMessage(conversationInfo, message).catch(() => undefined)
    }
    // 3. å‘送请求:按会话类型分发到不同接口;成功后 ackMessage æ›´æ–°ä¸º NORMAL,失败更新为 FAILED
    try {
      if (conversation.type === ImConversationType.PRIVATE) {
        const data = await apiSendPrivateMessage({
          clientMessageId,
          receiverId: realTarget,
          type,
          content
        })
        void messageStore
          .ackMessage(conversation.type, realTarget, clientMessageId, {
            id: data.id,
            sendTime: new Date(data.sendTime).getTime(),
            status: data.status,
            receiptStatus: data.receiptStatus,
            content: data.content
          })
          .catch(() => undefined)
      } else if (conversation.type === ImConversationType.GROUP) {
        const data = await apiSendGroupMessage({
          clientMessageId,
          groupId: realTarget,
          type,
          content,
          atUserIds: options?.atUserIds,
          receipt: options?.receipt
        })
        void messageStore
          .ackMessage(conversation.type, realTarget, clientMessageId, {
            id: data.id,
            sendTime: new Date(data.sendTime).getTime(),
            status: data.status,
            receiptStatus: data.receiptStatus,
            readCount: data.readCount,
            content: data.content
          })
          .catch(() => undefined)
      }
      return true
    } catch (error) {
      console.error('[IM] æ¶ˆæ¯å‘送失败', { type, realTarget, clientMessageId }, error)
      void messageStore
        .ackMessage(conversation.type, realTarget, clientMessageId, {
          status: ImMessageStatus.FAILED
        })
        .catch(() => undefined)
      return false
    }
  }
  /**
   * å‘送文本消息(最常用的快捷入口):message-input.vue æ–‡æœ¬å›žè½¦èµ°è¿™é‡Œ
   * è¿”回值:成功 true / å¤±è´¥ false / ç©ºæ–‡æœ¬ false(与 sendRaw å¯¹é½ï¼Œè½¬å‘场景按返回值判断)
   */
  const send = async (text: string, options?: SendExtOptions): Promise<boolean> => {
    if (!text.trim()) {
      return false
    }
    const payload = withQuotePayload<TextMessage>({ content: text }, options?.quote)
    return sendRaw(ImContentType.TEXT, serializeMessage(payload), options)
  }
  /**
   * æ’¤å›žæŸæ¡æ¶ˆæ¯
   * 1. æœåŠ¡ç«¯ä¼šé€šè¿‡ WebSocket RECALL äº‹ä»¶å›žä¼ ï¼Œæœ¬ç«¯ UI ç”± websocketStore ç»Ÿä¸€æ›´æ–°
   * 2. æ­¤å¤„不做乐观撤回,避免网络失败后状态不可回退
   */
  const recall = async (message: Message) => {
    // å‚数校验:本地占位消息不能撤回
    if (!message.id) {
      return
    }
    const conversation = conversationStore.activeConversation
    if (!conversation) {
      return
    }
    // ç§èŠ / ç¾¤èŠæŽ¥å£ç­¾åä¸€è‡´ï¼ŒæŒ‰ä¼šè¯ç±»åž‹åˆ†å‘
    const isPrivate = conversation.type === ImConversationType.PRIVATE
    try {
      await (isPrivate ? apiRecallPrivateMessage(message.id) : apiRecallGroupMessage(message.id))
    } catch (error) {
      console.error('[IM] æ’¤å›žå¤±è´¥', { messageId: message.id, type: conversation.type }, error)
    }
  }
  /**
   * è§¦å‘当前会话的已读上报(切会话 / è¿›å…¥é¡µé¢æ—¶è°ƒç”¨ï¼‰
   * 1. æœ¬ç«¯ç«‹åˆ»æ¸…未读数并推进读位置
   * 2. å·²è¯»ä½ç½®å–已加载消息和会话末条消息的最大服务端 id
   */
  const readActive = async () => {
    const conversation = conversationStore.activeConversation
    if (!conversation) {
      return
    }
    let loadedMaxMessageId = 0
    for (const message of messageStore.getMessages(
      getClientConversationId(conversation.type, conversation.targetId)
    )) {
      if (message.id && message.id > loadedMaxMessageId) {
        loadedMaxMessageId = message.id
      }
    }
    const maxMessageId = Math.max(loadedMaxMessageId, conversation.lastMessageId || 0)
    const readReported = conversationStore.isReportedReadPositionCovered(
      conversation.type,
      conversation.targetId,
      maxMessageId
    )
    if (readReported) {
      conversationStore.markConversationRead(conversation.type, conversation.targetId)
      return
    }
    const isPrivate = conversation.type === ImConversationType.PRIVATE
    const isGroup = conversation.type === ImConversationType.GROUP
    const isChannel = conversation.type === ImConversationType.CHANNEL
    // æœ¬åœ°æ ‡è®°å·²è¯»ï¼šæœªè¯»æ•°æ¸…零(UI ç«‹åˆ»å“åº”)
    conversationStore.markConversationRead(conversation.type, conversation.targetId, maxMessageId)
    if (!maxMessageId) {
      return
    }
    // æŽ¥å£è°ƒç”¨ï¼šæŒ‰ä¼šè¯ç±»åž‹åˆ†å‘,并按对应已读开关控制
    if (!isPrivate && !isGroup && !isChannel) {
      return
    }
    if (isPrivate && !MESSAGE_PRIVATE_READ_ENABLED) {
      return
    }
    if (isGroup && !MESSAGE_GROUP_READ_ENABLED) {
      return
    }
    try {
      if (isPrivate) {
        await apiReadPrivateMessages(conversation.targetId, maxMessageId)
      } else if (isGroup) {
        await apiReadGroupMessages(conversation.targetId, maxMessageId)
      } else {
        await apiReadChannelMessages(conversation.targetId, maxMessageId)
      }
      conversationStore.markConversationReadReported(
        conversation.type,
        conversation.targetId,
        maxMessageId
      )
    } catch (error) {
      console.error(
        '[IM] æ ‡è®°å·²è¯»å¤±è´¥',
        { type: conversation.type, targetId: conversation.targetId, maxMessageId },
        error
      )
    }
  }
  /**
   * æ‹‰å–「对方已读到我哪条消息」并补齐本地状态
   *
   * 1. å¼¥è¡¥ç¦»çº¿ / å¤šç«¯æœŸé—´é”™è¿‡çš„ RECEIPT æŽ¨é€ï¼šè¿›å…¥ç§èŠä¼šè¯æˆ–断线重连后调一次,
   *    æŠŠå¯¹æ–¹ maxReadId åŒæ­¥åˆ°æœ¬åœ°æ¶ˆæ¯ status,避免对方明明读了、本端却仍显示未读
   * 2. ä»…私聊使用:群聊已读位置在每条消息的 readCount / receiptStatus å­—段,离线拉取自带回
   */
  const syncPrivateReadStatus = async (peerId: number) => {
    if (!peerId) {
      return
    }
    // ç§èŠå·²è¯»å…³é—­ï¼šè·³è¿‡å¯¹æ–¹å·²è¯»ä½ç½®åŒæ­¥ï¼Œé¿å…æ— è°“接口调用
    if (!MESSAGE_PRIVATE_READ_ENABLED) {
      return
    }
    const cachedMaxReadId = messageStore.getPrivateReadMaxId(peerId)
    if (cachedMaxReadId !== undefined) {
      if (cachedMaxReadId > 0) {
        messageStore.applyMessageReadReceipt({
          conversationType: ImConversationType.PRIVATE,
          targetId: peerId,
          privateReadMaxId: cachedMaxReadId
        })
      }
      return
    }
    try {
      // æ‹‰å–对方已读到的最大消息 id
      const maxReadId = await apiGetPrivateMaxReadMessageId(peerId)
      messageStore.updatePrivateReadMaxId(peerId, maxReadId)
      if (!maxReadId) {
        return
      }
      // applyMessageReadReceipt å†…部把 â‰¤ maxReadId çš„æœ¬ç«¯æ¶ˆæ¯å›žæ‰§æ›´æ–°ä¸º DONE
      messageStore.applyMessageReadReceipt({
        conversationType: ImConversationType.PRIVATE,
        targetId: peerId,
        privateReadMaxId: maxReadId
      })
    } catch (error) {
      console.warn('[IM] æ‹‰å–对方已读位置失败', { peerId }, error)
    }
  }
  return { send, sendRaw, recall, readActive, syncPrivateReadStatus }
}