Skip to content

@studnicky/resilience

Circuit breaker, token bucket rate limiter, and bounded dead-letter queue. Each primitive is independently usable and composable.

Install

bash
pnpm add @studnicky/resilience

Requires @studnicky:registry=https://npm.pkg.github.com in .npmrc.

Construct runtime primitives through the package root. Schema-backed data declarations live at @studnicky/resilience/entities, and type-only contracts live at @studnicky/resilience/interfaces.

Usage

CircuitBreaker

Tracks failures and opens the circuit after a threshold, then probes with limited calls after a timeout.

ts
import { CircuitBreaker, CircuitBreakerOpenError } from '../src/index.js';

// Deterministic clock so tests are instant with no real waits.
let now = 0;
class Clock {
  static now(): number { const result = now + 0; return result; }
}

const breaker = CircuitBreaker.create({
  'clock': Clock.now,
  'failureThreshold': 3,
  'name': 'test-service',
  'resetTimeoutMs': 1_000,
  'successThreshold': 2
});

// --- CLOSED: successful calls pass through ---
await breaker.execute(() => { const result = Promise.resolve('ok'); return result; });
console.log('State after success:', breaker.state);

// --- Trip the breaker: 3 consecutive failures ---
class Fail {
  static boom(): Promise<never> { throw RuntimeError.create('boom'); }
}
for (let i = 0; i < 3; i++) {
  await breaker.execute(Fail.boom).catch(() => { /* expected */ });
}
console.log('State after 3 failures:', breaker.state);

// --- OPEN: next call is fast-rejected with CircuitBreakerOpenError ---
await breaker.execute(() => { const result = Promise.resolve('should not run'); return result; }).catch((error) => {
  console.log('Open-circuit rejection:', error instanceof CircuitBreakerOpenError ? 'CircuitBreakerOpenError' : 'other');
});

// --- Advance past resetTimeoutMs → halfOpen on next call ---
now = 1_001;
await breaker.execute(() => { const result = Promise.resolve('probe 1'); return result; });
console.log('State after probe 1:', breaker.state);

// --- 2 successes in halfOpen close the circuit (successThreshold: 2) ---
await breaker.execute(() => { const result = Promise.resolve('probe 2'); return result; });
console.log('State after probe 2:', breaker.state);

// --- forceOpen / reset ---
breaker.forceOpen();
console.log('State after forceOpen:', breaker.state);
breaker.reset();
console.log('State after reset:', breaker.state);

TokenBucket

Token-bucket rate limiter; consume throws immediately when exhausted, waitForToken blocks until tokens refill.

ts
import { TokenBucket, TokenBucketExhaustedError } from '../src/index.js';

// Deterministic clock: start at t=0, advance manually.
let now = 0;
class Clock {
  static now(): number { const result = now + 0; return result; }
}

// 2 tokens/s, burst capacity 3 → starts full (3 tokens).
const bucket = TokenBucket.create({ 'burstSize': 3, 'clock': Clock.now, 'requestsPerSecond': 2 });
console.log('Initial available tokens:', bucket.available);

// --- consume() drains tokens ---
bucket.consume();
bucket.consume();
bucket.consume();
console.log('Available after 3 consumes:', bucket.available);

// --- Advance 500 ms → 1 new token (2 tokens/s × 0.5 s = 1) ---
now = 500;
console.log('Available at t=500ms:', bucket.available);
bucket.consume();

// --- Advance to 1500 ms → 2 tokens refilled, capped at burstSize=3 ---
now = 1_500;
console.log('Available at t=1500ms:', bucket.available);

// --- waitForToken with abort signal cancels correctly ---
// Drain remaining tokens so waitForToken has to block, then abort concurrently.
bucket.consume();
bucket.consume();
const controller = new AbortController();
const abortError = RuntimeError.create('cancelled');
// Abort after a microtask so waitForToken is already suspended in the Promise.
const waitPromise = bucket.waitForToken({ 'signal': controller.signal, 'tokens': 1 });
queueMicrotask(() => { controller.abort(abortError); });
const abortRejected = await waitPromise.then(() => { const result = false; return result; }).catch(() => { const result = true; return result; });
console.log('waitForToken aborted:', abortRejected);

DeadLetterQueue

Bounded FIFO queue for items that failed processing. Drain via async generator.

DeadLetterQueueRetryGenerator: timed re-delivery

ts
import {
  DeadLetterQueue,
  DeadLetterQueueClosedError,
  DeadLetterQueueFullError,
  DeadLetterQueueRetryGenerator
} from '../src/index.js';

await (async function runDeadLetterQueueExample(): Promise<void> {
  // --- Basic enqueue and drain ---
  const deadLetterQueue = DeadLetterQueue.create<string>({ 'capacity': 5 });

  deadLetterQueue.enqueue('job-1', 'timeout');
  deadLetterQueue.enqueue('job-2', 'network error', RuntimeError.create('ECONNREFUSED'));
  console.log('Queue size after 2 enqueues:', deadLetterQueue.size);

  // Close before draining so the generator terminates instead of waiting.
  deadLetterQueue.close();

  const collected: string[] = [];
  for await (const entry of deadLetterQueue.drain()) {
    collected.push(entry.item);
  }
  console.log('Drained items:', collected);
  console.log('Queue size after drain:', deadLetterQueue.size);

  // --- Capacity enforcement ---
  const bounded = DeadLetterQueue.create<number>({ 'capacity': 2 });
  bounded.enqueue(1, 'err');
  bounded.enqueue(2, 'err');
  console.log('Bounded queue size:', bounded.size);

  // --- DeadLetterQueueRetryGenerator re-yields entries with a pause ---
  const retryDeadLetterQueue = DeadLetterQueue.create<string>();
  retryDeadLetterQueue.enqueue('retry-job-1', 'failed');
  retryDeadLetterQueue.enqueue('retry-job-2', 'failed');
  retryDeadLetterQueue.close();

  const generator = DeadLetterQueueRetryGenerator.create({ 'deadLetterQueue': retryDeadLetterQueue, 'intervalMs': 0 });
  const retried: string[] = [];
  for await (const entry of generator.generate()) {
    retried.push(entry.item);
  }
  console.log('Retried items:', retried);

  // --- AbortSignal aborts the queue on construction ---
  const controller = new AbortController();
  controller.abort();
  const abortedDeadLetterQueue = DeadLetterQueue.create<string>({ 'signal': controller.signal });
  const abortedEntries: string[] = [];
  for await (const entry of abortedDeadLetterQueue.drain()) {
    abortedEntries.push(entry.item);
  }
  console.log('Aborted drain count:', abortedEntries.length);

  assert.deepEqual(collected, ['job-1', 'job-2']);
  assert.equal(deadLetterQueue.size, 0);
  assert.throws(() => { deadLetterQueue.enqueue('job-3', 'late'); }, DeadLetterQueueClosedError);
  assert.throws(() => { bounded.enqueue(3, 'overflow'); }, DeadLetterQueueFullError);
  assert.deepEqual(retried, ['retry-job-1', 'retry-job-2']);
  assert.equal(abortedEntries.length, 0);

  console.log('dead-letter-queue: all assertions passed');
})();

Observability hooks

Subclass any primitive and override protected hooks to add logging, metrics, or tracing without coupling the core to any observability library.

CircuitBreaker hooks

HookWhen it firesArgs
onSuccess()After fn() resolves in any state
onFailure(error)After fn() throws in any stateerror: unknown
onTrip()When failure threshold is reached and state transitions closed → open
onOpen()Every time state becomes open (threshold trip or halfOpen → open on failure)
onHalfOpen()When state transitions open → halfOpen after resetTimeoutMs
onClose()When state becomes closed (success threshold reached in halfOpen or manual reset)
onReject()When a call is short-circuited because the circuit is open

TokenBucket hooks

HookWhen it firesArgs
onTokenAcquired(count)After consume() or waitForToken() successfully deducts tokenscount: number
onTokenDepleted()When consume() finds insufficient tokens (before throwing)
onRefill(added)When the internal refill adds tokens due to elapsed timeadded: number

DeadLetterQueue hooks

HookWhen it firesArgs
onEnqueue(item)After an item is added to the queueitem: T
onDequeue(item)After an item is shifted from the queue during drainitem: T
onOverflow()When enqueue() is called on a full queue (before throwing)
onClose()At the end of close()
onAbort()At the end of abort()

DeadLetterQueueRetryGenerator hooks

HookWhen it firesArgs
onYield(entry)Immediately before each entry is yielded from generate()entry: DeadLetterQueueEntryInterface<T>
onWait(intervalMs)Before each inter-entry delayintervalMs: number
onDone()When the generator finishes (DLQ closed or aborted)

CircuitBreaker, DeadLetterQueue, DeadLetterQueueRetryGenerator, and TokenBucket each use an owner-bound, instance-local hook recorder. A lifecycle hook that throws or rejects adds a HookInvocationError to that instance's protected hookErrors array; the entry's cause is the exact thrown or rejected value. Subclasses can inspect hookErrors, and the classes expose no public hook-error getter. Hook failures do not replace the primitive's canonical result or error.

ts
import {
  CircuitBreaker,
  CircuitBreakerOpenError,
  type CircuitBreakerOptionsInterface,
  DeadLetterQueue,
  type DeadLetterQueueOptionsInterface
} from '../src/index.js';

// --- Observed CircuitBreaker ---
class TracedBreaker extends CircuitBreaker {
  constructor(options: CircuitBreakerOptionsInterface) { super(options); }
  protected override onSuccess(): void { console.log('[resilience:cb] onSuccess — circuit closed and call succeeded'); }
  protected override onFailure(error: Error): void { console.log(`[resilience:cb] onFailure — error=${error.message}`); }
  protected override onTrip(): void { console.log('[resilience:cb] onTrip — failure threshold reached, circuit OPEN'); }
  protected override onOpen(): void { console.log('[resilience:cb] onOpen — circuit is now open'); }
  protected override onHalfOpen(): void { console.log('[resilience:cb] onHalfOpen — probing after timeout'); }
  protected override onClose(): void { console.log('[resilience:cb] onClose — circuit CLOSED, service healthy'); }
  protected override onReject(): void { console.log('[resilience:cb] onReject — call short-circuited (circuit open)'); }
}

// --- Observed DeadLetterQueue ---
class TracedDeadLetterQueue<T> extends DeadLetterQueue<T> {
  constructor(options?: DeadLetterQueueOptionsInterface) { super(options); }
  protected override onEnqueue(item: T): void { console.log(`[resilience:dlq] onEnqueue — item=${String(item)}`); }
  protected override onDequeue(item: T): void { console.log(`[resilience:dlq] onDequeue — item=${String(item)}`); }
  protected override onOverflow(): void { console.log('[resilience:dlq] onOverflow — queue full, item dropped'); }
  protected override onClose(): void { console.log('[resilience:dlq] onClose — queue sealed'); }
  protected override onAbort(): void { console.log('[resilience:dlq] onAbort — queue aborted'); }
}

const events: string[] = [];

// Scenario: breaker trips open → half-open → closes; rejected calls go to DLQ
let time = 0;
class Clock {
  static now(): number { const result = time + 0; return result; }
}

const circuitBreaker = new TracedBreaker({ 'clock': Clock.now, 'failureThreshold': 2, 'resetTimeoutMs': 100, 'successThreshold': 1 });
const deadLetterQueue = new TracedDeadLetterQueue<string>({ 'capacity': 10 });

// Trip the breaker open with 2 failures
console.log('\n--- Tripping open ---');
try { await circuitBreaker.execute(() => { throw RuntimeError.create('service down'); }); } catch { deadLetterQueue.enqueue('msg-1', 'service-failure'); events.push('failure-1'); }
try { await circuitBreaker.execute(() => { throw RuntimeError.create('service down'); }); } catch { deadLetterQueue.enqueue('msg-2', 'service-failure'); events.push('failure-2'); }

// Call while open — should be rejected and sent to DLQ
console.log('\n--- Rejecting while open ---');
try {
  await circuitBreaker.execute(() => { const result = Promise.resolve('ok'); return result; });
} catch (error) {
  if (error instanceof CircuitBreakerOpenError) {
    deadLetterQueue.enqueue('msg-3', 'circuit-open');
    events.push('rejected');
  }
}

// Advance clock past resetTimeoutMs → half-open
console.log('\n--- Half-open probe ---');
time = 100;
await circuitBreaker.execute(() => { const result = Promise.resolve('ok'); return result; }); // half-open → closed
events.push('recovered');

// Drain the DLQ
console.log('\n--- Draining DLQ ---');
deadLetterQueue.close();
for await (const entry of deadLetterQueue.drain()) {
  events.push(`drain:${entry.item}`);
}

The base class never calls any logger or metrics library. All hooks are no-ops by default.

Try it

The hooks demo subclasses both CircuitBreaker and DeadLetterQueue and overrides their lifecycle hooks. Watch the full scenario: two failures trigger onFailure, onTrip, and onOpen; a rejected call triggers onReject; advancing the virtual clock into half-open triggers onHalfOpen, onSuccess, and onClose; and DLQ drain emits onDequeue for every item recovered from the queue.

Loading example…

Exports

SymbolPurposeImport path
CircuitBreakerThree-state async circuit breaker.@studnicky/resilience
CircuitBreakerOpenErrorSignals a call rejected by an open circuit.@studnicky/resilience
CircuitBreakerOptionsInterfaceCaller-supplied circuit-breaker options, including clock and error classifier.@studnicky/resilience
DeadLetterQueue<T>Bounded FIFO queue with async-generator drain.@studnicky/resilience
DeadLetterQueueAbortedErrorSignals enqueue after queue abort.@studnicky/resilience
DeadLetterQueueClosedErrorSignals enqueue after queue close.@studnicky/resilience
DeadLetterQueueFullErrorSignals enqueue at queue capacity.@studnicky/resilience
DeadLetterQueueOptionsInterfaceCaller-supplied queue options, including clock and abort signal.@studnicky/resilience
DeadLetterQueueRetryGenerator<T>Re-yields queue entries after a configurable pause.@studnicky/resilience
DeadLetterQueueRetryGeneratorOptionsInterface<T>Caller-supplied retry-generator options with a live queue.@studnicky/resilience
ResilienceConfigErrorSignals invalid resilience configuration.@studnicky/resilience
ResilienceErrorBase error for the package.@studnicky/resilience
TokenBucketToken-bucket rate limiter.@studnicky/resilience
TokenBucketExhaustedErrorSignals insufficient available tokens.@studnicky/resilience
TokenBucketOptionsInterfaceCaller-supplied token-bucket options, including clock.@studnicky/resilience

Entities

@studnicky/resilience/entities exports all schema-backed configuration, state, event, and effect declarations. Each entity namespace provides Schema, Type, and validate.

typescript
import { CircuitBreakerOptionsEntity } from '@studnicky/resilience/entities';

Interfaces

@studnicky/resilience/interfaces exports type-only event, effect, queue-entry, and option contracts. Option interfaces that callers pass to public factories are also available from the package root.

typescript
import type { DeadLetterQueueEntryInterface } from '@studnicky/resilience/interfaces';

CircuitBreaker

MemberSignatureDescription
execute<T>(fn: () => Promise<T>) => Promise<T>Runs fn; throws CircuitBreakerOpenError when open
stateget state(): CircuitStateEntity.TypeCurrent circuit state
reset() => voidRestores the closed state and clears failure counters
forceOpen() => voidForces circuit open

TokenBucket

MemberSignatureDescription
consume(tokens?: number) => voidConsumes tokens; throws TokenBucketExhaustedError if insufficient
waitForToken(options?: { tokens?: number; signal?: AbortSignal }) => Promise<void>Waits until tokens are available, then consumes
availablenumberCurrent token count (triggers a refill calculation)

DeadLetterQueue<T>

MemberSignatureDescription
enqueue(item, reason, error?) => voidAdds item; throws on full, closed, or aborted
drain() => AsyncGenerator<DeadLetterQueueEntryInterface<T>>Yields all entries; suspends when queue is empty
close() => voidSignals drain to stop after the current entries
abort() => voidImmediately stops drain
sizeget size(): numberCurrent entry count

Source on GitHub