You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
295 lines
9.6 KiB
295 lines
9.6 KiB
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<string, string>) => 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<CharacterChatEngineResult> {
|
|
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<Uint8Array>,
|
|
sessionId: string,
|
|
serverAbortController: AbortController,
|
|
requestId: string,
|
|
): ReadableStream<Uint8Array> {
|
|
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();
|
|
},
|
|
});
|
|
}
|
|
|