114 lines
3.5 KiB
TypeScript
114 lines
3.5 KiB
TypeScript
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 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 }) => {
|
|
return sendMessageViaHttp(input.characterId, input.content);
|
|
});
|
|
|
|
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();
|
|
|
|
queue.setConsumer(async (content) => {
|
|
sendBack({ type: "ChatQueuedSendStarted" });
|
|
try {
|
|
const output = await sendMessageViaHttp(characterId, content);
|
|
sendBack({ type: "ChatQueuedHttpDone", output });
|
|
} catch (error) {
|
|
const errorMessage = ExceptionHandler.message(error);
|
|
log.error("[chat-machine] message queue send failed", {
|
|
error,
|
|
});
|
|
sendBack({
|
|
type: "ChatQueuedSendError",
|
|
content,
|
|
errorMessage,
|
|
});
|
|
}
|
|
});
|
|
|
|
receive((event) => {
|
|
if (event.type !== "ChatSendMessage") return;
|
|
const content = event.content.trim();
|
|
if (content.length === 0) return;
|
|
queue.enqueue(content);
|
|
});
|
|
|
|
return () => queue.dispose();
|
|
}
|
|
|
|
async function sendMessageViaHttp(characterId: string, content: string): Promise<{
|
|
response: ChatSendResponse;
|
|
reply: UiMessage | null;
|
|
}> {
|
|
const chatRepo = await loadChatRepository();
|
|
const cacheIdentityResult = await resolveChatConversationKey(characterId);
|
|
const result = await chatRepo.sendMessage(characterId, content);
|
|
if (Result.isErr(result)) {
|
|
log.error("[chat-machine] sendMessageHttpActor failed", {
|
|
error: result.error,
|
|
});
|
|
throw result.error;
|
|
}
|
|
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
|
|
);
|
|
}
|