import { streamText, type LanguageModel, type ModelMessage } from "ai"; import { resolveModelForUser, resolveModelAny, toLanguageModel } from "#server/service/llm/model-resolver"; import { buildCharacterSystemPrompt } from "./persona"; import { getMessagesByCharacterSession, getMaxCharacterSortOrder, saveCharacterMessage, touchCharacterSession, countCharacterMessagesByRole, updateCharacterSession, } from "./session"; import { registerCharacterAbort, unregisterCharacterAbort } from "./abort"; import { checkCharacterRateLimit, incrementCharacterRateLimit } from "./rate-limit"; import { createCharacterStreamBuffer, appendCharacterChunk, markCharacterBufferDone, } from "./stream-buffer"; import type { CharacterCardRow, CharacterSessionRow } from "./types"; import log4js from "logger"; const logger = log4js.getLogger("APP"); export interface CharacterChatEngineParams { card: CharacterCardRow; session: CharacterSessionRow; user: { id: number } | null; content: string; preferredModelId?: number | null; globalDefaultModelId?: number | null; ip: string; setRateLimitHeaders?: (headers: Record) => void; requestId?: string; onClientClose?: (cb: () => void) => void; } export interface CharacterChatEngineResult { response: Response; userMessageId: string | null; } function truncateFirstLine(text: string): string { const firstLine = text.split("\n")[0]?.trim() ?? ""; if (firstLine.length <= 50) return firstLine; return firstLine.slice(0, 50) + "…"; } /** * 角色对话引擎(纯文本流式)。 * * 与 agent 的 chat-engine 完全解耦:不支持工具调用、审批、思考模式、 * 持久记忆注入、多步推理。只做一件事 —— 用组装好的人设上下文 * (buildCharacterSystemPrompt)做角色扮演式流式回复。 */ export async function executeCharacterChat( params: CharacterChatEngineParams, ): Promise { const { card, session, user, content, ip, requestId = "-" } = params; const sessionId = session.id; const textContent = content.trim(); // 0. 访客限流(仅未登录用户生效) if (!user) { const rateLimit = checkCharacterRateLimit(card.slug, sessionId, ip); if (rateLimit.blocked) { params.setRateLimitHeaders?.({ "X-RateLimit-Session-Remaining": String(rateLimit.sessionRemaining), "X-RateLimit-Ip-Remaining": String(rateLimit.ipRemaining), }); throw createError({ statusCode: 429, statusMessage: rateLimit.sessionRemaining === 0 ? "当前会话已达回复上限,请登录后继续" : "今日回复次数已达上限,请登录或明日再试", }); } params.setRateLimitHeaders?.({ "X-RateLimit-Session-Remaining": String(rateLimit.sessionRemaining), "X-RateLimit-Ip-Remaining": String(rateLimit.ipRemaining), }); } // 1. 模型解析(复用 llm 原语):session > 本次请求覆盖 > card > 全局默认 const targetModelId = session.modelId ?? params.preferredModelId ?? card.defaultModelId ?? params.globalDefaultModelId ?? null; if (!targetModelId) { throw createError({ statusCode: 400, statusMessage: "未配置模型,请先在设置中选择模型" }); } const resolvedModel = user ? await resolveModelForUser(targetModelId, user.id) : await resolveModelAny(targetModelId); if (!resolvedModel) { throw createError({ statusCode: 404, statusMessage: "模型不存在或无权访问" }); } if (!resolvedModel.provider.apiKey) { throw createError({ statusCode: 400, statusMessage: "供应商未配置 API Key" }); } const languageModel: LanguageModel = toLanguageModel(resolvedModel); // 2. 上下文自主设计:由角色卡组装 system prompt const systemPrompt = buildCharacterSystemPrompt(card); // 3. 持久化用户消息 const sortOrder = (await getMaxCharacterSortOrder(sessionId)) + 1; const userMessage = await saveCharacterMessage({ sessionId, role: "user", content: textContent, sortOrder, }); // 首条用户消息 → 用首行作为会话标题(无需额外 LLM 调用) const userCount = await countCharacterMessagesByRole(sessionId, "user"); if (userCount === 1) { const title = truncateFirstLine(textContent); if (title) { await updateCharacterSession(sessionId, { title }); } } // 4. 组装历史消息(纯文本 user/assistant) const history = await getMessagesByCharacterSession(sessionId, { limit: 100, latest: true }); const messages: ModelMessage[] = history .filter((m) => m.id !== userMessage.id && (m.role === "user" || m.role === "assistant")) .map((m) => ({ role: m.role as "user" | "assistant", content: m.content } as ModelMessage)); messages.push({ role: "user", content: textContent } as ModelMessage); // 5. 流式生成 const controller = new AbortController(); registerCharacterAbort(sessionId, controller); params.onClientClose?.(() => { if (!controller.signal.aborted) { controller.abort(); } }); createCharacterStreamBuffer(sessionId, { modelId: resolvedModel.dbId, userMessageId: userMessage.id, }); let streamedText = ""; logger.info( "[%s] [CHAR-ENGINE] card=%s sessionId=%s model=%d", requestId, card.slug, sessionId, resolvedModel.dbId, ); const result = streamText({ model: languageModel, system: systemPrompt, messages, abortSignal: controller.signal, onChunk: ({ chunk }) => { if (chunk.type === "text-delta") { streamedText += chunk.text; } }, onFinish: async ({ text }) => { try { if (controller.signal.aborted) { unregisterCharacterAbort(sessionId); markCharacterBufferDone(sessionId); return; } const finalText = text.trim(); if (finalText) { const assistantSortOrder = (await getMaxCharacterSortOrder(sessionId)) + 1; await saveCharacterMessage({ sessionId, role: "assistant", content: finalText, sortOrder: assistantSortOrder, }); } await touchCharacterSession(sessionId); if (!user) { incrementCharacterRateLimit(card.slug, sessionId, ip); const rateLimit = checkCharacterRateLimit(card.slug, sessionId, ip); params.setRateLimitHeaders?.({ "X-RateLimit-Session-Remaining": String(rateLimit.sessionRemaining), "X-RateLimit-Ip-Remaining": String(rateLimit.ipRemaining), }); } unregisterCharacterAbort(sessionId); markCharacterBufferDone(sessionId); } catch (e) { logger.error("[%s] [CHAR-ENGINE] onFinish error: %s", requestId, e instanceof Error ? e.message : String(e)); unregisterCharacterAbort(sessionId); markCharacterBufferDone(sessionId); } }, onAbort: async () => { try { const partial = streamedText.trim(); if (partial) { const assistantSortOrder = (await getMaxCharacterSortOrder(sessionId)) + 1; await saveCharacterMessage({ sessionId, role: "assistant", content: partial, sortOrder: assistantSortOrder, }); } await touchCharacterSession(sessionId); if (!user) { incrementCharacterRateLimit(card.slug, sessionId, ip); } unregisterCharacterAbort(sessionId); markCharacterBufferDone(sessionId); } catch (e) { logger.error("[%s] [CHAR-ENGINE] onAbort error: %s", requestId, e instanceof Error ? e.message : String(e)); unregisterCharacterAbort(sessionId); markCharacterBufferDone(sessionId); } }, }); const response = result.toUIMessageStreamResponse({ headers: { "X-User-Message-Id": userMessage.id }, }); if (response.body) { const bufferedBody = wrapStreamWithBuffer(response.body, sessionId, controller, requestId); return { response: new Response(bufferedBody, { status: response.status, statusText: response.statusText, headers: response.headers, }), userMessageId: userMessage.id, }; } return { response, userMessageId: userMessage.id }; } function wrapStreamWithBuffer( originalBody: ReadableStream, sessionId: string, serverAbortController: AbortController, requestId: string, ): ReadableStream { return new ReadableStream({ start(controller) { const reader = originalBody.getReader(); let clientDisconnected = false; function pump() { reader.read().then(({ done, value }) => { if (done) { markCharacterBufferDone(sessionId); if (!clientDisconnected) { try { controller.close(); } catch { /* already closed */ } } return; } if (value) { appendCharacterChunk(sessionId, value); if (!clientDisconnected) { try { controller.enqueue(value); } catch { clientDisconnected = true; } } } pump(); }).catch((err) => { if (serverAbortController.signal.aborted) { logger.info("[%s] [CHAR-ENGINE] stream reader aborted by server abort", requestId); } else { logger.error("[%s] [CHAR-ENGINE] buffer stream read error: %s", requestId, err instanceof Error ? err.message : String(err)); } markCharacterBufferDone(sessionId); if (!clientDisconnected) { try { controller.close(); } catch { /* already closed */ } } }); } pump(); }, }); }