refactor: enhance realtime message handling and conversation state management

This commit is contained in:
mlogclub
2026-04-25 17:24:26 +08:00
parent 0b22c003b0
commit 96fde7aaab
7 changed files with 496 additions and 126 deletions
+2
View File
@@ -2,6 +2,7 @@ package services
import (
"cs-agent/internal/pkg/dto"
"cs-agent/internal/pkg/dto/response"
"cs-agent/internal/pkg/enums"
"cs-agent/internal/pkg/openidentity"
"encoding/json"
@@ -134,6 +135,7 @@ func (e RealtimeResyncRequiredEvent) EventPayload() RealtimeEventPayload {
type RealtimeMessageCreatedPayload struct {
ConversationID int64 `json:"conversationId,omitempty"`
MessageID int64 `json:"messageId,omitempty"`
Message response.MessageResponse `json:"message,omitempty"`
Status enums.IMConversationStatus `json:"status,omitempty"`
CurrentAssigneeID int64 `json:"currentAssigneeId,omitempty"`
SenderType enums.IMSenderType `json:"senderType,omitempty"`
+82
View File
@@ -3,6 +3,7 @@ package services
import (
"cs-agent/internal/models"
"cs-agent/internal/pkg/dto"
"cs-agent/internal/pkg/dto/response"
"cs-agent/internal/pkg/enums"
"cs-agent/internal/pkg/errorsx"
"cs-agent/internal/pkg/openidentity"
@@ -275,6 +276,7 @@ func (s *wsService) PublishMessageCreated(conversation *models.Conversation, mes
Payload: RealtimeMessageCreatedPayload{
ConversationID: conversation.ID,
MessageID: message.ID,
Message: s.buildRealtimeMessage(message),
Status: conversation.Status,
CurrentAssigneeID: conversation.CurrentAssigneeID,
SenderType: message.SenderType,
@@ -290,6 +292,86 @@ func (s *wsService) PublishMessageCreated(conversation *models.Conversation, mes
s.PublishToTopics(s.routeConversationTopics(conversation), event)
}
func (s *wsService) buildRealtimeMessage(item *models.Message) response.MessageResponse {
if item == nil {
return response.MessageResponse{}
}
agentReadState, customerReadState := ConversationReadStateService.GetConversationReadStates(item.ConversationID)
content, payload := utils.BuildRenderableMessage(item)
ret := response.MessageResponse{
ID: item.ID,
ConversationID: item.ConversationID,
ClientMsgID: item.ClientMsgID,
SenderType: item.SenderType,
SenderID: item.SenderID,
MessageType: item.MessageType,
Content: content,
Payload: payload,
SeqNo: item.SeqNo,
SendStatus: item.SendStatus,
SentAt: utils.FormatTimePtr(item.SentAt),
DeliveredAt: utils.FormatTimePtr(item.DeliveredAt),
ReadAt: utils.FormatTimePtr(item.ReadAt),
CustomerRead: isRealtimeMessageRead(item, customerReadState),
CustomerReadAt: realtimeReadMessageAt(item, customerReadState),
AgentRead: isRealtimeMessageRead(item, agentReadState),
AgentReadAt: realtimeReadMessageAt(item, agentReadState),
RecalledAt: utils.FormatTimePtr(item.RecalledAt),
QuotedMessageID: item.QuotedMessageID,
}
s.fillRealtimeMessageSender(&ret, item)
return ret
}
func (s *wsService) fillRealtimeMessageSender(ret *response.MessageResponse, item *models.Message) {
if ret == nil || item == nil || item.SenderID <= 0 {
return
}
switch item.SenderType {
case enums.IMSenderTypeAI:
if aiAgent := AIAgentService.Get(item.SenderID); aiAgent != nil {
ret.SenderName = aiAgent.Name
}
case enums.IMSenderTypeAgent:
if profile := AgentProfileService.GetByUserID(item.SenderID); profile != nil {
if displayName := strings.TrimSpace(profile.DisplayName); displayName != "" {
ret.SenderName = displayName
}
if avatar := strings.TrimSpace(profile.Avatar); avatar != "" {
ret.SenderAvatar = avatar
}
}
if ret.SenderName == "" {
s.fillRealtimeMessageUserName(ret, item.SenderID)
}
default:
s.fillRealtimeMessageUserName(ret, item.SenderID)
}
}
func (s *wsService) fillRealtimeMessageUserName(ret *response.MessageResponse, userID int64) {
if ret == nil || userID <= 0 {
return
}
if user := UserService.Get(userID); user != nil {
ret.SenderName = user.Nickname
if ret.SenderName == "" {
ret.SenderName = user.Username
}
}
}
func isRealtimeMessageRead(item *models.Message, state *models.ConversationReadState) bool {
return item != nil && state != nil && state.LastReadSeqNo >= item.SeqNo
}
func realtimeReadMessageAt(item *models.Message, state *models.ConversationReadState) string {
if !isRealtimeMessageRead(item, state) {
return ""
}
return utils.FormatTimePtr(state.LastReadAt)
}
func (s *wsService) PublishMessageRecalled(conversation *models.Conversation, message *models.Message) {
if conversation == nil || message == nil {
return
+63 -65
View File
@@ -4,23 +4,33 @@ import { useEffect, useRef } from "react"
import { toast } from "sonner"
import { createAdminWebSocketUrl } from "@/lib/api/admin"
import { type AgentMessage } from "@/lib/api/agent"
import { readSession } from "@/lib/auth"
import {
normalizeRealtimeMessage,
type RealtimeConversationPatch,
type RealtimeMessageCreatedPayload,
} from "@/lib/im-realtime-state"
import { createRealtimeConnectionManager } from "@/lib/realtime-connection"
import { getNotificationBody, showNotification } from "@/lib/services/notification"
import { useAgentConversationsStore } from "@/lib/stores/agent-conversations"
type AgentRealtimeConnection = ReturnType<typeof createRealtimeConnectionManager>
type AgentRealtimeEnvelope = {
eventId?: string
type?: string
data?: RealtimeMessageCreatedPayload<AgentMessage> &
RealtimeConversationPatch & {
messageId?: number
recalledAt?: string
sendStatus?: number
}
}
export function useAgentConversationRealtime() {
const selectedConversationId = useAgentConversationsStore(
(state) => state.selectedConversationId
)
const loadConversations = useAgentConversationsStore(
(state) => state.loadConversations
)
const syncLatestMessages = useAgentConversationsStore(
(state) => state.syncLatestMessages
)
const setRealtimeStatus = useAgentConversationsStore(
(state) => state.setRealtimeStatus
)
@@ -57,22 +67,11 @@ export function useAgentConversationRealtime() {
},
onMessage: (event, socket) => {
try {
const payload = JSON.parse(event.data) as {
eventId?: string
type?: string
data?: {
conversationId?: number
messageId?: number
status?: number
currentAssigneeId?: number
senderType?: string
messageType?: string
content?: string
}
}
const eventType = payload.type ?? ""
const conversationId = payload.data?.conversationId ?? 0
const eventId = payload.eventId?.trim() ?? ""
const envelope = JSON.parse(event.data) as AgentRealtimeEnvelope
const eventType = envelope.type ?? ""
const payload = envelope.data
const conversationId = payload?.conversationId ?? 0
const eventId = envelope.eventId?.trim() ?? ""
if (
eventType === "" ||
@@ -88,51 +87,50 @@ export function useAgentConversationRealtime() {
socket.send(JSON.stringify({ type: "ack", eventId }))
}
if (eventType === "message.created" && conversationId > 0) {
const senderType = payload.data?.senderType ?? ""
const status = payload.data?.status ?? 0
const currentAssigneeId = payload.data?.currentAssigneeId ?? 0
void loadConversations()
.then(() => {
const store = useAgentConversationsStore.getState()
const shouldNotify =
senderType === "customer" &&
status === 2 &&
currentAssigneeId > 0 &&
currentAssigneeId === currentUserIdRef.current &&
typeof document !== "undefined" &&
document.visibilityState !== "visible"
if (!shouldNotify) {
return
}
showNotification(
"新消息",
getNotificationBody({
messageType: payload.data?.messageType ?? "",
content: payload.data?.content ?? "",
}),
() => {
void store.selectConversation(conversationId)
}
)
})
.catch((error) => {
toast.error(error instanceof Error ? error.message : "加载消息失败")
})
} else {
void loadConversations().catch((error) => {
toast.error(error instanceof Error ? error.message : "加载会话列表失败")
const store = useAgentConversationsStore.getState()
if (eventType === "resync.required") {
void store.resyncRealtimeData(conversationId).catch((error) => {
toast.error(error instanceof Error ? error.message : "同步会话数据失败")
})
return
}
if (
conversationId > 0 &&
selectedConversationIdRef.current === conversationId
) {
void syncLatestMessages(conversationId)
if (eventType === "message.created") {
const message = normalizeRealtimeMessage<AgentMessage>(payload)
if (!message) {
void store.resyncRealtimeData(conversationId).catch((error) => {
toast.error(error instanceof Error ? error.message : "同步消息失败")
})
return
}
store.applyRealtimeMessageCreated(message)
const shouldNotify =
message.senderType === "customer" &&
payload?.status === 2 &&
(payload.currentAssigneeId ?? 0) > 0 &&
payload.currentAssigneeId === currentUserIdRef.current &&
typeof document !== "undefined" &&
document.visibilityState !== "visible"
if (shouldNotify) {
showNotification("新消息", getNotificationBody(message), () => {
void store.selectConversation(message.conversationId)
})
}
return
}
if (eventType === "message.recalled" && payload?.messageId) {
store.applyRealtimeMessageRecalled(payload.messageId, {
sendStatus: payload.sendStatus,
recalledAt: payload.recalledAt,
})
return
}
if (eventType.startsWith("conversation.") && payload) {
store.applyRealtimeConversationChanged(payload)
}
} catch {
// ignore invalid ws payload
@@ -167,7 +165,7 @@ export function useAgentConversationRealtime() {
realtime.disconnect()
subscribedConversationIdRef.current = null
}
}, [loadConversations, setRealtimeStatus, syncLatestMessages])
}, [setRealtimeStatus])
useEffect(() => {
const socket = realtimeRef.current?.getSocket()
+174
View File
@@ -0,0 +1,174 @@
import type { AgentConversation, AgentMessage } from "@/lib/api/agent"
import type { ImConversation, ImMessage } from "@/lib/api/im"
import { mergeImMessagesByIdAsc } from "@/lib/im-message-merge"
import { summarizeIMMessage } from "@/lib/im-message"
export type RealtimeMessage = AgentMessage | ImMessage
export type RealtimeConversation = AgentConversation | ImConversation
export type RealtimeMessageCreatedPayload<TMessage extends RealtimeMessage> = {
conversationId?: number
messageId?: number
message?: TMessage
senderType?: string
senderId?: number
senderName?: string
senderAvatar?: string
messageType?: string
content?: string
payload?: string
seqNo?: number
sendStatus?: number
sentAt?: string
}
export type RealtimeConversationPatch = Partial<RealtimeConversation> & {
conversationId?: number
}
export function normalizeRealtimeMessage<TMessage extends RealtimeMessage>(
payload: RealtimeMessageCreatedPayload<TMessage> | null | undefined
): TMessage | null {
if (!payload) {
return null
}
if (payload.message?.id) {
return payload.message
}
const id = payload.messageId ?? 0
const conversationId = payload.conversationId ?? 0
if (id <= 0 || conversationId <= 0) {
return null
}
return {
id,
conversationId,
senderType: payload.senderType ?? "",
senderId: payload.senderId ?? 0,
senderName: payload.senderName,
senderAvatar: payload.senderAvatar,
messageType: payload.messageType ?? "",
content: payload.content ?? "",
payload: payload.payload,
seqNo: payload.seqNo ?? 0,
sendStatus: payload.sendStatus ?? 0,
sentAt: payload.sentAt,
customerRead: false,
agentRead: false,
} as TMessage
}
export function mergeRealtimeMessage<TMessage extends RealtimeMessage>(
messages: TMessage[],
message: TMessage | null | undefined
): TMessage[] {
if (!message) {
return messages
}
return mergeImMessagesByIdAsc(messages, [message])
}
export function patchConversation<TConversation extends RealtimeConversation>(
conversation: TConversation | null,
patch: RealtimeConversationPatch | null | undefined
): TConversation | null {
if (!conversation || !patch) {
return conversation
}
const id = patch.id ?? patch.conversationId
if (!id || conversation.id !== id) {
return conversation
}
const fields = { ...patch }
delete fields.conversationId
return {
...conversation,
...fields,
} as TConversation
}
export function patchConversationList<TConversation extends RealtimeConversation>(
conversations: TConversation[],
patch: RealtimeConversationPatch | null | undefined
): TConversation[] {
if (!patch) {
return conversations
}
const id = patch.id ?? patch.conversationId
if (!id) {
return conversations
}
let changed = false
const next = conversations.map((item) => {
const patched = patchConversation(item, patch)
if (patched !== item) {
changed = true
}
return patched ?? item
})
return changed ? next : conversations
}
export function patchConversationWithMessage<
TConversation extends RealtimeConversation,
TMessage extends RealtimeMessage,
>(conversation: TConversation | null, message: TMessage | null | undefined) {
if (!conversation || !message || conversation.id !== message.conversationId) {
return conversation
}
return {
...conversation,
lastMessageId: message.id,
lastMessageAt: message.sentAt ?? conversation.lastMessageAt,
lastActiveAt: message.sentAt ?? conversation.lastActiveAt,
lastMessageSummary: summarizeIMMessage(message),
} as TConversation
}
export function patchConversationListWithMessage<
TConversation extends RealtimeConversation,
TMessage extends RealtimeMessage,
>(conversations: TConversation[], message: TMessage | null | undefined) {
if (!message) {
return conversations
}
let changed = false
const next = conversations.map((item) => {
const patched = patchConversationWithMessage(item, message)
if (patched !== item) {
changed = true
}
return patched ?? item
})
return changed ? next : conversations
}
export function markMessagesReadToSeqNo<TMessage extends RealtimeMessage>(
messages: TMessage[],
seqNo: number,
reader: "agent" | "customer",
readAt?: string
): TMessage[] {
if (seqNo <= 0) {
return messages
}
let changed = false
const next = messages.map((message) => {
if (message.seqNo > seqNo) {
return message
}
if (reader === "agent") {
if (message.agentRead && message.agentReadAt === readAt) {
return message
}
changed = true
return { ...message, agentRead: true, agentReadAt: readAt } as TMessage
}
if (message.customerRead && message.customerReadAt === readAt) {
return message
}
changed = true
return { ...message, customerRead: true, customerReadAt: readAt } as TMessage
})
return changed ? next : messages
}
+7 -9
View File
@@ -1,18 +1,16 @@
import { createWebSocketBaseUrl } from "@/lib/api/websocket"
import { getImVisitorId } from "@/lib/api/im"
import { getImVisitorId, type ImMessage } from "@/lib/api/im"
import { readKefuWidgetConfig } from "@/lib/kefu-widget-config"
import type {
RealtimeConversationPatch,
RealtimeMessageCreatedPayload,
} from "@/lib/im-realtime-state"
export type ImRealtimeEnvelope = {
type: string
topic?: string
data?: {
conversationId?: number
messageId?: number
}
payload?: {
conversationId?: number
messageId?: number
}
data?: RealtimeMessageCreatedPayload<ImMessage> & RealtimeConversationPatch
payload?: RealtimeMessageCreatedPayload<ImMessage> & RealtimeConversationPatch
}
export function createImRealtimeConnection() {
+103 -29
View File
@@ -21,6 +21,12 @@ import {
mergeImMessagesByIdAsc,
parseImMessageCursorId,
} from "@/lib/im-message-merge"
import {
markMessagesReadToSeqNo,
patchConversationList,
patchConversationListWithMessage,
type RealtimeConversationPatch,
} from "@/lib/im-realtime-state"
import { summarizeIMMessage } from "@/lib/im-message"
import { generateUUID } from "@/lib/utils"
@@ -89,6 +95,10 @@ type AgentConversationsStore = {
uploadImage: (file: File) => Promise<AgentAsset | null>
sendAttachment: (file: File) => Promise<AgentMessage | null>
recallMessage: (messageId: number) => Promise<AgentMessage | null>
applyRealtimeMessageCreated: (message: AgentMessage) => void
applyRealtimeConversationChanged: (patch: RealtimeConversationPatch) => void
applyRealtimeMessageRecalled: (messageId: number, patch: Partial<AgentMessage>) => void
resyncRealtimeData: (conversationId?: number) => Promise<void>
}
let conversationsRequestSeq = 0
@@ -391,6 +401,77 @@ export const useAgentConversationsStore = create<AgentConversationsStore>((set,
}
},
applyRealtimeMessageCreated: (message) => {
set((state) => {
const isSelected = state.selectedConversationId === message.conversationId
const nextMessages = isSelected
? mergeImMessagesByIdAsc(state.messages, [message])
: state.messages
return {
messages: nextMessages,
conversations: patchConversationListWithMessage(
state.conversations,
message
),
}
})
},
applyRealtimeConversationChanged: (patch) => {
set((state) => {
const conversationId = patch.id ?? patch.conversationId ?? 0
let nextMessages = state.messages
if (
conversationId > 0 &&
state.selectedConversationId === conversationId
) {
if ((patch.agentLastReadSeqNo ?? 0) > 0) {
nextMessages = markMessagesReadToSeqNo(
nextMessages,
patch.agentLastReadSeqNo ?? 0,
"agent",
patch.agentLastReadAt
)
}
if ((patch.customerLastReadSeqNo ?? 0) > 0) {
nextMessages = markMessagesReadToSeqNo(
nextMessages,
patch.customerLastReadSeqNo ?? 0,
"customer",
patch.customerLastReadAt
)
}
}
return {
messages: nextMessages,
conversations: patchConversationList(state.conversations, patch),
}
})
},
applyRealtimeMessageRecalled: (messageId, patch) => {
if (messageId <= 0) {
return
}
set((state) => ({
messages: state.messages.map((item) =>
item.id === messageId ? { ...item, ...patch, id: item.id } : item
),
}))
},
resyncRealtimeData: async (conversationId) => {
await get().loadConversations()
const selectedConversationId = get().selectedConversationId
const targetConversationId = conversationId ?? selectedConversationId
if (targetConversationId && selectedConversationId === targetConversationId) {
await get().loadMessages(targetConversationId, {
forceLoading: false,
reset: false,
})
}
},
sendMessage: async (html) => {
const trimmedContent = html.trim()
const { selectedConversationId, sending } = get()
@@ -412,22 +493,17 @@ export const useAgentConversationsStore = create<AgentConversationsStore>((set,
messages: current.messages.some((m) => m.id === message.id)
? current.messages.map((m) => (m.id === message.id ? message : m))
: [...current.messages, message],
conversations: current.conversations.map((item) =>
item.id === selectedConversationId
? {
...item,
lastMessageAt: message.sentAt,
lastActiveAt: message.sentAt,
lastMessageSummary: summarizeIMMessage({
messageType: "html",
content: trimmedContent,
}),
agentUnreadCount: 0,
customerUnreadCount: (item.customerUnreadCount ?? 0) + 1,
agentLastReadMessageId: message.id,
agentLastReadSeqNo: message.seqNo,
}
: item
conversations: patchConversationList(
patchConversationListWithMessage(current.conversations, message),
{
conversationId: selectedConversationId,
agentUnreadCount: 0,
customerUnreadCount:
(current.conversations.find((item) => item.id === selectedConversationId)
?.customerUnreadCount ?? 0) + 1,
agentLastReadMessageId: message.id,
agentLastReadSeqNo: message.seqNo,
}
),
}))
}
@@ -474,19 +550,17 @@ export const useAgentConversationsStore = create<AgentConversationsStore>((set,
messages: current.messages.some((m) => m.id === message.id)
? current.messages.map((m) => (m.id === message.id ? message : m))
: [...current.messages, message],
conversations: current.conversations.map((item) =>
item.id === selectedConversationId
? {
...item,
lastMessageAt: message.sentAt,
lastActiveAt: message.sentAt,
lastMessageSummary: summarizeIMMessage(message),
agentUnreadCount: 0,
customerUnreadCount: (item.customerUnreadCount ?? 0) + 1,
agentLastReadMessageId: message.id,
agentLastReadSeqNo: message.seqNo,
}
: item
conversations: patchConversationList(
patchConversationListWithMessage(current.conversations, message),
{
conversationId: selectedConversationId,
agentUnreadCount: 0,
customerUnreadCount:
(current.conversations.find((item) => item.id === selectedConversationId)
?.customerUnreadCount ?? 0) + 1,
agentLastReadMessageId: message.id,
agentLastReadSeqNo: message.seqNo,
}
),
}))
}
+65 -23
View File
@@ -26,6 +26,12 @@ import {
mergeImMessagesByIdAsc,
parseImMessageCursorId,
} from "@/lib/im-message-merge"
import {
markMessagesReadToSeqNo,
normalizeRealtimeMessage,
patchConversation,
patchConversationWithMessage,
} from "@/lib/im-realtime-state"
import { summarizeIMMessage } from "@/lib/im-message"
import { createRealtimeConnectionManager } from "@/lib/realtime-connection"
import { generateUUID } from "@/lib/utils"
@@ -70,6 +76,30 @@ function ensureMessageList(value: ImMessage[] | null | undefined): ImMessage[] {
return Array.isArray(value) ? value : []
}
function markConversationReadMessages(
messages: ImMessage[],
payload: ImRealtimeEnvelope["data"] | ImRealtimeEnvelope["payload"]
) {
let next = messages
if ((payload?.agentLastReadSeqNo ?? 0) > 0) {
next = markMessagesReadToSeqNo(
next,
payload?.agentLastReadSeqNo ?? 0,
"agent",
payload?.agentLastReadAt
)
}
if ((payload?.customerLastReadSeqNo ?? 0) > 0) {
next = markMessagesReadToSeqNo(
next,
payload?.customerLastReadSeqNo ?? 0,
"customer",
payload?.customerLastReadAt
)
}
return next
}
export type KefuChatStore = {
title: string
subtitle: string
@@ -133,31 +163,43 @@ export const useKefuChatStore = create<KefuChatStore>((set, get) => {
return
}
const payload = event.data ?? event.payload
const needsRefresh =
event.type === "message.created" ||
event.type?.startsWith("conversation.")
if (event.type === "resync.required") {
void get().refreshMessages()
return
}
if (payload?.conversationId !== conversationId) {
return
}
if (needsRefresh && payload?.conversationId === conversationId) {
void get()
.syncLatestMessages()
.then(() => {
if (event.type !== "message.created") {
return
}
const state = get()
const lastMessage = state.messages.at(-1)
if (
lastMessage &&
lastMessage.senderType !== "customer" &&
typeof document !== "undefined" &&
document.visibilityState !== "visible"
) {
showNotification("新消息", getNotificationBody(lastMessage), () => {
state.setIsOpen(true)
state.setIsVisible(true)
})
}
if (event.type === "message.created") {
const message = normalizeRealtimeMessage<ImMessage>(payload)
if (!message) {
void get().syncLatestMessages()
return
}
set((state) => ({
messages: mergeImMessagesByIdAsc(state.messages, [message]),
conversation: patchConversationWithMessage(state.conversation, message),
}))
if (
message.senderType !== "customer" &&
typeof document !== "undefined" &&
document.visibilityState !== "visible"
) {
const state = get()
showNotification("新消息", getNotificationBody(message), () => {
state.setIsOpen(true)
state.setIsVisible(true)
})
}
return
}
if (event.type?.startsWith("conversation.")) {
set((state) => ({
conversation: patchConversation(state.conversation, payload),
messages: markConversationReadMessages(state.messages, payload),
}))
}
},
})