@claudewerk/bun-mq
v0.2.0
Published
Transport-agnostic in-process message queue for Bun: per-message QoS, idempotency/dedupe, opt-in durable SQLite storage, dead-letter queue, named queues. Zero required runtime deps.
Maintainers
Readme
bun-mq
A transport-agnostic, in-process message queue for Bun,
built on bun:sqlite. Per-message QoS, id-based dedupe, opt-in durable
storage (crash-recoverable), a dead-letter queue, and named queues.
- Zero required runtime dependencies. Just
bun+bun:sqlite.pinois an optional peer dependency used only if you ask for the pino logger adapter. - Transport-agnostic. No network code. Producers
enqueue(value), consumersdequeue()/subscribe(). Values are opaque; the queue never inspects them. - Durability is opt-in.
memoryStore()for a fast in-process queue, orsqliteStore(path)for a WAL-backed store that survives a crash. - Storage is pluggable.
storagetakes anyStore. To keep queues inside a database your app already owns, borrow its handle:sqliteStore(appDb, { tablePrefix: "mq_" }). The tables sit next to yours under the prefix,atomic()joins your own writes in one transaction, andclose()leaves your handle open. The handle must be opened{ strict: true }.
import { MQ, sqliteStore, memoryStore } from "@claudewerk/bun-mq";
const mq = new MQ({ storage: sqliteStore("./data.db") }); // or memoryStore()
const q = mq.queue("orders", {
qos: "at-least-once",
maxDeliveries: 5,
ttlMs: 86_400_000,
visibilityTimeoutMs: 30_000,
});
const id = await q.enqueue(new Uint8Array([1, 2, 3])); // durable-committed before resolve
const msg = await q.dequeue(); // { id, value, receipt, deliveries, ... }
if (msg) {
try {
// ... process msg.value ...
await q.ack(msg.receipt); // settle THIS delivery (fenced; see below)
} catch {
await q.nack(msg.receipt, { requeue: true }); // retry (counts a delivery)
}
}
// Or subscribe: auto-acks on return, auto-nacks (requeue) on throw.
const unsubscribe = q.subscribe(async (m) => { await handle(m.value); });
// Dead-letter queue
const dead = await q.dlq.list();
await q.dlq.requeue(deadId);
await q.dlq.purge(deadId); // or purge() to clear allInstall
bun add @claudewerk/bun-mq
# optional, only for the pino logger adapter:
bun add pinoRequires Bun (uses bun:sqlite and node:v8). Not a Node.js package.
API
| Call | Description |
| --- | --- |
| new MQ({ storage, logger?, clock? }) | Engine over one storage backend. |
| mq.queue(name, config?) | Get/create a named queue. Config fixed on first reference. |
| mq.close() | Close the storage backend. |
| q.enqueue(value, { id?, qos?, ttlMs? }) | Enqueue an opaque value; returns the id. Durable stores commit before resolving. |
| q.dequeue({ visibilityTimeoutMs? }) | Take the next eligible message (with a per-delivery receipt), or undefined. |
| q.ack(receipt) | Mark this delivery processed; removes permanently. Throws StaleReceiptError if the receipt isn't the current delivery. |
| q.nack(receipt, { requeue? }) | requeue:true (default) redelivers; false dead-letters (rejected). Throws StaleReceiptError on a stale receipt. |
| q.subscribe(handler, { visibilityTimeoutMs?, pollMs? }) | Poll loop; acks on return, nacks on throw. Returns an unsubscribe fn. |
| q.size() | Count of live (non-dead) messages. |
| q.dlq.list() / get(id) / count() | Inspect dead-lettered messages. |
| q.dlq.requeue(id) | Restore a dead message (fresh delivery count, tail of queue). |
| q.dlq.purge(id?) | Delete one dead message, or all if id omitted. |
Queue config
| Option | Default | Meaning |
| --- | --- | --- |
| qos | "at-least-once" | Default QoS for messages that don't override it. |
| maxDeliveries | 5 | Deliveries before a message is dead-lettered. |
| ttlMs | none | Time-to-live; expired messages are dead-lettered. |
| visibilityTimeoutMs | 30_000 | How long a delivered at-least-once message stays invisible before redelivery. |
| dedupeWindowMs | 7 days | Window in which a re-enqueue of a known id collapses. |
Guarantees
| Guarantee | Holds | Exact semantics |
| --- | --- | --- |
| at-most-once | ✅ | On dequeue the message is removed and handed over exactly once. No ack, no redelivery. If the consumer dies mid-process, the message is lost. ttl still dead-letters it if it expires before delivery. |
| at-least-once (default) | ✅ | Delivered and made invisible for visibilityTimeoutMs. Redelivered on nack{requeue:true} or when the timeout lapses without an ack. Each delivery increments deliveries. Delivered 1+ times until acked; never lost while un-acked (durable store: survives crash). |
| Delivery fencing | ✅ | Each delivery carries an opaque receipt distinct from id. ack/nack take the receipt and reject a stale one (StaleReceiptError) — a delivery that already timed out and was redelivered can't settle the message a later delivery now owns. No silent stale-ack no-op. See DESIGN-FEEDBACK.md §1. |
| exactly-once | ❌ (by design) | Not promised. Use at-least-once + idempotent consumers. Dedupe collapses duplicate enqueues, fencing prevents stale-ack loss — but neither prevents a duplicate delivery from running your handler twice. |
| Idempotency / dedupe | ✅ | A re-enqueue of a known id is collapsed (returns the original id, stores nothing new) when a live copy exists or a dedupe record for that id is within dedupeWindowMs. Scope is per-queue + id. On durable stores the dedupe window persists across restarts. |
| Durability | ✅ (opt-in) | With sqliteStore, enqueue commits to SQLite (WAL) before the promise resolves — the durability point. Un-acked messages replay after a process crash + reopen. memoryStore gives the same API with no persistence. |
| Dead-letter queue | ✅ | A message is moved to the DLQ (never silently dropped) when it reaches maxDeliveries, exceeds ttlMs, or is nacked with requeue:false. Inspect / requeue / purge via q.dlq. Reasons: max-deliveries, ttl, rejected. |
| Ordering | ⚠️ FIFO among visible | FIFO per named queue by a monotonic sequence. A redelivered message resumes its original sequence position, so it is picked before higher-sequence visible messages. But because delivery only makes a message invisible (not removed) until acked, a later message can be dequeued/acked while an earlier one is mid-redelivery — so strict enqueue-order observation does not survive redelivery. See DESIGN-FEEDBACK.md §2. |
| Named-queue isolation | ✅ | Independent named queues in one engine; each has its own config, messages, and DLQ. |
Delivery-count semantics
maxDeliveries: N means a message is delivered up to N times; on the
attempt that would exceed N it is dead-lettered instead of delivered. nacking
a message already at the cap dead-letters it immediately.
Receipts & fencing
dequeue returns a Message with a receipt — an opaque, per-delivery token
(distinct from id). Always settle with the receipt:
const msg = await q.dequeue();
if (msg) await q.ack(msg.receipt); // NOT msg.idIf the delivery you hold has already timed out and been redelivered (to you or
another consumer), ack(receipt) / nack(receipt) throw StaleReceiptError
instead of silently deleting a message a newer delivery now owns. This is the
fence that makes at-least-once safe under competing consumers. subscribe
handles it for you: a superseded settle is treated as a no-op (another consumer
owns the message). If you settle by hand, catch it where a stale settle is
expected — e.g. an idempotent acker on the source side of a bridge:
import { StaleReceiptError } from "@claudewerk/bun-mq";
try {
await q.ack(receipt);
} catch (e) {
if (!(e instanceof StaleReceiptError)) throw e;
// already settled/superseded — safe to ignore in an idempotent path
}Sharp edges
- Everything is
async.enqueue,dequeue,ack,nack,size, and everydlq.*method return Promises — even onmemoryStore, where the work is synchronous. You mustawaitthem. In particularq.size() === 3compares a Promise to a number and is always false; write(await q.size()) === 3. This uniform async surface is deliberate somemoryStoreandsqliteStoreare drop-in interchangeable. - Settle by
receipt, notid(see above) —ack(msg.id)will throwStaleReceiptErrorbecause a bare id is not a valid receipt. ackis not idempotent. A receipt settles exactly once; re-acking it (or acking a superseded delivery) throws. CatchStaleReceiptErrorif your caller may retry a settle.
Values & serialization
Values are opaque. memoryStore keeps a structuredClone snapshot (so mutating
the original after enqueue doesn't change the stored copy). sqliteStore
serializes with node:v8 (serialize/deserialize) to a BLOB, giving full
structured-clone fidelity — Uint8Array is byte-exact, and Map/Date/
nested objects round-trip. Anything structured-cloneable works on both stores.
Custom stores
Implement the Store interface in src/store/types.ts. Every method is
synchronous and single-op-atomic. Besides the queue/DLQ primitives, a store
provides seen(queue, id, dedupeWindowMs, now): a read-only dedupe probe
that reports whether insert of that id would collapse right now (a live row
exists, or a dedupe record is within the window) without committing anything.
A protocol layer uses it to re-acknowledge a duplicate before validating a
frame, and to validate a fresh frame before committing it.
Two more primitives serve callers that keep state next to their queues:
meta(getMeta/setMeta/deleteMeta/listMeta(prefix)): a durable key/value side table with the same durability as the rows. Persist a counter, a frontier or a registry here instead of inventing a second store.atomic(fn): run every store call insidefnas one unit. OnsqliteStoreit is a transaction (nested calls become savepoints, so wrapping aninsertis fine); onmemoryStoreit runsfninline with no rollback.
Three primitives exist for paths that run per dequeue or per inbound message. A custom store may satisfy their signatures with a scan, but it will be slow in exactly the place it matters, so they are worth implementing properly:
size(queue)counts live rows. Do not implement it aslistLive().length— that decodes every stored value to produce an integer.sweepExpired(queue, now, deadAt)dead-letters every row past itsexpiresAtin one atomic unit and returns how many moved. It runs on every dequeue, so index it on the expiry;sqliteStoreuses a partial index over(queue, expires_at)holding only rows that carry a TTL, which makes the common no-TTL case a probe on an empty index. Do not skip it when the queue has nottlMs: a single message enqueued with its ownttlMswould then never expire.admit(queue, id, windowMs, now)records(queue, id)and returns whether it was the first sighting insidewindowMs, in one statement. This is the dedupe-ledger primitive: probe-then-write is two commits to learn one boolean, with a crash window between them.sqliteStoredoes it withINSERT ... ON CONFLICT DO UPDATE ... WHERE ... RETURNING.
Group commit
Enqueues issued in the same tick are committed in one transaction — roughly 4.9x the throughput of one transaction each, since per-transaction overhead rather than row encoding is the cost of a small write.
await Promise.all([q.enqueue(a), q.enqueue(b), q.enqueue(c)]); // one transaction
await q.enqueue(a); await q.enqueue(b); // two, as askedThe durability point does not move. Each enqueue resolves only after the
batch's transaction has committed. A caller that awaits each enqueue in turn is
never batched with itself, because it demanded durability before issuing the
next one. A batch is all-or-nothing: if the transaction rolls back, every
enqueue in it rejects — which is the truth, since none of them is stored. FIFO
order within a batch is the order the enqueues were issued.
Logging
Inject any object satisfying the Logger interface
(debug/info/warn/error(msg, fields?)):
import { MQ, consoleLogger, pinoLogger } from "@claudewerk/bun-mq";
new MQ({ storage: memoryStore(), logger: consoleLogger() });
new MQ({ storage: memoryStore(), logger: await pinoLogger() }); // requires optional pinoLifecycle events log at debug; dead-letter moves at warn, with structured
fields (queue, id, deliveries, reason). With no logger injected, the
default is a silent shim — zero deps, zero output.
Clock and timers
Two injection points, both structural (this package imports nothing to accept them):
clock: () => number— epoch ms, used for visibility, TTL and dedupe windows. DefaultDate.now.timers: { setTimeout, clearTimeout }— whatsubscribe()'s poll loop sleeps on. Default: the host timers. Pass abun-schedulerClockorFakeClockto drive the loop deterministically, or a Durable Object alarm clock on Cloudflare.src/never callssetTimeoutdirectly outside the default; the workspace test inbun-schedulerenforces that.
Recipe: bridging two queues over an unreliable transport
bun-mq is transport-agnostic, so a common topology is a store-and-forward
bridge: a message hops from a queue on one host, across a link that can drop
and re-establish, into a queue on another host.
producer → q1 → [bridge A] ==transport (may drop)== [bridge B] → q2 → consumerThe bridge on host A is a consumer of q1 and a producer to q2. The
primitives below are exactly what makes this reliable, and bun-mq supports
every piece:
| Requirement | How bun-mq provides it |
| --- | --- |
| Source message isn't lost while in flight | dequeue makes it invisible, not removed; it stays until you ack. Un-acked messages redeliver on visibility timeout. |
| Survives a host crash mid-transfer | Durable sqliteStore: enqueue commits before it resolves, and un-acked messages replay on reopen. |
| Resends don't duplicate at the receiver | q2.enqueue(value, { id }) dedupes on the id within dedupeWindowMs — idempotent, commits durably. |
| A long outage doesn't dead-letter good messages | Tune the source queue's maxDeliveries high (or effectively-infinite) and its visibilityTimeoutMs to the round-trip. |
The three rules that make it correct
1. Commit before receipt. B must durably enqueue into q2 before
acking A, and A must remove from q1 only on B's ack. Reversed, a crash
between the ack and the commit loses the message.
2. The receiver always re-acks; dedupe gates only the enqueue. When B's ack
is lost, A resends and B sees a known id again. B must still ack (A is still
waiting) while not re-enqueuing. Ack and dedupe are independent: a duplicate
means "don't enqueue again, but do ack again." Silence on a duplicate strands
the message in q1 until it dead-letters. So the bridge-B handler is branchless:
async function onReceive(msg: { id: string; value: Uint8Array }) {
await q2.enqueue(msg.value, { id: msg.id }); // idempotent: dup collapses, commits durably
sendAckBackToA(msg.id); // ALWAYS ack, duplicate or not
}And on host A:
const m = await q1.dequeue({ visibilityTimeoutMs: 60_000 });
if (m) {
sendOverTransport(m); // to B (carry m.id end-to-end for q2's dedupe)
onAckFromB(m.id, async () => {
try {
await q1.ack(m.receipt); // remove ONLY after B confirms durable receipt
} catch (e) {
// A duplicate B-ack, or a redelivery already superseded this one:
// the fence throws instead of deleting the wrong delivery. Ignore it.
if (!(e instanceof StaleReceiptError)) throw e;
}
});
// no ack from B → visibility lapses → q1 redelivers → resend → q2 dedupes → B re-acks → q1.ack
}3. Receiver memory ≥ sender persistence. q2's dedupeWindowMs must exceed
the maximum time A can keep retrying (q1's ttlMs, or
maxDeliveries × visibilityTimeoutMs). If the sender can outlast the receiver's
dedupe memory, a late resend is re-enqueued as a genuine duplicate. Size the
window from the sender's worst-case outage, not a round number.
This yields effectively-once delivery (at-least-once + an idempotent receiver). True exactly-once is impossible across an unreliable link (two-generals), so the receiving queue's dedupe is load-bearing, not optional — even with a single well-behaved bridge. See DESIGN-FEEDBACK.md §1 and §4 for the protocol-level version of these rules (and the fencing-token gap they don't fully close).
Notes on scale
This is an in-process engine tuned for correctness and small-to-medium
in-memory / single-file workloads. TTL expiry is swept on each dequeue by
scanning live rows, which is fine for typical queue depths but O(n) per dequeue
on very large backlogs.
Deviations from the original SPEC
Greenfield build — the SPEC was a starting contract, improved where it was weak:
- Added
DeadReason: "rejected"fornack{requeue:false}(SPEC listed onlymax-deliveries/ttl) so a manual reject is honestly distinguishable from a policy-triggered dead-letter. - Durable values use
node:v8structured-clone serialization (BLOB) rather than JSON, for exact fidelity of binary and non-JSON values. - Added
q.size(),dlq.count(), anddlq.purge()(purge-all) as small, obviously-useful surface not spelled out in the SPEC. - Per-delivery
receipt+ fencing.ack/nacktake an opaque per-delivery receipt (not the bareid) and reject a stale one withStaleReceiptError, closing the unfenced-ack hazard where a slow consumer deletes a message a newer delivery owns. Implements this package's own DESIGN-FEEDBACK §1.
License
MIT
