diff --git a/.codegraph/daemon.pid b/.codegraph/daemon.pid index f5b924b..2149e9e 100644 --- a/.codegraph/daemon.pid +++ b/.codegraph/daemon.pid @@ -1,6 +1,6 @@ { - "pid": 385650, + "pid": 6400, "version": "0.9.7", "socketPath": "/home/dash/code/nuxt-app/.codegraph/daemon.sock", - "startedAt": 1786124270447 + "startedAt": 1786156990118 } diff --git a/app/composables/useAgentChat.ts b/app/composables/useAgentChat.ts index 5824418..07b3a0b 100644 --- a/app/composables/useAgentChat.ts +++ b/app/composables/useAgentChat.ts @@ -99,7 +99,12 @@ export function useAgentChat(options: UseAgentChatOptions) { async function tryResumeStream(sid: string) { const lastMsg = messages.value[messages.value.length - 1]; - if (!lastMsg || lastMsg.role !== "user") return; + if (!lastMsg) return; + + const isApprovalResume = lastMsg.role === "assistant" && + lastMsg.parts?.some((p) => p.state === "approval-responded"); + const isNormalResume = lastMsg.role === "user"; + if (!isApprovalResume && !isNormalResume) return; let res: Response; try { @@ -118,18 +123,36 @@ export function useAgentChat(options: UseAgentChatOptions) { const modelIdHeader = res.headers.get("X-Stream-Model-Id"); const userMessageIdHeader = res.headers.get("X-Stream-User-Message-Id"); - const assistantMsg: AgentMessage = { - id: generateId(), - role: "assistant", - content: "", - parts: [], - modelId: modelIdHeader ? Number(modelIdHeader) : null, - }; - messages.value.push(assistantMsg); - const assistantIdx = messages.value.length - 1; + let assistantIdx: number; + let assistantMsg: AgentMessage; - if (userMessageIdHeader && lastMsg) { - lastMsg.id = userMessageIdHeader; + if (isApprovalResume) { + assistantMsg = lastMsg; + assistantIdx = messages.value.length - 1; + if (modelIdHeader) { + assistantMsg.modelId = Number(modelIdHeader); + } + staleApprovalToolCallIds = new Set( + messages.value.flatMap((m) => + (m.parts ?? []) + .filter((p) => p.state === "approval-requested" && !p.isAutomaticApproval) + .map((p) => p.toolCallId ?? ""), + ), + ); + } else { + assistantMsg = { + id: generateId(), + role: "assistant", + content: "", + parts: [], + modelId: modelIdHeader ? Number(modelIdHeader) : null, + }; + messages.value.push(assistantMsg); + assistantIdx = messages.value.length - 1; + + if (userMessageIdHeader && lastMsg) { + lastMsg.id = userMessageIdHeader; + } } isLoading.value = true; @@ -137,7 +160,7 @@ export function useAgentChat(options: UseAgentChatOptions) { try { await processStream(res, assistantIdx); - validateAssistantContent(assistantIdx); + validateAssistantContent(assistantIdx, isApprovalResume); onStreamComplete?.(); } catch (err: any) { if (err.name !== "AbortError") { diff --git a/packages/drizzle-pkg/db.sqlite b/packages/drizzle-pkg/db.sqlite index 31f118b..a21cabd 100644 Binary files a/packages/drizzle-pkg/db.sqlite and b/packages/drizzle-pkg/db.sqlite differ diff --git a/server/api/agent/chat/index.post.ts b/server/api/agent/chat/index.post.ts index 8032432..56d74de 100644 --- a/server/api/agent/chat/index.post.ts +++ b/server/api/agent/chat/index.post.ts @@ -414,6 +414,21 @@ export default defineEventHandler(async (event) => { userMessageId: userMessage?.id ?? null, }); + let approvalPrefixChunk: Uint8Array | null = null; + if (isApprovalContinue && Object.keys(approvalToolResults).length > 0) { + const prefixChunks = Object.entries(approvalToolResults).map( + ([toolCallId, output]) => + `data: ${JSON.stringify({ + type: "tool-output-available", + toolCallId, + output, + })}\n\n`, + ); + approvalPrefixChunk = new TextEncoder().encode(prefixChunks.join("")); + appendChunk(sessionId, approvalPrefixChunk); + logger.info("[%s] [APPROVAL-PREFIX] written %d bytes to buffer before stream start", event.context.requestId ?? "-", approvalPrefixChunk.length); + } + const result = streamText({ model: languageModel, system: systemPrompt || undefined, @@ -602,19 +617,10 @@ export default defineEventHandler(async (event) => { }); } - if (continueAfterApproval && Object.keys(approvalToolResults).length > 0) { + if (continueAfterApproval && approvalPrefixChunk && 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("")); - appendChunk(sessionId, prefix); + const prefix = approvalPrefixChunk; const transformed = new ReadableStream({ async start(controller) { diff --git a/server/service/agent/stream-buffer.ts b/server/service/agent/stream-buffer.ts index 0b20479..cb606b6 100644 --- a/server/service/agent/stream-buffer.ts +++ b/server/service/agent/stream-buffer.ts @@ -43,6 +43,13 @@ export function createStreamBuffer( sessionId: string, meta: { modelId?: number | null; userMessageId?: string | null }, ): StreamBuffer { + const existing = buffers.get(sessionId); + if (existing && !existing.done) { + existing.modelId = meta.modelId ?? existing.modelId; + existing.userMessageId = meta.userMessageId ?? existing.userMessageId; + logger.info("[STREAM-BUFFER] reused existing sessionId=%s chunks=%d", sessionId, existing.chunks.length); + return existing; + } const buf: StreamBuffer = { chunks: [], done: false,