a3117c2842
- Relocated `AsyncDrain` class from `coding-agent` to `utils` package. - Updated `HistoryStorage` to import `AsyncDrain` from shared utilities. - Centralized the utility to allow reuse across the codebase.
92 lines
2.5 KiB
TypeScript
92 lines
2.5 KiB
TypeScript
/**
|
|
* Wrap a promise with a timeout and optional abort signal.
|
|
* Rejects with the given message if the timeout fires first.
|
|
* Cleans up all listeners on settlement.
|
|
*/
|
|
export function withTimeout<T>(promise: Promise<T>, ms: number, message: string, signal?: AbortSignal): Promise<T> {
|
|
if (signal?.aborted) {
|
|
const reason = signal.reason instanceof Error ? signal.reason : new Error("Aborted");
|
|
return Promise.reject(reason);
|
|
}
|
|
|
|
const { promise: wrapped, resolve, reject } = Promise.withResolvers<T>();
|
|
let settled = false;
|
|
const timeoutId = setTimeout(() => {
|
|
if (settled) return;
|
|
settled = true;
|
|
if (signal) signal.removeEventListener("abort", onAbort);
|
|
reject(new Error(message));
|
|
}, ms);
|
|
|
|
const onAbort = () => {
|
|
if (settled) return;
|
|
settled = true;
|
|
clearTimeout(timeoutId);
|
|
reject(signal?.reason instanceof Error ? signal.reason : new Error("Aborted"));
|
|
};
|
|
|
|
if (signal) {
|
|
signal.addEventListener("abort", onAbort, { once: true });
|
|
}
|
|
|
|
promise.then(
|
|
value => {
|
|
if (settled) return;
|
|
settled = true;
|
|
clearTimeout(timeoutId);
|
|
if (signal) signal.removeEventListener("abort", onAbort);
|
|
resolve(value);
|
|
},
|
|
err => {
|
|
if (settled) return;
|
|
settled = true;
|
|
clearTimeout(timeoutId);
|
|
if (signal) signal.removeEventListener("abort", onAbort);
|
|
reject(err);
|
|
},
|
|
);
|
|
|
|
return wrapped;
|
|
}
|
|
|
|
/**
|
|
* Coalesces rapid-fire writes into one deferred batch. `push` queues a value
|
|
* and returns a promise for the batch flush; the first push of a batch arms a
|
|
* timer (`delayMs`, or a microtask at 0), and every push before it fires joins
|
|
* the same batch and shares the same promise. Used to keep hot paths off
|
|
* synchronous storage (prompt history, model perf).
|
|
*/
|
|
export class AsyncDrain<T> {
|
|
#queue?: T[];
|
|
#promise = Promise.resolve();
|
|
|
|
constructor(readonly delayMs: number = 0) {}
|
|
|
|
/** Queue `value`; `hnd` receives the whole batch when the window closes. */
|
|
push(value: T, hnd: (values: T[]) => Promise<void> | void): Promise<void> {
|
|
let queue = this.#queue;
|
|
if (!queue) {
|
|
this.#queue = queue = [];
|
|
const { promise, resolve, reject } = Promise.withResolvers<void>();
|
|
const exec = (): void => {
|
|
try {
|
|
if (this.#queue === queue) {
|
|
this.#queue = undefined;
|
|
}
|
|
resolve(hnd(queue!));
|
|
} catch (error) {
|
|
reject(error);
|
|
}
|
|
};
|
|
if (this.delayMs > 0) {
|
|
setTimeout(exec, this.delayMs);
|
|
} else {
|
|
queueMicrotask(exec);
|
|
}
|
|
this.#promise = promise;
|
|
}
|
|
queue.push(value);
|
|
return this.#promise;
|
|
}
|
|
}
|