Merge branch 'composed-push-iterator' into 'development'
Composed push iterator See merge request GeneralProtocols/xo/utils!14
This commit is contained in:
60
package-lock.json
generated
60
package-lock.json
generated
@@ -1105,9 +1105,9 @@
|
|||||||
"license": "MIT"
|
"license": "MIT"
|
||||||
},
|
},
|
||||||
"node_modules/@eslint/config-array/node_modules/brace-expansion": {
|
"node_modules/@eslint/config-array/node_modules/brace-expansion": {
|
||||||
"version": "1.1.15",
|
"version": "1.1.16",
|
||||||
"resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.15.tgz",
|
"resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.16.tgz",
|
||||||
"integrity": "sha512-EwOCDEex4quD37XhqM3omwtMoJjr//isUZz1JopUNWms+4Z2ViyM/k1YIRePpoVNnQhENnxtFjLaxNHrT7xIUg==",
|
"integrity": "sha512-IDw48K2/2kRkg9LdJxurvq3lV3aBgq0REY89duEqFRthjlPdXHKMj7EnQOXVckxzgisinf3nHfrcE2FufFLXMw==",
|
||||||
"dev": true,
|
"dev": true,
|
||||||
"license": "MIT",
|
"license": "MIT",
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
@@ -1186,9 +1186,9 @@
|
|||||||
"license": "MIT"
|
"license": "MIT"
|
||||||
},
|
},
|
||||||
"node_modules/@eslint/eslintrc/node_modules/brace-expansion": {
|
"node_modules/@eslint/eslintrc/node_modules/brace-expansion": {
|
||||||
"version": "1.1.15",
|
"version": "1.1.16",
|
||||||
"resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.15.tgz",
|
"resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.16.tgz",
|
||||||
"integrity": "sha512-EwOCDEex4quD37XhqM3omwtMoJjr//isUZz1JopUNWms+4Z2ViyM/k1YIRePpoVNnQhENnxtFjLaxNHrT7xIUg==",
|
"integrity": "sha512-IDw48K2/2kRkg9LdJxurvq3lV3aBgq0REY89duEqFRthjlPdXHKMj7EnQOXVckxzgisinf3nHfrcE2FufFLXMw==",
|
||||||
"dev": true,
|
"dev": true,
|
||||||
"license": "MIT",
|
"license": "MIT",
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
@@ -3435,9 +3435,9 @@
|
|||||||
}
|
}
|
||||||
},
|
},
|
||||||
"node_modules/body-parser": {
|
"node_modules/body-parser": {
|
||||||
"version": "1.20.5",
|
"version": "1.20.6",
|
||||||
"resolved": "https://registry.npmjs.org/body-parser/-/body-parser-1.20.5.tgz",
|
"resolved": "https://registry.npmjs.org/body-parser/-/body-parser-1.20.6.tgz",
|
||||||
"integrity": "sha512-3grm+/2tUOvu2cjJkvsIxrv/wVpfXQW4PsQHYm7yk4vfpu7Ekl6nEsYBoJUL6qDwZUx8wUhQ8tR2qz+ad9c9OA==",
|
"integrity": "sha512-p5tAzS57i5MV9fZFDj9LeIiTZEufbSe2eDozP+ElheSUq1m74CRq1jI4mYNDdVs9vQztXFLuk/Gd6BWTdwRJ5g==",
|
||||||
"dev": true,
|
"dev": true,
|
||||||
"license": "MIT",
|
"license": "MIT",
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
@@ -3503,9 +3503,9 @@
|
|||||||
}
|
}
|
||||||
},
|
},
|
||||||
"node_modules/brace-expansion": {
|
"node_modules/brace-expansion": {
|
||||||
"version": "5.0.6",
|
"version": "5.0.7",
|
||||||
"resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-5.0.6.tgz",
|
"resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-5.0.7.tgz",
|
||||||
"integrity": "sha512-kLpxurY4Z4r9sgMsyG0Z9uzsBlgiU/EFKhj/h91/8yHu0edo7XuixOIH3VcJ8kkxs6/jPzoI6U9Vj3WqbMQ94g==",
|
"integrity": "sha512-7oFy703dxfY3/NLxC1fh2SUCQ0H9rmAY+5EpDVfXjUTTs+HEwR2nYaqLv+GWcTsumwxPfiz6CzCNkwXwBUwqCA==",
|
||||||
"dev": true,
|
"dev": true,
|
||||||
"license": "MIT",
|
"license": "MIT",
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
@@ -5292,9 +5292,9 @@
|
|||||||
"peer": true
|
"peer": true
|
||||||
},
|
},
|
||||||
"node_modules/eslint-plugin-import/node_modules/brace-expansion": {
|
"node_modules/eslint-plugin-import/node_modules/brace-expansion": {
|
||||||
"version": "1.1.15",
|
"version": "1.1.16",
|
||||||
"resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.15.tgz",
|
"resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.16.tgz",
|
||||||
"integrity": "sha512-EwOCDEex4quD37XhqM3omwtMoJjr//isUZz1JopUNWms+4Z2ViyM/k1YIRePpoVNnQhENnxtFjLaxNHrT7xIUg==",
|
"integrity": "sha512-IDw48K2/2kRkg9LdJxurvq3lV3aBgq0REY89duEqFRthjlPdXHKMj7EnQOXVckxzgisinf3nHfrcE2FufFLXMw==",
|
||||||
"dev": true,
|
"dev": true,
|
||||||
"license": "MIT",
|
"license": "MIT",
|
||||||
"peer": true,
|
"peer": true,
|
||||||
@@ -5379,9 +5379,9 @@
|
|||||||
"peer": true
|
"peer": true
|
||||||
},
|
},
|
||||||
"node_modules/eslint-plugin-jsx-a11y/node_modules/brace-expansion": {
|
"node_modules/eslint-plugin-jsx-a11y/node_modules/brace-expansion": {
|
||||||
"version": "1.1.15",
|
"version": "1.1.16",
|
||||||
"resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.15.tgz",
|
"resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.16.tgz",
|
||||||
"integrity": "sha512-EwOCDEex4quD37XhqM3omwtMoJjr//isUZz1JopUNWms+4Z2ViyM/k1YIRePpoVNnQhENnxtFjLaxNHrT7xIUg==",
|
"integrity": "sha512-IDw48K2/2kRkg9LdJxurvq3lV3aBgq0REY89duEqFRthjlPdXHKMj7EnQOXVckxzgisinf3nHfrcE2FufFLXMw==",
|
||||||
"dev": true,
|
"dev": true,
|
||||||
"license": "MIT",
|
"license": "MIT",
|
||||||
"peer": true,
|
"peer": true,
|
||||||
@@ -5468,9 +5468,9 @@
|
|||||||
"peer": true
|
"peer": true
|
||||||
},
|
},
|
||||||
"node_modules/eslint-plugin-react/node_modules/brace-expansion": {
|
"node_modules/eslint-plugin-react/node_modules/brace-expansion": {
|
||||||
"version": "1.1.15",
|
"version": "1.1.16",
|
||||||
"resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.15.tgz",
|
"resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.16.tgz",
|
||||||
"integrity": "sha512-EwOCDEex4quD37XhqM3omwtMoJjr//isUZz1JopUNWms+4Z2ViyM/k1YIRePpoVNnQhENnxtFjLaxNHrT7xIUg==",
|
"integrity": "sha512-IDw48K2/2kRkg9LdJxurvq3lV3aBgq0REY89duEqFRthjlPdXHKMj7EnQOXVckxzgisinf3nHfrcE2FufFLXMw==",
|
||||||
"dev": true,
|
"dev": true,
|
||||||
"license": "MIT",
|
"license": "MIT",
|
||||||
"peer": true,
|
"peer": true,
|
||||||
@@ -5542,9 +5542,9 @@
|
|||||||
"license": "MIT"
|
"license": "MIT"
|
||||||
},
|
},
|
||||||
"node_modules/eslint/node_modules/brace-expansion": {
|
"node_modules/eslint/node_modules/brace-expansion": {
|
||||||
"version": "1.1.15",
|
"version": "1.1.16",
|
||||||
"resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.15.tgz",
|
"resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-1.1.16.tgz",
|
||||||
"integrity": "sha512-EwOCDEex4quD37XhqM3omwtMoJjr//isUZz1JopUNWms+4Z2ViyM/k1YIRePpoVNnQhENnxtFjLaxNHrT7xIUg==",
|
"integrity": "sha512-IDw48K2/2kRkg9LdJxurvq3lV3aBgq0REY89duEqFRthjlPdXHKMj7EnQOXVckxzgisinf3nHfrcE2FufFLXMw==",
|
||||||
"dev": true,
|
"dev": true,
|
||||||
"license": "MIT",
|
"license": "MIT",
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
@@ -7129,9 +7129,9 @@
|
|||||||
"license": "MIT"
|
"license": "MIT"
|
||||||
},
|
},
|
||||||
"node_modules/js-yaml": {
|
"node_modules/js-yaml": {
|
||||||
"version": "4.2.0",
|
"version": "4.3.0",
|
||||||
"resolved": "https://registry.npmjs.org/js-yaml/-/js-yaml-4.2.0.tgz",
|
"resolved": "https://registry.npmjs.org/js-yaml/-/js-yaml-4.3.0.tgz",
|
||||||
"integrity": "sha512-ePWsvanv0DWuDRsW8dnt+R4jQ31SCRCQ7hhNcPXZPsoBZiemuZNYGf7adZdqX2D86j6rvKp3RpCxVTSb8WQlOw==",
|
"integrity": "sha512-1td788aAnnZ5qs7V2QIRl1owjtYpbKt749Y3xauqQgwIIGF/xXWz1wMTEBx5O3LK3lXLVuqXPdPxj2BoFHaW9Q==",
|
||||||
"dev": true,
|
"dev": true,
|
||||||
"funding": [
|
"funding": [
|
||||||
{
|
{
|
||||||
@@ -7524,9 +7524,9 @@
|
|||||||
}
|
}
|
||||||
},
|
},
|
||||||
"node_modules/linkify-it": {
|
"node_modules/linkify-it": {
|
||||||
"version": "5.0.1",
|
"version": "5.0.2",
|
||||||
"resolved": "https://registry.npmjs.org/linkify-it/-/linkify-it-5.0.1.tgz",
|
"resolved": "https://registry.npmjs.org/linkify-it/-/linkify-it-5.0.2.tgz",
|
||||||
"integrity": "sha512-wVoTjP4Q6R0NW5hiZkVJaFZPWgtXfoGF+6LucL3/FtiNjmcHhYjEr5f1Kqjirc1nBW07J/ZuRFumqr2oqccEWg==",
|
"integrity": "sha512-ONTm2jCMAVZjgQa/Fy1kScXsuOoF5NPTsoFBdE1KVIZ2vAh/r9+Bqo+0jINCBYnavTPQZz38QzFTme79ENoN3Q==",
|
||||||
"dev": true,
|
"dev": true,
|
||||||
"funding": [
|
"funding": [
|
||||||
{
|
{
|
||||||
|
|||||||
105
source/sse-session/async-push-iterator.ts
Normal file
105
source/sse-session/async-push-iterator.ts
Normal file
@@ -0,0 +1,105 @@
|
|||||||
|
/**
|
||||||
|
* An async iterable queue that bridges push-based producers and pull-based consumers.
|
||||||
|
*
|
||||||
|
* Composes an internal {@link ReadableStream} instead of extending it, so producers
|
||||||
|
* call {@link push} while consumers use standard async iteration (`for await...of`).
|
||||||
|
*
|
||||||
|
* ```ts
|
||||||
|
* const messages = new AsyncPushIterator<SSEvent>();
|
||||||
|
*
|
||||||
|
* // Producer (elsewhere)
|
||||||
|
* messages.push(event);
|
||||||
|
*
|
||||||
|
* // Consumer
|
||||||
|
* for await (const event of messages) {
|
||||||
|
* handle(event);
|
||||||
|
* }
|
||||||
|
* ```
|
||||||
|
*
|
||||||
|
* {@link Symbol.asyncIterator} returns `stream.values({ preventCancel: true })` so
|
||||||
|
* breaking out of `for await...of` does not cancel the underlying stream. That
|
||||||
|
* matters for long-lived sessions where the producer keeps pushing after a consumer
|
||||||
|
* stops reading early (for example, test helpers that only collect a fixed count).
|
||||||
|
*/
|
||||||
|
export class AsyncPushIterator<T> {
|
||||||
|
/** ReadableStream backing the async iterator returned from {@link Symbol.asyncIterator}. */
|
||||||
|
#stream: ReadableStream<T>;
|
||||||
|
|
||||||
|
/** Controller used to enqueue values and close the stream from {@link push} and {@link close}. */
|
||||||
|
#controller: ReadableStreamDefaultController<T> | undefined;
|
||||||
|
|
||||||
|
/** When true, no more values are accepted and iteration eventually completes. */
|
||||||
|
#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({
|
||||||
|
start: (controller: ReadableStreamDefaultController<T>): void => {
|
||||||
|
this.#controller = controller;
|
||||||
|
},
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Flag indicating if the iterator is closed.
|
||||||
|
*/
|
||||||
|
public get closed(): boolean {
|
||||||
|
return this.#closed;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Enqueues a value for the consumer.
|
||||||
|
*
|
||||||
|
* After {@link close}, pushes are silently dropped.
|
||||||
|
*
|
||||||
|
* @param value - The next value to yield from the iterator.
|
||||||
|
*/
|
||||||
|
push(value: T): void {
|
||||||
|
if (this.#closed) return;
|
||||||
|
|
||||||
|
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);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Ends the stream.
|
||||||
|
*
|
||||||
|
* Marks the iterator closed so future {@link push} calls are ignored.
|
||||||
|
* Buffered values are still yielded before iteration completes.
|
||||||
|
*/
|
||||||
|
close(): void {
|
||||||
|
this.#closed = true;
|
||||||
|
|
||||||
|
try {
|
||||||
|
this.#controller?.close();
|
||||||
|
} catch {
|
||||||
|
// The reader may already have released or cancelled the stream.
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Returns an async iterator over the composed stream.
|
||||||
|
*
|
||||||
|
* 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<T> {
|
||||||
|
return this.#stream.values({ preventCancel: true });
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -1 +1,2 @@
|
|||||||
|
export * from './async-push-iterator.ts';
|
||||||
export * from './sse-event-parser.ts';
|
export * from './sse-event-parser.ts';
|
||||||
|
|||||||
214
test/sse-session/async-push-iterator.test.ts
Normal file
214
test/sse-session/async-push-iterator.test.ts
Normal file
@@ -0,0 +1,214 @@
|
|||||||
|
import { expect, test, vi } from 'vitest';
|
||||||
|
|
||||||
|
import { AsyncPushIterator } from '../../source/sse-session/async-push-iterator.ts';
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Collects every value from the iterator into an array.
|
||||||
|
*
|
||||||
|
* @param iterator - Iterator under test.
|
||||||
|
*/
|
||||||
|
const collectAll = async <T>(iterator: AsyncPushIterator<T>): Promise<T[]> => {
|
||||||
|
const results: T[] = [];
|
||||||
|
|
||||||
|
for await (const value of iterator) {
|
||||||
|
results.push(value);
|
||||||
|
}
|
||||||
|
|
||||||
|
return results;
|
||||||
|
};
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Tests that values pushed while a consumer is already waiting are delivered in order.
|
||||||
|
*/
|
||||||
|
const testPushComposedPushAndConsume = async (): Promise<void> => {
|
||||||
|
const iterator = new AsyncPushIterator<number>();
|
||||||
|
|
||||||
|
const result = (): Promise<number[]> => collectAll(iterator);
|
||||||
|
|
||||||
|
iterator.push(1);
|
||||||
|
iterator.push(2);
|
||||||
|
iterator.push(3);
|
||||||
|
iterator.close();
|
||||||
|
|
||||||
|
await expect(result()).resolves.toEqual([ 1, 2, 3 ]);
|
||||||
|
};
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Tests that values pushed before `for await...of` starts are buffered and yielded
|
||||||
|
* once the consumer begins reading.
|
||||||
|
*/
|
||||||
|
const testPushComposedBuffersValuesPushedBeforeLoopStarts = async (): Promise<void> => {
|
||||||
|
const iterator = new AsyncPushIterator<number>();
|
||||||
|
|
||||||
|
iterator.push(1);
|
||||||
|
iterator.push(2);
|
||||||
|
iterator.push(3);
|
||||||
|
|
||||||
|
const result = collectAll(iterator);
|
||||||
|
|
||||||
|
iterator.close();
|
||||||
|
|
||||||
|
await expect(result).resolves.toEqual([ 1, 2, 3 ]);
|
||||||
|
};
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Tests that the iterator completes with no values when nothing was pushed.
|
||||||
|
*/
|
||||||
|
const testPushComposedResolvesWithNoValues = async (): Promise<void> => {
|
||||||
|
const iterator = new AsyncPushIterator<number>();
|
||||||
|
|
||||||
|
const result = async (): Promise<number[]> => collectAll(iterator);
|
||||||
|
|
||||||
|
iterator.close();
|
||||||
|
|
||||||
|
await expect(result()).resolves.toEqual([]);
|
||||||
|
};
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Tests that values pushed after {@link AsyncPushIterator.close} are ignored.
|
||||||
|
*/
|
||||||
|
const testPushComposedIgnoresValuesAfterClose = async (): Promise<void> => {
|
||||||
|
const iterator = new AsyncPushIterator<number>();
|
||||||
|
|
||||||
|
const result = async (): Promise<number[]> => collectAll(iterator);
|
||||||
|
|
||||||
|
iterator.push(1);
|
||||||
|
iterator.push(2);
|
||||||
|
iterator.push(3);
|
||||||
|
iterator.close();
|
||||||
|
iterator.push(4);
|
||||||
|
|
||||||
|
await expect(result()).resolves.toEqual([ 1, 2, 3 ]);
|
||||||
|
};
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Tests that only one async consumer can read from the composed ReadableStream at a time.
|
||||||
|
*
|
||||||
|
* Unlike the hand-rolled async-push-iterator, the second consumer fails with a
|
||||||
|
* stream lock error rather than TooManyAsyncIteratorsError.
|
||||||
|
*/
|
||||||
|
const testPushComposedRejectsMultipleConsumers = async (): Promise<void> => {
|
||||||
|
const iterator = new AsyncPushIterator<number>();
|
||||||
|
|
||||||
|
const failureFlag = vi.fn();
|
||||||
|
|
||||||
|
const successfulIterator = (): Promise<number[]> => collectAll(iterator);
|
||||||
|
|
||||||
|
const failedIterator = async (): Promise<void> => {
|
||||||
|
try {
|
||||||
|
/* eslint-disable-next-line */
|
||||||
|
for await (const _value of iterator) {
|
||||||
|
}
|
||||||
|
} catch (error) {
|
||||||
|
failureFlag();
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
const promises = [ successfulIterator(), failedIterator() ];
|
||||||
|
|
||||||
|
iterator.close();
|
||||||
|
|
||||||
|
await Promise.all(promises);
|
||||||
|
|
||||||
|
expect(failureFlag).toHaveBeenCalledOnce();
|
||||||
|
};
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Tests that closing before iteration starts lets the loop finish immediately.
|
||||||
|
*/
|
||||||
|
const testPushComposedResolvesWhenClosedBeforeLoop = async (): Promise<void> => {
|
||||||
|
const iterator = new AsyncPushIterator<number>();
|
||||||
|
|
||||||
|
iterator.close();
|
||||||
|
|
||||||
|
await expect(collectAll(iterator)).resolves.toEqual([]);
|
||||||
|
};
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Tests that breaking out of `for await...of` early does not cancel the stream.
|
||||||
|
*
|
||||||
|
* {@link AsyncPushIterator} uses `preventCancel: true` so producers can keep pushing
|
||||||
|
* and a later consumer can read the remaining values.
|
||||||
|
*/
|
||||||
|
const testPushComposedAllowsPushingAfterEarlyBreak = async (): Promise<void> => {
|
||||||
|
const iterator = new AsyncPushIterator<number>();
|
||||||
|
|
||||||
|
iterator.push(1);
|
||||||
|
|
||||||
|
const firstPass: number[] = [];
|
||||||
|
|
||||||
|
for await (const value of iterator) {
|
||||||
|
firstPass.push(value);
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
|
||||||
|
iterator.push(2);
|
||||||
|
iterator.push(3);
|
||||||
|
iterator.close();
|
||||||
|
|
||||||
|
const secondPass = await collectAll(iterator);
|
||||||
|
|
||||||
|
expect(firstPass).toEqual([ 1 ]);
|
||||||
|
expect(secondPass).toEqual([ 2, 3 ]);
|
||||||
|
};
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Tests that the iterator rejects after {@link AsyncPushIterator.error} is called.
|
||||||
|
*/
|
||||||
|
const testPushIteratorRejectsAfterError = async (): Promise<void> => {
|
||||||
|
const iterator = new AsyncPushIterator<number>();
|
||||||
|
|
||||||
|
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<void> => {
|
||||||
|
const iterator = new AsyncPushIterator<number>();
|
||||||
|
|
||||||
|
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<void> => {
|
||||||
|
const iterator = new AsyncPushIterator<number>();
|
||||||
|
|
||||||
|
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<void> => {
|
||||||
|
const iterator = new AsyncPushIterator<number>();
|
||||||
|
|
||||||
|
expect(iterator.closed).toBe(false);
|
||||||
|
|
||||||
|
iterator.close();
|
||||||
|
expect(iterator.closed).toBe(true);
|
||||||
|
};
|
||||||
|
|
||||||
|
const runTests = async (): Promise<void> => {
|
||||||
|
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();
|
||||||
Reference in New Issue
Block a user