diff --git a/.changeset/clean-up-web-chat-streams.md b/.changeset/clean-up-web-chat-streams.md new file mode 100644 index 0000000000..8f10370e28 --- /dev/null +++ b/.changeset/clean-up-web-chat-streams.md @@ -0,0 +1,5 @@ +--- +"eve": patch +--- + +Generated Web Chat apps now include their new-session and resumable-session routes. Abandoned browser session streams also release their local workflow listeners instead of accumulating them across navigation. diff --git a/apps/docs/registry.json b/apps/docs/registry.json index 9bd080207c..c70a4dc440 100644 --- a/apps/docs/registry.json +++ b/apps/docs/registry.json @@ -1491,6 +1491,16 @@ "type": "registry:file", "target": "app/page.tsx" }, + { + "path": "registry/channel/web/app/s/[sessionId]/page.tsx", + "type": "registry:file", + "target": "app/s/[sessionId]/page.tsx" + }, + { + "path": "registry/channel/web/app/s/page.tsx", + "type": "registry:file", + "target": "app/s/page.tsx" + }, { "path": "registry/channel/web/components.json", "type": "registry:file", diff --git a/packages/eve/src/public/channels/eve.test.ts b/packages/eve/src/public/channels/eve.test.ts index 9616a07c5c..eb956c9afa 100644 --- a/packages/eve/src/public/channels/eve.test.ts +++ b/packages/eve/src/public/channels/eve.test.ts @@ -344,13 +344,13 @@ function createEveStreamHandler(input: EveChannelInput) { return { getEventStream, getStreamTailIndex, - async fetch(url: string) { + async fetch(url: string, init?: RequestInit) { const args: RouteHandlerArgs = { ...createRouteArgs(), attachSession: () => session, params: { sessionId: "test-session-id" }, }; - return (streamRoute as any).handler(new Request(url), args); + return (streamRoute as any).handler(new Request(url, init), args); }, }; } @@ -533,6 +533,26 @@ describe("eveChannel — stream cursor", () => { expect(new TextDecoder().decode(firstChunk.value)).toBe("\n"); }); + it("cancels the durable stream when the request aborts", async () => { + const handler = createEveStreamHandler({ auth: none() }); + const cancelled = vi.fn(); + handler.getEventStream.mockResolvedValueOnce( + new ReadableStream({ + cancel: cancelled, + }), + ); + const abort = new AbortController(); + + const response = await handler.fetch("https://eve.test/eve/v1/session/test-session-id/stream", { + signal: abort.signal, + }); + const reader = response.body!.getReader(); + await reader.read(); + abort.abort(); + + await vi.waitFor(() => expect(cancelled).toHaveBeenCalledOnce()); + }); + it("forwards negative tail-relative start indices", async () => { const handler = createEveStreamHandler({ auth: none() }); diff --git a/packages/eve/src/public/channels/eve.ts b/packages/eve/src/public/channels/eve.ts index ce38d0c697..29bda462da 100644 --- a/packages/eve/src/public/channels/eve.ts +++ b/packages/eve/src/public/channels/eve.ts @@ -1068,7 +1068,7 @@ async function createSessionStreamResponse(request: Request, session: Session): if (tailIndex !== undefined) { headers.set(EVE_STREAM_TAIL_INDEX_HEADER, String(tailIndex)); } - return new Response(serializeAsNdjson(events), { + return new Response(serializeAsNdjson(events, request.signal), { headers, }); } catch { @@ -1352,16 +1352,19 @@ function parseStartIndex(request: Request): number | undefined | Response { return parsed; } -function serializeAsNdjson(events: ReadableStream): ReadableStream { +function serializeAsNdjson( + events: ReadableStream, + signal: AbortSignal, +): ReadableStream { const encoder = new TextEncoder(); - return events.pipeThrough( - new TransformStream({ - start(controller) { - controller.enqueue(encoder.encode("\n")); - }, - transform(event, controller) { - controller.enqueue(encoder.encode(`${JSON.stringify(event)}\n`)); - }, - }), - ); + const transform = new TransformStream({ + start(controller) { + controller.enqueue(encoder.encode("\n")); + }, + transform(event, controller) { + controller.enqueue(encoder.encode(`${JSON.stringify(event)}\n`)); + }, + }); + void events.pipeTo(transform.writable, { signal }).catch(() => {}); + return transform.readable; } diff --git a/packages/eve/src/setup/scaffold/index.integration.test.ts b/packages/eve/src/setup/scaffold/index.integration.test.ts index 0c889e2014..7fa1e01fa4 100644 --- a/packages/eve/src/setup/scaffold/index.integration.test.ts +++ b/packages/eve/src/setup/scaffold/index.integration.test.ts @@ -228,6 +228,12 @@ describe("ensureChannel", () => { await expect(readFile(join(projectRoot, "app/page.tsx"), "utf8")).resolves.toContain( "AgentChat", ); + await expect(readFile(join(projectRoot, "app/s/page.tsx"), "utf8")).resolves.toContain( + "", + ); + await expect( + readFile(join(projectRoot, "app/s/[sessionId]/page.tsx"), "utf8"), + ).resolves.toContain(""); await expect( readFile(join(projectRoot, "agent/tools/randomize.ts"), "utf8"), ).rejects.toMatchObject({