From fa9a77c3914874f2457d83f60057760d47b43620 Mon Sep 17 00:00:00 2001 From: Thomas Hart Date: Sun, 23 Aug 2026 19:02:24 +0000 Subject: [PATCH] feat: Add Per-Key Ordering Guarantees KeyedWorkQueue keeps FIFO order within a message key and lets different keys run in parallel on competing consumers. --- README.md | 4 + src/index.ts | 10 ++ src/keyed-queue.ts | 236 +++++++++++++++++++++++++++++++++++++++ test/keyed-queue.test.ts | 189 +++++++++++++++++++++++++++++++ 4 files changed, 439 insertions(+) create mode 100644 src/keyed-queue.ts create mode 100644 test/keyed-queue.test.ts diff --git a/README.md b/README.md index b63e153..bf16bf0 100644 --- a/README.md +++ b/README.md @@ -22,6 +22,8 @@ Message brokers hide a lot of machinery behind `publish` and `subscribe`. This r - **Bounded buffers** with a finite ready-queue `capacity` and **reject-on-full** (`QueueFullError` / `tryEnqueue`) - **High/low watermark backpressure** (hysteresis / Schmitt trigger) so producers pause before the wall and resume after the queue drains, without flapping in the band - **Producer-facing occupancy**: capacity is ready depth. Prefetch already caps in-flight. Redelivery of accepted work is allowed to sit above capacity so a nack cannot drop a message the queue already took. +- **Per-key FIFO ordering**: at most one in-flight delivery per key (SQS FIFO message group / Kafka key), a partial order across keys, and head-of-line blocking only within a key + ## 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). @@ -34,6 +36,8 @@ Message brokers hide a lot of machinery behind `publish` and `subscribe`. This r - **High/low watermark backpressure** (hysteresis / Schmitt trigger) so producers pause before the wall and resume after the queue drains, without flapping in the band - **Producer-facing occupancy**: capacity is ready depth. Prefetch already caps in-flight. Redelivery of accepted work is allowed to sit above capacity so a nack cannot drop a message the queue already took. - **Bounded queues with backpressure signaling.** Give `WorkQueue` a finite `capacity`. New publishes that would grow the ready backlog past that bound are rejected (`enqueue` throws `QueueFullError`, `tryEnqueue` returns `{ accepted: false }`). A `WatermarkGate` watches ready depth: occupancy at or above `highWatermark` emits `paused`, occupancy at or below `lowWatermark` emits `open`, and values in between hold the last state. Subscribe with `onBackpressure` or poll `backpressure()`. Defaults: high equals capacity, low is half the capacity (or `high - 1` when high is small). Unbounded queues stay the default so existing callers do not change. +- **Per-key FIFO ordering**: at most one in-flight delivery per key (SQS FIFO message group / Kafka key), a partial order across keys, and head-of-line blocking only within a key +- **Per-key ordering guarantees.** A `KeyedWorkQueue` where each message carries a group key. The head of a key is exclusive: the tail stays queued until that head is acked, nacked-and-dropped, or cancelled. Different keys can be in flight on competing consumers at the same time. A nack requeues at the head of that key so a later message of the same key cannot overtake. Prefetch can hold several keys, never two messages of one key. ## Usage ```ts diff --git a/src/index.ts b/src/index.ts index 3d5bf16..fdfcd93 100644 --- a/src/index.ts +++ b/src/index.ts @@ -47,3 +47,13 @@ export type { QueueBounds, QueueBoundOptions, } from './backpressure.js' + +export { KeyedWorkQueue } from './keyed-queue.js' + +export type { + KeyedMessage, + KeyedDelivery, + KeyedHandler, + KeyedWorkQueueOptions, + Unsubscribe as KeyedUnsubscribe, +} from './keyed-queue.js' diff --git a/src/keyed-queue.ts b/src/keyed-queue.ts new file mode 100644 index 0000000..6f02132 --- /dev/null +++ b/src/keyed-queue.ts @@ -0,0 +1,236 @@ +export interface KeyedMessage { + readonly id: number + readonly key: string + readonly payload: T + readonly deliveryCount: number +} + +export interface KeyedDelivery { + readonly message: KeyedMessage + readonly deliveryTag: number + ack(): void + nack(options?: { requeue?: boolean }): void +} + +export type KeyedHandler = (delivery: KeyedDelivery) => void +export type Unsubscribe = () => void + +export interface KeyedWorkQueueOptions { + maxDeliveryCount?: number +} + +interface Pending { + id: number + key: string + payload: T + deliveryCount: number +} + +interface Group { + pending: Pending[] + inFlight: boolean +} + +interface InFlightEntry { + pending: Pending + consumerId: number + settled: boolean +} + +interface ConsumerState { + id: number + handler: KeyedHandler + prefetch: number + inFlight: number + active: boolean +} + +const DEFAULT_MAX_DELIVERY_COUNT = 10 + +export class KeyedWorkQueue { + private readonly groups = new Map>() + private readonly runnable: string[] = [] + private readonly queuedKeys = new Set() + private readonly consumers: ConsumerState[] = [] + private readonly inFlight = new Map>() + private readonly maxDeliveryCount: number + private nextMessageId = 1 + private nextDeliveryTag = 1 + private nextConsumerId = 1 + private nextConsumerIndex = 0 + private pumping = false + + constructor(options?: KeyedWorkQueueOptions) { + const max = options?.maxDeliveryCount ?? DEFAULT_MAX_DELIVERY_COUNT + if (!Number.isInteger(max) || max < 1) { + throw new Error(`maxDeliveryCount must be a positive integer, got ${max}`) + } + this.maxDeliveryCount = max + } + + enqueue(key: string, payload: T): number { + if (key.length === 0) throw new Error('key must be a non-empty string') + const id = this.nextMessageId++ + let group = this.groups.get(key) + if (!group) { + group = { pending: [], inFlight: false } + this.groups.set(key, group) + } + group.pending.push({ id, key, payload, deliveryCount: 0 }) + this.schedule(key) + this.pump() + return id + } + + consume(handler: KeyedHandler, options?: { prefetch?: number }): Unsubscribe { + const prefetch = options?.prefetch ?? 1 + if (!Number.isInteger(prefetch) || prefetch < 1) { + throw new Error(`prefetch must be a positive integer, got ${prefetch}`) + } + const consumer: ConsumerState = { + id: this.nextConsumerId++, + handler, + prefetch, + inFlight: 0, + active: true, + } + this.consumers.push(consumer) + this.pump() + return () => { + if (!consumer.active) return + consumer.active = false + for (const [tag, entry] of this.inFlight) { + if (entry.consumerId !== consumer.id || entry.settled) continue + entry.settled = true + this.inFlight.delete(tag) + const group = this.groups.get(entry.pending.key) + if (!group) continue + group.inFlight = false + group.pending.unshift(entry.pending) + this.schedule(entry.pending.key) + } + consumer.inFlight = 0 + const i = this.consumers.indexOf(consumer) + if (i !== -1) this.consumers.splice(i, 1) + this.pump() + } + } + + readyCount(): number { + let n = 0 + for (const group of this.groups.values()) n += group.pending.length + return n + } + + inFlightCount(): number { + return this.inFlight.size + } + + consumerCount(): number { + return this.consumers.length + } + + private schedule(key: string): void { + const group = this.groups.get(key) + if (!group || group.inFlight || group.pending.length === 0) return + if (this.queuedKeys.has(key)) return + this.queuedKeys.add(key) + this.runnable.push(key) + } + + private pickConsumer(): { consumer: ConsumerState; index: number } | undefined { + const n = this.consumers.length + if (n === 0) return undefined + const start = this.nextConsumerIndex % n + for (let offset = 0; offset < n; offset++) { + const i = (start + offset) % n + const consumer = this.consumers[i] + if (!consumer || !consumer.active || consumer.inFlight >= consumer.prefetch) continue + return { consumer, index: i } + } + return undefined + } + + private pump(): void { + if (this.pumping) return + this.pumping = true + try { + while (this.runnable.length > 0) { + const deliveredThisRound = new Set() + let deliveredAny = false + while (this.runnable.length > 0) { + const key = this.runnable[0] + if (!key || deliveredThisRound.has(key)) break + const picked = this.pickConsumer() + if (!picked) break + this.runnable.shift() + this.queuedKeys.delete(key) + const group = this.groups.get(key) + if (!group || group.inFlight || group.pending.length === 0) continue + const pending = group.pending.shift() + if (!pending) continue + group.inFlight = true + deliveredThisRound.add(key) + this.deliver(picked.consumer, pending) + const len = this.consumers.length + this.nextConsumerIndex = len > 0 ? (picked.index + 1) % len : 0 + deliveredAny = true + } + if (!deliveredAny) break + } + } finally { + this.pumping = false + } + } + + private deliver(consumer: ConsumerState, pending: Pending): void { + pending.deliveryCount += 1 + const deliveryTag = this.nextDeliveryTag++ + const entry: InFlightEntry = { + pending, + consumerId: consumer.id, + settled: false, + } + this.inFlight.set(deliveryTag, entry) + consumer.inFlight += 1 + const delivery: KeyedDelivery = { + message: { + id: pending.id, + key: pending.key, + payload: pending.payload, + deliveryCount: pending.deliveryCount, + }, + deliveryTag, + ack: () => this.settle(deliveryTag, false), + nack: (options) => this.settle(deliveryTag, options?.requeue !== false), + } + try { + consumer.handler(delivery) + } catch { + if (!entry.settled) this.settle(deliveryTag, true) + } + } + + private settle(deliveryTag: number, requeue: boolean): void { + const entry = this.inFlight.get(deliveryTag) + if (!entry || entry.settled) return + entry.settled = true + this.inFlight.delete(deliveryTag) + const consumer = this.consumers.find((c) => c.id === entry.consumerId) + if (consumer?.active) consumer.inFlight -= 1 + const key = entry.pending.key + const group = this.groups.get(key) + if (group) group.inFlight = false + if (requeue && entry.pending.deliveryCount < this.maxDeliveryCount) { + // Head of this key so a later message cannot overtake a nack. + group?.pending.unshift(entry.pending) + } + if (group && group.pending.length === 0 && !group.inFlight) { + this.groups.delete(key) + this.queuedKeys.delete(key) + } else if (group) { + this.schedule(key) + } + this.pump() + } +} diff --git a/test/keyed-queue.test.ts b/test/keyed-queue.test.ts new file mode 100644 index 0000000..8bab250 --- /dev/null +++ b/test/keyed-queue.test.ts @@ -0,0 +1,189 @@ +import { describe, it, expect } from 'vitest' +import { KeyedWorkQueue, type KeyedDelivery } from '../src/index.js' + +describe('KeyedWorkQueue per-key order', () => { + it('throws on empty key and invalid options', () => { + const queue = new KeyedWorkQueue() + expect(() => queue.enqueue('', 'x')).toThrow(/non-empty/) + expect(() => new KeyedWorkQueue({ maxDeliveryCount: 0 })).toThrow(/maxDeliveryCount/) + expect(() => queue.consume(() => {}, { prefetch: 0 })).toThrow(/prefetch/) + }) + + it('delivers the same key in enqueue order with monotonic ids', () => { + const queue = new KeyedWorkQueue() + const seen: { id: number; payload: number }[] = [] + for (let i = 0; i < 2; i++) { + queue.consume((d) => { + seen.push({ id: d.message.id, payload: d.message.payload }) + d.ack() + }) + } + queue.enqueue('order-1', 1) + queue.enqueue('order-1', 2) + queue.enqueue('order-1', 3) + expect(seen.map((s) => s.payload)).toEqual([1, 2, 3]) + expect(seen.map((s) => s.id)).toEqual([1, 2, 3]) + }) + + it('blocks a key tail until the head acks, while another key is in flight', () => { + const queue = new KeyedWorkQueue() + const held: KeyedDelivery[] = [] + queue.consume((d) => held.push(d)) + queue.consume((d) => held.push(d)) + queue.enqueue('A', 'A1') + queue.enqueue('A', 'A2') + queue.enqueue('B', 'B1') + expect(held.map((d) => `${d.message.key}:${d.message.payload}`).sort()).toEqual([ + 'A:A1', + 'B:B1', + ]) + expect(queue.readyCount()).toBe(1) + expect(queue.inFlightCount()).toBe(2) + held.find((d) => d.message.payload === 'A1')?.ack() + expect(held.map((d) => d.message.payload)).toContain('A2') + expect(queue.readyCount()).toBe(0) + expect(queue.inFlightCount()).toBe(2) + }) + + it('prefetch cannot pull two messages of the same key', () => { + const queue = new KeyedWorkQueue() + const held: KeyedDelivery[] = [] + queue.consume((d) => held.push(d), { prefetch: 4 }) + queue.enqueue('A', '1') + queue.enqueue('A', '2') + queue.enqueue('B', '1') + expect(held.map((d) => d.message.key).sort()).toEqual(['A', 'B']) + expect(queue.readyCount()).toBe(1) + }) + + it('buffers with no consumer, then drains per key', () => { + const queue = new KeyedWorkQueue() + queue.enqueue('A', 1) + queue.enqueue('B', 10) + queue.enqueue('A', 2) + expect(queue.readyCount()).toBe(3) + const seen: string[] = [] + queue.consume((d) => { + seen.push(`${d.message.key}:${d.message.payload}`) + d.ack() + }) + expect(seen).toEqual(['A:1', 'B:10', 'A:2']) + }) +}) + +describe('nack, drop, and cancel keep per-key order', () => { + it('redelivers a nacked head before the next message of that key', () => { + const queue = new KeyedWorkQueue() + const seen: string[] = [] + queue.enqueue('A', 'A1') + queue.enqueue('A', 'A2') + queue.consume((d) => { + seen.push(`${d.message.payload}#${d.message.deliveryCount}`) + if (d.message.payload === 'A1' && d.message.deliveryCount === 1) d.nack() + else d.ack() + }) + expect(seen).toEqual(['A1#1', 'A1#2', 'A2#1']) + }) + + it('nack({ requeue: false }) drops the head and releases the tail', () => { + const queue = new KeyedWorkQueue() + const seen: string[] = [] + queue.enqueue('A', 'A1') + queue.enqueue('A', 'A2') + queue.consume((d) => { + seen.push(d.message.payload) + if (d.message.payload === 'A1') d.nack({ requeue: false }) + else d.ack() + }) + expect(seen).toEqual(['A1', 'A2']) + expect(queue.readyCount()).toBe(0) + }) + + it('handler throw requeues the head of that key', () => { + const queue = new KeyedWorkQueue() + let throws = true + const seen: string[] = [] + queue.consume((d) => { + seen.push(d.message.payload) + if (throws && d.message.payload === 'A1') { + throws = false + throw new Error('boom') + } + d.ack() + }) + queue.enqueue('A', 'A1') + queue.enqueue('A', 'A2') + expect(seen).toEqual(['A1', 'A1', 'A2']) + }) + + it('cancel requeues the in-flight head ahead of the same-key tail', () => { + const queue = new KeyedWorkQueue() + let stolen: KeyedDelivery | undefined + const off = queue.consume((d) => { + stolen = d + }) + queue.enqueue('A', 'A1') + queue.enqueue('A', 'A2') + expect(queue.inFlightCount()).toBe(1) + off() + expect(() => off()).not.toThrow() + expect(queue.consumerCount()).toBe(0) + expect(queue.readyCount()).toBe(2) + const recovered: string[] = [] + queue.consume((d) => { + recovered.push(d.message.payload) + d.ack() + }) + expect(recovered).toEqual(['A1', 'A2']) + expect(() => stolen?.ack()).not.toThrow() + }) + + it('double settle is a no-op', () => { + const queue = new KeyedWorkQueue() + let first: KeyedDelivery | undefined + let count = 0 + queue.consume((d) => { + count += 1 + if (count === 1) { + first = d + d.nack() + return + } + d.ack() + }) + queue.enqueue('k', 'x') + expect(() => first?.ack()).not.toThrow() + expect(() => first?.nack()).not.toThrow() + expect(count).toBe(2) + }) + + it('always-nack stays bounded and does not starve another key', () => { + const maxDeliveryCount = 4 + const queue = new KeyedWorkQueue({ maxDeliveryCount }) + const seen: string[] = [] + queue.enqueue('poison', 'spin') + queue.enqueue('ok', 'live') + queue.consume((d) => { + seen.push(d.message.payload) + if (d.message.key === 'poison') d.nack() + else d.ack() + }) + expect(seen.filter((p) => p === 'spin')).toHaveLength(maxDeliveryCount) + expect(seen.indexOf('live')).toBeLessThan(seen.lastIndexOf('spin')) + expect(queue.readyCount()).toBe(0) + expect(queue.enqueue('ok', 'after')).toBe(3) + }) + + it('enqueue from a handler preserves that key order', () => { + const queue = new KeyedWorkQueue() + const seen: string[] = [] + queue.consume((d) => { + seen.push(d.message.payload) + if (d.message.payload === 'seed') queue.enqueue('A', 'child') + d.ack() + }) + queue.enqueue('A', 'seed') + queue.enqueue('A', 'after') + expect(seen).toEqual(['seed', 'child', 'after']) + }) +})