From bfb03bc029f029fd0c9a48c6e1cb40b9467799f1 Mon Sep 17 00:00:00 2001 From: carl chen Date: Mon, 3 Aug 2026 12:54:49 +0800 Subject: [PATCH] fix(x-sdk): release XStream reader on early exit --- .../src/x-stream/__tests__/index.test.ts | 81 +++++++++++++++++++ packages/x-sdk/src/x-stream/index.ts | 28 +++++-- 2 files changed, 103 insertions(+), 6 deletions(-) create mode 100644 packages/x-sdk/src/x-stream/__tests__/index.test.ts diff --git a/packages/x-sdk/src/x-stream/__tests__/index.test.ts b/packages/x-sdk/src/x-stream/__tests__/index.test.ts new file mode 100644 index 0000000..9e5bdfd --- /dev/null +++ b/packages/x-sdk/src/x-stream/__tests__/index.test.ts @@ -0,0 +1,81 @@ +import { describe, expect, it, vi } from "vite-plus/test"; + +import XStream from "../index"; + +const encoder = new TextEncoder(); + +describe("XStream", () => { + it("cancels the source and releases the reader lock on early exit", async () => { + const cancel = vi.fn(); + const stream = XStream({ + readableStream: new ReadableStream({ + start(controller) { + controller.enqueue(encoder.encode("data: first\n\n")); + }, + cancel, + }), + }); + + for await (const value of stream) { + expect(value).toEqual({ data: "first" }); + break; + } + + expect(stream.locked).toBe(false); + await vi.waitFor(() => expect(cancel).toHaveBeenCalledOnce()); + + const reader = stream.getReader(); + await expect(reader.read()).resolves.toEqual({ + done: true, + value: undefined, + }); + reader.releaseLock(); + }); + + it("preserves the consumer error when cancellation fails", async () => { + const consumerError = new Error("consumer failed"); + const cancelError = new Error("cancel failed"); + const cancel = vi.fn(() => Promise.reject(cancelError)); + const stream = XStream({ + readableStream: new ReadableStream({ + start(controller) { + controller.enqueue(encoder.encode("data: first\n\n")); + }, + cancel, + }), + }); + + const consume = async () => { + for await (const _value of stream) { + throw consumerError; + } + }; + + await expect(consume()).rejects.toBe(consumerError); + expect(cancel).toHaveBeenCalledOnce(); + expect(stream.locked).toBe(false); + }); + + it("releases the reader lock without cancellation after normal completion", async () => { + const cancel = vi.fn(); + const stream = XStream({ + readableStream: new ReadableStream({ + start(controller) { + controller.enqueue(encoder.encode("data: first\n\n")); + controller.enqueue(encoder.encode("data: second\n\n")); + controller.close(); + }, + cancel, + }), + }); + const values = []; + + for await (const value of stream) { + values.push(value); + } + + expect(values).toEqual([{ data: "first" }, { data: "second" }]); + expect(cancel).not.toHaveBeenCalled(); + expect(stream.locked).toBe(false); + }); +}); diff --git a/packages/x-sdk/src/x-stream/index.ts b/packages/x-sdk/src/x-stream/index.ts index 21130e3..f11274d 100644 --- a/packages/x-sdk/src/x-stream/index.ts +++ b/packages/x-sdk/src/x-stream/index.ts @@ -213,16 +213,32 @@ function XStream(options: XStreamOptions) { // @ts-ignore ReadableStream async iterator signature conflicts with AsyncGenerator in this cross-runtime typing. stream[Symbol.asyncIterator] = async function* () { const reader = this.getReader(); + let completed = false; - while (true) { - const { done, value } = await reader.read(); + try { + while (true) { + const { done, value } = await reader.read(); - if (done) break; + if (done) { + completed = true; + break; + } + + if (!value) continue; - if (!value) continue; + // Transformed data through all transform pipes + yield value; + } + } finally { + if (!completed) { + try { + await reader.cancel(); + } catch { + // Cancellation must not replace the consumer's original error. + } + } - // Transformed data through all transform pipes - yield value; + reader.releaseLock(); } };