diff --git a/apps/web-antd/src/views/im/home/components/friend/friend-add-dialog.vue b/apps/web-antd/src/views/im/home/components/friend/friend-add-dialog.vue index a994440fd..94adfdb37 100644 --- a/apps/web-antd/src/views/im/home/components/friend/friend-add-dialog.vue +++ b/apps/web-antd/src/views/im/home/components/friend/friend-add-dialog.vue @@ -67,6 +67,8 @@ const dialogTitle = computed(() => const presetMode = computed(() => !!presetUser.value); function resetAll() { + loading.value = false; + submitting.value = false; keyword.value = ''; users.value = []; searched.value = false; @@ -102,14 +104,17 @@ function buildPresetApplyContent(): string { /** 按昵称搜索用户:空关键字直接清空结果 */ async function handleSearch() { searched.value = true; - if (!keyword.value.trim()) { + const query = keyword.value.trim(); + if (!query) { users.value = []; + loading.value = false; return; } loading.value = true; try { - users.value = - (await getSimpleUserListByNickname(keyword.value.trim())) || []; + users.value = (await getSimpleUserListByNickname(query)) || []; + } catch (error) { + console.warn('[IM FriendAddDialog] 搜索用户失败', error); } finally { loading.value = false; } @@ -132,25 +137,27 @@ function backToSearch() { /** 提交好友申请:返回 requestId 走「等待验证」;返回 null 表示后端命中「单向好友静默重启」分支,已直接成为好友 */ async function handleSubmitApply() { - if (!targetUser.value) { + const target = targetUser.value; + if (!target) { return; } // 预校验:不能加自己(搜索列表已过滤,这里兜底 presetUser / 名片入口等场景) - if (targetUser.value.id === currentUserId.value) { + if (target.id === currentUserId.value) { message.warning('不能添加自己为好友'); return; } + const payload = { + toUserId: target.id, + applyContent: applyContent.value.trim() || undefined, + displayName: displayName.value.trim() || undefined, + addSource: addSource.value, + }; submitting.value = true; try { - const requestId = await friendStore.applyFriendRequest({ - toUserId: targetUser.value.id, - applyContent: applyContent.value.trim() || undefined, - displayName: displayName.value.trim() || undefined, - addSource: addSource.value, - }); + const requestId = await friendStore.applyFriendRequest(payload); // silent 分支(已是单向好友被静默重启):主动 fetchFriendInfo 入库,不依赖 WS FRIEND_ADD 推送,避免丢推时列表看不到 if (requestId === null) { - await friendStore.fetchFriendInfo(targetUser.value.id); + await friendStore.fetchFriendInfo(target.id); } message.success(requestId ? '申请已发送,等待对方验证' : '已添加为好友'); visible.value = false; diff --git a/apps/web-antd/src/views/im/home/components/group/group-info-card.vue b/apps/web-antd/src/views/im/home/components/group/group-info-card.vue index cb292cad3..35026d59d 100644 --- a/apps/web-antd/src/views/im/home/components/group/group-info-card.vue +++ b/apps/web-antd/src/views/im/home/components/group/group-info-card.vue @@ -66,6 +66,7 @@ function handleChat(group: GroupLite) { /** 加入群聊:先关浮层(避免与 prompt 的 mask 互相遮挡)→ 弹申请理由(可选)→ applyJoinGroup */ async function handleApply(group: GroupLite) { + const groupId = group.id; handleClose(); let applyContent: string; try { @@ -86,12 +87,16 @@ async function handleApply(group: GroupLite) { } catch { return; } - await applyJoinGroup({ - groupId: group.id, - applyContent: applyContent || undefined, - addSource: ImGroupAddSource.SHARE_LINK, - }); - message.success('加群申请已发送'); + try { + await applyJoinGroup({ + groupId, + applyContent: applyContent || undefined, + addSource: ImGroupAddSource.SHARE_LINK, + }); + message.success('加群申请已发送'); + } catch (error) { + console.warn('[IM GroupInfoCard] 申请加群失败', error); + } } diff --git a/apps/web-antd/src/views/im/home/components/group/group-info.vue b/apps/web-antd/src/views/im/home/components/group/group-info.vue index acc52cf41..f245c3267 100644 --- a/apps/web-antd/src/views/im/home/components/group/group-info.vue +++ b/apps/web-antd/src/views/im/home/components/group/group-info.vue @@ -78,7 +78,7 @@ const memberCountText = computed(() => { return count ? `${count} 位成员` : ''; }); -/** member 切群 / 首挂:拉成员;竞态用 group.id 比对丢弃陈旧响应避免上一条群成员错位 */ +/** member 切群 / 首挂:拉取群成员 */ watch( () => [props.group?.id, isMember.value] as const, async ([id, member]) => { @@ -86,13 +86,14 @@ watch( if (!id || !member) { return; } - const list = await groupStore.fetchGroupMemberList(id, true); - if (props.group?.id !== id) { - return; + try { + const list = await groupStore.fetchGroupMemberList(id, true); + members.value = list.map((m) => + convertGroupMemberLite(m, friendStore.getFriend(m.userId)), + ); + } catch (error) { + console.warn('[IM GroupInfo] 群成员加载失败', { groupId: id }, error); } - members.value = list.map((m) => - convertGroupMemberLite(m, friendStore.getFriend(m.userId)), - ); }, { immediate: true }, ); diff --git a/apps/web-antd/src/views/im/home/components/group/group-owner-transfer-dialog.vue b/apps/web-antd/src/views/im/home/components/group/group-owner-transfer-dialog.vue index 624132a51..b18363fe0 100644 --- a/apps/web-antd/src/views/im/home/components/group/group-owner-transfer-dialog.vue +++ b/apps/web-antd/src/views/im/home/components/group/group-owner-transfer-dialog.vue @@ -52,12 +52,15 @@ const newOwner = computed(() => { /** 二次确认转让:转让后旧群主降为普通成员,无法撤销 */ async function handleOk() { - if (!groupId.value || !newOwner.value) { + const targetGroupId = groupId.value; + const newOwnerUserId = newOwner.value?.userId; + const newOwnerName = newOwner.value?.showName; + if (!targetGroupId || !newOwnerUserId) { return; } try { await confirm( - `确定将群主转让给 ${newOwner.value.showName}?转让后你将变为普通成员,无法撤销。`, + `确定将群主转让给 ${newOwnerName}?转让后你将变为普通成员,无法撤销。`, '确认转让群主', ); } catch { @@ -66,12 +69,14 @@ async function handleOk() { submitting.value = true; try { await transferGroupOwner({ - id: groupId.value, - newOwnerUserId: newOwner.value.userId, + id: targetGroupId, + newOwnerUserId, }); message.success('群主转让成功'); emit('reload'); visible.value = false; + } catch (error) { + console.warn('[IM GroupOwnerTransferDialog] 转让群主失败', error); } finally { submitting.value = false; } diff --git a/apps/web-antd/src/views/im/home/components/rtc/rtc-group-call-banner.vue b/apps/web-antd/src/views/im/home/components/rtc/rtc-group-call-banner.vue index 8cad20ed7..5f63264ee 100644 --- a/apps/web-antd/src/views/im/home/components/rtc/rtc-group-call-banner.vue +++ b/apps/web-antd/src/views/im/home/components/rtc/rtc-group-call-banner.vue @@ -7,7 +7,7 @@ import { IconifyIcon as Icon } from '@vben/icons'; import { message, Popover } from 'ant-design-vue'; -import { getActiveCall, joinCall } from '#/api/im/rtc'; +import { getActiveCall, joinCall, leaveCall } from '#/api/im/rtc'; import { getCurrentUserId } from '#/views/im/utils/auth'; import { useGroupCallMembers } from '../../composables/useGroupCallMembers'; @@ -25,6 +25,7 @@ const rtcStore = useRtcStore(); const groupStore = useGroupStore(); const popoverVisible = ref(false); +const joining = ref(false); /** 当前群的活跃通话;rtcStore 维护,参与者加入 / 离开通知增删 joinedUserIds,通话结束移除 */ const activeCall = computed(() => rtcStore.getGroupCall(props.groupId)); @@ -51,10 +52,14 @@ watch( activeCall.value?.room, groupStore.isGroupActiveCallExpired(props.groupId), ] as const, - async ([groupId, room], oldValues) => { + async ([groupId, room], oldValues, onCleanup) => { if (!groupId) { return; } + let obsolete = false; + onCleanup(() => { + obsolete = true; + }); if (!activeCall.value) { if (!groupStore.isGroupActiveCallExpired(groupId)) { @@ -62,6 +67,9 @@ watch( } try { const data = await getActiveCall(groupId); + if (obsolete) { + return; + } if (data) { rtcStore.setGroupCall(data, true); } else { @@ -94,6 +102,9 @@ watch( // 拉最新参与者写回 store;接口返回空 → 该群已无活跃通话,移除本地缓存 try { const data = await getActiveCall(groupId); + if (obsolete) { + return; + } if (data) { rtcStore.setGroupCall(data, true); } else { @@ -124,10 +135,11 @@ const serverSaysJoined = computed(() => { }); /** 加入按钮禁用:仅在本端实际持有 LiveKit 连接时禁用 */ -const joinDisabled = computed(() => isInThisCall.value); +const joinDisabled = computed(() => isInThisCall.value || joining.value); /** 加入按钮文案;本端连着 → 已在通话中;服务端还残留我但本端断了 → 重新加入;其它 → 加入 */ const joinLabel = computed(() => { + if (joining.value) return '加入中...'; if (isInThisCall.value) return '已在通话中'; if (serverSaysJoined.value) return '重新加入'; return '加入'; @@ -136,16 +148,35 @@ const joinLabel = computed(() => { /** 主动加入:调 invite 命中已有 call 拿 token;rtcStore 按 status 自动进 RUNNING */ async function handleJoin() { const call = activeCall.value; - if (!call || joinDisabled.value) { + if (!call || joinDisabled.value || joining.value) { return; } if (rtcStore.isActive) { message.warning('您正在通话中'); return; } + // 加入结果会获取物理房间凭证,必须阻断旧用户或旧群结果发布 popoverVisible.value = false; - const data = await joinCall(call.room); - rtcStore.startInviting(data); + joining.value = true; + const userId = getCurrentUserId(); + try { + const data = await joinCall(call.room); + if ( + getCurrentUserId() !== userId || + activeCall.value?.room !== call.room || + (rtcStore.isActive && rtcStore.call?.room !== data.room) + ) { + if (getCurrentUserId() === userId) { + await leaveCall(data.room || call.room).catch(() => undefined); + } + return; + } + rtcStore.startInviting(data); + } catch (error) { + console.warn('[IM RtcGroupCallBanner] 加入群通话失败', error); + } finally { + joining.value = false; + } } diff --git a/apps/web-antd/src/views/im/home/components/user/user-info.vue b/apps/web-antd/src/views/im/home/components/user/user-info.vue index 0132e872a..e9c4e7b73 100644 --- a/apps/web-antd/src/views/im/home/components/user/user-info.vue +++ b/apps/web-antd/src/views/im/home/components/user/user-info.vue @@ -105,7 +105,6 @@ const remarkInputRef = ref void; select?: () => void }>( * user.id 变化的统一处理: * 1. 起手用 prop 兜底首屏(full = props.user),再 getSimpleUser 命中后合并替换 * 2. 顺便复位备注编辑态,避免上一个用户的脏输入泄漏到下一个 - * 3. 竞态用 id 比对丢弃陈旧响应 */ watch( () => props.user?.id, @@ -116,9 +115,6 @@ watch( return; } const data = (await getSimpleUser(id)) as User; - if (props.user?.id !== id) { - return; - } full.value = { ...props.user, ...data }; }, { immediate: true }, @@ -235,10 +231,11 @@ async function handleBlock() { /** 移出黑名单:操作温和不弹 confirm;后端 FRIEND_UNBLOCK 推到时由 dispatcher 同步多端 */ async function handleUnblock() { - if (!props.user?.id) { + const targetId = props.user?.id; + if (!targetId) { return; } - await friendStore.unblockFriend(props.user.id); + await friendStore.unblockFriend(targetId); message.success('已移出黑名单'); } @@ -272,7 +269,12 @@ async function handleDeleteFriend() { } catch { return; } - await friendStore.deleteFriend(target.id, clearConversation.value); + try { + await friendStore.deleteFriend(target.id, clearConversation.value); + } catch (error) { + console.warn('[IM UserInfo] 删除好友失败', error); + return; + } message.success('已删除好友'); emit('deleted', target); } diff --git a/apps/web-antd/src/views/im/home/pages/contact/friend-request-detail.vue b/apps/web-antd/src/views/im/home/pages/contact/friend-request-detail.vue index 53c2f9282..2a9464ec0 100644 --- a/apps/web-antd/src/views/im/home/pages/contact/friend-request-detail.vue +++ b/apps/web-antd/src/views/im/home/pages/contact/friend-request-detail.vue @@ -72,6 +72,7 @@ const processing = ref(false); /** 同意申请:互斥锁 + 状态二次校验,避免并发 / 服务端已处理后再次提交 */ async function handleAgree() { + const requestId = props.request.id; if (processing.value) { return; } @@ -81,8 +82,10 @@ async function handleAgree() { processing.value = true; agreeing.value = true; try { - await friendStore.agreeFriendRequest(props.request.id); + await friendStore.agreeFriendRequest(requestId); message.success('已同意好友申请'); + } catch (error) { + console.warn('[IM FriendRequestDetail] 同意好友申请失败', error); } finally { agreeing.value = false; processing.value = false; @@ -91,6 +94,7 @@ async function handleAgree() { /** 拒绝申请:弹 prompt 收集可选拒绝理由(点取消则中止),随后调 store 落库 + 提示 */ async function handleRefuse() { + const requestId = props.request.id; if (processing.value) { return; } @@ -138,8 +142,10 @@ async function handleRefuse() { processing.value = true; refusing.value = true; try { - await friendStore.refuseFriendRequest(props.request.id, handleContent); + await friendStore.refuseFriendRequest(requestId, handleContent); message.success('已拒绝好友申请'); + } catch (error) { + console.warn('[IM FriendRequestDetail] 拒绝好友申请失败', error); } finally { refusing.value = false; processing.value = false; diff --git a/apps/web-antd/src/views/im/home/pages/contact/index.vue b/apps/web-antd/src/views/im/home/pages/contact/index.vue index 5deb944bb..4b238c84c 100644 --- a/apps/web-antd/src/views/im/home/pages/contact/index.vue +++ b/apps/web-antd/src/views/im/home/pages/contact/index.vue @@ -220,18 +220,26 @@ function handleChatGroup(group: GroupLite) { /** 删除好友:二次确认 → store 落库 → 清空当前选中 */ async function handleDeleteFriend(friend: FriendLite) { + const friendId = friend.id; try { await confirm(`确定删除好友「${friend.nickname}」吗?`, '删除联系人'); + } catch { + return; + } + try { // friendStore.deleteFriend 内部已经级联清理对应私聊会话 - await friendStore.deleteFriend(friend.id); - if ( - selection.value?.type === 'friend' && - selection.value.friend.id === friend.id - ) { - selection.value = null; - } - message.success('已删除好友'); - } catch {} + await friendStore.deleteFriend(friendId); + } catch (error) { + console.warn('[IM contact] 删除好友失败', error); + return; + } + if ( + selection.value?.type === 'friend' && + selection.value.friend.id === friendId + ) { + selection.value = null; + } + message.success('已删除好友'); } /** 备注已保存:UserInfo 内部已经走完 friendStore 落库 + 提示,本侧只负责同步 selection 持的旧 FriendLite 副本 */ diff --git a/apps/web-antd/src/views/im/home/pages/conversation/components/conversation/conversation-item.vue b/apps/web-antd/src/views/im/home/pages/conversation/components/conversation/conversation-item.vue index 8193fd428..b91eb03ff 100644 --- a/apps/web-antd/src/views/im/home/pages/conversation/components/conversation/conversation-item.vue +++ b/apps/web-antd/src/views/im/home/pages/conversation/components/conversation/conversation-item.vue @@ -171,24 +171,28 @@ function handleMuted() { type === ImConversationType.PRIVATE ? friendStore.setFriendSilent(targetId, next) : groupStore.setGroupSilent(targetId, next); - sync.catch((error) => { + void sync.catch((error) => { console.error('[IM] 切换免打扰失败', error); - conversationStore.setConversationSilent(type, targetId, !next); + const conversation = conversationStore.getConversation(type, targetId); + if (conversation?.silent === next) { + conversationStore.setConversationSilent(type, targetId, !next); + } }); } /** 删除会话:二次确认后软删 */ async function handleDelete() { + const { type, targetId, name } = props.conversation; try { - await confirm( - `确定删除与「${props.conversation.name}」的会话吗?`, - '删除会话', - ); - conversationStore.removeConversation( - props.conversation.type, - props.conversation.targetId, - ); - } catch {} + await confirm(`确定删除与「${name}」的会话吗?`, '删除会话'); + } catch { + return; + } + try { + await conversationStore.removeConversation(type, targetId); + } catch (error) { + console.warn('[IM ConversationItem] 删除会话失败', error); + } } /** 右键菜单:置顶 / 免打扰 / 删除 */ diff --git a/apps/web-antd/src/views/im/home/pages/conversation/components/input/face-picker.vue b/apps/web-antd/src/views/im/home/pages/conversation/components/input/face-picker.vue index a070f3c6d..3a96d3742 100644 --- a/apps/web-antd/src/views/im/home/pages/conversation/components/input/face-picker.vue +++ b/apps/web-antd/src/views/im/home/pages/conversation/components/input/face-picker.vue @@ -124,6 +124,9 @@ async function onUploadPicked(e: Event) { } const payload = { url, width: size.width, height: size.height }; await faceStore.addFaceUserItem(payload); + } catch (error) { + console.warn('[IM FacePicker] 上传个人表情失败', error); + message.error('上传失败'); } finally { uploading.value = false; } diff --git a/apps/web-antd/src/views/im/home/pages/conversation/components/input/message-multi-select-bar.vue b/apps/web-antd/src/views/im/home/pages/conversation/components/input/message-multi-select-bar.vue index d45237b35..6b794e03c 100644 --- a/apps/web-antd/src/views/im/home/pages/conversation/components/input/message-multi-select-bar.vue +++ b/apps/web-antd/src/views/im/home/pages/conversation/components/input/message-multi-select-bar.vue @@ -6,6 +6,8 @@ import { computed, inject } from 'vue'; import { confirm } from '@vben/common-ui'; import { IconifyIcon as Icon } from '@vben/icons'; +import { message } from 'ant-design-vue'; + import { useMessageMultiSelect } from '#/views/im/home/composables/useMessageMultiSelect'; import { useConversationStore } from '#/views/im/home/store/conversationStore'; import { useMessageStore } from '#/views/im/home/store/messageStore'; @@ -83,6 +85,7 @@ async function handleDelete() { if (!conversation) { return; } + const { type, targetId } = conversation; const messages = getSelectedMessages(); if (messages.length === 0) { return; @@ -97,13 +100,21 @@ async function handleDelete() { } catch { return; } - for (const m of messages) { - messageStore.removeMessage(conversation.type, conversation.targetId, { - id: m.id, - clientMessageId: m.clientMessageId, - }); + try { + await Promise.all( + messages.map((item) => + messageStore.removeMessage(type, targetId, { + id: item.id, + clientMessageId: item.clientMessageId, + }), + ), + ); + } catch (error) { + console.warn('[IM MessageMultiSelectBar] 批量删除消息失败', error); + message.error('删除失败,请重试'); + } finally { + multiSelect.exit(); } - multiSelect.exit(); } /** 取消多选 */ diff --git a/apps/web-antd/src/views/im/home/pages/conversation/components/message/message-panel.vue b/apps/web-antd/src/views/im/home/pages/conversation/components/message/message-panel.vue index 02dfe1476..caeb57daf 100644 --- a/apps/web-antd/src/views/im/home/pages/conversation/components/message/message-panel.vue +++ b/apps/web-antd/src/views/im/home/pages/conversation/components/message/message-panel.vue @@ -8,7 +8,7 @@ import { IconifyIcon as Icon } from '@vben/icons'; import { message, Popover, Tooltip } from 'ant-design-vue'; -import { createCall } from '#/api/im/rtc'; +import { createCall, leaveCall } from '#/api/im/rtc'; import { getCurrentUserId } from '#/views/im/utils/auth'; import { ImConversationType, @@ -275,8 +275,12 @@ function reloadGroupData() { if (!conversation || conversation.type !== ImConversationType.GROUP) { return; } - groupStore.fetchGroupInfo(conversation.targetId, true); - groupStore.fetchGroupMemberList(conversation.targetId, true); + void groupStore.fetchGroupInfo(conversation.targetId, true); + void groupStore + .fetchGroupMemberList(conversation.targetId, true) + .catch((error) => { + console.warn('[IM MessagePanel] 强制刷新群成员失败', error); + }); } const historyDialogRef = ref>(); // 历史消息抽屉 ref:「聊天历史」icon / 抽屉「查找聊天内容」入口都调 open() 触发 @@ -359,8 +363,18 @@ async function doInvite(reqVO: { return; } callInviting.value = true; + const userId = getCurrentUserId(); try { const data = await createCall(reqVO); + if ( + getCurrentUserId() !== userId || + (rtcStore.isActive && rtcStore.call?.room !== data.room) + ) { + if (getCurrentUserId() === userId) { + await leaveCall(data.room).catch(() => undefined); + } + return; + } // 后端已 INSERT + 立即 end(如忙线):toast 提示,不进 INVITING 阶段;chat tip 由 RTC_CALL_END 推送写入消息流 if (data.status === ImRtcCallStatus.ENDED) { message.warning(resolveCallEndReasonText(data.endReason)); @@ -368,6 +382,8 @@ async function doInvite(reqVO: { } // 正常进入 INVITING 阶段:走 store 逻辑发起通话,后续状态更新 / 消息流更新由 RTC 模块监听推送处理 rtcStore.startInviting(data); + } catch (error) { + console.warn('[IM MessagePanel] 发起通话失败', error); } finally { callInviting.value = false; } @@ -497,16 +513,6 @@ async function handleLocateMention() { conversation.type, conversation.targetId, ); - const isActive = () => { - const activeConversation = conversationStore.activeConversation; - return ( - !!activeConversation && - getClientConversationId( - activeConversation.type, - activeConversation.targetId, - ) === clientConversationId - ); - }; for (let guard = 0; guard < 50; guard++) { const loadedMessages = messageStore.getMessages(clientConversationId); if (loadedMessages.some((item) => item.id === messageId)) { @@ -516,9 +522,6 @@ async function handleLocateMention() { clientConversationId, 50, ); - if (!isActive()) { - return; - } if ( messageStore .getMessages(clientConversationId) @@ -530,10 +533,7 @@ async function handleLocateMention() { break; } } - if (!isActive()) { - return; - } - await handleLocate(messageId, isActive); + await handleLocate(messageId); } /** @@ -545,14 +545,11 @@ async function handleLocateMention() { * 4. 加 --highlight class 短暂高亮,提示用户"就是这条" * 5. 找不到 wrapper(原消息已分页出去)时弹 warning 提示,与微信"消息已不在窗口"观感一致 */ -async function handleLocate(messageId: number, isActive?: () => boolean) { +async function handleLocate(messageId: number) { if (!messageId) { return; } await nextTick(); - if (isActive && !isActive()) { - return; - } if (!listRef.value) { return; } diff --git a/apps/web-antd/src/views/im/home/pages/conversation/components/message/message-read-status.vue b/apps/web-antd/src/views/im/home/pages/conversation/components/message/message-read-status.vue index 842ba1ce9..62e9b9ce4 100644 --- a/apps/web-antd/src/views/im/home/pages/conversation/components/message/message-read-status.vue +++ b/apps/web-antd/src/views/im/home/pages/conversation/components/message/message-read-status.vue @@ -109,7 +109,7 @@ async function loadReadUsers() { // 全可见成员都已读 → 更新为 DONE,让外面 label 直接命中「全部已读」分支; // 否则只更新 readCount,receiptStatus 维持不变(PENDING) const allRead = readCount > 0 && readCount >= visibleMembers.value.length; - messageStore.applyMessageReadReceipt({ + await messageStore.applyMessageReadReceipt({ conversationType: ImConversationType.GROUP, targetId: props.groupId, groupMessageId: props.message.id, diff --git a/apps/web-antd/src/views/im/utils/messageSync.ts b/apps/web-antd/src/views/im/utils/messageSync.ts new file mode 100644 index 000000000..03f3780d6 --- /dev/null +++ b/apps/web-antd/src/views/im/utils/messageSync.ts @@ -0,0 +1,159 @@ +/** 消息状态优先级;高优先级终态不可被普通消息覆盖 */ +export enum MessageTerminalPriority { + NORMAL = 0, // 普通消息 + CONFIRMED = 1, // 服务端已确认消息 + RECALL = 2 // 撤回终态 +} + +/** 会话写 lane 与全量屏障 */ +interface ConversationWriteState { + barrierTail: Promise // 当前全量屏障尾部 + tails: Map> // 各会话写入尾部 +} + +interface RelationState { + terminated: boolean + messageId: number + localTerminationPending: boolean +} + +const writeState: ConversationWriteState = { + // 当前运行时的会话写状态 + barrierTail: Promise.resolve(), + tails: new Map() +} +const relationStates = new Map() // 群关系消息终态 + +/** 同一会话串行执行消息与会话终态写入 */ +export async function enqueueConversationWrite( + clientConversationId: string, + operation: () => Promise +): Promise { + return enqueueConversationWrites([clientConversationId], operation) +} + +/** 一次写入原子占用全部会话 lane,避免嵌套获取与屏障互锁 */ +export function enqueueConversationWrites( + clientConversationIds: string[], + operation: () => Promise +): Promise { + // 1. 等待全量屏障和所有目标会话的前驱写入 + const keys = Array.from(new Set(clientConversationIds)).sort() + const predecessors = [ + writeState.barrierTail, + ...keys.map((key) => writeState.tails.get(key) || Promise.resolve()) + ] + const current = Promise.all(predecessors.map((task) => task.catch(() => undefined))).then( + operation + ) + const settled = current.then( + () => undefined, + () => undefined + ) + // 2. 先发布 recovery tail;完成时仅清理仍指向本任务的 lane + keys.forEach((key) => writeState.tails.set(key, settled)) + return current.finally(() => { + keys.forEach((key) => { + if (writeState.tails.get(key) === settled) { + writeState.tails.delete(key) + } + }) + }) +} + +/** 独占全部会话写入;只用于全量快照重建,常规写仍按会话并行 */ +export function enqueueConversationBarrier(operation: () => Promise): Promise { + // 1. 同步发布 gate,阻止后续会话写越过本次全量操作 + const previousBarrier = writeState.barrierTail + const existingWrites = Array.from(writeState.tails.values()) + let release!: () => void + const gate = new Promise((resolve) => { + release = resolve + }) + writeState.barrierTail = previousBarrier.catch(() => undefined).then(() => gate) + return (async () => { + // 2. 排空封门前的屏障和会话写,再独占执行全量操作 + await previousBarrier.catch(() => undefined) + await Promise.all(existingWrites.map((task) => task.catch(() => undefined))) + return await operation() + })().finally(release) +} + +/** 终态优先;相同优先级使用后到达状态 */ +export function reduceMessageState( + current: { priority: MessageTerminalPriority; value?: T } | undefined, + incoming: { priority: MessageTerminalPriority; value?: T } +) { + return current && current.priority > incoming.priority ? current : incoming +} + +/** 在会话写 lane 内记录关系终态;本地主动操作等待服务端终态消息后才允许重开 */ +export function markRelationTerminated(clientConversationId: string, messageId?: number): boolean { + return applyRelationState(clientConversationId, true, messageId) +} + +/** 显式重新加入后清除关系终态;旧通知不得重开新终态 */ +export function reopenRelation(clientConversationId: string, messageId?: number): boolean { + return applyRelationState(clientConversationId, false, messageId) +} + +export function isRelationTerminated(clientConversationId: string): boolean { + return relationStates.get(clientConversationId)?.terminated === true +} + +/** 排空当前 IM 运行时的消息写入,并在调用方仍允许时清理关系终态 */ +export async function clearMessageSyncState(shouldClear: () => boolean): Promise { + const barrier = writeState.barrierTail + const tails = Array.from(writeState.tails.entries()) + await Promise.all([ + barrier.catch(() => undefined), + ...tails.map(([, task]) => task.catch(() => undefined)) + ]) + if (!shouldClear()) { + return + } + relationStates.clear() + if (writeState.barrierTail === barrier) { + writeState.barrierTail = Promise.resolve() + } + tails.forEach(([key, task]) => { + if (writeState.tails.get(key) === task) { + writeState.tails.delete(key) + } + }) +} + +/** 按服务端关系消息编号单调归约群关系;本地主动终止在服务端终态确认前阻止旧成员消息重开 */ +function applyRelationState( + clientConversationId: string, + terminated: boolean, + messageId?: number +): boolean { + const current = relationStates.get(clientConversationId) + if (messageId === undefined) { + if (!terminated && current?.localTerminationPending) { + return false + } + if (terminated && current?.terminated) { + return true + } + relationStates.set(clientConversationId, { + terminated, + messageId: current?.messageId ?? 0, + localTerminationPending: terminated + }) + return true + } + if (!terminated && current?.localTerminationPending) { + return false + } + if (current && messageId <= current.messageId) { + return false + } + relationStates.set(clientConversationId, { + terminated, + messageId, + localTerminationPending: false + }) + return true +} diff --git a/apps/web-antd/src/views/im/utils/resourceRequest.ts b/apps/web-antd/src/views/im/utils/resourceRequest.ts new file mode 100644 index 000000000..2e5382a2a --- /dev/null +++ b/apps/web-antd/src/views/im/utils/resourceRequest.ts @@ -0,0 +1,108 @@ +/** 可合并请求的固定资源 */ +export enum ResourceRequestKey { + FACE_PACKS = 'facePacks', // 系统表情包 + FACE_USER_ITEMS = 'faceUserItems', // 用户表情 + FRIEND_LIST = 'friendList', // 好友列表 + GROUP_LIST = 'groupList', // 群列表 + CHANNEL_LIST = 'channelList', // 频道列表 + GROUP_REQUEST_UNHANDLED = 'groupRequestUnhandled' // 未处理加群申请 +} + +/** 固定资源请求的 task 生命周期模式 */ +export enum ResourceRequestMode { + CACHE_SUCCESS = 'cache-success', // 成功后持续复用,清理运行时状态时失效 + SINGLE_FLIGHT = 'single-flight' // 仅合并当前在途请求 +} + +/** 固定资源请求策略 */ +type ResourceRequestPolicy = + | { + mode: ResourceRequestMode.CACHE_SUCCESS + } + | { + mode: ResourceRequestMode.SINGLE_FLIGHT + refreshAfterPending?: boolean + } + +/** 单个固定资源当前发布的请求状态 */ +interface ResourceRequestEntry { + mode: ResourceRequestMode // task 生命周期模式 + task: Promise // 当前请求 task + trailingExecute?: () => Promise // 合并后的尾随刷新 +} + +const resourceRequests = new Map() // 每个 key 仅发布一个当前 entry + +/** 运行固定资源请求 */ +export function runResourceRequest( + key: ResourceRequestKey, + execute: () => Promise, + policy: ResourceRequestPolicy +): Promise { + const existing = resourceRequests.get(key) + // 1. 复用 task;force 只覆盖为一个最新尾随执行器 + if (existing) { + if (existing.mode !== policy.mode) { + return Promise.reject(new Error(`IM resource policy mismatch: ${key}`)) + } + if ( + existing.mode === ResourceRequestMode.SINGLE_FLIGHT && + policy.mode === ResourceRequestMode.SINGLE_FLIGHT && + policy.refreshAfterPending + ) { + existing.trailingExecute = execute + } + return existing.task as Promise + } + const task = Promise.resolve().then(execute) + const entry: ResourceRequestEntry = { mode: policy.mode, task } + resourceRequests.set(key, entry) + void task.then( + () => finishResourceRequest(key, entry, true), + () => finishResourceRequest(key, entry, false) + ) + return task +} + +/** 完成请求并按策略释放或补刷 */ +function finishResourceRequest( + key: ResourceRequestKey, + entry: ResourceRequestEntry, + succeeded: boolean +): void { + // 1. 旧 finalizer 不能修改已经替换的新 entry + if (resourceRequests.get(key) !== entry) { + return + } + // 2. once 成功保留;其余情况先释放当前 entry + if (entry.mode === ResourceRequestMode.CACHE_SUCCESS && succeeded) { + return + } + resourceRequests.delete(key) + // 3. single-flight 的多次 force 合并为一次后台尾随刷新 + if (entry.trailingExecute) { + void runResourceRequest(key, entry.trailingExecute, { + mode: ResourceRequestMode.SINGLE_FLIGHT + }).catch((error) => console.warn(`[IM] 尾随刷新 ${key} 失败`, error)) + } +} + +/** 排空并清理固定资源请求状态 */ +export async function clearResourceRequests(): Promise { + const entries = Array.from(resourceRequests.entries()) + entries.forEach(([, entry]) => { + entry.trailingExecute = undefined + }) + await Promise.all(entries.map(([, entry]) => entry.task.catch(() => undefined))) + entries.forEach(([key, entry]) => { + if (resourceRequests.get(key) === entry) { + resourceRequests.delete(key) + } + }) +} + +/** 判断固定资源当前是否有请求在途 */ +export function isResourceRequestPending(key: ResourceRequestKey): boolean { + const entry = resourceRequests.get(key) + return entry?.mode === ResourceRequestMode.SINGLE_FLIGHT +} diff --git a/apps/web-antdv-next/src/views/im/utils/messageSync.ts b/apps/web-antdv-next/src/views/im/utils/messageSync.ts new file mode 100644 index 000000000..03f3780d6 --- /dev/null +++ b/apps/web-antdv-next/src/views/im/utils/messageSync.ts @@ -0,0 +1,159 @@ +/** 消息状态优先级;高优先级终态不可被普通消息覆盖 */ +export enum MessageTerminalPriority { + NORMAL = 0, // 普通消息 + CONFIRMED = 1, // 服务端已确认消息 + RECALL = 2 // 撤回终态 +} + +/** 会话写 lane 与全量屏障 */ +interface ConversationWriteState { + barrierTail: Promise // 当前全量屏障尾部 + tails: Map> // 各会话写入尾部 +} + +interface RelationState { + terminated: boolean + messageId: number + localTerminationPending: boolean +} + +const writeState: ConversationWriteState = { + // 当前运行时的会话写状态 + barrierTail: Promise.resolve(), + tails: new Map() +} +const relationStates = new Map() // 群关系消息终态 + +/** 同一会话串行执行消息与会话终态写入 */ +export async function enqueueConversationWrite( + clientConversationId: string, + operation: () => Promise +): Promise { + return enqueueConversationWrites([clientConversationId], operation) +} + +/** 一次写入原子占用全部会话 lane,避免嵌套获取与屏障互锁 */ +export function enqueueConversationWrites( + clientConversationIds: string[], + operation: () => Promise +): Promise { + // 1. 等待全量屏障和所有目标会话的前驱写入 + const keys = Array.from(new Set(clientConversationIds)).sort() + const predecessors = [ + writeState.barrierTail, + ...keys.map((key) => writeState.tails.get(key) || Promise.resolve()) + ] + const current = Promise.all(predecessors.map((task) => task.catch(() => undefined))).then( + operation + ) + const settled = current.then( + () => undefined, + () => undefined + ) + // 2. 先发布 recovery tail;完成时仅清理仍指向本任务的 lane + keys.forEach((key) => writeState.tails.set(key, settled)) + return current.finally(() => { + keys.forEach((key) => { + if (writeState.tails.get(key) === settled) { + writeState.tails.delete(key) + } + }) + }) +} + +/** 独占全部会话写入;只用于全量快照重建,常规写仍按会话并行 */ +export function enqueueConversationBarrier(operation: () => Promise): Promise { + // 1. 同步发布 gate,阻止后续会话写越过本次全量操作 + const previousBarrier = writeState.barrierTail + const existingWrites = Array.from(writeState.tails.values()) + let release!: () => void + const gate = new Promise((resolve) => { + release = resolve + }) + writeState.barrierTail = previousBarrier.catch(() => undefined).then(() => gate) + return (async () => { + // 2. 排空封门前的屏障和会话写,再独占执行全量操作 + await previousBarrier.catch(() => undefined) + await Promise.all(existingWrites.map((task) => task.catch(() => undefined))) + return await operation() + })().finally(release) +} + +/** 终态优先;相同优先级使用后到达状态 */ +export function reduceMessageState( + current: { priority: MessageTerminalPriority; value?: T } | undefined, + incoming: { priority: MessageTerminalPriority; value?: T } +) { + return current && current.priority > incoming.priority ? current : incoming +} + +/** 在会话写 lane 内记录关系终态;本地主动操作等待服务端终态消息后才允许重开 */ +export function markRelationTerminated(clientConversationId: string, messageId?: number): boolean { + return applyRelationState(clientConversationId, true, messageId) +} + +/** 显式重新加入后清除关系终态;旧通知不得重开新终态 */ +export function reopenRelation(clientConversationId: string, messageId?: number): boolean { + return applyRelationState(clientConversationId, false, messageId) +} + +export function isRelationTerminated(clientConversationId: string): boolean { + return relationStates.get(clientConversationId)?.terminated === true +} + +/** 排空当前 IM 运行时的消息写入,并在调用方仍允许时清理关系终态 */ +export async function clearMessageSyncState(shouldClear: () => boolean): Promise { + const barrier = writeState.barrierTail + const tails = Array.from(writeState.tails.entries()) + await Promise.all([ + barrier.catch(() => undefined), + ...tails.map(([, task]) => task.catch(() => undefined)) + ]) + if (!shouldClear()) { + return + } + relationStates.clear() + if (writeState.barrierTail === barrier) { + writeState.barrierTail = Promise.resolve() + } + tails.forEach(([key, task]) => { + if (writeState.tails.get(key) === task) { + writeState.tails.delete(key) + } + }) +} + +/** 按服务端关系消息编号单调归约群关系;本地主动终止在服务端终态确认前阻止旧成员消息重开 */ +function applyRelationState( + clientConversationId: string, + terminated: boolean, + messageId?: number +): boolean { + const current = relationStates.get(clientConversationId) + if (messageId === undefined) { + if (!terminated && current?.localTerminationPending) { + return false + } + if (terminated && current?.terminated) { + return true + } + relationStates.set(clientConversationId, { + terminated, + messageId: current?.messageId ?? 0, + localTerminationPending: terminated + }) + return true + } + if (!terminated && current?.localTerminationPending) { + return false + } + if (current && messageId <= current.messageId) { + return false + } + relationStates.set(clientConversationId, { + terminated, + messageId, + localTerminationPending: false + }) + return true +} diff --git a/apps/web-antdv-next/src/views/im/utils/resourceRequest.ts b/apps/web-antdv-next/src/views/im/utils/resourceRequest.ts new file mode 100644 index 000000000..2e5382a2a --- /dev/null +++ b/apps/web-antdv-next/src/views/im/utils/resourceRequest.ts @@ -0,0 +1,108 @@ +/** 可合并请求的固定资源 */ +export enum ResourceRequestKey { + FACE_PACKS = 'facePacks', // 系统表情包 + FACE_USER_ITEMS = 'faceUserItems', // 用户表情 + FRIEND_LIST = 'friendList', // 好友列表 + GROUP_LIST = 'groupList', // 群列表 + CHANNEL_LIST = 'channelList', // 频道列表 + GROUP_REQUEST_UNHANDLED = 'groupRequestUnhandled' // 未处理加群申请 +} + +/** 固定资源请求的 task 生命周期模式 */ +export enum ResourceRequestMode { + CACHE_SUCCESS = 'cache-success', // 成功后持续复用,清理运行时状态时失效 + SINGLE_FLIGHT = 'single-flight' // 仅合并当前在途请求 +} + +/** 固定资源请求策略 */ +type ResourceRequestPolicy = + | { + mode: ResourceRequestMode.CACHE_SUCCESS + } + | { + mode: ResourceRequestMode.SINGLE_FLIGHT + refreshAfterPending?: boolean + } + +/** 单个固定资源当前发布的请求状态 */ +interface ResourceRequestEntry { + mode: ResourceRequestMode // task 生命周期模式 + task: Promise // 当前请求 task + trailingExecute?: () => Promise // 合并后的尾随刷新 +} + +const resourceRequests = new Map() // 每个 key 仅发布一个当前 entry + +/** 运行固定资源请求 */ +export function runResourceRequest( + key: ResourceRequestKey, + execute: () => Promise, + policy: ResourceRequestPolicy +): Promise { + const existing = resourceRequests.get(key) + // 1. 复用 task;force 只覆盖为一个最新尾随执行器 + if (existing) { + if (existing.mode !== policy.mode) { + return Promise.reject(new Error(`IM resource policy mismatch: ${key}`)) + } + if ( + existing.mode === ResourceRequestMode.SINGLE_FLIGHT && + policy.mode === ResourceRequestMode.SINGLE_FLIGHT && + policy.refreshAfterPending + ) { + existing.trailingExecute = execute + } + return existing.task as Promise + } + const task = Promise.resolve().then(execute) + const entry: ResourceRequestEntry = { mode: policy.mode, task } + resourceRequests.set(key, entry) + void task.then( + () => finishResourceRequest(key, entry, true), + () => finishResourceRequest(key, entry, false) + ) + return task +} + +/** 完成请求并按策略释放或补刷 */ +function finishResourceRequest( + key: ResourceRequestKey, + entry: ResourceRequestEntry, + succeeded: boolean +): void { + // 1. 旧 finalizer 不能修改已经替换的新 entry + if (resourceRequests.get(key) !== entry) { + return + } + // 2. once 成功保留;其余情况先释放当前 entry + if (entry.mode === ResourceRequestMode.CACHE_SUCCESS && succeeded) { + return + } + resourceRequests.delete(key) + // 3. single-flight 的多次 force 合并为一次后台尾随刷新 + if (entry.trailingExecute) { + void runResourceRequest(key, entry.trailingExecute, { + mode: ResourceRequestMode.SINGLE_FLIGHT + }).catch((error) => console.warn(`[IM] 尾随刷新 ${key} 失败`, error)) + } +} + +/** 排空并清理固定资源请求状态 */ +export async function clearResourceRequests(): Promise { + const entries = Array.from(resourceRequests.entries()) + entries.forEach(([, entry]) => { + entry.trailingExecute = undefined + }) + await Promise.all(entries.map(([, entry]) => entry.task.catch(() => undefined))) + entries.forEach(([key, entry]) => { + if (resourceRequests.get(key) === entry) { + resourceRequests.delete(key) + } + }) +} + +/** 判断固定资源当前是否有请求在途 */ +export function isResourceRequestPending(key: ResourceRequestKey): boolean { + const entry = resourceRequests.get(key) + return entry?.mode === ResourceRequestMode.SINGLE_FLIGHT +} diff --git a/apps/web-ele/src/views/im/utils/messageSync.ts b/apps/web-ele/src/views/im/utils/messageSync.ts new file mode 100644 index 000000000..03f3780d6 --- /dev/null +++ b/apps/web-ele/src/views/im/utils/messageSync.ts @@ -0,0 +1,159 @@ +/** 消息状态优先级;高优先级终态不可被普通消息覆盖 */ +export enum MessageTerminalPriority { + NORMAL = 0, // 普通消息 + CONFIRMED = 1, // 服务端已确认消息 + RECALL = 2 // 撤回终态 +} + +/** 会话写 lane 与全量屏障 */ +interface ConversationWriteState { + barrierTail: Promise // 当前全量屏障尾部 + tails: Map> // 各会话写入尾部 +} + +interface RelationState { + terminated: boolean + messageId: number + localTerminationPending: boolean +} + +const writeState: ConversationWriteState = { + // 当前运行时的会话写状态 + barrierTail: Promise.resolve(), + tails: new Map() +} +const relationStates = new Map() // 群关系消息终态 + +/** 同一会话串行执行消息与会话终态写入 */ +export async function enqueueConversationWrite( + clientConversationId: string, + operation: () => Promise +): Promise { + return enqueueConversationWrites([clientConversationId], operation) +} + +/** 一次写入原子占用全部会话 lane,避免嵌套获取与屏障互锁 */ +export function enqueueConversationWrites( + clientConversationIds: string[], + operation: () => Promise +): Promise { + // 1. 等待全量屏障和所有目标会话的前驱写入 + const keys = Array.from(new Set(clientConversationIds)).sort() + const predecessors = [ + writeState.barrierTail, + ...keys.map((key) => writeState.tails.get(key) || Promise.resolve()) + ] + const current = Promise.all(predecessors.map((task) => task.catch(() => undefined))).then( + operation + ) + const settled = current.then( + () => undefined, + () => undefined + ) + // 2. 先发布 recovery tail;完成时仅清理仍指向本任务的 lane + keys.forEach((key) => writeState.tails.set(key, settled)) + return current.finally(() => { + keys.forEach((key) => { + if (writeState.tails.get(key) === settled) { + writeState.tails.delete(key) + } + }) + }) +} + +/** 独占全部会话写入;只用于全量快照重建,常规写仍按会话并行 */ +export function enqueueConversationBarrier(operation: () => Promise): Promise { + // 1. 同步发布 gate,阻止后续会话写越过本次全量操作 + const previousBarrier = writeState.barrierTail + const existingWrites = Array.from(writeState.tails.values()) + let release!: () => void + const gate = new Promise((resolve) => { + release = resolve + }) + writeState.barrierTail = previousBarrier.catch(() => undefined).then(() => gate) + return (async () => { + // 2. 排空封门前的屏障和会话写,再独占执行全量操作 + await previousBarrier.catch(() => undefined) + await Promise.all(existingWrites.map((task) => task.catch(() => undefined))) + return await operation() + })().finally(release) +} + +/** 终态优先;相同优先级使用后到达状态 */ +export function reduceMessageState( + current: { priority: MessageTerminalPriority; value?: T } | undefined, + incoming: { priority: MessageTerminalPriority; value?: T } +) { + return current && current.priority > incoming.priority ? current : incoming +} + +/** 在会话写 lane 内记录关系终态;本地主动操作等待服务端终态消息后才允许重开 */ +export function markRelationTerminated(clientConversationId: string, messageId?: number): boolean { + return applyRelationState(clientConversationId, true, messageId) +} + +/** 显式重新加入后清除关系终态;旧通知不得重开新终态 */ +export function reopenRelation(clientConversationId: string, messageId?: number): boolean { + return applyRelationState(clientConversationId, false, messageId) +} + +export function isRelationTerminated(clientConversationId: string): boolean { + return relationStates.get(clientConversationId)?.terminated === true +} + +/** 排空当前 IM 运行时的消息写入,并在调用方仍允许时清理关系终态 */ +export async function clearMessageSyncState(shouldClear: () => boolean): Promise { + const barrier = writeState.barrierTail + const tails = Array.from(writeState.tails.entries()) + await Promise.all([ + barrier.catch(() => undefined), + ...tails.map(([, task]) => task.catch(() => undefined)) + ]) + if (!shouldClear()) { + return + } + relationStates.clear() + if (writeState.barrierTail === barrier) { + writeState.barrierTail = Promise.resolve() + } + tails.forEach(([key, task]) => { + if (writeState.tails.get(key) === task) { + writeState.tails.delete(key) + } + }) +} + +/** 按服务端关系消息编号单调归约群关系;本地主动终止在服务端终态确认前阻止旧成员消息重开 */ +function applyRelationState( + clientConversationId: string, + terminated: boolean, + messageId?: number +): boolean { + const current = relationStates.get(clientConversationId) + if (messageId === undefined) { + if (!terminated && current?.localTerminationPending) { + return false + } + if (terminated && current?.terminated) { + return true + } + relationStates.set(clientConversationId, { + terminated, + messageId: current?.messageId ?? 0, + localTerminationPending: terminated + }) + return true + } + if (!terminated && current?.localTerminationPending) { + return false + } + if (current && messageId <= current.messageId) { + return false + } + relationStates.set(clientConversationId, { + terminated, + messageId, + localTerminationPending: false + }) + return true +} diff --git a/apps/web-ele/src/views/im/utils/resourceRequest.ts b/apps/web-ele/src/views/im/utils/resourceRequest.ts new file mode 100644 index 000000000..2e5382a2a --- /dev/null +++ b/apps/web-ele/src/views/im/utils/resourceRequest.ts @@ -0,0 +1,108 @@ +/** 可合并请求的固定资源 */ +export enum ResourceRequestKey { + FACE_PACKS = 'facePacks', // 系统表情包 + FACE_USER_ITEMS = 'faceUserItems', // 用户表情 + FRIEND_LIST = 'friendList', // 好友列表 + GROUP_LIST = 'groupList', // 群列表 + CHANNEL_LIST = 'channelList', // 频道列表 + GROUP_REQUEST_UNHANDLED = 'groupRequestUnhandled' // 未处理加群申请 +} + +/** 固定资源请求的 task 生命周期模式 */ +export enum ResourceRequestMode { + CACHE_SUCCESS = 'cache-success', // 成功后持续复用,清理运行时状态时失效 + SINGLE_FLIGHT = 'single-flight' // 仅合并当前在途请求 +} + +/** 固定资源请求策略 */ +type ResourceRequestPolicy = + | { + mode: ResourceRequestMode.CACHE_SUCCESS + } + | { + mode: ResourceRequestMode.SINGLE_FLIGHT + refreshAfterPending?: boolean + } + +/** 单个固定资源当前发布的请求状态 */ +interface ResourceRequestEntry { + mode: ResourceRequestMode // task 生命周期模式 + task: Promise // 当前请求 task + trailingExecute?: () => Promise // 合并后的尾随刷新 +} + +const resourceRequests = new Map() // 每个 key 仅发布一个当前 entry + +/** 运行固定资源请求 */ +export function runResourceRequest( + key: ResourceRequestKey, + execute: () => Promise, + policy: ResourceRequestPolicy +): Promise { + const existing = resourceRequests.get(key) + // 1. 复用 task;force 只覆盖为一个最新尾随执行器 + if (existing) { + if (existing.mode !== policy.mode) { + return Promise.reject(new Error(`IM resource policy mismatch: ${key}`)) + } + if ( + existing.mode === ResourceRequestMode.SINGLE_FLIGHT && + policy.mode === ResourceRequestMode.SINGLE_FLIGHT && + policy.refreshAfterPending + ) { + existing.trailingExecute = execute + } + return existing.task as Promise + } + const task = Promise.resolve().then(execute) + const entry: ResourceRequestEntry = { mode: policy.mode, task } + resourceRequests.set(key, entry) + void task.then( + () => finishResourceRequest(key, entry, true), + () => finishResourceRequest(key, entry, false) + ) + return task +} + +/** 完成请求并按策略释放或补刷 */ +function finishResourceRequest( + key: ResourceRequestKey, + entry: ResourceRequestEntry, + succeeded: boolean +): void { + // 1. 旧 finalizer 不能修改已经替换的新 entry + if (resourceRequests.get(key) !== entry) { + return + } + // 2. once 成功保留;其余情况先释放当前 entry + if (entry.mode === ResourceRequestMode.CACHE_SUCCESS && succeeded) { + return + } + resourceRequests.delete(key) + // 3. single-flight 的多次 force 合并为一次后台尾随刷新 + if (entry.trailingExecute) { + void runResourceRequest(key, entry.trailingExecute, { + mode: ResourceRequestMode.SINGLE_FLIGHT + }).catch((error) => console.warn(`[IM] 尾随刷新 ${key} 失败`, error)) + } +} + +/** 排空并清理固定资源请求状态 */ +export async function clearResourceRequests(): Promise { + const entries = Array.from(resourceRequests.entries()) + entries.forEach(([, entry]) => { + entry.trailingExecute = undefined + }) + await Promise.all(entries.map(([, entry]) => entry.task.catch(() => undefined))) + entries.forEach(([key, entry]) => { + if (resourceRequests.get(key) === entry) { + resourceRequests.delete(key) + } + }) +} + +/** 判断固定资源当前是否有请求在途 */ +export function isResourceRequestPending(key: ResourceRequestKey): boolean { + const entry = resourceRequests.get(key) + return entry?.mode === ResourceRequestMode.SINGLE_FLIGHT +}