Merge branch 'development' into sse-and-backoff

This commit is contained in:
2026-07-26 04:00:35 +00:00
4 changed files with 137 additions and 87 deletions

60
package-lock.json generated
View File

@@ -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": [
{ {

View File

@@ -23,24 +23,31 @@
*/ */
export class AsyncPushIterator<T> { export class AsyncPushIterator<T> {
/** ReadableStream backing the async iterator returned from {@link Symbol.asyncIterator}. */ /** ReadableStream backing the async iterator returned from {@link Symbol.asyncIterator}. */
private readonly stream: ReadableStream<T>; #stream: ReadableStream<T>;
/** Controller used to enqueue values and close the stream from {@link push} and {@link close}. */ /** Controller used to enqueue values and close the stream from {@link push} and {@link close}. */
private controller: ReadableStreamDefaultController<T> | undefined; #controller: ReadableStreamDefaultController<T> | undefined;
/** When true, no more values are accepted and iteration eventually completes. */ /** When true, no more values are accepted and iteration eventually completes. */
public closed = false; #closed = false;
public constructor() { public constructor() {
// `start`'s `this` is the underlying source object when using a plain method. // `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. // An arrow function captures the class instance so the controller is stored here.
this.stream = new ReadableStream({ this.#stream = new ReadableStream({
start: (controller: ReadableStreamDefaultController<T>): void => { start: (controller: ReadableStreamDefaultController<T>): 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. * Enqueues a value for the consumer.
* *
@@ -49,9 +56,22 @@ export class AsyncPushIterator<T> {
* @param value - The next value to yield from the iterator. * @param value - The next value to yield from the iterator.
*/ */
push(value: T): void { 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<T> {
* Buffered values are still yielded before iteration completes. * Buffered values are still yielded before iteration completes.
*/ */
close(): void { close(): void {
this.closed = true; this.#closed = true;
try { try {
this.controller?.close(); this.#controller?.close();
} catch { } catch {
// The reader may already have released or cancelled the stream. // The reader may already have released or cancelled the stream.
} }
@@ -75,8 +95,11 @@ export class AsyncPushIterator<T> {
* *
* Uses `preventCancel: true` so early `break` from `for await...of` does not * Uses `preventCancel: true` so early `break` from `for await...of` does not
* close the stream and block later pushes. * 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> { [Symbol.asyncIterator](): AsyncIterableIterator<T> {
return this.stream.values({ preventCancel: true }); return this.#stream.values({ preventCancel: true });
} }
} }

View File

@@ -1,3 +1,4 @@
export * from './async-push-iterator.ts';
export * from './types.ts'; export * from './types.ts';
export * from './sse-session.ts'; export * from './sse-session.ts';
export * from './sse-event-parser.ts'; export * from './sse-event-parser.ts';

View File

@@ -23,18 +23,14 @@ const collectAll = async <T>(iterator: AsyncPushIterator<T>): Promise<T[]> => {
const testPushComposedPushAndConsume = async (): Promise<void> => { const testPushComposedPushAndConsume = async (): Promise<void> => {
const iterator = new AsyncPushIterator<number>(); const iterator = new AsyncPushIterator<number>();
const result = new Promise<number[]>((resolve) => { const result = (): Promise<number[]> => collectAll(iterator);
void (async (): Promise<void> => {
resolve(await collectAll(iterator));
})();
});
iterator.push(1); iterator.push(1);
iterator.push(2); iterator.push(2);
iterator.push(3); iterator.push(3);
iterator.close(); 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<vo
const testPushComposedResolvesWithNoValues = async (): Promise<void> => { const testPushComposedResolvesWithNoValues = async (): Promise<void> => {
const iterator = new AsyncPushIterator<number>(); const iterator = new AsyncPushIterator<number>();
const result = new Promise<number[]>((resolve) => { const result = async (): Promise<number[]> => collectAll(iterator);
void (async (): Promise<void> => {
resolve(await collectAll(iterator));
})();
});
iterator.close(); iterator.close();
await expect(result).resolves.toEqual([]); await expect(result()).resolves.toEqual([]);
}; };
/** /**
@@ -78,11 +70,7 @@ const testPushComposedResolvesWithNoValues = async (): Promise<void> => {
const testPushComposedIgnoresValuesAfterClose = async (): Promise<void> => { const testPushComposedIgnoresValuesAfterClose = async (): Promise<void> => {
const iterator = new AsyncPushIterator<number>(); const iterator = new AsyncPushIterator<number>();
const result = new Promise<number[]>((resolve) => { const result = async (): Promise<number[]> => collectAll(iterator);
void (async (): Promise<void> => {
resolve(await collectAll(iterator));
})();
});
iterator.push(1); iterator.push(1);
iterator.push(2); iterator.push(2);
@@ -90,7 +78,7 @@ const testPushComposedIgnoresValuesAfterClose = async (): Promise<void> => {
iterator.close(); iterator.close();
iterator.push(4); 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<void> => {
const failureFlag = vi.fn(); const failureFlag = vi.fn();
const successfulIterator = (): Promise<number[]> => const successfulIterator = (): Promise<number[]> => collectAll(iterator);
new Promise((resolve) => {
void (async (): Promise<void> => {
resolve(await collectAll(iterator));
})();
});
const failedIterator = (): Promise<void> => const failedIterator = async (): Promise<void> => {
new Promise((resolve, reject) => { try {
void (async (): Promise<void> => { /* eslint-disable-next-line */
try { for await (const _value of iterator) {
/* eslint-disable-next-line */ }
for await (const _value of iterator) { } catch (error) {
} failureFlag();
} catch (error) { }
failureFlag(); };
reject(error);
}
resolve(); const promises = [ successfulIterator(), failedIterator() ];
})();
});
const promises = [ successfulIterator(), failedIterator().catch(() => {}) ];
iterator.close(); iterator.close();
@@ -175,14 +152,63 @@ const testPushComposedAllowsPushingAfterEarlyBreak = async (): Promise<void> =>
expect(secondPass).toEqual([ 2, 3 ]); 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> => { const runTests = async (): Promise<void> => {
test('AsyncPushIterator (composed): pushes and consumes values', testPushComposedPushAndConsume); test('AsyncPushIterator: pushes and consumes values', testPushComposedPushAndConsume);
test('AsyncPushIterator (composed): buffers values pushed before the for-await loop starts', testPushComposedBuffersValuesPushedBeforeLoopStarts); test('AsyncPushIterator: buffers values pushed before the for-await loop starts', testPushComposedBuffersValuesPushedBeforeLoopStarts);
test('AsyncPushIterator (composed): resolves with no values when nothing was pushed', testPushComposedResolvesWithNoValues); test('AsyncPushIterator: resolves with no values when nothing was pushed', testPushComposedResolvesWithNoValues);
test('AsyncPushIterator (composed): ignores values pushed after close', testPushComposedIgnoresValuesAfterClose); test('AsyncPushIterator: ignores values pushed after close', testPushComposedIgnoresValuesAfterClose);
test('AsyncPushIterator (composed): rejects multiple consumers', testPushComposedRejectsMultipleConsumers); test('AsyncPushIterator: rejects multiple consumers', testPushComposedRejectsMultipleConsumers);
test('AsyncPushIterator (composed): resolves immediately when closed before the loop starts', testPushComposedResolvesWhenClosedBeforeLoop); test('AsyncPushIterator: resolves immediately when closed before the loop starts', testPushComposedResolvesWhenClosedBeforeLoop);
test('AsyncPushIterator (composed): keeps the stream open after an early break', testPushComposedAllowsPushingAfterEarlyBreak); 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(); await runTests();