import { fromCallback, fromPromise } from "xstate"; import { ExceptionHandler } from "@/core/errors"; import { MessageQueue } from "@/core/net/message-queue"; import type { ChatSendResponse } from "@/data/schemas/chat"; import { loadChatRepository } from "@/data/repositories/chat_repository_loader"; import { resolveChatConversationKey } from "@/data/repositories/chat_cache_identity"; import { Logger } from "@/utils/logger"; import { Result } from "@/utils/result"; import { isAbortError } from "@/utils/abort"; import { getCharacterErrorCode } from "@/data/services/api"; import type { ChatEvent } from "../../chat-events"; import { sendResponseToUiMessage } from "../../helper/message-mappers"; const log = new Logger("StoresChatChatSendFlow"); type UiMessage = import("@/stores/chat/ui-message").UiMessage; export const sendMessageHttpActor = fromPromise< { response: ChatSendResponse; reply: UiMessage | null }, { characterId: string; content: string } >(async ({ input, signal }) => { return sendMessageViaHttp(input.characterId, input.content, signal); }); export const httpMessageQueueActor = fromCallback< ChatEvent, { characterId: string } >( ({ input, sendBack, receive }) => { return createMessageQueueActor(input.characterId, sendBack, receive); }, ); function createMessageQueueActor( characterId: string, sendBack: (event: ChatEvent) => void, receive: (listener: (event: ChatEvent) => void) => void, ): () => void { const queue = new MessageQueue(); let activeController: AbortController | null = null; queue.setConsumer(async (content) => { const controller = new AbortController(); activeController = controller; sendBack({ type: "ChatQueuedSendStarted" }); try { const output = await sendMessageViaHttp( characterId, content, controller.signal, ); sendBack({ type: "ChatQueuedHttpDone", output }); } catch (error) { if (isAbortError(error)) return; const errorMessage = ExceptionHandler.message(error); log.error("[chat-machine] message queue send failed", { error, }); sendBack({ type: "ChatQueuedSendError", content, errorMessage, ...(getCharacterErrorCode(error) ? { characterErrorCode: getCharacterErrorCode(error) ?? undefined } : {}), }); } finally { if (activeController === controller) activeController = null; } }); receive((event) => { if (event.type !== "ChatSendMessage") return; const content = event.content.trim(); if (content.length === 0) return; queue.enqueue(content); }); return () => { activeController?.abort(); queue.dispose(); }; } async function sendMessageViaHttp( characterId: string, content: string, signal?: AbortSignal, ): Promise<{ response: ChatSendResponse; reply: UiMessage | null; }> { const chatRepo = await loadChatRepository(); const cacheIdentityResult = await resolveChatConversationKey(characterId); const result = await chatRepo.sendMessage(characterId, content, { signal }); if (Result.isErr(result)) { if (isAbortError(result.error)) throw result.error; log.error("[chat-machine] sendMessageHttpActor failed", { error: result.error, }); throw result.error; } signal?.throwIfAborted(); if (Result.isOk(cacheIdentityResult)) { void chatRepo.prefetchMediaForSendResponse( result.data, characterId, cacheIdentityResult.data, ); } const isInsufficientCredits = result.data.canSendMessage === false && !hasRenderableSendResponse(result.data); return { response: result.data, reply: isInsufficientCredits ? null : sendResponseToUiMessage(result.data), }; } function hasRenderableSendResponse(response: ChatSendResponse): boolean { const isLockedPaidMessage = response.lockDetail.locked && (response.lockDetail.reason === "private_message" || response.lockDetail.reason === "voice_message" || response.lockDetail.reason === "image_paywall" || response.lockDetail.reason === "image"); return ( response.reply.trim().length > 0 || (!response.lockDetail.locked && response.audioUrl.trim().length > 0) || Boolean(response.image.url) || isLockedPaidMessage ); }