refactor: implement centralized realtime connection manager and enhance message handling logic

This commit is contained in:
mlogclub
2026-04-25 16:27:58 +08:00
parent e109d19ea9
commit 0b22c003b0
5 changed files with 358 additions and 341 deletions
+45
View File
@@ -39,6 +39,51 @@ export function mergeImMessagesByIdAsc<T extends MergeableImMessage>(
return Array.from(byId.values()).sort((x, y) => x.id - y.id)
}
export function parseImMessageCursorId(cursor: string): number {
const value = Number.parseInt(cursor, 10)
return Number.isFinite(value) && value > 0 ? value : 0
}
export function cursorFromLoadedImMessages<T extends Pick<MergeableImMessage, "id">>(
messages: T[]
): string {
if (messages.length === 0) {
return ""
}
return String(Math.min(...messages.map((message) => message.id)))
}
export function hasMoreAfterLatestImMessageMerge<
T extends Pick<MergeableImMessage, "id">,
>(args: {
previousMessages: T[]
previousHasMore: boolean
merged: T[]
apiHasMore: boolean
}): boolean {
const prevMin = minImMessageId(args.previousMessages)
const mergedMin = minImMessageId(args.merged)
if (mergedMin === null) {
return Boolean(args.apiHasMore)
}
if (!args.previousHasMore && prevMin !== null && mergedMin >= prevMin) {
return false
}
return args.previousHasMore || Boolean(args.apiHasMore)
}
function minImMessageId<T extends Pick<MergeableImMessage, "id">>(
messages: T[]
): number | null {
if (messages.length === 0) {
return null
}
return Math.min(...messages.map((message) => message.id))
}
export function mergeImMessage<T extends MergeableImMessage>(
existing: T,
incoming: T
+190
View File
@@ -0,0 +1,190 @@
"use client"
export type RealtimeConnectionStatus = "connecting" | "connected" | "disconnected"
type RealtimeConnectionManagerOptions = {
createSocket: () => WebSocket
canReconnect?: () => boolean
onStatusChange?: (status: RealtimeConnectionStatus) => void
onSocketChange?: (socket: WebSocket | null) => void
onOpen?: (socket: WebSocket) => void
onMessage?: (event: MessageEvent, socket: WebSocket) => void
onClose?: (event: CloseEvent, socket: WebSocket) => void
onError?: (event: Event, socket: WebSocket) => void
onConnectError?: (error: unknown) => void
buildPingMessage?: () => string
pingIntervalMs?: number
reconnectBaseDelayMs?: number
reconnectMaxDelayMs?: number
}
type DisconnectOptions = {
reconnect?: boolean
updateStatus?: boolean
}
const DEFAULT_PING_INTERVAL_MS = 20000
const DEFAULT_RECONNECT_BASE_DELAY_MS = 2000
const DEFAULT_RECONNECT_MAX_DELAY_MS = 30000
export function createRealtimeConnectionManager(
options: RealtimeConnectionManagerOptions
) {
let socket: WebSocket | null = null
let reconnectTimer: number | null = null
let pingTimer: number | null = null
let reconnectAttempt = 0
let reconnectEnabled = false
const clearReconnectTimer = () => {
if (reconnectTimer !== null) {
window.clearTimeout(reconnectTimer)
reconnectTimer = null
}
}
const clearPingTimer = () => {
if (pingTimer !== null) {
window.clearInterval(pingTimer)
pingTimer = null
}
}
const clearTimers = () => {
clearReconnectTimer()
clearPingTimer()
}
const setSocket = (nextSocket: WebSocket | null) => {
socket = nextSocket
options.onSocketChange?.(nextSocket)
}
const canReconnect = () => {
return reconnectEnabled && (options.canReconnect?.() ?? true)
}
const scheduleReconnect = () => {
if (!canReconnect() || reconnectTimer !== null) {
return
}
const delay = Math.min(
(options.reconnectBaseDelayMs ?? DEFAULT_RECONNECT_BASE_DELAY_MS) *
2 ** reconnectAttempt,
options.reconnectMaxDelayMs ?? DEFAULT_RECONNECT_MAX_DELAY_MS
)
options.onStatusChange?.("connecting")
reconnectTimer = window.setTimeout(() => {
reconnectTimer = null
reconnectAttempt += 1
if (canReconnect()) {
connect()
}
}, delay)
}
const connect = () => {
reconnectEnabled = true
if (!(options.canReconnect?.() ?? true)) {
return
}
disconnect({ reconnect: true, updateStatus: false })
reconnectEnabled = true
options.onStatusChange?.("connecting")
let nextSocket: WebSocket
try {
nextSocket = options.createSocket()
} catch (error) {
options.onStatusChange?.("disconnected")
options.onConnectError?.(error)
scheduleReconnect()
return
}
setSocket(nextSocket)
nextSocket.addEventListener("open", () => {
if (socket !== nextSocket) {
return
}
clearReconnectTimer()
clearPingTimer()
reconnectAttempt = 0
options.onStatusChange?.("connected")
const buildPingMessage =
options.buildPingMessage ?? (() => JSON.stringify({ type: "ping" }))
pingTimer = window.setInterval(() => {
if (nextSocket.readyState === WebSocket.OPEN) {
nextSocket.send(buildPingMessage())
}
}, options.pingIntervalMs ?? DEFAULT_PING_INTERVAL_MS)
options.onOpen?.(nextSocket)
})
nextSocket.addEventListener("message", (event) => {
if (socket === nextSocket) {
options.onMessage?.(event, nextSocket)
}
})
nextSocket.addEventListener("close", (event) => {
clearPingTimer()
const isCurrentSocket = socket === nextSocket
if (isCurrentSocket) {
setSocket(null)
}
options.onClose?.(event, nextSocket)
if (!isCurrentSocket) {
return
}
if (canReconnect()) {
scheduleReconnect()
} else {
options.onStatusChange?.("disconnected")
}
})
nextSocket.addEventListener("error", (event) => {
if (socket !== nextSocket) {
return
}
options.onError?.(event, nextSocket)
options.onStatusChange?.("disconnected")
scheduleReconnect()
})
}
const disconnect = (disconnectOptions?: DisconnectOptions) => {
reconnectEnabled = disconnectOptions?.reconnect ?? false
clearTimers()
if (!reconnectEnabled) {
reconnectAttempt = 0
}
const currentSocket = socket
setSocket(null)
if (
currentSocket &&
(currentSocket.readyState === WebSocket.OPEN ||
currentSocket.readyState === WebSocket.CONNECTING)
) {
currentSocket.close()
}
if (disconnectOptions?.updateStatus ?? !reconnectEnabled) {
options.onStatusChange?.("disconnected")
}
}
return {
connect,
disconnect,
reconnect: () => {
reconnectEnabled = true
connect()
},
getSocket: () => socket,
}
}
+11 -54
View File
@@ -15,7 +15,12 @@ import {
type AgentMessage,
} from "@/lib/api/agent"
import type { RealtimeConnectionStatusValue } from "@/components/realtime-connection-status"
import { mergeImMessagesByIdAsc } from "@/lib/im-message-merge"
import {
cursorFromLoadedImMessages,
hasMoreAfterLatestImMessageMerge,
mergeImMessagesByIdAsc,
parseImMessageCursorId,
} from "@/lib/im-message-merge"
import { summarizeIMMessage } from "@/lib/im-message"
import { generateUUID } from "@/lib/utils"
@@ -49,54 +54,6 @@ function ensureArray<T>(value: T[] | null | undefined): T[] {
return Array.isArray(value) ? value : []
}
function parseCursorId(cursor: string): number {
const n = Number.parseInt(cursor, 10)
return Number.isFinite(n) && n > 0 ? n : 0
}
/** 下一页「更旧」请求应传入的游标:当前已加载列表中的最小 message id(后端用 id < cursor */
function cursorFromLoadedMessages(messages: AgentMessage[]): string {
if (messages.length === 0) {
return ""
}
return String(Math.min(...messages.map((m) => m.id)))
}
function minMessageId(messages: AgentMessage[]): number | null {
if (messages.length === 0) {
return null
}
return Math.min(...messages.map((m) => m.id))
}
/**
* 拉「最新一页」做增量合并后,是否仍显示「还有更旧」。
* 若本地已确认没有更旧,且合并后最早一条 id 没有变小,则不能用接口对「最新一页」的 hasMore 再次打开(满页会误报)。
*/
function hasMoreAfterLatestSyncMerge(args: {
previousMessages: AgentMessage[]
previousHasMore: boolean
merged: AgentMessage[]
apiHasMore: boolean
}): boolean {
const prevMin = minMessageId(args.previousMessages)
const mergedMin = minMessageId(args.merged)
if (mergedMin === null) {
return Boolean(args.apiHasMore)
}
if (
!args.previousHasMore &&
prevMin !== null &&
mergedMin >= prevMin
) {
return false
}
return args.previousHasMore || Boolean(args.apiHasMore)
}
type AgentConversationsStore = {
searchKeyword: string
conversationFilter: AgentConversationFilterKey
@@ -298,7 +255,7 @@ export const useAgentConversationsStore = create<AgentConversationsStore>((set,
messagesLoading: false,
messagesLoadedConversationId: conversationId,
messagesCursor:
cursorFromLoadedMessages(list) || (data.cursor ?? ""),
cursorFromLoadedImMessages(list) || (data.cursor ?? ""),
messagesHasMore: Boolean(data.hasMore),
})
} catch (error) {
@@ -314,7 +271,7 @@ export const useAgentConversationsStore = create<AgentConversationsStore>((set,
if (!conversationId || get().messagesLoadingMore || !get().messagesHasMore) {
return
}
const cursorId = parseCursorId(get().messagesCursor)
const cursorId = parseImMessageCursorId(get().messagesCursor)
if (cursorId <= 0) {
return
}
@@ -335,7 +292,7 @@ export const useAgentConversationsStore = create<AgentConversationsStore>((set,
return {
messages: merged,
messagesCursor:
cursorFromLoadedMessages(merged) ||
cursorFromLoadedImMessages(merged) ||
(data.cursor ?? state.messagesCursor),
messagesHasMore: Boolean(data.hasMore),
messagesLoadingMore: false,
@@ -368,9 +325,9 @@ export const useAgentConversationsStore = create<AgentConversationsStore>((set,
return {
messages: merged,
messagesCursor:
cursorFromLoadedMessages(merged) ||
cursorFromLoadedImMessages(merged) ||
(data.cursor ?? state.messagesCursor),
messagesHasMore: hasMoreAfterLatestSyncMerge({
messagesHasMore: hasMoreAfterLatestImMessageMerge({
previousMessages: state.messages,
previousHasMore: state.messagesHasMore,
merged,
+46 -157
View File
@@ -20,14 +20,18 @@ import {
createImRealtimeConnection,
type ImRealtimeEnvelope,
} from "@/lib/im-realtime"
import { mergeImMessagesByIdAsc } from "@/lib/im-message-merge"
import {
cursorFromLoadedImMessages,
hasMoreAfterLatestImMessageMerge,
mergeImMessagesByIdAsc,
parseImMessageCursorId,
} from "@/lib/im-message-merge"
import { summarizeIMMessage } from "@/lib/im-message"
import { createRealtimeConnectionManager } from "@/lib/realtime-connection"
import { generateUUID } from "@/lib/utils"
type ChatStatus = "connecting" | "connected" | "disconnected"
const RECONNECT_BASE_DELAY = 2000
const RECONNECT_MAX_DELAY = 30000
const DEFAULT_PAGE_LIMIT = 50
function getNotificationBody(message: ImMessage): string {
@@ -66,45 +70,6 @@ function ensureMessageList(value: ImMessage[] | null | undefined): ImMessage[] {
return Array.isArray(value) ? value : []
}
function parseCursorId(cursor: string): number {
const value = Number.parseInt(cursor, 10)
return Number.isFinite(value) && value > 0 ? value : 0
}
function cursorFromLoadedMessages(messages: ImMessage[]): string {
if (messages.length === 0) {
return ""
}
return String(Math.min(...messages.map((message) => message.id)))
}
function minMessageId(messages: ImMessage[]): number | null {
if (messages.length === 0) {
return null
}
return Math.min(...messages.map((message) => message.id))
}
function hasMoreAfterLatestSyncMerge(args: {
previousMessages: ImMessage[]
previousHasMore: boolean
merged: ImMessage[]
apiHasMore: boolean
}): boolean {
const prevMin = minMessageId(args.previousMessages)
const mergedMin = minMessageId(args.merged)
if (mergedMin === null) {
return Boolean(args.apiHasMore)
}
if (!args.previousHasMore && prevMin !== null && mergedMin >= prevMin) {
return false
}
return args.previousHasMore || Boolean(args.apiHasMore)
}
export type KefuChatStore = {
title: string
subtitle: string
@@ -144,74 +109,29 @@ export type KefuChatStore = {
let bootstrapToken = 0
export const useKefuChatStore = create<KefuChatStore>((set, get) => {
let reconnectTimer: number | null = null
let pingTimer: number | null = null
let reconnectAttempt = 0
let shouldReconnect = false
const clearRealtimeTimers = () => {
if (reconnectTimer !== null) {
window.clearTimeout(reconnectTimer)
reconnectTimer = null
}
if (pingTimer !== null) {
window.clearInterval(pingTimer)
pingTimer = null
}
}
const scheduleReconnect = () => {
if (!shouldReconnect || reconnectTimer !== null) {
return
}
const delay = Math.min(
RECONNECT_BASE_DELAY * 2 ** reconnectAttempt,
RECONNECT_MAX_DELAY
)
set({ status: "connecting" })
reconnectTimer = window.setTimeout(() => {
reconnectTimer = null
reconnectAttempt += 1
if (!shouldReconnect || !get().isOpen) {
const realtime = createRealtimeConnectionManager({
createSocket: createImRealtimeConnection,
canReconnect: () => Boolean(get().isOpen && get().conversation?.id),
onStatusChange: (status) => {
if (get().isOpen || status === "disconnected") {
set({ status })
}
},
onSocketChange: (socket) => {
set({ socket })
},
onMessage: (messageEvent) => {
let event: ImRealtimeEnvelope
try {
event = JSON.parse(messageEvent.data) as ImRealtimeEnvelope
} catch {
return
}
connectSocket()
}, delay)
}
const closeSocket = (options?: { reconnect?: boolean }) => {
shouldReconnect = options?.reconnect ?? false
clearRealtimeTimers()
if (!shouldReconnect) {
reconnectAttempt = 0
}
const socket = get().socket
if (
socket &&
(socket.readyState === WebSocket.OPEN ||
socket.readyState === WebSocket.CONNECTING)
) {
socket.close()
}
set({ socket: null })
}
const connectSocket = () => {
const conversationId = get().conversation?.id
if (!conversationId) {
return
}
closeSocket({ reconnect: false })
shouldReconnect = true
const socket = createImRealtimeConnection()
set({ socket })
const handleRealtimeEvent = (event: ImRealtimeEnvelope) => {
const conversationId = get().conversation?.id
if (!conversationId) {
return
}
const payload = event.data ?? event.payload
const needsRefresh =
event.type === "message.created" ||
@@ -239,51 +159,21 @@ export const useKefuChatStore = create<KefuChatStore>((set, get) => {
}
})
}
},
})
const closeSocket = (options?: { reconnect?: boolean }) => {
realtime.disconnect({
reconnect: options?.reconnect ?? false,
updateStatus: true,
})
}
const connectSocket = () => {
if (!get().conversation?.id) {
return
}
socket.addEventListener("message", (event) => {
try {
handleRealtimeEvent(JSON.parse(event.data) as ImRealtimeEnvelope)
} catch {
return
}
})
socket.addEventListener("open", () => {
clearRealtimeTimers()
reconnectAttempt = 0
pingTimer = window.setInterval(() => {
if (socket.readyState === WebSocket.OPEN) {
socket.send(JSON.stringify({ type: "ping" }))
}
}, 20000)
if (get().isOpen && get().socket === socket) {
set({ status: "connected" })
}
})
socket.addEventListener("error", () => {
if (get().socket === socket) {
scheduleReconnect()
}
})
socket.addEventListener("close", () => {
if (pingTimer !== null) {
window.clearInterval(pingTimer)
pingTimer = null
}
if (get().socket === socket) {
set({ socket: null })
}
if (get().isOpen) {
if (shouldReconnect) {
scheduleReconnect()
} else {
set({ status: "disconnected" })
}
}
})
realtime.connect()
}
return {
@@ -388,7 +278,7 @@ export const useKefuChatStore = create<KefuChatStore>((set, get) => {
const results = ensureMessageList(page.results)
set({
messages: results,
messagesCursor: cursorFromLoadedMessages(results) || page.cursor || "",
messagesCursor: cursorFromLoadedImMessages(results) || page.cursor || "",
messagesHasMore: Boolean(page.hasMore) || results.length >= DEFAULT_PAGE_LIMIT,
})
} catch (error) {
@@ -418,8 +308,8 @@ export const useKefuChatStore = create<KefuChatStore>((set, get) => {
const merged = mergeImMessagesByIdAsc(state.messages, batch)
return {
messages: merged,
messagesCursor: cursorFromLoadedMessages(merged) || page.cursor || "",
messagesHasMore: hasMoreAfterLatestSyncMerge({
messagesCursor: cursorFromLoadedImMessages(merged) || page.cursor || "",
messagesHasMore: hasMoreAfterLatestImMessageMerge({
previousMessages: state.messages,
previousHasMore: state.messagesHasMore,
merged,
@@ -444,7 +334,7 @@ export const useKefuChatStore = create<KefuChatStore>((set, get) => {
return
}
const cursorId = parseCursorId(get().messagesCursor)
const cursorId = parseImMessageCursorId(get().messagesCursor)
if (cursorId <= 0) {
return
}
@@ -464,7 +354,7 @@ export const useKefuChatStore = create<KefuChatStore>((set, get) => {
)
return {
messages: merged,
messagesCursor: cursorFromLoadedMessages(merged) || page.cursor || "",
messagesCursor: cursorFromLoadedImMessages(merged) || page.cursor || "",
messagesHasMore: Boolean(page.hasMore) || results.length >= DEFAULT_PAGE_LIMIT,
messagesLoadingMore: false,
}
@@ -667,7 +557,6 @@ export const useKefuChatStore = create<KefuChatStore>((set, get) => {
try {
await get().refreshMessages()
if (get().isOpen) {
shouldReconnect = true
connectSocket()
}
} catch (error) {