diff --git a/README.md b/README.md index bb91cc0..a86d9a3 100644 --- a/README.md +++ b/README.md @@ -1,3 +1,29 @@ # steadfast -Zero-dependency async resilience: retry with exponential backoff + jitter, timeout, and concurrency limiting. \ No newline at end of file +Zero-dependency async resilience primitives for agents making unreliable calls. + +```js +const { retry, timeout, pLimit } = require('@community/steadfast'); + +// Retry a flaky call: 5 attempts, exponential backoff + jitter, only on 5xx. +const data = await retry(() => fetchThing(), { + attempts: 5, + shouldRetry: (err) => err.status >= 500, + onRetry: (err, n, delay) => console.warn(`attempt ${n} failed, retrying in ${delay|0}ms`), +}); + +// Bound any promise with a timeout. +const res = await timeout(slowCall(), 2000); + +// Run many tasks with bounded concurrency. +const limit = pLimit(4); +await Promise.all(urls.map((u) => limit(() => fetch(u)))); +``` + +## API +- `retry(fn, opts)` — `attempts`, `minDelay`, `maxDelay`, `factor`, `jitter`, `shouldRetry(err,attempt)`, `onRetry(err,attempt,delay)`, `signal` (AbortSignal). +- `timeout(promise, ms, message?)` — rejects if not settled in time; cleans up its timer. +- `pLimit(concurrency)` — returns `run(fn)`; at most `concurrency` run at once. +- `sleep(ms, signal?)` — abortable delay. + +MIT. diff --git a/index.js b/index.js new file mode 100644 index 0000000..6a59fc1 --- /dev/null +++ b/index.js @@ -0,0 +1,63 @@ +'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} 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 }; diff --git a/package.json b/package.json new file mode 100644 index 0000000..135c964 --- /dev/null +++ b/package.json @@ -0,0 +1,8 @@ +{ + "name": "@community/steadfast", + "version": "1.0.0", + "description": "Zero-dependency async resilience: retry with exponential backoff + jitter, timeout, and concurrency limiting.", + "main": "index.js", + "keywords": ["retry", "backoff", "jitter", "timeout", "concurrency", "p-limit", "resilience", "async"], + "license": "MIT" +} diff --git a/test.sh b/test.sh new file mode 100755 index 0000000..bf2cbc7 --- /dev/null +++ b/test.sh @@ -0,0 +1,34 @@ +#!/bin/bash +set -e +node -e ' +const assert = require("assert"); +const { retry, timeout, pLimit, sleep } = require("./index.js"); +(async () => { + // retry: succeeds on the 3rd attempt + let calls = 0; + const v = await retry(async () => { calls++; if (calls < 3) throw new Error("fail"); return "ok"; }, { minDelay: 1, jitter: false }); + assert.strictEqual(v, "ok"); assert.strictEqual(calls, 3); + + // retry: gives up after `attempts`, calling exactly that many times + let n = 0; + await assert.rejects(retry(async () => { n++; throw new Error("boom"); }, { attempts: 4, minDelay: 1 })); + assert.strictEqual(n, 4); + + // retry: shouldRetry=false → fails fast (1 call) + let m = 0; + await assert.rejects(retry(async () => { m++; throw new Error("nope"); }, { attempts: 5, minDelay: 1, shouldRetry: () => false })); + assert.strictEqual(m, 1); + + // timeout: fast resolves, slow rejects + assert.strictEqual(await timeout(Promise.resolve(42), 50), 42); + await assert.rejects(timeout(sleep(100), 10), /timed out/); + + // pLimit: never exceeds the concurrency cap + const limit = pLimit(2); + let active = 0, peak = 0; + await Promise.all(Array.from({ length: 6 }, () => limit(async () => { active++; peak = Math.max(peak, active); await sleep(15); active--; }))); + assert.ok(peak <= 2, "peak concurrency " + peak + " exceeded 2"); + + console.log("all tests passed (retry, timeout, pLimit)"); +})().catch((e) => { console.error("TEST FAILED:", e.message); process.exit(1); }); +'