diff --git a/package-lock.json b/package-lock.json index 2d1a9ee..813ede6 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1105,9 +1105,9 @@ "license": "MIT" }, "node_modules/@eslint/config-array/node_modules/brace-expansion": { - "version": "1.1.15", - "resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.15.tgz", - "integrity": "sha512-EwOCDEex4quD37XhqM3omwtMoJjr//isUZz1JopUNWms+4Z2ViyM/k1YIRePpoVNnQhENnxtFjLaxNHrT7xIUg==", + "version": "1.1.16", + "resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.16.tgz", + "integrity": "sha512-IDw48K2/2kRkg9LdJxurvq3lV3aBgq0REY89duEqFRthjlPdXHKMj7EnQOXVckxzgisinf3nHfrcE2FufFLXMw==", "dev": true, "license": "MIT", "dependencies": { @@ -1186,9 +1186,9 @@ "license": "MIT" }, "node_modules/@eslint/eslintrc/node_modules/brace-expansion": { - "version": "1.1.15", - "resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.15.tgz", - "integrity": "sha512-EwOCDEex4quD37XhqM3omwtMoJjr//isUZz1JopUNWms+4Z2ViyM/k1YIRePpoVNnQhENnxtFjLaxNHrT7xIUg==", + "version": "1.1.16", + "resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.16.tgz", + "integrity": "sha512-IDw48K2/2kRkg9LdJxurvq3lV3aBgq0REY89duEqFRthjlPdXHKMj7EnQOXVckxzgisinf3nHfrcE2FufFLXMw==", "dev": true, "license": "MIT", "dependencies": { @@ -3435,9 +3435,9 @@ } }, "node_modules/body-parser": { - "version": "1.20.5", - "resolved": "https://registry.npmjs.org/body-parser/-/body-parser-1.20.5.tgz", - "integrity": "sha512-3grm+/2tUOvu2cjJkvsIxrv/wVpfXQW4PsQHYm7yk4vfpu7Ekl6nEsYBoJUL6qDwZUx8wUhQ8tR2qz+ad9c9OA==", + "version": "1.20.6", + "resolved": "https://registry.npmjs.org/body-parser/-/body-parser-1.20.6.tgz", + "integrity": "sha512-p5tAzS57i5MV9fZFDj9LeIiTZEufbSe2eDozP+ElheSUq1m74CRq1jI4mYNDdVs9vQztXFLuk/Gd6BWTdwRJ5g==", "dev": true, "license": "MIT", "dependencies": { @@ -3503,9 +3503,9 @@ } }, "node_modules/brace-expansion": { - "version": "5.0.6", - "resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-5.0.6.tgz", - "integrity": "sha512-kLpxurY4Z4r9sgMsyG0Z9uzsBlgiU/EFKhj/h91/8yHu0edo7XuixOIH3VcJ8kkxs6/jPzoI6U9Vj3WqbMQ94g==", + "version": "5.0.7", + "resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-5.0.7.tgz", + "integrity": "sha512-7oFy703dxfY3/NLxC1fh2SUCQ0H9rmAY+5EpDVfXjUTTs+HEwR2nYaqLv+GWcTsumwxPfiz6CzCNkwXwBUwqCA==", "dev": true, "license": "MIT", "dependencies": { @@ -5292,9 +5292,9 @@ "peer": true }, "node_modules/eslint-plugin-import/node_modules/brace-expansion": { - "version": "1.1.15", - "resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.15.tgz", - "integrity": "sha512-EwOCDEex4quD37XhqM3omwtMoJjr//isUZz1JopUNWms+4Z2ViyM/k1YIRePpoVNnQhENnxtFjLaxNHrT7xIUg==", + "version": "1.1.16", + "resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.16.tgz", + "integrity": "sha512-IDw48K2/2kRkg9LdJxurvq3lV3aBgq0REY89duEqFRthjlPdXHKMj7EnQOXVckxzgisinf3nHfrcE2FufFLXMw==", "dev": true, "license": "MIT", "peer": true, @@ -5379,9 +5379,9 @@ "peer": true }, "node_modules/eslint-plugin-jsx-a11y/node_modules/brace-expansion": { - "version": "1.1.15", - "resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.15.tgz", - "integrity": "sha512-EwOCDEex4quD37XhqM3omwtMoJjr//isUZz1JopUNWms+4Z2ViyM/k1YIRePpoVNnQhENnxtFjLaxNHrT7xIUg==", + "version": "1.1.16", + "resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.16.tgz", + "integrity": "sha512-IDw48K2/2kRkg9LdJxurvq3lV3aBgq0REY89duEqFRthjlPdXHKMj7EnQOXVckxzgisinf3nHfrcE2FufFLXMw==", "dev": true, "license": "MIT", "peer": true, @@ -5468,9 +5468,9 @@ "peer": true }, "node_modules/eslint-plugin-react/node_modules/brace-expansion": { - "version": "1.1.15", - "resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.15.tgz", - "integrity": "sha512-EwOCDEex4quD37XhqM3omwtMoJjr//isUZz1JopUNWms+4Z2ViyM/k1YIRePpoVNnQhENnxtFjLaxNHrT7xIUg==", + "version": "1.1.16", + "resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.16.tgz", + "integrity": "sha512-IDw48K2/2kRkg9LdJxurvq3lV3aBgq0REY89duEqFRthjlPdXHKMj7EnQOXVckxzgisinf3nHfrcE2FufFLXMw==", "dev": true, "license": "MIT", "peer": true, @@ -5542,9 +5542,9 @@ "license": "MIT" }, "node_modules/eslint/node_modules/brace-expansion": { - "version": "1.1.15", - "resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.15.tgz", - "integrity": "sha512-EwOCDEex4quD37XhqM3omwtMoJjr//isUZz1JopUNWms+4Z2ViyM/k1YIRePpoVNnQhENnxtFjLaxNHrT7xIUg==", + "version": "1.1.16", + "resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.16.tgz", + "integrity": "sha512-IDw48K2/2kRkg9LdJxurvq3lV3aBgq0REY89duEqFRthjlPdXHKMj7EnQOXVckxzgisinf3nHfrcE2FufFLXMw==", "dev": true, "license": "MIT", "dependencies": { @@ -7129,9 +7129,9 @@ "license": "MIT" }, "node_modules/js-yaml": { - "version": "4.2.0", - "resolved": "https://registry.npmjs.org/js-yaml/-/js-yaml-4.2.0.tgz", - "integrity": "sha512-ePWsvanv0DWuDRsW8dnt+R4jQ31SCRCQ7hhNcPXZPsoBZiemuZNYGf7adZdqX2D86j6rvKp3RpCxVTSb8WQlOw==", + "version": "4.3.0", + "resolved": "https://registry.npmjs.org/js-yaml/-/js-yaml-4.3.0.tgz", + "integrity": "sha512-1td788aAnnZ5qs7V2QIRl1owjtYpbKt749Y3xauqQgwIIGF/xXWz1wMTEBx5O3LK3lXLVuqXPdPxj2BoFHaW9Q==", "dev": true, "funding": [ { @@ -7524,9 +7524,9 @@ } }, "node_modules/linkify-it": { - "version": "5.0.1", - "resolved": "https://registry.npmjs.org/linkify-it/-/linkify-it-5.0.1.tgz", - "integrity": "sha512-wVoTjP4Q6R0NW5hiZkVJaFZPWgtXfoGF+6LucL3/FtiNjmcHhYjEr5f1Kqjirc1nBW07J/ZuRFumqr2oqccEWg==", + "version": "5.0.2", + "resolved": "https://registry.npmjs.org/linkify-it/-/linkify-it-5.0.2.tgz", + "integrity": "sha512-ONTm2jCMAVZjgQa/Fy1kScXsuOoF5NPTsoFBdE1KVIZ2vAh/r9+Bqo+0jINCBYnavTPQZz38QzFTme79ENoN3Q==", "dev": true, "funding": [ { diff --git a/source/sse-session/async-push-iterator.ts b/source/sse-session/async-push-iterator.ts index 29736ea..949b3d9 100644 --- a/source/sse-session/async-push-iterator.ts +++ b/source/sse-session/async-push-iterator.ts @@ -23,24 +23,31 @@ */ export class AsyncPushIterator { /** ReadableStream backing the async iterator returned from {@link Symbol.asyncIterator}. */ - private readonly stream: ReadableStream; + #stream: ReadableStream; /** Controller used to enqueue values and close the stream from {@link push} and {@link close}. */ - private controller: ReadableStreamDefaultController | undefined; + #controller: ReadableStreamDefaultController | undefined; /** When true, no more values are accepted and iteration eventually completes. */ - public closed = false; + #closed = false; public constructor() { // `start`'s `this` is the underlying source object when using a plain method. // An arrow function captures the class instance so the controller is stored here. - this.stream = new ReadableStream({ + this.#stream = new ReadableStream({ start: (controller: ReadableStreamDefaultController): void => { - this.controller = controller; + this.#controller = controller; }, }); } + /** + * Flag indicating if the iterator is closed. + */ + public get closed(): boolean { + return this.#closed; + } + /** * Enqueues a value for the consumer. * @@ -49,9 +56,22 @@ export class AsyncPushIterator { * @param value - The next value to yield from the iterator. */ push(value: T): void { - if (this.closed) return; + if (this.#closed) return; - this.controller?.enqueue(value); + this.#controller?.enqueue(value); + } + + /** + * Causes any future interactions with the associated stream to error with {@link error}. + * Calling this will also clear the pending values immediately, so iterators that were listening will not receive them. + * + * @param error - The error to throw from the stream. + */ + error(error: Error): void { + if (this.#closed) return; + + this.#closed = true; + this.#controller?.error(error); } /** @@ -61,10 +81,10 @@ export class AsyncPushIterator { * Buffered values are still yielded before iteration completes. */ close(): void { - this.closed = true; + this.#closed = true; try { - this.controller?.close(); + this.#controller?.close(); } catch { // The reader may already have released or cancelled the stream. } @@ -75,8 +95,11 @@ export class AsyncPushIterator { * * Uses `preventCancel: true` so early `break` from `for await...of` does not * close the stream and block later pushes. + * + * Because values are discarded after being read, only a single consumer is supported. + * Additional consumers will receive a stream lock error. Unread values will be preserved until {@link close} is called. */ [Symbol.asyncIterator](): AsyncIterableIterator { - return this.stream.values({ preventCancel: true }); + return this.#stream.values({ preventCancel: true }); } } diff --git a/source/sse-session/index.ts b/source/sse-session/index.ts index 3e963df..bb0dee9 100644 --- a/source/sse-session/index.ts +++ b/source/sse-session/index.ts @@ -1,3 +1,4 @@ +export * from './async-push-iterator.ts'; export * from './types.ts'; export * from './sse-session.ts'; export * from './sse-event-parser.ts'; diff --git a/test/sse-session/async-push-iterator.test.ts b/test/sse-session/async-push-iterator.test.ts index e110b5c..82a1940 100644 --- a/test/sse-session/async-push-iterator.test.ts +++ b/test/sse-session/async-push-iterator.test.ts @@ -23,18 +23,14 @@ const collectAll = async (iterator: AsyncPushIterator): Promise => { const testPushComposedPushAndConsume = async (): Promise => { const iterator = new AsyncPushIterator(); - const result = new Promise((resolve) => { - void (async (): Promise => { - resolve(await collectAll(iterator)); - })(); - }); + const result = (): Promise => collectAll(iterator); iterator.push(1); iterator.push(2); iterator.push(3); iterator.close(); - await expect(result).resolves.toEqual([ 1, 2, 3 ]); + await expect(result()).resolves.toEqual([ 1, 2, 3 ]); }; /** @@ -61,15 +57,11 @@ const testPushComposedBuffersValuesPushedBeforeLoopStarts = async (): Promise => { const iterator = new AsyncPushIterator(); - const result = new Promise((resolve) => { - void (async (): Promise => { - resolve(await collectAll(iterator)); - })(); - }); + const result = async (): Promise => collectAll(iterator); iterator.close(); - await expect(result).resolves.toEqual([]); + await expect(result()).resolves.toEqual([]); }; /** @@ -78,11 +70,7 @@ const testPushComposedResolvesWithNoValues = async (): Promise => { const testPushComposedIgnoresValuesAfterClose = async (): Promise => { const iterator = new AsyncPushIterator(); - const result = new Promise((resolve) => { - void (async (): Promise => { - resolve(await collectAll(iterator)); - })(); - }); + const result = async (): Promise => collectAll(iterator); iterator.push(1); iterator.push(2); @@ -90,7 +78,7 @@ const testPushComposedIgnoresValuesAfterClose = async (): Promise => { iterator.close(); iterator.push(4); - await expect(result).resolves.toEqual([ 1, 2, 3 ]); + await expect(result()).resolves.toEqual([ 1, 2, 3 ]); }; /** @@ -104,30 +92,19 @@ const testPushComposedRejectsMultipleConsumers = async (): Promise => { const failureFlag = vi.fn(); - const successfulIterator = (): Promise => - new Promise((resolve) => { - void (async (): Promise => { - resolve(await collectAll(iterator)); - })(); - }); + const successfulIterator = (): Promise => collectAll(iterator); - const failedIterator = (): Promise => - new Promise((resolve, reject) => { - void (async (): Promise => { - try { - /* eslint-disable-next-line */ - for await (const _value of iterator) { - } - } catch (error) { - failureFlag(); - reject(error); - } + const failedIterator = async (): Promise => { + try { + /* eslint-disable-next-line */ + for await (const _value of iterator) { + } + } catch (error) { + failureFlag(); + } + }; - resolve(); - })(); - }); - - const promises = [ successfulIterator(), failedIterator().catch(() => {}) ]; + const promises = [ successfulIterator(), failedIterator() ]; iterator.close(); @@ -175,14 +152,63 @@ const testPushComposedAllowsPushingAfterEarlyBreak = async (): Promise => expect(secondPass).toEqual([ 2, 3 ]); }; +/** + * Tests that the iterator rejects after {@link AsyncPushIterator.error} is called. + */ +const testPushIteratorRejectsAfterError = async (): Promise => { + const iterator = new AsyncPushIterator(); + + iterator.error(new Error('Stream has been closed for a test')); + + await expect(collectAll(iterator)).rejects.toThrow('Stream has been closed for a test'); +}; + +/** + * Tests that the iterator closes the stream when error() is called. + */ +const testPushIteratorClosesWhenErrorIsCalled = async (): Promise => { + const iterator = new AsyncPushIterator(); + + iterator.error(new Error('Stream has been closed for a test')); + expect(iterator.closed).toBe(true); +}; + +/** + * Tests that subsequent calls to error() are ignored. + */ +const testPushIteratorIgnoresSubsequentErrorCalls = async (): Promise => { + const iterator = new AsyncPushIterator(); + + iterator.error(new Error('Stream has been closed for a test')); + expect(iterator.closed).toBe(true); + + expect(iterator.error(new Error('Second error'))).toBe(undefined); +}; + +/** + * Tests that a consumer can check if the iterator is closed. + */ +const testPushIteratorCanCheckIfClosed = async (): Promise => { + const iterator = new AsyncPushIterator(); + + expect(iterator.closed).toBe(false); + + iterator.close(); + expect(iterator.closed).toBe(true); +}; + const runTests = async (): Promise => { - test('AsyncPushIterator (composed): pushes and consumes values', testPushComposedPushAndConsume); - test('AsyncPushIterator (composed): buffers values pushed before the for-await loop starts', testPushComposedBuffersValuesPushedBeforeLoopStarts); - test('AsyncPushIterator (composed): resolves with no values when nothing was pushed', testPushComposedResolvesWithNoValues); - test('AsyncPushIterator (composed): ignores values pushed after close', testPushComposedIgnoresValuesAfterClose); - test('AsyncPushIterator (composed): rejects multiple consumers', testPushComposedRejectsMultipleConsumers); - test('AsyncPushIterator (composed): resolves immediately when closed before the loop starts', testPushComposedResolvesWhenClosedBeforeLoop); - test('AsyncPushIterator (composed): keeps the stream open after an early break', testPushComposedAllowsPushingAfterEarlyBreak); + test('AsyncPushIterator: pushes and consumes values', testPushComposedPushAndConsume); + test('AsyncPushIterator: buffers values pushed before the for-await loop starts', testPushComposedBuffersValuesPushedBeforeLoopStarts); + test('AsyncPushIterator: resolves with no values when nothing was pushed', testPushComposedResolvesWithNoValues); + test('AsyncPushIterator: ignores values pushed after close', testPushComposedIgnoresValuesAfterClose); + test('AsyncPushIterator: rejects multiple consumers', testPushComposedRejectsMultipleConsumers); + test('AsyncPushIterator: resolves immediately when closed before the loop starts', testPushComposedResolvesWhenClosedBeforeLoop); + test('AsyncPushIterator: keeps the stream open after an early break', testPushComposedAllowsPushingAfterEarlyBreak); + test('AsyncPushIterator: rejects after error', testPushIteratorRejectsAfterError); + test('AsyncPushIterator: closes the stream when error() is called', testPushIteratorClosesWhenErrorIsCalled); + test('AsyncPushIterator: ignores subsequent error() calls', testPushIteratorIgnoresSubsequentErrorCalls); + test('AsyncPushIterator: can check if closed', testPushIteratorCanCheckIfClosed); }; await runTests();