From 37372f39dde3d1674f181f0aabac19b4bdbedd2f Mon Sep 17 00:00:00 2001 From: BigSimmo <87357024+BigSimmo@users.noreply.github.com> Date: Sun, 16 Aug 2026 19:54:01 +0800 Subject: [PATCH 1/2] fix(api): avoid universal NDJSON stream errors and unify stream error propagation --- src/app/api/search/universal/route.ts | 9 +++++++-- src/lib/universal-search-stream.ts | 9 +++++++-- tests/universal-search-stream.test.ts | 7 +++++++ tests/universal-search.test.ts | 18 ++++++++++++++++++ 4 files changed, 39 insertions(+), 4 deletions(-) diff --git a/src/app/api/search/universal/route.ts b/src/app/api/search/universal/route.ts index 32d4d1702c..ed17b2a63d 100644 --- a/src/app/api/search/universal/route.ts +++ b/src/app/api/search/universal/route.ts @@ -95,7 +95,7 @@ function universalStreamResponse( enqueue({ type: "complete", response: { ...response, ...decoration } }); controller.close(); }) - .catch((error: unknown) => { + .catch(() => { if (searchController.signal.aborted) { try { controller.close(); @@ -104,7 +104,12 @@ function universalStreamResponse( } return; } - controller.error(error); + // Headers have already been returned, so this failure cannot become a + // JSON 500. End the NDJSON protocol with a fixed, non-sensitive error + // event instead of erroring the response stream (which Next reports as + // an unhandled server request error and can leave consumers waiting). + enqueue({ type: "error", code: "universal_search_failed" }); + controller.close(); }) .finally(() => request.signal.removeEventListener("abort", abortFromRequest)); }, diff --git a/src/lib/universal-search-stream.ts b/src/lib/universal-search-stream.ts index c3ba49156a..c96d446324 100644 --- a/src/lib/universal-search-stream.ts +++ b/src/lib/universal-search-stream.ts @@ -7,7 +7,8 @@ export type UniversalSearchStreamResponse = UniversalSearchResponse & { export type UniversalSearchStreamEvent = | { type: "group"; query: string; group: UniversalSearchGroup } - | { type: "complete"; response: UniversalSearchStreamResponse }; + | { type: "complete"; response: UniversalSearchStreamResponse } + | { type: "error"; code: "universal_search_failed" }; function abortReason(signal: AbortSignal): Error { return signal.reason instanceof Error ? signal.reason : new DOMException("The operation was aborted.", "AbortError"); @@ -25,6 +26,9 @@ function parseEvent(line: string): UniversalSearchStreamEvent { if (parsed.type === "complete" && parsed.response) { return parsed as Extract; } + if (parsed.type === "error" && parsed.code === "universal_search_failed") { + return parsed as Extract; + } throw new Error("Invalid universal-search NDJSON event."); } @@ -52,7 +56,8 @@ export async function consumeUniversalSearchNdjson( if (!line.trim()) return; const event = parseEvent(line); if (event.type === "group") await options.onGroup?.(event.group, event.query); - else complete = event.response; + else if (event.type === "complete") complete = event.response; + else throw new Error("Universal search failed."); }; try { diff --git a/tests/universal-search-stream.test.ts b/tests/universal-search-stream.test.ts index 338be2f889..a34bd9469d 100644 --- a/tests/universal-search-stream.test.ts +++ b/tests/universal-search-stream.test.ts @@ -100,4 +100,11 @@ describe("consumeUniversalSearchNdjson", () => { await expect(consumeUniversalSearchNdjson(response)).rejects.toThrow("complete event"); }); + + it("rejects a redacted server error event without requiring a complete event", async () => { + const { consumeUniversalSearchNdjson } = await import("../src/lib/universal-search-stream"); + const response = responseFromChunks([`${JSON.stringify({ type: "error", code: "universal_search_failed" })}\n`]); + + await expect(consumeUniversalSearchNdjson(response)).rejects.toThrow("Universal search failed."); + }); }); diff --git a/tests/universal-search.test.ts b/tests/universal-search.test.ts index e2f1a79ab5..6d96e0c5a3 100644 --- a/tests/universal-search.test.ts +++ b/tests/universal-search.test.ts @@ -789,4 +789,22 @@ describe("GET /api/search/universal (live public/owner path)", () => { expect(args.demo).toBe(false); expect(args.ownerId).toBe(userId); }); + + it("ends a rejected NDJSON search with a redacted error event", async () => { + const client = createSupabaseMock(); + const runUniversalSearch = createRunMock(); + runUniversalSearch.mockRejectedValue(new Error("private dependency failure")); + mockRuntime(client, runUniversalSearch); + const { GET } = await import("../src/app/api/search/universal/route"); + + const response = await GET(new Request("http://localhost/api/search/universal?q=clozapine&stream=ndjson")); + const events = (await response.text()) + .trim() + .split("\n") + .map((line) => JSON.parse(line) as Record); + + expect(response.status).toBe(200); + expect(events).toEqual([{ type: "error", code: "universal_search_failed" }]); + expect(JSON.stringify(events)).not.toContain("private dependency failure"); + }); }); From 3028cfaa5e2556dfb41a62cc4645b9c6f9de27d8 Mon Sep 17 00:00:00 2001 From: BigSimmo <87357024+BigSimmo@users.noreply.github.com> Date: Sun, 16 Aug 2026 20:14:32 +0800 Subject: [PATCH 2/2] fix(search): cancel open NDJSON body on stream errors --- ...4a753cf9a1522a6298f9286d8835adfd.record.md | 1 + src/lib/universal-search-stream.ts | 6 +++++- tests/universal-search-stream.test.ts | 19 +++++++++++++++++-- 3 files changed, 23 insertions(+), 3 deletions(-) create mode 100644 docs/branch-review-records/2f102c68fe61b45d4dbfd9b43a1ca1d14a753cf9a1522a6298f9286d8835adfd.record.md diff --git a/docs/branch-review-records/2f102c68fe61b45d4dbfd9b43a1ca1d14a753cf9a1522a6298f9286d8835adfd.record.md b/docs/branch-review-records/2f102c68fe61b45d4dbfd9b43a1ca1d14a753cf9a1522a6298f9286d8835adfd.record.md new file mode 100644 index 0000000000..e3d39922ce --- /dev/null +++ b/docs/branch-review-records/2f102c68fe61b45d4dbfd9b43a1ca1d14a753cf9a1522a6298f9286d8835adfd.record.md @@ -0,0 +1 @@ +| 2026-08-16 | codex/chat-search-error-342-search-error-342 | 019281ae9c8a6883737691438feba82117335e14 | PR #2002 universal NDJSON stream error propagation review | Fixed reproducible P2: structured error rejection now cancels an open response body so upstream stream work is released; no P0/P1 findings | Manual adversarial pass; consumer Web Streams harness 4/4; server stream harness 3/3; pre-fix exact-head unit coverage, static checks, SAST and secret scan passed | diff --git a/src/lib/universal-search-stream.ts b/src/lib/universal-search-stream.ts index c96d446324..33af25025f 100644 --- a/src/lib/universal-search-stream.ts +++ b/src/lib/universal-search-stream.ts @@ -57,7 +57,11 @@ export async function consumeUniversalSearchNdjson( const event = parseEvent(line); if (event.type === "group") await options.onGroup?.(event.group, event.query); else if (event.type === "complete") complete = event.response; - else throw new Error("Universal search failed."); + else { + const error = new Error("Universal search failed."); + void reader.cancel(error).catch(() => undefined); + throw error; + } }; try { diff --git a/tests/universal-search-stream.test.ts b/tests/universal-search-stream.test.ts index a34bd9469d..197e40ac91 100644 --- a/tests/universal-search-stream.test.ts +++ b/tests/universal-search-stream.test.ts @@ -101,10 +101,25 @@ describe("consumeUniversalSearchNdjson", () => { await expect(consumeUniversalSearchNdjson(response)).rejects.toThrow("complete event"); }); - it("rejects a redacted server error event without requiring a complete event", async () => { + it("rejects a redacted server error and cancels the open response body", async () => { const { consumeUniversalSearchNdjson } = await import("../src/lib/universal-search-stream"); - const response = responseFromChunks([`${JSON.stringify({ type: "error", code: "universal_search_failed" })}\n`]); + const cancelled = vi.fn(); + const response = new Response( + new ReadableStream({ + start(controller) { + controller.enqueue( + encoder.encode(`${JSON.stringify({ type: "error", code: "universal_search_failed" })}\n`), + ); + }, + cancel(reason) { + cancelled(reason); + }, + }), + { headers: { "Content-Type": "application/x-ndjson; charset=utf-8" } }, + ); await expect(consumeUniversalSearchNdjson(response)).rejects.toThrow("Universal search failed."); + expect(cancelled).toHaveBeenCalledOnce(); + expect(cancelled).toHaveBeenCalledWith(expect.objectContaining({ message: "Universal search failed." })); }); });