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.
593 lines
22 KiB
593 lines
22 KiB
import { defineEventHandler, getQuery, getRouterParam, readBody, setResponseHeaders } from "h3";
|
|
import { R } from "#server/utils/response";
|
|
import { getCurrentUser, getConfigGlobal, getConfigUser } from "#server/utils/context";
|
|
import { createOpenAICompatible } from "@ai-sdk/openai-compatible";
|
|
import { type LanguageModel, type ModelMessage, streamText, stepCountIs } from "ai";
|
|
import { getSessionByIdAndUser, getMessagesBySession, getMaxSortOrder, saveMessage, truncateMessagesAfter, deleteMessage, getMessageById, updateMessageContent, updateMessageParts, touchSession, countAssistantMessages } from "#server/service/agent/session";
|
|
import { getAgentToolsForChat, getAgentToolsBySlugs, executeAgentTool } from "#server/service/agent-tool";
|
|
import { getTempTokenFromCookie } from "#server/service/agent/temp-token";
|
|
import { checkRateLimit, incrementRateLimit } from "#server/service/agent/rate-limit";
|
|
import { getModelWithProviderById, getModelWithProviderByIdAny } from "#server/service/llm";
|
|
import { generateSessionTitle } from "#server/service/agent/title";
|
|
import { type StoredPart } from "./types";
|
|
import log4js from "logger";
|
|
|
|
const logger = log4js.getLogger("APP");
|
|
|
|
function resolveModel(
|
|
provider: {
|
|
name: string;
|
|
apiKey: string | null;
|
|
baseUrl: string | null;
|
|
parseMode: string;
|
|
},
|
|
modelId: string,
|
|
): LanguageModel {
|
|
const baseUrl = provider.baseUrl?.replace(/\/+$/, "") || undefined;
|
|
if (provider.parseMode === "anthropic") {
|
|
throw createError({
|
|
statusCode: 400,
|
|
statusMessage: "Anthropic 解析模式暂不支持流式对话,请使用 OpenAI 兼容模式",
|
|
});
|
|
}
|
|
const openaiCompatible = createOpenAICompatible({
|
|
name: provider.name,
|
|
apiKey: provider.apiKey || undefined,
|
|
baseURL: baseUrl || "https://api.openai.com/v1",
|
|
});
|
|
return openaiCompatible(modelId) as LanguageModel;
|
|
}
|
|
|
|
function buildModelMessages(
|
|
historyMessages: { id: string; role: string; content: string; parts: string | null }[],
|
|
opts: {
|
|
skipMessageId?: string;
|
|
currentUserContent?: string;
|
|
approvalToolResults?: Record<string, unknown>;
|
|
},
|
|
): ModelMessage[] {
|
|
const messages: ModelMessage[] = [];
|
|
|
|
for (const m of historyMessages) {
|
|
if (opts.skipMessageId && m.id === opts.skipMessageId) continue;
|
|
|
|
if (m.role === "assistant" && m.parts) {
|
|
let parts: StoredPart[];
|
|
try {
|
|
parts = JSON.parse(m.parts);
|
|
} catch {
|
|
parts = [];
|
|
}
|
|
|
|
const assistantContent: Array<Record<string, unknown>> = [];
|
|
|
|
if (m.content) {
|
|
assistantContent.push({ type: "text", text: m.content });
|
|
}
|
|
|
|
for (const p of parts) {
|
|
if (p.type === "tool-call" && p.toolCallId && p.toolName) {
|
|
if (p.state === "approval-requested" && p.approvalId) {
|
|
// Tool still waiting for approval - only generate tool-approval-request
|
|
// (SDK filters it out), do NOT generate tool-call to avoid MissingToolResultsError
|
|
assistantContent.push({
|
|
type: "tool-approval-request",
|
|
approvalId: p.approvalId,
|
|
toolCallId: p.toolCallId,
|
|
toolName: p.toolName,
|
|
input: p.args ?? {},
|
|
});
|
|
} else if (p.state === "approval-responded" && p.approvalId) {
|
|
const hasResult = opts.approvalToolResults?.[p.toolCallId] !== undefined;
|
|
if (p.approved === true && hasResult) {
|
|
assistantContent.push({
|
|
type: "tool-call",
|
|
toolCallId: p.toolCallId,
|
|
toolName: p.toolName,
|
|
input: p.args ?? {},
|
|
});
|
|
} else {
|
|
// Denied: generate BOTH tool-call AND tool-approval-request.
|
|
// collectToolApprovals (SDK) scans original messages for tool-approval-response
|
|
// and requires a matching tool-call, otherwise throws ToolCallNotFoundForApprovalError.
|
|
// convertToLanguageModelPrompt adds toolCallId to approvedToolCallIds when
|
|
// tool-approval-response exists, preventing MissingToolResultsError.
|
|
assistantContent.push({
|
|
type: "tool-call",
|
|
toolCallId: p.toolCallId,
|
|
toolName: p.toolName,
|
|
input: p.args ?? {},
|
|
});
|
|
assistantContent.push({
|
|
type: "tool-approval-request",
|
|
approvalId: p.approvalId,
|
|
toolCallId: p.toolCallId,
|
|
toolName: p.toolName,
|
|
input: p.args ?? {},
|
|
});
|
|
}
|
|
} else {
|
|
const toolCallPart: Record<string, unknown> = {
|
|
type: "tool-call",
|
|
toolCallId: p.toolCallId,
|
|
toolName: p.toolName,
|
|
input: p.args ?? {},
|
|
};
|
|
if (p.state === "result" && p.result !== undefined) {
|
|
toolCallPart.output = p.result;
|
|
}
|
|
assistantContent.push(toolCallPart);
|
|
}
|
|
} else if (p.type === "reasoning" && p.text) {
|
|
assistantContent.push({ type: "reasoning", text: p.text });
|
|
} else if (p.type === "text" && p.text && !m.content) {
|
|
assistantContent.push({ type: "text", text: p.text });
|
|
}
|
|
}
|
|
|
|
if (assistantContent.length > 0) {
|
|
messages.push({ role: "assistant", content: assistantContent } as ModelMessage);
|
|
}
|
|
|
|
const toolResults: Array<Record<string, unknown>> = [];
|
|
const approvalResponses: Array<Record<string, unknown>> = [];
|
|
for (const p of parts) {
|
|
if (p.type === "tool-call" && p.toolCallId && p.toolName && p.state === "result" && p.result !== undefined) {
|
|
toolResults.push({
|
|
type: "tool-result",
|
|
toolCallId: p.toolCallId,
|
|
toolName: p.toolName,
|
|
output: { type: "json" as const, value: p.result as unknown },
|
|
});
|
|
}
|
|
if (p.type === "tool-call" && p.state === "approval-responded" && p.approvalId) {
|
|
const hasResult = opts.approvalToolResults?.[p.toolCallId] !== undefined;
|
|
if (p.approved === true && hasResult) {
|
|
toolResults.push({
|
|
type: "tool-result",
|
|
toolCallId: p.toolCallId,
|
|
toolName: p.toolName,
|
|
output: { type: "json" as const, value: opts.approvalToolResults![p.toolCallId] as unknown },
|
|
});
|
|
} else {
|
|
approvalResponses.push({
|
|
type: "tool-approval-response",
|
|
approvalId: p.approvalId,
|
|
approved: p.approved ?? false,
|
|
reason: p.approvalReason,
|
|
});
|
|
}
|
|
}
|
|
}
|
|
|
|
if (toolResults.length > 0 || approvalResponses.length > 0) {
|
|
messages.push({ role: "tool", content: [...toolResults, ...approvalResponses] } as ModelMessage);
|
|
}
|
|
} else if (m.role === "user") {
|
|
messages.push({ role: "user", content: m.content });
|
|
}
|
|
}
|
|
|
|
if (opts.currentUserContent !== undefined) {
|
|
messages.push({ role: "user", content: opts.currentUserContent });
|
|
}
|
|
|
|
return messages;
|
|
}
|
|
|
|
export default defineEventHandler(async (event) => {
|
|
const user = await getCurrentUser(event);
|
|
const tempToken = getTempTokenFromCookie(event);
|
|
|
|
if (!user && !tempToken) {
|
|
throw createError({ statusCode: 401, statusMessage: "请先登录或创建临时会话" });
|
|
}
|
|
|
|
const body = await readBody(event);
|
|
const { sessionId, content, editMessageId, regenerate, continueAfterApproval, approvalToolCallId, approved, approvalReason } = body as {
|
|
sessionId: string;
|
|
content: string;
|
|
editMessageId?: string;
|
|
regenerate?: boolean;
|
|
continueAfterApproval?: boolean;
|
|
approvalToolCallId?: string;
|
|
approved?: boolean;
|
|
approvalReason?: string;
|
|
};
|
|
|
|
const isApprovalContinue = continueAfterApproval === true && approvalToolCallId;
|
|
|
|
if (!sessionId) {
|
|
throw createError({ statusCode: 400, statusMessage: "参数无效" });
|
|
}
|
|
|
|
if (!isApprovalContinue && (!content || typeof content !== "string" || content.trim().length === 0)) {
|
|
throw createError({ statusCode: 400, statusMessage: "参数无效" });
|
|
}
|
|
|
|
const session = await getSessionByIdAndUser(sessionId, user?.id ?? null, tempToken);
|
|
if (!session) {
|
|
throw createError({ statusCode: 404, statusMessage: "会话不存在" });
|
|
}
|
|
|
|
const ip = getRequestIP(event, { xForwardedFor: true }) || "unknown";
|
|
|
|
if (!user) {
|
|
const rateLimit = checkRateLimit(sessionId, ip);
|
|
if (rateLimit.blocked) {
|
|
setResponseHeaders(event, {
|
|
"X-RateLimit-Session-Remaining": String(rateLimit.sessionRemaining),
|
|
"X-RateLimit-Ip-Remaining": String(rateLimit.ipRemaining),
|
|
});
|
|
throw createError({
|
|
statusCode: 429,
|
|
statusMessage: rateLimit.sessionRemaining === 0 ? "当前会话已达回复上限,请登录后继续" : "今日回复次数已达上限,请登录或明日再试",
|
|
});
|
|
}
|
|
setResponseHeaders(event, {
|
|
"X-RateLimit-Session-Remaining": String(rateLimit.sessionRemaining),
|
|
"X-RateLimit-Ip-Remaining": String(rateLimit.ipRemaining),
|
|
});
|
|
}
|
|
|
|
let modelRow;
|
|
if (user) {
|
|
const preferredModelId = await getConfigUser(event, "preferredLlmModelId");
|
|
const targetModelId = session.modelId ?? preferredModelId;
|
|
if (!targetModelId) {
|
|
throw createError({ statusCode: 400, statusMessage: "未配置模型,请先在设置中选择模型" });
|
|
}
|
|
modelRow = await getModelWithProviderById(targetModelId, user.id);
|
|
} else {
|
|
if (!session.modelId) {
|
|
throw createError({ statusCode: 400, statusMessage: "系统未配置默认模型" });
|
|
}
|
|
modelRow = await getModelWithProviderByIdAny(session.modelId);
|
|
}
|
|
|
|
if (!modelRow) {
|
|
throw createError({ statusCode: 404, statusMessage: "模型不存在或无权访问" });
|
|
}
|
|
|
|
const { model, provider } = modelRow;
|
|
if (!provider.apiKey) {
|
|
throw createError({ statusCode: 400, statusMessage: "供应商未配置 API Key" });
|
|
}
|
|
|
|
const systemPrompt = await getConfigGlobal("agentSystemPrompt");
|
|
const titleModelId = await getConfigGlobal("agentTitleModelId");
|
|
|
|
let publicToolSlugs: string[] = [];
|
|
if (!user) {
|
|
const raw = await getConfigGlobal("agentPublicToolSlugs");
|
|
if (raw && Array.isArray(raw)) {
|
|
publicToolSlugs = raw as string[];
|
|
}
|
|
}
|
|
|
|
const { tools } = await getAgentToolsForChat({
|
|
userId: user?.id ?? null,
|
|
userRole: user?.role ?? null,
|
|
enableTools: session.enableTools === 1,
|
|
publicToolSlugs,
|
|
});
|
|
|
|
let historyMessages = await getMessagesBySession(sessionId, { limit: 100, latest: true });
|
|
|
|
let userMessage;
|
|
let shouldGenerateTitle = false;
|
|
|
|
if (isApprovalContinue) {
|
|
userMessage = historyMessages.filter((m) => m.role === "user").pop() ?? historyMessages[0];
|
|
const assistantCount = await countAssistantMessages(sessionId);
|
|
if (assistantCount <= 1) {
|
|
shouldGenerateTitle = true;
|
|
}
|
|
} else if (editMessageId) {
|
|
const editMsg = await getMessageById(editMessageId);
|
|
if (!editMsg || editMsg.sessionId !== sessionId) {
|
|
throw createError({ statusCode: 400, statusMessage: "编辑的消息不存在" });
|
|
}
|
|
await truncateMessagesAfter(sessionId, editMsg.sortOrder);
|
|
await updateMessageContent(editMessageId, content.trim());
|
|
historyMessages = await getMessagesBySession(sessionId, { limit: 100, latest: true });
|
|
userMessage = editMsg;
|
|
} else if (regenerate) {
|
|
const lastAssistant = historyMessages.filter((m) => m.role === "assistant").pop();
|
|
if (lastAssistant) {
|
|
await deleteMessage(lastAssistant.id);
|
|
historyMessages = await getMessagesBySession(sessionId, { limit: 100, latest: true });
|
|
}
|
|
const lastUser = historyMessages.filter((m) => m.role === "user").pop();
|
|
if (lastUser) {
|
|
await updateMessageContent(lastUser.id, content.trim());
|
|
userMessage = lastUser;
|
|
} else {
|
|
const sortOrder = (await getMaxSortOrder(sessionId)) + 1;
|
|
userMessage = await saveMessage({
|
|
sessionId,
|
|
role: "user",
|
|
content: content.trim(),
|
|
sortOrder,
|
|
});
|
|
shouldGenerateTitle = true;
|
|
}
|
|
} else {
|
|
const sortOrder = (await getMaxSortOrder(sessionId)) + 1;
|
|
userMessage = await saveMessage({
|
|
sessionId,
|
|
role: "user",
|
|
content: content.trim(),
|
|
sortOrder,
|
|
});
|
|
const assistantCount = await countAssistantMessages(sessionId);
|
|
if (assistantCount === 0) {
|
|
shouldGenerateTitle = true;
|
|
}
|
|
}
|
|
|
|
const approvalToolResults: Record<string, unknown> = {};
|
|
if (isApprovalContinue) {
|
|
const pendingSlugs = new Set<string>();
|
|
for (const m of historyMessages) {
|
|
if (!m.parts) continue;
|
|
let parts: StoredPart[];
|
|
try {
|
|
parts = JSON.parse(m.parts);
|
|
} catch {
|
|
continue;
|
|
}
|
|
for (const p of parts) {
|
|
if (p.type !== "tool-call" || !p.toolCallId || !p.toolName) continue;
|
|
if (p.state !== "approval-responded" || p.approved !== true) continue;
|
|
if (approvalToolResults[p.toolCallId] !== undefined) continue;
|
|
pendingSlugs.add(p.toolName);
|
|
}
|
|
}
|
|
const toolMap = await getAgentToolsBySlugs([...pendingSlugs]);
|
|
|
|
for (const m of historyMessages) {
|
|
if (!m.parts) continue;
|
|
let parts: StoredPart[];
|
|
try {
|
|
parts = JSON.parse(m.parts);
|
|
} catch {
|
|
continue;
|
|
}
|
|
for (const p of parts) {
|
|
if (p.type !== "tool-call" || !p.toolCallId || !p.toolName) continue;
|
|
if (p.state !== "approval-responded" || p.approved !== true) continue;
|
|
if (approvalToolResults[p.toolCallId] !== undefined) continue;
|
|
const agentTool = toolMap.get(p.toolName);
|
|
if (!agentTool) {
|
|
approvalToolResults[p.toolCallId] = { error: `Tool ${p.toolName} not found` };
|
|
continue;
|
|
}
|
|
try {
|
|
const result = await executeAgentTool(agentTool.id, p.args, user?.id ?? null);
|
|
approvalToolResults[p.toolCallId] = result;
|
|
logger.info("[%s] [APPROVAL-EXEC] tool=%s result=%s", event.context.requestId ?? "-", p.toolName, JSON.stringify(result).slice(0, 200));
|
|
} catch (err) {
|
|
approvalToolResults[p.toolCallId] = { error: err instanceof Error ? err.message : String(err) };
|
|
logger.error("[%s] [APPROVAL-EXEC] tool=%s error=%s", event.context.requestId ?? "-", p.toolName, approvalToolResults[p.toolCallId]);
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
const modelMessages = buildModelMessages(historyMessages, {
|
|
skipMessageId: isApprovalContinue ? undefined : userMessage?.id,
|
|
currentUserContent: isApprovalContinue ? undefined : content.trim(),
|
|
approvalToolResults,
|
|
});
|
|
|
|
const languageModel = resolveModel(provider, model.modelId);
|
|
|
|
logger.info(
|
|
"[%s] [AGENT-CHAT] userId=%s sessionId=%s modelId=%d thinking=%s tools=%s edit=%s regen=%s approval=%s",
|
|
event.context.requestId ?? "-",
|
|
user?.id ?? "temp",
|
|
sessionId,
|
|
model.id,
|
|
session.enableThinking === 1 ? "on" : "off",
|
|
Object.keys(tools).length > 0 ? "on" : "off",
|
|
editMessageId ? "yes" : "no",
|
|
regenerate ? "yes" : "no",
|
|
isApprovalContinue ? "yes" : "no",
|
|
);
|
|
|
|
const result = streamText({
|
|
model: languageModel,
|
|
system: systemPrompt || undefined,
|
|
messages: modelMessages,
|
|
maxOutputTokens: model.maxTokens || undefined,
|
|
...(Object.keys(tools).length > 0
|
|
? { tools, stopWhen: stepCountIs(8) }
|
|
: {}),
|
|
...(session.enableThinking === 1
|
|
? {
|
|
providerOptions: {
|
|
openaiCompatible: { reasoningEffort: "high" },
|
|
},
|
|
}
|
|
: {}),
|
|
onError: (errorData) => {
|
|
const errMsg = errorData?.error instanceof Error
|
|
? errorData.error.message
|
|
: String(errorData?.error ?? "未知错误");
|
|
logger.error("[%s] [AGENT-CHAT] streamText error: %s", event.context.requestId ?? "-", errMsg);
|
|
},
|
|
onFinish: async ({ finishReason, usage, steps, text: assistantText }) => {
|
|
logger.info(
|
|
"[%s] [AGENT-CHAT] finished: reason=%s steps=%d inputTokens=%d outputTokens=%d",
|
|
event.context.requestId ?? "-",
|
|
finishReason,
|
|
steps.length,
|
|
usage?.inputTokens ?? 0,
|
|
usage?.outputTokens ?? 0,
|
|
);
|
|
|
|
const parts: Array<Record<string, unknown>> = [];
|
|
for (const step of steps) {
|
|
for (const r of step.reasoning) {
|
|
if (r.type === "reasoning" && r.text) {
|
|
parts.push({ id: `p_${parts.length}`, type: "reasoning", text: r.text });
|
|
}
|
|
}
|
|
const toolResultMap = new Map(step.toolResults.map((tr) => [tr.toolCallId, tr.output]));
|
|
const approvalRequestMap = new Map<string, string>();
|
|
for (const cp of step.content as Array<Record<string, unknown>>) {
|
|
if (cp.type === "tool-approval-request" && cp.approvalId && cp.toolCall) {
|
|
const tc = cp.toolCall as Record<string, unknown>;
|
|
if (tc.toolCallId) {
|
|
approvalRequestMap.set(tc.toolCallId as string, cp.approvalId as string);
|
|
}
|
|
}
|
|
}
|
|
for (const tc of step.toolCalls) {
|
|
const approvalId = approvalRequestMap.get(tc.toolCallId);
|
|
if (approvalId) {
|
|
parts.push({
|
|
id: `p_${parts.length}`,
|
|
type: "tool-call",
|
|
toolCallId: tc.toolCallId,
|
|
toolName: tc.toolName,
|
|
args: tc.input,
|
|
state: "approval-requested",
|
|
approvalId,
|
|
});
|
|
} else {
|
|
parts.push({
|
|
id: `p_${parts.length}`,
|
|
type: "tool-call",
|
|
toolCallId: tc.toolCallId,
|
|
toolName: tc.toolName,
|
|
args: tc.input,
|
|
result: toolResultMap.get(tc.toolCallId),
|
|
state: "result",
|
|
});
|
|
}
|
|
}
|
|
if (step.text) {
|
|
parts.push({ id: `p_${parts.length}`, type: "text", text: step.text });
|
|
}
|
|
}
|
|
|
|
if (isApprovalContinue) {
|
|
const lastAssistant = historyMessages.filter((m) => m.role === "assistant").pop();
|
|
if (lastAssistant) {
|
|
let existingParts: Array<Record<string, unknown>> = [];
|
|
try {
|
|
existingParts = lastAssistant.parts ? JSON.parse(lastAssistant.parts) : [];
|
|
} catch {
|
|
existingParts = [];
|
|
}
|
|
logger.info("[%s] [APPROVAL-CONTINUE] merged %d existing parts with %d new parts", event.context.requestId ?? "-", existingParts.length, parts.length);
|
|
for (const p of existingParts) {
|
|
if (p.type === "tool-call" && p.state === "approval-responded") {
|
|
if (p.approved === true && approvalToolResults[p.toolCallId] !== undefined) {
|
|
p.result = approvalToolResults[p.toolCallId];
|
|
p.state = "result";
|
|
} else if (p.approved === false) {
|
|
p.result = "工具执行被拒绝";
|
|
p.state = "result";
|
|
}
|
|
}
|
|
}
|
|
const newParts = parts.filter(
|
|
(np) => np.type !== "tool-call" || !existingParts.some((ep) => ep.toolCallId === np.toolCallId),
|
|
);
|
|
const idOffset = existingParts.length;
|
|
for (let i = 0; i < newParts.length; i++) {
|
|
newParts[i].id = `p_${idOffset + i}`;
|
|
}
|
|
const mergedParts = [...existingParts, ...newParts];
|
|
const mergedContent = (lastAssistant.content || "") + assistantText;
|
|
await updateMessageContent(lastAssistant.id, mergedContent);
|
|
await updateMessageParts(lastAssistant.id, JSON.stringify(mergedParts));
|
|
}
|
|
} else {
|
|
const assistantSortOrder = (await getMaxSortOrder(sessionId)) + 1;
|
|
await saveMessage({
|
|
sessionId,
|
|
role: "assistant",
|
|
content: assistantText,
|
|
parts: parts.length > 0 ? JSON.stringify(parts) : null,
|
|
modelId: model.id,
|
|
inputTokens: usage?.inputTokens ?? null,
|
|
outputTokens: usage?.outputTokens ?? null,
|
|
sortOrder: assistantSortOrder,
|
|
});
|
|
}
|
|
|
|
await touchSession(sessionId);
|
|
|
|
if (!user) {
|
|
incrementRateLimit(sessionId, ip);
|
|
const rateLimit = checkRateLimit(sessionId, ip);
|
|
setResponseHeaders(event, {
|
|
"X-RateLimit-Session-Remaining": String(rateLimit.sessionRemaining),
|
|
"X-RateLimit-Ip-Remaining": String(rateLimit.ipRemaining),
|
|
});
|
|
}
|
|
|
|
if (shouldGenerateTitle && titleModelId && assistantText.trim().length > 0) {
|
|
generateSessionTitle(sessionId, content.trim(), assistantText, titleModelId, user?.id ?? null).catch((e) => { logger.error("[AGENT-CHAT] title gen failed: %s", e?.message ?? e); });
|
|
}
|
|
},
|
|
});
|
|
|
|
const responseHeaders: Record<string, string> = {
|
|
"X-User-Message-Id": userMessage?.id ?? "",
|
|
};
|
|
|
|
const response = result.toUIMessageStreamResponse({
|
|
sendReasoning: true,
|
|
headers: responseHeaders,
|
|
});
|
|
|
|
if (continueAfterApproval && Object.keys(approvalToolResults).length > 0) {
|
|
const originalBody = response.body;
|
|
if (originalBody) {
|
|
const prefixChunks = Object.entries(approvalToolResults).map(
|
|
([toolCallId, output]) =>
|
|
`data: ${JSON.stringify({
|
|
type: "tool-output-available",
|
|
toolCallId,
|
|
output,
|
|
})}\n\n`,
|
|
);
|
|
const prefix = new TextEncoder().encode(prefixChunks.join(""));
|
|
|
|
const transformed = new ReadableStream({
|
|
async start(controller) {
|
|
controller.enqueue(prefix);
|
|
const reader = originalBody.getReader();
|
|
try {
|
|
while (true) {
|
|
const { done, value } = await reader.read();
|
|
if (done) break;
|
|
controller.enqueue(value);
|
|
}
|
|
} catch (err) {
|
|
logger.error("[%s] [AGENT-CHAT] stream read error: %s", event.context.requestId ?? "-", err instanceof Error ? err.message : String(err));
|
|
const errorChunk = new TextEncoder().encode(
|
|
`data: ${JSON.stringify({ type: "error", errorText: "流式响应中断" })}\n\n`,
|
|
);
|
|
controller.enqueue(errorChunk);
|
|
} finally {
|
|
controller.close();
|
|
}
|
|
},
|
|
});
|
|
|
|
return new Response(transformed, {
|
|
status: response.status,
|
|
statusText: response.statusText,
|
|
headers: response.headers,
|
|
});
|
|
}
|
|
}
|
|
|
|
return response;
|
|
});
|
|
|