Files
plainleaf/plug-api/lib/async.ts
T

176 lines
4.1 KiB
TypeScript

export function throttle(
func: () => void,
limit: number,
): (() => void) & { flush(): void } {
let timer: any = null;
const throttled = () => {
if (!timer) {
timer = setTimeout(() => {
func();
timer = null;
}, limit);
}
};
// Immediately execute any pending call and cancel the timer
throttled.flush = () => {
if (timer) {
clearTimeout(timer);
timer = null;
func();
}
};
return throttled;
}
export function throttleImmediately(
func: () => void,
limit: number,
): () => void {
let timer: any = null;
return () => {
if (!timer) {
func();
timer = setTimeout(() => {
timer = null;
}, limit);
}
};
}
// race for promises returns first promise that resolves
export function race<T>(promises: Promise<T>[]): Promise<T> {
return new Promise((resolve, reject) => {
for (const p of promises) {
p.then(resolve, reject);
}
});
}
export function timeout(ms: number): Promise<never> {
return new Promise((_resolve, reject) =>
setTimeout(() => {
reject(new Error("timeout"));
}, ms),
);
}
export function sleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
}
export class PromiseQueue {
private queue: {
fn: () => Promise<any>;
resolve: (value: any) => void;
reject: (error: any) => void;
}[] = [];
private processing = false;
runInQueue(fn: () => Promise<any>): Promise<any> {
return new Promise((resolve, reject) => {
this.queue.push({ fn, resolve, reject });
if (!this.processing) {
void this.process();
}
});
}
private async process(): Promise<void> {
if (this.queue.length === 0) {
this.processing = false;
return;
}
this.processing = true;
const { fn, resolve, reject } = this.queue.shift()!;
try {
const result = await fn();
resolve(result);
} catch (error) {
reject(error);
}
void this.process(); // Continue processing the next promise in the queue
}
}
/**
* Batches up values, and processes in batches of batchSize in parallel
* then merges the results in the appropriate order.
* @param values - The values to batch.
* @param fn - The function to run on each batch.
* @param batchSize - The size of each batch.
*/
export async function batchRequests<I, O>(
values: I[],
fn: (batch: I[]) => Promise<O[]>,
batchSize: number,
): Promise<O[]> {
const results: O[] = [];
// Split values into batches of batchSize
const batches: I[][] = [];
for (let i = 0; i < values.length; i += batchSize) {
batches.push(values.slice(i, i + batchSize));
}
// Run fn on them in parallel
const batchResults = await Promise.all(batches.map(fn));
// Flatten the results
for (const batchResult of batchResults) {
if (Array.isArray(batchResult)) {
// If fn returns an array, collect them
results.push(...batchResult);
}
}
return results;
}
/**
* Processes items in parallel with a specified concurrency limit.
* @param items - The items to process.
* @param handler - The function to run on each item.
* @param concurrency - The maximum number of concurrent operations.
*/
export async function processWithConcurrency<I, O>(
items: I[],
handler: (item: I) => Promise<O>,
concurrency: number,
): Promise<O[]> {
const results: O[] = [];
let idx = 0;
async function worker() {
while (idx < items.length) {
const currentIdx = idx++;
const item = items[currentIdx];
const result = await handler(item);
results[currentIdx] = result;
}
}
const workers = [];
for (let i = 0; i < Math.min(concurrency, items.length); i++) {
workers.push(worker());
}
await Promise.all(workers);
return results;
}
/**
* Runs a function safely by catching any errors and logging them to the console.
* @param fn - The function to run.
*/
export function safeRun(fn: () => Promise<void>): void {
fn().catch((e) => {
console.error(e);
});
}
/**
* Generates a random delay between 0 and 1000 milliseconds.
*/
export function jitter(maxLength = 1000): number {
return Math.floor(Math.random() * maxLength);
}