71 lines
2.8 KiB
TypeScript
71 lines
2.8 KiB
TypeScript
import { createUIMessageStream, createUIMessageStreamResponse } from "ai";
|
|
import { NextRequest, NextResponse } from "next/server";
|
|
import { backendFetch, readBackendError } from "@/lib/backend";
|
|
import { messagesForBackend } from "@/lib/cad-messages";
|
|
import { backendEventToUiChunk } from "@/lib/cad-stream";
|
|
import type { CadUIMessage } from "@/lib/cad-types";
|
|
|
|
export const runtime = "nodejs";
|
|
|
|
type BackendEvent = { event: string; data: Record<string, unknown> };
|
|
|
|
async function* parseSse(response: Response): AsyncGenerator<BackendEvent> {
|
|
if (!response.body) return;
|
|
const reader = response.body.getReader();
|
|
const decoder = new TextDecoder();
|
|
let buffer = "";
|
|
try {
|
|
for (;;) {
|
|
const { value, done } = await reader.read();
|
|
buffer += decoder.decode(value || new Uint8Array(), { stream: !done });
|
|
const blocks = buffer.split(/\r?\n\r?\n/);
|
|
buffer = done ? "" : blocks.pop() || "";
|
|
for (const block of blocks) {
|
|
const eventName = /^event:\s*(.+)$/m.exec(block)?.[1]?.trim() || "message";
|
|
const dataText = /^data:\s*(.+)$/m.exec(block)?.[1]?.trim() || "{}";
|
|
try {
|
|
yield { event: eventName, data: JSON.parse(dataText) as Record<string, unknown> };
|
|
} catch {
|
|
yield { event: "cad_error", data: { stage: "stream", message: "Invalid backend event." } };
|
|
}
|
|
}
|
|
if (done) break;
|
|
}
|
|
} finally {
|
|
reader.releaseLock();
|
|
}
|
|
}
|
|
|
|
export async function POST(request: NextRequest) {
|
|
const body = await request.json().catch(() => ({}));
|
|
const upstream = await backendFetch("/v1/chat/stream", {
|
|
method: "POST",
|
|
headers: { "Content-Type": "application/json", Accept: "text/event-stream" },
|
|
body: JSON.stringify({
|
|
conversation_id: body.conversationId || null,
|
|
selected_task_id: body.selectedTaskId || null,
|
|
provider_id: body.providerId || null,
|
|
model_id: body.modelId || null,
|
|
messages: messagesForBackend((Array.isArray(body.messages) ? body.messages : []) as CadUIMessage[]),
|
|
}),
|
|
signal: request.signal,
|
|
});
|
|
if (!upstream.ok) {
|
|
return NextResponse.json({ error: await readBackendError(upstream) }, { status: upstream.status });
|
|
}
|
|
const stream = createUIMessageStream({
|
|
execute: async ({ writer }) => {
|
|
const textId = `assistant_${Date.now()}`;
|
|
writer.write({ type: "start", messageId: textId });
|
|
writer.write({ type: "text-start", id: textId });
|
|
for await (const item of parseSse(upstream)) {
|
|
const chunk = backendEventToUiChunk(item, textId);
|
|
if (chunk) writer.write(chunk);
|
|
}
|
|
writer.write({ type: "text-end", id: textId });
|
|
writer.write({ type: "finish", finishReason: "stop" });
|
|
},
|
|
});
|
|
return createUIMessageStreamResponse({ stream, headers: { "Cache-Control": "no-store" } });
|
|
}
|