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/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/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..14a039b 100644 --- a/src/work-queue.ts +++ b/src/work-queue.ts @@ -7,6 +7,7 @@ import { 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' @@ -31,12 +32,20 @@ export type EnqueueResult = { readonly accepted: true; readonly id: number } | { export interface WorkQueueOptions extends QueueBoundOptions { 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 { @@ -65,6 +74,10 @@ export class WorkQueue { 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 private nextMessageId = 1 private nextDeliveryTag = 1 private nextConsumerId = 1 @@ -85,6 +98,8 @@ export class WorkQueue { 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 { @@ -96,7 +111,7 @@ export class WorkQueue { 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 } @@ -158,6 +173,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,6 +205,7 @@ export class WorkQueue { if (this.pumping) return 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()) { const deliveredThisRound = new Set() @@ -216,6 +245,7 @@ export class WorkQueue { } finally { this.pumping = false } + this.armRetryTimer() } private deliver(consumer: ConsumerState, pending: PendingMessage): void { @@ -259,12 +289,50 @@ export class WorkQueue { // 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 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 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) + } + 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) + } + private applyFlow(): void { if (this.pumping || !this.gate) return const next = this.gate.observe(this.ready.length) 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-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 diff --git a/test/work-queue.test.ts b/test/work-queue.test.ts index 37c8aa7..7698060 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,184 @@ describe('consumer lifecycle', () => { expect(counts).toEqual([3, 3, 3]) }) }) + +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({ + 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('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({ + 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']) + }) +}) +