steadfast/index.js

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