Files
oh-my-pi/packages/utils/src/async.ts
T
can1357 a3117c2842 feat(agent): migrated async-drain utility to shared package for reuse
- 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.
2026-07-12 01:06:17 +02:00

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;
}
}