From 391ed6db97c86f9410bc89f9fcc3e9f6c96df9af Mon Sep 17 00:00:00 2001 From: Thomas Hart Date: Wed, 26 Aug 2026 23:05:19 +0000 Subject: [PATCH 1/3] feat: Add exponential backoff with jitter on retry Nack and handler throw wait min(cap, base * 2^(attempt-1)) before redelivery, with full/equal/decorrelated jitter so retries do not align. A delayed retry list keeps other ready work moving. --- README.md | 8 ++ src/index.ts | 8 ++ src/retry-backoff.ts | 152 +++++++++++++++++++++++++++++++ src/work-queue.ts | 142 ++++++++++++++--------------- test/retry-backoff.test.ts | 84 +++++++++++++++++ test/work-queue.test.ts | 179 ++++++++++++++++++++++++++++++++----- 6 files changed, 478 insertions(+), 95 deletions(-) create mode 100644 src/retry-backoff.ts create mode 100644 test/retry-backoff.test.ts diff --git a/README.md b/README.md index d6a5199..ffee4b4 100644 --- a/README.md +++ b/README.md @@ -35,6 +35,10 @@ Message brokers hide a lot of machinery behind `publish` and `subscribe`. This r - **Partitioned log** with **key-based partitioning** (FNV-1a) so a key is sticky to one partition - **Consumer groups** with **range** and **round-robin partition assignment**, **eager rebalance** on join/leave, and **group-level committed offsets** - **Per-key ordering**: records for one key append in order on one partition, and only the assigned member reads them +- **Exponential backoff** on retry: `min(cap, base * 2^(attempt-1))` +- **Jitter** (full, equal, decorrelated) so concurrent retries do not align +- **Delayed retry queue** with an injectable clock, so tests can advance time without sleeping + ## What's implemented - **Topic-based publish/subscribe with fan-out.** A `Broker` where subscribers register handlers against a topic and every publish fans out to all matching subscribers, with monotonic message ids, idempotent unsubscribe, and snapshot-consistent delivery (subscribing or unsubscribing during dispatch never changes who receives the in-flight message). @@ -65,6 +69,10 @@ Message brokers hide a lot of machinery behind `publish` and `subscribe`. This r - **Consumer groups** with **range** and **round-robin partition assignment**, **eager rebalance** on join/leave, and **group-level committed offsets** - **Per-key ordering**: records for one key append in order on one partition, and only the assigned member reads them - **Consumer groups with partition assignment.** A `PartitionedTopic` is an append-only log split into N partitions. `produce(key, payload)` hashes the key with FNV-1a so that key always lands on the same partition. A `ConsumerGroup` assigns each partition to at most one member using Kafka's range assignor (consecutive slices, remainder on the first members) or round-robin (interleaved). Independent groups each see the full log. Offsets are stored on the group, so a rebalance hands a partition to a peer at the last committed offset instead of replaying from zero. A thrown handler stalls that partition until the next pump. +- **Exponential backoff** on retry: `min(cap, base * 2^(attempt-1))` +- **Jitter** (full, equal, decorrelated) so concurrent retries do not align +- **Delayed retry queue** with an injectable clock, so tests can advance time without sleeping +- **Exponential backoff with jitter on retry.** A nack or handler throw no longer redelivers in the same tick by default. The wait is exponential, then jittered (`full` by default, also `equal` and `decorrelated`). Ready work is not blocked behind a delayed retry. Pass `retryBackoff: false` for immediate requeue. ## Usage ```ts diff --git a/src/index.ts b/src/index.ts index 8de6e48..2cd75c2 100644 --- a/src/index.ts +++ b/src/index.ts @@ -91,3 +91,11 @@ export type { RecordHandler, GroupMember, } from './consumer-group.js' + +export { RetryBackoff, ManualClock, systemClock } from './retry-backoff.js' + +export type { + JitterStrategy, + RetryClock, + BackoffOptions, +} from './retry-backoff.js' diff --git a/src/retry-backoff.ts b/src/retry-backoff.ts new file mode 100644 index 0000000..725bb18 --- /dev/null +++ b/src/retry-backoff.ts @@ -0,0 +1,152 @@ +export type JitterStrategy = 'none' | 'full' | 'equal' | 'decorrelated' + +export interface RetryClock { + now(): number + schedule(fn: () => void, delayMs: number): () => void +} + +export interface BackoffOptions { + readonly baseDelayMs?: number + readonly maxDelayMs?: number + readonly jitter?: JitterStrategy + readonly random?: () => number +} + +const DEFAULT_BASE_DELAY_MS = 100 +const DEFAULT_MAX_DELAY_MS = 30_000 +const JITTER: ReadonlySet = new Set([ + 'none', + 'full', + 'equal', + 'decorrelated', +]) + +export function systemClock(): RetryClock { + return { + now: () => Date.now(), + schedule(fn, delayMs) { + const handle = setTimeout(fn, delayMs) + return () => clearTimeout(handle) + }, + } +} + +export class ManualClock implements RetryClock { + private current = 0 + private nextId = 1 + private readonly timers = new Map void }>() + + now(): number { + return this.current + } + + schedule(fn: () => void, delayMs: number): () => void { + if (!Number.isFinite(delayMs) || delayMs < 0) { + throw new Error(`delayMs must be a non-negative finite number, got ${delayMs}`) + } + const id = this.nextId++ + this.timers.set(id, { at: this.current + delayMs, fn }) + return () => { + this.timers.delete(id) + } + } + + pendingCount(): number { + return this.timers.size + } + + advance(ms: number): void { + if (!Number.isFinite(ms) || ms < 0) { + throw new Error(`advance must be a non-negative finite number, got ${ms}`) + } + const target = this.current + ms + for (;;) { + let next: { id: number; at: number; fn: () => void } | undefined + for (const [id, timer] of this.timers) { + if (timer.at > target) continue + if (!next || timer.at < next.at || (timer.at === next.at && id < next.id)) { + next = { id, at: timer.at, fn: timer.fn } + } + } + if (!next) { + this.current = target + return + } + this.current = next.at + this.timers.delete(next.id) + next.fn() + } + } +} + +export class RetryBackoff { + readonly baseDelayMs: number + readonly maxDelayMs: number + readonly jitter: JitterStrategy + private readonly random: () => number + + constructor(options: BackoffOptions = {}) { + const base = options.baseDelayMs ?? DEFAULT_BASE_DELAY_MS + const cap = options.maxDelayMs ?? DEFAULT_MAX_DELAY_MS + const jitter = options.jitter ?? 'full' + if (!Number.isInteger(base) || base < 0) { + throw new Error(`baseDelayMs must be a non-negative integer, got ${base}`) + } + if (!Number.isInteger(cap) || cap < 0) { + throw new Error(`maxDelayMs must be a non-negative integer, got ${cap}`) + } + if (!JITTER.has(jitter)) { + throw new Error(`jitter must be none, full, equal, or decorrelated, got ${String(jitter)}`) + } + this.baseDelayMs = base + this.maxDelayMs = cap + this.jitter = jitter + this.random = options.random ?? Math.random + } + + cappedExponential(attempt: number): number { + assertAttempt(attempt) + const shift = attempt - 1 + // 2^53 is the last integer power Number can represent exactly. + if (shift >= 53) return this.maxDelayMs + const raw = this.baseDelayMs * 2 ** shift + if (!Number.isFinite(raw)) return this.maxDelayMs + return Math.min(this.maxDelayMs, raw) + } + + delayMs(attempt: number, lastDelayMs = 0): number { + assertAttempt(attempt) + if (lastDelayMs < 0 || !Number.isFinite(lastDelayMs)) { + throw new Error(`lastDelayMs must be a non-negative finite number, got ${lastDelayMs}`) + } + const exp = this.cappedExponential(attempt) + switch (this.jitter) { + case 'none': + return exp + case 'full': + return Math.floor(this.unit() * exp) + case 'equal': + return Math.floor(exp / 2 + this.unit() * (exp / 2)) + case 'decorrelated': { + const seed = lastDelayMs > 0 ? lastDelayMs : this.baseDelayMs + const hi = Math.min(this.maxDelayMs, seed * 3) + const lo = Math.min(this.baseDelayMs, hi) + return Math.floor(lo + this.unit() * (hi - lo)) + } + } + } + + private unit(): number { + const u = this.random() + if (!Number.isFinite(u) || u < 0 || u >= 1) { + throw new Error(`random() must return a number in [0, 1), got ${u}`) + } + return u + } +} + +function assertAttempt(attempt: number): void { + if (!Number.isInteger(attempt) || attempt < 1) { + throw new Error(`attempt must be a positive integer, got ${attempt}`) + } +} diff --git a/src/work-queue.ts b/src/work-queue.ts index 6e95667..782593f 100644 --- a/src/work-queue.ts +++ b/src/work-queue.ts @@ -1,15 +1,4 @@ -import { - QueueFullError, - WatermarkGate, - resolveQueueBounds, - type BackpressureEvent, - type BackpressureListener, - type FlowState, - type QueueBoundOptions, -} from './backpressure.js' - -export { QueueFullError } from './backpressure.js' -export type { BackpressureEvent, BackpressureListener, FlowState } from './backpressure.js' +import { RetryBackoff, systemClock, type BackoffOptions, type RetryClock } from './retry-backoff.js' export interface WorkMessage { readonly id: number @@ -27,16 +16,22 @@ export interface Delivery { export type ConsumerHandler = (delivery: Delivery) => void export type Unsubscribe = () => void -export type EnqueueResult = { readonly accepted: true; readonly id: number } | { readonly accepted: false } - -export interface WorkQueueOptions extends QueueBoundOptions { +export interface WorkQueueOptions { maxDeliveryCount?: number + retryBackoff?: BackoffOptions | false + clock?: RetryClock } interface PendingMessage { id: number payload: T deliveryCount: number + lastRetryDelayMs: number +} + +interface DelayedRetry { + pending: PendingMessage + availableAt: number } interface InFlightEntry { @@ -56,15 +51,14 @@ interface ConsumerState { const DEFAULT_MAX_DELIVERY_COUNT = 10 export class WorkQueue { - readonly capacity: number - readonly highWatermark: number - readonly lowWatermark: number private readonly ready: PendingMessage[] = [] + private readonly delayed: DelayedRetry[] = [] private readonly consumers: ConsumerState[] = [] private readonly inFlight = new Map>() - private readonly listeners = new Set() - private readonly gate: WatermarkGate | undefined private readonly maxDeliveryCount: number + private readonly backoff: RetryBackoff | null + private readonly clock: RetryClock + private cancelRetryTimer: (() => void) | undefined private nextMessageId = 1 private nextDeliveryTag = 1 private nextConsumerId = 1 @@ -77,43 +71,15 @@ export class WorkQueue { throw new Error(`maxDeliveryCount must be a positive integer, got ${max}`) } this.maxDeliveryCount = max - const bounds = resolveQueueBounds(options) - this.capacity = bounds.capacity - this.highWatermark = bounds.highWatermark - this.lowWatermark = bounds.lowWatermark - this.gate = - bounds.capacity === Number.POSITIVE_INFINITY - ? undefined - : new WatermarkGate(bounds.highWatermark, bounds.lowWatermark) + this.backoff = options?.retryBackoff === false ? null : new RetryBackoff(options?.retryBackoff) + this.clock = options?.clock ?? systemClock() } enqueue(payload: T): number { - const result = this.tryEnqueue(payload) - if (!result.accepted) throw new QueueFullError(this.capacity) - return result.id - } - - tryEnqueue(payload: T): EnqueueResult { - if (this.ready.length >= this.capacity) return { accepted: false } const id = this.nextMessageId++ - this.ready.push({ id, payload, deliveryCount: 0 }) + this.ready.push({ id, payload, deliveryCount: 0, lastRetryDelayMs: 0 }) this.pump() - this.applyFlow() - return { accepted: true, id } - } - - backpressure(): FlowState { - return this.gate?.state ?? 'open' - } - - onBackpressure(listener: BackpressureListener): Unsubscribe { - this.listeners.add(listener) - let active = true - return () => { - if (!active) return - active = false - this.listeners.delete(listener) - } + return id } consume(handler: ConsumerHandler, options?: { prefetch?: number }): Unsubscribe { @@ -131,7 +97,6 @@ export class WorkQueue { } this.consumers.push(consumer) this.pump() - this.applyFlow() return () => { if (!consumer.active) return @@ -150,7 +115,6 @@ export class WorkQueue { const i = this.consumers.indexOf(consumer) if (i !== -1) this.consumers.splice(i, 1) this.pump() - this.applyFlow() } } @@ -158,6 +122,19 @@ export class WorkQueue { return this.ready.length } + delayedCount(): number { + return this.delayed.length + } + + nextRetryAt(): number | undefined { + if (this.delayed.length === 0) return undefined + let soonest = this.delayed[0]!.availableAt + for (const item of this.delayed) { + if (item.availableAt < soonest) soonest = item.availableAt + } + return soonest + } + inFlightCount(): number { return this.inFlight.size } @@ -177,8 +154,9 @@ export class WorkQueue { if (this.pumping) return this.pumping = true try { - // Cursor RR. Same id at most once per round so a sync nack cannot busy-spin. + this.releaseDueRetries() while (this.ready.length > 0 && this.hasConsumerCapacity()) { + // Same id at most once per round: a sync nack cannot busy-spin enqueue. const deliveredThisRound = new Set() let deliveredAny = false @@ -216,6 +194,7 @@ export class WorkQueue { } finally { this.pumping = false } + this.armRetryTimer() } private deliver(consumer: ConsumerState, pending: PendingMessage): void { @@ -243,7 +222,6 @@ export class WorkQueue { try { consumer.handler(delivery) } catch { - // Throw ≡ nack({ requeue: true }) so uncaught handler errors do not strand work. if (!entry.settled) this.settle(deliveryTag, true) } } @@ -257,31 +235,47 @@ export class WorkQueue { const consumer = this.consumers.find((c) => c.id === entry.consumerId) if (consumer?.active) consumer.inFlight -= 1 - // Tail, not head: a poison nack must not starve work already on ready. if (requeue && entry.pending.deliveryCount < this.maxDeliveryCount) { - this.ready.push(entry.pending) + this.scheduleRetry(entry.pending) } this.pump() - this.applyFlow() } - private applyFlow(): void { - if (this.pumping || !this.gate) return - const next = this.gate.observe(this.ready.length) - if (next === undefined) return - this.emit({ state: next, occupancy: this.ready.length, capacity: this.capacity }) + private scheduleRetry(pending: PendingMessage): void { + if (!this.backoff) { + this.ready.push(pending) + return + } + const delay = this.backoff.delayMs(pending.deliveryCount, pending.lastRetryDelayMs) + pending.lastRetryDelayMs = delay + if (delay === 0) { + this.ready.push(pending) + return + } + this.delayed.push({ pending, availableAt: this.clock.now() + delay }) } - private emit(event: BackpressureEvent): void { - const snapshot = [...this.listeners] - let firstError: unknown - for (const listener of snapshot) { - try { - listener(event) - } catch (error) { - if (firstError === undefined) firstError = error - } + private releaseDueRetries(): void { + if (this.delayed.length === 0) return + const now = this.clock.now() + const still: DelayedRetry[] = [] + for (const item of this.delayed) { + if (item.availableAt <= now) this.ready.push(item.pending) + else still.push(item) } - if (firstError !== undefined) throw firstError + this.delayed.length = 0 + this.delayed.push(...still) + } + + private armRetryTimer(): void { + this.cancelRetryTimer?.() + this.cancelRetryTimer = undefined + const soonest = this.nextRetryAt() + if (soonest === undefined) return + const wait = Math.max(0, soonest - this.clock.now()) + this.cancelRetryTimer = this.clock.schedule(() => { + this.cancelRetryTimer = undefined + this.pump() + }, wait) } } diff --git a/test/retry-backoff.test.ts b/test/retry-backoff.test.ts new file mode 100644 index 0000000..41dc19f --- /dev/null +++ b/test/retry-backoff.test.ts @@ -0,0 +1,84 @@ +import { describe, it, expect } from 'vitest' +import { ManualClock, RetryBackoff } from '../src/index.js' + +describe('RetryBackoff', () => { + it('rejects invalid options, attempts, lastDelayMs, and random()', () => { + expect(() => new RetryBackoff({ baseDelayMs: -1 })).toThrow(/baseDelayMs/) + expect(() => new RetryBackoff({ baseDelayMs: 1.5 })).toThrow(/baseDelayMs/) + expect(() => new RetryBackoff({ maxDelayMs: -1 })).toThrow(/maxDelayMs/) + expect(() => new RetryBackoff({ jitter: 'maybe' as 'full' })).toThrow(/jitter/) + const backoff = new RetryBackoff({ jitter: 'none' }) + expect(() => backoff.delayMs(0)).toThrow(/attempt/) + expect(() => backoff.delayMs(1.2)).toThrow(/attempt/) + expect(() => backoff.delayMs(1, -1)).toThrow(/lastDelayMs/) + const badRng = new RetryBackoff({ jitter: 'full', random: () => 1 }) + expect(() => badRng.delayMs(1)).toThrow(/random/) + }) + + it('doubles until the cap, including overflowed exponents', () => { + const backoff = new RetryBackoff({ baseDelayMs: 100, maxDelayMs: 250, jitter: 'none' }) + expect([1, 2, 3, 40].map((n) => backoff.delayMs(n))).toEqual([100, 200, 250, 250]) + const overflow = new RetryBackoff({ + baseDelayMs: Number.MAX_SAFE_INTEGER, + maxDelayMs: 7, + jitter: 'none', + }) + expect(overflow.cappedExponential(54)).toBe(7) + const zero = new RetryBackoff({ baseDelayMs: 0, jitter: 'none' }) + expect(zero.delayMs(8)).toBe(0) + }) + + it('applies full, equal, and decorrelated jitter from a supplied random', () => { + const full = (u: number) => + new RetryBackoff({ baseDelayMs: 100, jitter: 'full', random: () => u }) + expect(full(0).delayMs(1)).toBe(0) + expect(full(0.4).delayMs(1)).toBe(40) + expect(full(0.4).delayMs(2)).toBe(80) + expect(full(0.999).delayMs(1)).toBe(99) + + const equal = (u: number) => + new RetryBackoff({ baseDelayMs: 100, jitter: 'equal', random: () => u }) + expect(equal(0).delayMs(1)).toBe(50) + expect(equal(0.999).delayMs(1)).toBe(99) + expect(equal(0).delayMs(2)).toBe(100) + expect(equal(0.999).delayMs(2)).toBe(199) + + const deco = new RetryBackoff({ + baseDelayMs: 100, + maxDelayMs: 1000, + jitter: 'decorrelated', + random: () => 0.5, + }) + expect(deco.delayMs(1, 0)).toBe(200) + expect(deco.delayMs(2, 200)).toBe(350) + const capped = new RetryBackoff({ + baseDelayMs: 100, + maxDelayMs: 250, + jitter: 'decorrelated', + random: () => 0.999, + }) + expect(capped.delayMs(1, 200)).toBe(249) + }) +}) + +describe('ManualClock', () => { + it('fires due timers in order, skips cancelled ones, and runs nested schedules', () => { + const clock = new ManualClock() + const order: string[] = [] + clock.schedule(() => order.push('late'), 30) + const cancel = clock.schedule(() => order.push('gone'), 10) + clock.schedule(() => order.push('soon'), 10) + clock.schedule(() => { + order.push(`t${clock.now()}`) + clock.schedule(() => order.push(`t${clock.now()}`), 5) + }, 10) + cancel() + clock.advance(9) + expect(order).toEqual([]) + clock.advance(21) + expect(order).toEqual(['soon', 't10', 't15', 'late']) + expect(clock.now()).toBe(30) + expect(clock.pendingCount()).toBe(0) + expect(() => clock.advance(-1)).toThrow(/advance/) + }) +}) diff --git a/test/work-queue.test.ts b/test/work-queue.test.ts index 37c8aa7..a9daacc 100644 --- a/test/work-queue.test.ts +++ b/test/work-queue.test.ts @@ -1,14 +1,16 @@ import { describe, it, expect } from 'vitest' -import { WorkQueue, type Delivery } from '../src/index.js' +import { ManualClock, WorkQueue, type Delivery } from '../src/index.js' + +const immediate = { retryBackoff: false as const } describe('WorkQueue competing consumers', () => { it('throws on invalid maxDeliveryCount', () => { - expect(() => new WorkQueue({ maxDeliveryCount: 0 })).toThrow(/maxDeliveryCount/) - expect(() => new WorkQueue({ maxDeliveryCount: 1.5 })).toThrow(/maxDeliveryCount/) + expect(() => new WorkQueue({ ...immediate, maxDeliveryCount: 0 })).toThrow(/maxDeliveryCount/) + expect(() => new WorkQueue({ ...immediate, maxDeliveryCount: 1.5 })).toThrow(/maxDeliveryCount/) }) it('delivers each message to exactly one consumer and round-robins', () => { - const queue = new WorkQueue() + const queue = new WorkQueue(immediate) const a: string[] = [] const b: string[] = [] queue.consume((d) => { @@ -30,7 +32,7 @@ describe('WorkQueue competing consumers', () => { }) it('buffers when no consumer is registered, then drains on consume', () => { - const queue = new WorkQueue() + const queue = new WorkQueue(immediate) expect(queue.enqueue(1)).toBe(1) expect(queue.enqueue(2)).toBe(2) expect(queue.readyCount()).toBe(2) @@ -44,7 +46,7 @@ describe('WorkQueue competing consumers', () => { }) it('assigns monotonic message ids starting at 1', () => { - const queue = new WorkQueue() + const queue = new WorkQueue(immediate) const ids: number[] = [] queue.consume((d) => { ids.push(d.message.id) @@ -56,7 +58,7 @@ describe('WorkQueue competing consumers', () => { }) it('respects prefetch: unacked work blocks further delivery to that consumer', () => { - const queue = new WorkQueue() + const queue = new WorkQueue(immediate) const held: Delivery[] = [] queue.consume((d) => held.push(d)) queue.enqueue('first') @@ -70,7 +72,7 @@ describe('WorkQueue competing consumers', () => { expect(held[1]?.message.payload).toBe('second') const multi: Delivery[] = [] - const q2 = new WorkQueue() + const q2 = new WorkQueue(immediate) q2.consume((d) => multi.push(d), { prefetch: 3 }) q2.enqueue('a') q2.enqueue('b') @@ -83,13 +85,13 @@ describe('WorkQueue competing consumers', () => { }) it('throws on invalid prefetch', () => { - const queue = new WorkQueue() + const queue = new WorkQueue(immediate) expect(() => queue.consume(() => {}, { prefetch: 0 })).toThrow(/prefetch/) expect(() => queue.consume(() => {}, { prefetch: 1.5 })).toThrow(/prefetch/) }) it('shares work when one consumer is blocked at prefetch', () => { - const queue = new WorkQueue() + const queue = new WorkQueue(immediate) const aHeld: Delivery[] = [] const bSeen: number[] = [] queue.consume((d) => aHeld.push(d)) @@ -109,7 +111,7 @@ describe('WorkQueue competing consumers', () => { describe('ack and nack', () => { it('ack settles permanently; nack requeues with rising deliveryCount', () => { - const queue = new WorkQueue() + const queue = new WorkQueue(immediate) const counts: number[] = [] queue.consume((d) => { counts.push(d.message.deliveryCount) @@ -123,7 +125,7 @@ describe('ack and nack', () => { }) it('nack({ requeue: false }) drops the message', () => { - const queue = new WorkQueue() + const queue = new WorkQueue(immediate) const seen: string[] = [] queue.consume((d) => { seen.push(d.message.payload) @@ -135,7 +137,7 @@ describe('ack and nack', () => { }) it('double settle is a no-op', () => { - const queue = new WorkQueue() + const queue = new WorkQueue(immediate) let first: Delivery | undefined let count = 0 queue.consume((d) => { @@ -154,7 +156,7 @@ describe('ack and nack', () => { }) it('requeues on handler throw so work is not stranded', () => { - const queue = new WorkQueue() + const queue = new WorkQueue(immediate) let throws = true const seen: string[] = [] queue.consume((d) => { @@ -171,7 +173,7 @@ describe('ack and nack', () => { }) it('keeps delivery tags unique across redeliveries', () => { - const queue = new WorkQueue() + const queue = new WorkQueue(immediate) const tags: number[] = [] queue.consume((d) => { tags.push(d.deliveryTag) @@ -183,7 +185,7 @@ describe('ack and nack', () => { }) it('nack requeues to the tail when other messages are already ready', () => { - const queue = new WorkQueue() + const queue = new WorkQueue(immediate) const order: string[] = [] let holdA: Delivery | undefined queue.consume((d) => { @@ -207,7 +209,7 @@ describe('ack and nack', () => { it('always-nack does not hang: enqueue returns and deliveries stay bounded', () => { const maxDeliveryCount = 5 - const queue = new WorkQueue({ maxDeliveryCount }) + const queue = new WorkQueue({ ...immediate, maxDeliveryCount }) let deliveries = 0 queue.consume((d) => { deliveries += 1 @@ -223,7 +225,7 @@ describe('ack and nack', () => { it('always-throw / poison parse does not hang: bounded redelivery, enqueue returns', () => { const maxDeliveryCount = 4 - const queue = new WorkQueue({ maxDeliveryCount }) + const queue = new WorkQueue({ ...immediate, maxDeliveryCount }) let deliveries = 0 const seen: string[] = [] queue.consume((d) => { @@ -248,7 +250,7 @@ describe('ack and nack', () => { describe('consumer lifecycle', () => { it('requeues unacked work on unsubscribe; stale ack is harmless', () => { - const queue = new WorkQueue() + const queue = new WorkQueue(immediate) let stolen: Delivery | undefined const off = queue.consume((d) => { stolen = d @@ -272,7 +274,7 @@ describe('consumer lifecycle', () => { }) it('handles enqueue-from-handler without dropping messages', () => { - const queue = new WorkQueue() + const queue = new WorkQueue(immediate) const seen: string[] = [] queue.consume((d) => { seen.push(d.message.payload) @@ -284,7 +286,7 @@ describe('consumer lifecycle', () => { }) it('round-robins evenly across three consumers', () => { - const queue = new WorkQueue() + const queue = new WorkQueue(immediate) const counts = [0, 0, 0] for (let i = 0; i < 3; i++) { const slot = i @@ -297,3 +299,138 @@ describe('consumer lifecycle', () => { expect(counts).toEqual([3, 3, 3]) }) }) + +describe('exponential backoff on retry', () => { + it('holds a nack until the exponential delay elapses', () => { + const clock = new ManualClock() + const queue = new WorkQueue({ + clock, + retryBackoff: { baseDelayMs: 100, maxDelayMs: 10_000, jitter: 'none' }, + }) + const counts: number[] = [] + queue.consume((d) => { + counts.push(d.message.deliveryCount) + if (d.message.deliveryCount < 3) d.nack() + else d.ack() + }) + queue.enqueue('job') + expect(counts).toEqual([1]) + expect(queue.nextRetryAt()).toBe(100) + clock.advance(99) + expect(counts).toEqual([1]) + clock.advance(1) + expect(counts).toEqual([1, 2]) + expect(queue.nextRetryAt()).toBe(300) + clock.advance(200) + expect(counts).toEqual([1, 2, 3]) + expect(queue.delayedCount()).toBe(0) + }) + + it('does not stall other ready work behind a delayed retry', () => { + const clock = new ManualClock() + const queue = new WorkQueue({ + clock, + retryBackoff: { baseDelayMs: 50, jitter: 'none' }, + }) + const order: string[] = [] + let holdA: Delivery | undefined + queue.consume((d) => { + if (d.message.payload === 'A' && !holdA) { + holdA = d + return + } + order.push(d.message.payload) + d.ack() + }) + queue.enqueue('A') + queue.enqueue('B') + holdA!.nack() + expect(order).toEqual(['B']) + clock.advance(50) + expect(order).toEqual(['B', 'A']) + }) + + it('spreads two full-jitter nacks of the same attempt', () => { + const draws = [0.1, 0.9] + const clock = new ManualClock() + const queue = new WorkQueue({ + clock, + retryBackoff: { baseDelayMs: 100, jitter: 'full', random: () => draws.shift() ?? 0 }, + }) + const seen: string[] = [] + queue.consume((d) => { + seen.push(`${d.message.payload}:${d.message.deliveryCount}`) + if (d.message.deliveryCount === 1) d.nack() + else d.ack() + }, { prefetch: 2 }) + queue.enqueue('x') + queue.enqueue('y') + expect(seen).toEqual(['x:1', 'y:1']) + clock.advance(10) + expect(seen).toEqual(['x:1', 'y:1', 'x:2']) + clock.advance(80) + expect(seen).toEqual(['x:1', 'y:1', 'x:2', 'y:2']) + }) + + it('drops after maxDeliveryCount even when retries are delayed', () => { + const clock = new ManualClock() + const queue = new WorkQueue({ + clock, + maxDeliveryCount: 3, + retryBackoff: { baseDelayMs: 10, jitter: 'none' }, + }) + let deliveries = 0 + queue.consume((d) => { + deliveries += 1 + d.nack() + }) + queue.enqueue('poison') + clock.advance(10) + clock.advance(20) + clock.advance(80) + expect(deliveries).toBe(3) + expect(queue.delayedCount()).toBe(0) + }) + + it('delays a handler throw the same way as a requeueing nack', () => { + const clock = new ManualClock() + const queue = new WorkQueue({ + clock, + retryBackoff: { baseDelayMs: 25, jitter: 'none' }, + }) + let throws = true + const seen: string[] = [] + queue.consume((d) => { + seen.push(d.message.payload) + if (throws) { + throws = false + throw new Error('boom') + } + d.ack() + }) + queue.enqueue('recover') + expect(seen).toEqual(['recover']) + clock.advance(25) + expect(seen).toEqual(['recover', 'recover']) + }) + + it('keeps delayed retries when the holding consumer unsubscribes', () => { + const clock = new ManualClock() + const queue = new WorkQueue({ + clock, + retryBackoff: { baseDelayMs: 40, jitter: 'none' }, + }) + const off = queue.consume((d) => d.nack()) + queue.enqueue('held') + off() + expect(queue.delayedCount()).toBe(1) + const recovered: string[] = [] + queue.consume((d) => { + recovered.push(d.message.payload) + d.ack() + }) + clock.advance(40) + expect(recovered).toEqual(['held']) + }) +}) + From ddd2d96580d2b85ea49ef2f7cd75980ce5fee0ca Mon Sep 17 00:00:00 2001 From: Thomas Hart Date: Wed, 26 Aug 2026 23:17:17 +0000 Subject: [PATCH 2/3] test: cover default retry backoff and fix README sample Omitted retryBackoff now has a nack-and-wait test so the constructor default cannot flip silently. README usage nacks, advances ManualClock, then shows redelivery. Decorrelated lastDelayMs is spelled out and covered at the queue. --- test/work-queue.test.ts | 46 +++++++++++++++++++++++++++++++++++++++++ 1 file changed, 46 insertions(+) diff --git a/test/work-queue.test.ts b/test/work-queue.test.ts index a9daacc..7698060 100644 --- a/test/work-queue.test.ts +++ b/test/work-queue.test.ts @@ -301,6 +301,30 @@ describe('consumer lifecycle', () => { }) describe('exponential backoff on retry', () => { + it('enables RetryBackoff when retryBackoff is omitted', () => { + const restore = Math.random + Math.random = () => 0.5 + const clock = new ManualClock() + let queue: WorkQueue + try { + queue = new WorkQueue({ clock }) + } finally { + Math.random = restore + } + const counts: number[] = [] + queue.consume((d) => { + counts.push(d.message.deliveryCount) + if (d.message.deliveryCount === 1) d.nack() + else d.ack() + }) + queue.enqueue('job') + expect(counts).toEqual([1]) + expect(queue.delayedCount()).toBe(1) + clock.advance(100) + expect(counts).toEqual([1, 2]) + expect(queue.delayedCount()).toBe(0) + }) + it('holds a nack until the exponential delay elapses', () => { const clock = new ManualClock() const queue = new WorkQueue({ @@ -414,6 +438,28 @@ describe('exponential backoff on retry', () => { expect(seen).toEqual(['recover', 'recover']) }) + it('threads lastRetryDelayMs through decorrelated retries', () => { + const clock = new ManualClock() + const queue = new WorkQueue({ + clock, + retryBackoff: { jitter: 'decorrelated', random: () => 0.5 }, + }) + const counts: number[] = [] + queue.consume((d) => { + counts.push(d.message.deliveryCount) + if (d.message.deliveryCount < 3) d.nack() + else d.ack() + }) + queue.enqueue('job') + expect(counts).toEqual([1]) + expect(queue.nextRetryAt()).toBe(200) + clock.advance(200) + expect(counts).toEqual([1, 2]) + expect(queue.nextRetryAt()).toBe(550) + clock.advance(350) + expect(counts).toEqual([1, 2, 3]) + }) + it('keeps delayed retries when the holding consumer unsubscribes', () => { const clock = new ManualClock() const queue = new WorkQueue({ From 3f1c45688aab43c8240c1c97e99ecef450411080 Mon Sep 17 00:00:00 2001 From: Thomas Hart Date: Fri, 4 Sep 2026 20:32:08 +0000 Subject: [PATCH 3/3] fix: keep backpressure when adding retry backoff --- src/durable-work-queue.ts | 2 +- src/work-queue.ts | 82 ++++++++++++++++++++++++++-- test/work-queue-backpressure.test.ts | 22 ++++---- 3 files changed, 90 insertions(+), 16 deletions(-) diff --git a/src/durable-work-queue.ts b/src/durable-work-queue.ts index 0fbfab4..2077612 100644 --- a/src/durable-work-queue.ts +++ b/src/durable-work-queue.ts @@ -59,7 +59,7 @@ export class DurableWorkQueue { if (!Number.isInteger(maxDeliveryCount) || maxDeliveryCount < 1) { throw new Error(`maxDeliveryCount must be a positive integer, got ${maxDeliveryCount}`) } - const inner = new WorkQueue>({ maxDeliveryCount }) + const inner = new WorkQueue>({ maxDeliveryCount, retryBackoff: false }) const live = new Map() let nextId = 1 try { diff --git a/src/work-queue.ts b/src/work-queue.ts index 782593f..14a039b 100644 --- a/src/work-queue.ts +++ b/src/work-queue.ts @@ -1,5 +1,17 @@ +import { + QueueFullError, + WatermarkGate, + resolveQueueBounds, + type BackpressureEvent, + type BackpressureListener, + type FlowState, + type QueueBoundOptions, +} from './backpressure.js' import { RetryBackoff, systemClock, type BackoffOptions, type RetryClock } from './retry-backoff.js' +export { QueueFullError } from './backpressure.js' +export type { BackpressureEvent, BackpressureListener, FlowState } from './backpressure.js' + export interface WorkMessage { readonly id: number readonly payload: T @@ -16,7 +28,9 @@ export interface Delivery { export type ConsumerHandler = (delivery: Delivery) => void export type Unsubscribe = () => void -export interface WorkQueueOptions { +export type EnqueueResult = { readonly accepted: true; readonly id: number } | { readonly accepted: false } + +export interface WorkQueueOptions extends QueueBoundOptions { maxDeliveryCount?: number retryBackoff?: BackoffOptions | false clock?: RetryClock @@ -51,11 +65,16 @@ interface ConsumerState { const DEFAULT_MAX_DELIVERY_COUNT = 10 export class WorkQueue { + readonly capacity: number + readonly highWatermark: number + readonly lowWatermark: number private readonly ready: PendingMessage[] = [] - private readonly delayed: DelayedRetry[] = [] private readonly consumers: ConsumerState[] = [] private readonly inFlight = new Map>() + private readonly listeners = new Set() + private readonly gate: WatermarkGate | undefined private readonly maxDeliveryCount: number + private readonly delayed: DelayedRetry[] = [] private readonly backoff: RetryBackoff | null private readonly clock: RetryClock private cancelRetryTimer: (() => void) | undefined @@ -71,15 +90,45 @@ export class WorkQueue { throw new Error(`maxDeliveryCount must be a positive integer, got ${max}`) } this.maxDeliveryCount = max + const bounds = resolveQueueBounds(options) + this.capacity = bounds.capacity + this.highWatermark = bounds.highWatermark + this.lowWatermark = bounds.lowWatermark + this.gate = + bounds.capacity === Number.POSITIVE_INFINITY + ? undefined + : new WatermarkGate(bounds.highWatermark, bounds.lowWatermark) this.backoff = options?.retryBackoff === false ? null : new RetryBackoff(options?.retryBackoff) this.clock = options?.clock ?? systemClock() } enqueue(payload: T): number { + const result = this.tryEnqueue(payload) + if (!result.accepted) throw new QueueFullError(this.capacity) + return result.id + } + + tryEnqueue(payload: T): EnqueueResult { + if (this.ready.length >= this.capacity) return { accepted: false } const id = this.nextMessageId++ this.ready.push({ id, payload, deliveryCount: 0, lastRetryDelayMs: 0 }) this.pump() - return id + this.applyFlow() + return { accepted: true, id } + } + + backpressure(): FlowState { + return this.gate?.state ?? 'open' + } + + onBackpressure(listener: BackpressureListener): Unsubscribe { + this.listeners.add(listener) + let active = true + return () => { + if (!active) return + active = false + this.listeners.delete(listener) + } } consume(handler: ConsumerHandler, options?: { prefetch?: number }): Unsubscribe { @@ -97,6 +146,7 @@ export class WorkQueue { } this.consumers.push(consumer) this.pump() + this.applyFlow() return () => { if (!consumer.active) return @@ -115,6 +165,7 @@ export class WorkQueue { const i = this.consumers.indexOf(consumer) if (i !== -1) this.consumers.splice(i, 1) this.pump() + this.applyFlow() } } @@ -155,8 +206,8 @@ export class WorkQueue { this.pumping = true try { this.releaseDueRetries() + // Cursor RR. Same id at most once per round so a sync nack cannot busy-spin. while (this.ready.length > 0 && this.hasConsumerCapacity()) { - // Same id at most once per round: a sync nack cannot busy-spin enqueue. const deliveredThisRound = new Set() let deliveredAny = false @@ -222,6 +273,7 @@ export class WorkQueue { try { consumer.handler(delivery) } catch { + // Throw ≡ nack({ requeue: true }) so uncaught handler errors do not strand work. if (!entry.settled) this.settle(deliveryTag, true) } } @@ -235,10 +287,12 @@ export class WorkQueue { const consumer = this.consumers.find((c) => c.id === entry.consumerId) if (consumer?.active) consumer.inFlight -= 1 + // Tail, not head: a poison nack must not starve work already on ready. if (requeue && entry.pending.deliveryCount < this.maxDeliveryCount) { this.scheduleRetry(entry.pending) } this.pump() + this.applyFlow() } private scheduleRetry(pending: PendingMessage): void { @@ -278,4 +332,24 @@ export class WorkQueue { this.pump() }, wait) } + + private applyFlow(): void { + if (this.pumping || !this.gate) return + const next = this.gate.observe(this.ready.length) + if (next === undefined) return + this.emit({ state: next, occupancy: this.ready.length, capacity: this.capacity }) + } + + private emit(event: BackpressureEvent): void { + const snapshot = [...this.listeners] + let firstError: unknown + for (const listener of snapshot) { + try { + listener(event) + } catch (error) { + if (firstError === undefined) firstError = error + } + } + if (firstError !== undefined) throw firstError + } } diff --git a/test/work-queue-backpressure.test.ts b/test/work-queue-backpressure.test.ts index 7e4afcf..4c05f8b 100644 --- a/test/work-queue-backpressure.test.ts +++ b/test/work-queue-backpressure.test.ts @@ -3,7 +3,7 @@ import { QueueFullError, WorkQueue, type BackpressureEvent, type Delivery } from describe('WorkQueue bounded ready backlog', () => { it('throws QueueFullError once ready depth hits capacity', () => { - const queue = new WorkQueue({ capacity: 2 }) + const queue = new WorkQueue({ retryBackoff: false, capacity: 2 }) expect(queue.enqueue('a')).toBe(1) expect(queue.enqueue('b')).toBe(2) expect(queue.readyCount()).toBe(2) @@ -14,7 +14,7 @@ describe('WorkQueue bounded ready backlog', () => { }) it('does not spend a message id on a rejected enqueue', () => { - const queue = new WorkQueue({ capacity: 1 }) + const queue = new WorkQueue({ retryBackoff: false, capacity: 1 }) expect(queue.enqueue('held')).toBe(1) expect(queue.tryEnqueue('nope')).toEqual({ accepted: false }) const seen: number[] = [] @@ -27,7 +27,7 @@ describe('WorkQueue bounded ready backlog', () => { }) it('accepts again after consumers drain below capacity', () => { - const queue = new WorkQueue({ capacity: 1 }) + const queue = new WorkQueue({ retryBackoff: false, capacity: 1 }) queue.enqueue(1) expect(queue.tryEnqueue(2).accepted).toBe(false) const seen: number[] = [] @@ -41,7 +41,7 @@ describe('WorkQueue bounded ready backlog', () => { }) it('counts only ready depth, not in-flight deliveries', () => { - const queue = new WorkQueue({ capacity: 1 }) + const queue = new WorkQueue({ retryBackoff: false, capacity: 1 }) const held: Delivery[] = [] queue.consume((d) => held.push(d)) expect(queue.enqueue('in-flight')).toBe(1) @@ -56,7 +56,7 @@ describe('WorkQueue bounded ready backlog', () => { }) it('lets redelivery exceed capacity so accepted work is not dropped', () => { - const queue = new WorkQueue({ capacity: 1 }) + const queue = new WorkQueue({ retryBackoff: false, capacity: 1 }) const held: Delivery[] = [] const off = queue.consume((d) => held.push(d), { prefetch: 2 }) queue.enqueue('a') @@ -73,7 +73,7 @@ describe('WorkQueue bounded ready backlog', () => { describe('WorkQueue backpressure signaling', () => { it('stays open when a consumer drains as fast as we enqueue', () => { - const queue = new WorkQueue({ capacity: 2, highWatermark: 2, lowWatermark: 1 }) + const queue = new WorkQueue({ retryBackoff: false, capacity: 2, highWatermark: 2, lowWatermark: 1 }) const events: BackpressureEvent[] = [] queue.onBackpressure((event) => events.push(event)) queue.consume((d) => d.ack()) @@ -85,7 +85,7 @@ describe('WorkQueue backpressure signaling', () => { }) it('pauses at high watermark and resumes at or below low, without flapping in the band', () => { - const queue = new WorkQueue({ capacity: 4, highWatermark: 3, lowWatermark: 1 }) + const queue = new WorkQueue({ retryBackoff: false, capacity: 4, highWatermark: 3, lowWatermark: 1 }) const events: Pick[] = [] queue.onBackpressure((event) => events.push({ state: event.state, occupancy: event.occupancy })) @@ -120,7 +120,7 @@ describe('WorkQueue backpressure signaling', () => { }) it('unsubscribing a listener stops further events', () => { - const queue = new WorkQueue({ capacity: 2 }) + const queue = new WorkQueue({ retryBackoff: false, capacity: 2 }) const seen: string[] = [] const off = queue.onBackpressure((event) => seen.push(event.state)) queue.enqueue('a') @@ -136,7 +136,7 @@ describe('WorkQueue backpressure signaling', () => { }) it('delivers the first listener error after the rest of the snapshot runs', () => { - const queue = new WorkQueue({ capacity: 1 }) + const queue = new WorkQueue({ retryBackoff: false, capacity: 1 }) const order: string[] = [] queue.onBackpressure(() => { order.push('a') @@ -152,7 +152,7 @@ describe('WorkQueue backpressure signaling', () => { }) it('surfaces a listener error when consume drains a paused queue with sync ack', () => { - const queue = new WorkQueue({ capacity: 4, highWatermark: 3, lowWatermark: 1 }) + const queue = new WorkQueue({ retryBackoff: false, capacity: 4, highWatermark: 3, lowWatermark: 1 }) queue.enqueue('1') queue.enqueue('2') queue.enqueue('3') @@ -178,7 +178,7 @@ describe('WorkQueue backpressure signaling', () => { }) it('does not reject a nack requeue when ready is already at capacity', () => { - const queue = new WorkQueue({ capacity: 1 }) + const queue = new WorkQueue({ retryBackoff: false, capacity: 1 }) const events: string[] = [] queue.onBackpressure((event) => events.push(event.state)) let first: Delivery | undefined