import { defineEventHandler, getQuery, getRouterParam, readBody, setResponseHeaders } from "h3"; import { R } from "#server/utils/response"; import { getCurrentUser, getConfigGlobal, getConfigUser } from "#server/utils/context"; import { createOpenAI } from "@ai-sdk/openai"; 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 === "openai") { const openai = createOpenAI({ apiKey: provider.apiKey || undefined, baseURL: baseUrl, }); return openai(modelId) as unknown as LanguageModel; } 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; }, ): 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> = []; 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 = { 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> = []; const approvalResponses: Array> = []; 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, enableThinking: bodyEnableThinking, enableTools: bodyEnableTools } = body as { sessionId: string; content: string; editMessageId?: string; regenerate?: boolean; continueAfterApproval?: boolean; approvalToolCallId?: string; approved?: boolean; approvalReason?: string; enableThinking?: boolean; enableTools?: boolean; }; 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 toolsEnabled = bodyEnableTools !== undefined ? bodyEnableTools : session.enableTools === 1; const { tools } = await getAgentToolsForChat({ userId: user?.id ?? null, userRole: user?.role ?? null, enableTools: toolsEnabled, 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 = {}; if (isApprovalContinue) { const pendingSlugs = new Set(); 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); const thinkingEnabled = bodyEnableThinking !== undefined ? bodyEnableThinking : session.enableThinking === 1; const modelSupportsTools = (model.supportsTools ?? 1) === 1; const effectiveTools = modelSupportsTools ? tools : {}; logger.info( "[%s] [AGENT-CHAT] userId=%s sessionId=%s modelId=%d thinking=%s tools=%s supportsTools=%s edit=%s regen=%s approval=%s", event.context.requestId ?? "-", user?.id ?? "temp", sessionId, model.id, thinkingEnabled ? "on" : "off", Object.keys(effectiveTools).length > 0 ? "on" : "off", modelSupportsTools ? "yes" : "no", 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(effectiveTools).length > 0 ? { tools: effectiveTools, stopWhen: stepCountIs(8) } : {}), ...(thinkingEnabled ? { 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> = []; for (const step of steps) { const toolResultMap = new Map(step.toolResults.map((tr) => [tr.toolCallId, tr.output])); const contentArr = step.content as Array>; const approvalToolCallIds = new Set(); for (const cp of contentArr) { if (cp.type === "tool-approval-request" && (cp as any).approvalId && (cp as any).toolCall) { approvalToolCallIds.add((cp as any).toolCall.toolCallId); } } for (const cp of contentArr) { if (cp.type === "reasoning" && (cp as any).text) { parts.push({ id: `p_${parts.length}`, type: "reasoning", text: (cp as any).text }); } else if (cp.type === "text" && (cp as any).text) { parts.push({ id: `p_${parts.length}`, type: "text", text: (cp as any).text }); } else if (cp.type === "tool-approval-request" && (cp as any).approvalId && (cp as any).toolCall) { const tc = (cp as any).toolCall; const toolResult = toolResultMap.get(tc.toolCallId); if (toolResult === undefined) { parts.push({ id: `p_${parts.length}`, type: "tool-call", toolCallId: tc.toolCallId, toolName: tc.toolName, args: tc.input, state: "approval-requested", approvalId: (cp as any).approvalId, }); } else { parts.push({ id: `p_${parts.length}`, type: "tool-call", toolCallId: tc.toolCallId, toolName: tc.toolName, args: tc.input, result: toolResult, state: "result", approvalId: (cp as any).approvalId, }); } } else if (cp.type === "tool-call") { const tc = cp as any; if (approvalToolCallIds.has(tc.toolCallId)) continue; const toolResult = toolResultMap.get(tc.toolCallId); parts.push({ id: `p_${parts.length}`, type: "tool-call", toolCallId: tc.toolCallId, toolName: tc.toolName, args: tc.input, result: toolResult, state: "result", }); } } } if (isApprovalContinue) { const lastAssistant = historyMessages.filter((m) => m.role === "assistant").pop(); if (lastAssistant) { let existingParts: Array> = []; 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 = { "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; });