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/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..33af25025f 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,12 @@ 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 { + 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 338be2f889..f4777f19ba 100644 --- a/tests/universal-search-stream.test.ts +++ b/tests/universal-search-stream.test.ts @@ -100,4 +100,24 @@ describe("consumeUniversalSearchNdjson", () => { await expect(consumeUniversalSearchNdjson(response)).rejects.toThrow("complete event"); }); + + it("rejects a redacted server error and cancels the open response body", async () => { + const { consumeUniversalSearchNdjson } = await import("../src/lib/universal-search-stream"); + 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." })); + }); }); 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"); + }); });