Browse Source

feat: update stream handling and buffer management for approval responses

feat/ai-sdk-v6-upgrade
npmrun 1 day ago
parent
commit
5d5af1f57b
  1. 4
      .codegraph/daemon.pid
  2. 49
      app/composables/useAgentChat.ts
  3. BIN
      packages/drizzle-pkg/db.sqlite
  4. 28
      server/api/agent/chat/index.post.ts
  5. 7
      server/service/agent/stream-buffer.ts

4
.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
}

49
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") {

BIN
packages/drizzle-pkg/db.sqlite

Binary file not shown.

28
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) {

7
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,

Loading…
Cancel
Save