From abb7cc80e2aca1666023ab0aa7f5772529d689f7 Mon Sep 17 00:00:00 2001 From: Alaister Young <10985857+alaister@users.noreply.github.com> Date: Mon, 11 May 2026 15:10:32 +0800 Subject: [PATCH] fix(studio): teach toWebHandler to stream when res.writeHead is called MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The pages-router AI route uses result.pipeUIMessageStreamToResponse (res, ...) from the AI SDK, which writes to a Node ServerResponse in three stages: res.writeHead(status, headers), then async res.write(chunk) calls, then res.end(). The buffering shim only exposed status/setHeader/write/end and crashed at writeHead. Switch the response proxy into streaming mode the moment writeHead (or flushHeaders) is called: a Web ReadableStream is opened, any already-buffered chunks flush into it, subsequent writes enqueue into it, and res.end() closes it. finalize() then returns new Response(stream, …) instead of the buffered payload, so the client starts receiving tokens before the handler resolves. Also stub the writable/EventEmitter surface (on/once/off/emit/ headersSent/writableEnded) as no-ops so pipe helpers that attach drain/close/error listeners don't crash. Abort is still handled at the req level via the Request.signal subscription. --- apps/studio/compat/next/api.ts | 113 +++++++++++++++++++++++++++++++-- 1 file changed, 109 insertions(+), 4 deletions(-) diff --git a/apps/studio/compat/next/api.ts b/apps/studio/compat/next/api.ts index 57516c440de..13a75d368af 100644 --- a/apps/studio/compat/next/api.ts +++ b/apps/studio/compat/next/api.ts @@ -176,6 +176,34 @@ function buildResponse() { return encoder.encode(String(data)) } + // Streaming mode is entered as soon as the handler calls + // `res.writeHead(...)` — that's the Node `ServerResponse` signal that + // headers are committed and chunks should flush as they're written + // (matches what `result.pipeUIMessageStreamToResponse(res, …)` from + // the AI SDK does for token-by-token streaming). + let streamMode = false + let streamController: ReadableStreamDefaultController | null = null + let streamClosed = false + // Backpressure for upstream writers: drain when the controller's + // desired-size goes negative. We don't get a true "drain" event + // from a Web ReadableStream, but most handlers just check the + // return value of `res.write()`. + const stream = new ReadableStream({ + start(controller) { + streamController = controller + }, + cancel() { + streamClosed = true + }, + }) + const enterStreamMode = () => { + if (streamMode) return + streamMode = true + // Flush anything that was buffered before the streaming switch. + for (const c of chunks) streamController?.enqueue(c) + chunks.length = 0 + } + const res = {} as NextApiResponse & Record res.statusCode = 200 @@ -206,11 +234,46 @@ function buildResponse() { res.removeHeader = (name: string) => { responseHeaders.delete(name.toLowerCase()) } + // Node `ServerResponse.writeHead(statusCode, statusMessage?, headers?)`. + // We honour `statusCode` and `headers`; `statusMessage` is dropped + // (Fetch `Response` doesn't preserve a custom HTTP/1 reason phrase + // when running through TanStack's runtime). + ;(res as unknown as { writeHead: (...args: unknown[]) => unknown }).writeHead = ( + code: number, + headersOrMessage?: string | Record, + maybeHeaders?: Record + ) => { + statusCode = code + res.statusCode = code + const headers = + typeof headersOrMessage === 'object' && headersOrMessage !== null + ? headersOrMessage + : maybeHeaders + if (headers) { + for (const [name, value] of Object.entries(headers)) { + const key = name.toLowerCase() + if (Array.isArray(value)) { + responseHeaders.delete(key) + for (const v of value) responseHeaders.append(key, String(v)) + } else if (value !== undefined && value !== null) { + responseHeaders.set(key, String(value)) + } + } + } + enterStreamMode() + return res + } + ;(res as unknown as { flushHeaders: () => unknown }).flushHeaders = () => { + enterStreamMode() + return res + } res.json = (data: unknown) => { if (!responseHeaders.has('content-type')) { responseHeaders.set('content-type', 'application/json') } - chunks.push(encode(JSON.stringify(data))) + const payload = encode(JSON.stringify(data)) + if (streamMode) streamController?.enqueue(payload) + else chunks.push(payload) return res } res.send = (data: unknown) => { @@ -222,15 +285,31 @@ function buildResponse() { ) { return res.json(data) } - chunks.push(encode(data)) + const payload = encode(data) + if (streamMode) streamController?.enqueue(payload) + else chunks.push(payload) return res } res.write = (chunk: unknown) => { - chunks.push(encode(chunk)) + const payload = encode(chunk) + if (streamMode) streamController?.enqueue(payload) + else chunks.push(payload) return true } res.end = (data?: unknown) => { - if (data !== undefined) chunks.push(encode(data)) + if (streamMode) { + if (data !== undefined && data !== null) streamController?.enqueue(encode(data)) + if (!streamClosed) { + try { + streamController?.close() + } catch { + // Already closed (or never opened): nothing to do. + } + streamClosed = true + } + } else if (data !== undefined) { + chunks.push(encode(data)) + } return res } res.redirect = (...args: unknown[]) => { @@ -242,7 +321,33 @@ function buildResponse() { return res } + // Node `ServerResponse` is a Writable EventEmitter. AI SDK / other + // pipe-style helpers may attach listeners (drain/close/error) for + // backpressure or cleanup. Accept and no-op them so they don't + // crash; abort is already handled at the `req` level via the + // Request.signal subscription. + const noopEE = (...args: unknown[]) => { + void args + return res + } + ;(res as unknown as Record).on = noopEE + ;(res as unknown as Record).once = noopEE + ;(res as unknown as Record).off = noopEE + ;(res as unknown as Record).removeListener = noopEE + ;(res as unknown as Record).removeAllListeners = noopEE + ;(res as unknown as Record).emit = () => false + ;(res as unknown as Record).addListener = noopEE + ;(res as unknown as Record).prependListener = noopEE + ;(res as unknown as { headersSent: boolean }).headersSent = false + ;(res as unknown as { writableEnded: boolean }).writableEnded = false + const finalize = (): Response => { + if (streamMode) { + // The handler entered streaming mode (writeHead/flushHeaders) and + // may still be pushing chunks asynchronously after this return — + // the stream stays open until res.end() runs. + return new Response(stream, { status: statusCode, headers: responseHeaders }) + } const totalLen = chunks.reduce((n, c) => n + c.length, 0) const payload = new Uint8Array(totalLen) let offset = 0