fix(chat): queue outgoing messages

This commit is contained in:
2026-06-26 11:11:43 +08:00
parent 2858ceccdc
commit fe85af0cd7
8 changed files with 240 additions and 60 deletions
+80 -14
View File
@@ -1,10 +1,11 @@
/**
* Chat 状态机:4 个 XState actor
* Chat 状态机:history / send / websocket / queue actors
*/
import { fromPromise, fromCallback } from "xstate";
import { createChatWebSocket, ChatWebSocket } from "@/core/net/chat-websocket";
import { MessageQueue } from "@/core/net/message-queue";
import type { ChatSendResponse } from "@/data/dto/chat";
import { Result, Logger } from "@/utils";
@@ -74,19 +75,7 @@ export const sendMessageHttpActor = fromPromise<
{ response: ChatSendResponse; reply: UiMessage | null },
{ content: string }
>(async ({ input }) => {
const result = await chatRepo.sendMessage(input.content);
if (Result.isErr(result)) {
log.error("[chat-machine] sendMessageHttpActor failed", { error: result.error });
throw result.error;
}
const isDailyLimit =
result.data.lockDetail.locked &&
result.data.lockDetail.showUpgrade &&
result.data.lockDetail.reason === "daily_limit";
return {
response: result.data,
reply: isDailyLimit ? null : sendResponseToUiMessage(result.data),
};
return sendMessageViaHttp(input.content);
});
/** 翻历史(pagination)—— 不走 local-first,纯 server fetch */
@@ -233,5 +222,82 @@ export const sendMessageWsActor = fromPromise<void, { content: string }>(
},
);
export const httpMessageQueueActor = fromCallback<ChatEvent>(
({ sendBack, receive }) => {
return createMessageQueueActor("http", sendBack, receive);
},
);
export const wsPreferredMessageQueueActor = fromCallback<ChatEvent>(
({ sendBack, receive }) => {
return createMessageQueueActor("wsPreferred", sendBack, receive);
},
);
function createMessageQueueActor(
mode: "http" | "wsPreferred",
sendBack: (event: ChatEvent) => void,
receive: (listener: (event: ChatEvent) => void) => void,
): () => void {
const queue = new MessageQueue();
queue.setConsumer(async (content) => {
try {
if (mode === "wsPreferred") {
const ws = getActiveChatWebSocket();
if (ws?.sendMessage(content)) return true;
log.debug("[chat-machine] message queue fallback to HTTP", {
contentLength: content.length,
});
}
const output = await sendMessageViaHttp(content);
sendBack({ type: "ChatQueuedHttpDone", output });
return true;
} catch (error) {
const errorMessage =
error instanceof Error ? error.message : "Message send failed";
log.error("[chat-machine] message queue send failed", {
error,
mode,
});
sendBack({
type: "ChatQueuedSendError",
content,
errorMessage,
});
return true;
}
});
receive((event) => {
if (event.type !== "ChatSendMessage") return;
const content = event.content.trim();
if (content.length === 0) return;
queue.enqueue(content);
});
return () => queue.dispose();
}
// Re-export `UiMessage` type for the chat-machine.ts public API
type UiMessage = import("@/data/dto/chat").UiMessage;
async function sendMessageViaHttp(content: string): Promise<{
response: ChatSendResponse;
reply: UiMessage | null;
}> {
const result = await chatRepo.sendMessage(content);
if (Result.isErr(result)) {
log.error("[chat-machine] sendMessageHttpActor failed", { error: result.error });
throw result.error;
}
const isDailyLimit =
result.data.lockDetail.locked &&
result.data.lockDetail.showUpgrade &&
result.data.lockDetail.reason === "daily_limit";
return {
response: result.data,
reply: isDailyLimit ? null : sendResponseToUiMessage(result.data),
};
}