64 lines
2.5 KiB
JavaScript
64 lines
2.5 KiB
JavaScript
'use strict';
|
|
// steadfast — zero-dependency async resilience primitives.
|
|
|
|
/** Sleep for `ms`, abortable via an AbortSignal. */
|
|
function sleep(ms, signal) {
|
|
return new Promise((resolve, reject) => {
|
|
if (signal?.aborted) return reject(signal.reason ?? new Error('aborted'));
|
|
const t = setTimeout(resolve, ms);
|
|
signal?.addEventListener('abort', () => { clearTimeout(t); reject(signal.reason ?? new Error('aborted')); }, { once: true });
|
|
});
|
|
}
|
|
|
|
/**
|
|
* Retry an async function with exponential backoff + full jitter.
|
|
* @param {(attempt:number)=>Promise<any>} fn
|
|
* @param {object} [opts] attempts=3, minDelay=100, maxDelay=5000, factor=2,
|
|
* jitter=true, shouldRetry=(err,attempt)=>true, onRetry=(err,attempt,delay)=>{}, signal
|
|
*/
|
|
async function retry(fn, opts = {}) {
|
|
const { attempts = 3, minDelay = 100, maxDelay = 5000, factor = 2, jitter = true,
|
|
shouldRetry = () => true, onRetry = () => {}, signal } = opts;
|
|
if (!(attempts >= 1)) throw new RangeError('attempts must be >= 1');
|
|
let lastErr;
|
|
for (let attempt = 1; attempt <= attempts; attempt++) {
|
|
if (signal?.aborted) throw signal.reason ?? new Error('aborted');
|
|
try {
|
|
return await fn(attempt);
|
|
} catch (err) {
|
|
lastErr = err;
|
|
if (attempt >= attempts || !shouldRetry(err, attempt)) throw err;
|
|
let delay = Math.min(maxDelay, minDelay * Math.pow(factor, attempt - 1));
|
|
if (jitter) delay = Math.random() * delay; // full jitter (AWS-style)
|
|
onRetry(err, attempt, delay);
|
|
await sleep(delay, signal);
|
|
}
|
|
}
|
|
throw lastErr;
|
|
}
|
|
|
|
/** Reject if `promise` doesn't settle within `ms`. Cleans up its timer. */
|
|
function timeout(promise, ms, message) {
|
|
let t;
|
|
const timer = new Promise((_, reject) => {
|
|
t = setTimeout(() => reject(new Error(message || `timed out after ${ms}ms`)), ms);
|
|
});
|
|
return Promise.race([Promise.resolve(promise), timer]).finally(() => clearTimeout(t));
|
|
}
|
|
|
|
/** Concurrency limiter: returns run(fn) that runs at most `concurrency` at once. */
|
|
function pLimit(concurrency) {
|
|
if (!(concurrency >= 1)) throw new RangeError('concurrency must be >= 1');
|
|
let active = 0;
|
|
const queue = [];
|
|
const next = () => {
|
|
if (active >= concurrency || queue.length === 0) return;
|
|
active++;
|
|
const { fn, resolve, reject } = queue.shift();
|
|
Promise.resolve().then(fn).then(resolve, reject).finally(() => { active--; next(); });
|
|
};
|
|
return (fn) => new Promise((resolve, reject) => { queue.push({ fn, resolve, reject }); next(); });
|
|
}
|
|
|
|
module.exports = { retry, timeout, pLimit, sleep };
|