@bumail/queue
v0.4.3
Published
The outbound mail queue of a mail server: every recipient's state, retries with back-off, delivery status notifications (RFC 3464), leases for several workers, and a memory, a bun:sqlite, a PostgreSQL (Bun.sql) and a Redis (Bun.redis) store, delivering th
Readme
@bumail/queue
The outbound queue of a mail server: it keeps every message your server
sends to another one, delivers it through @bumail/smtp/client, retries
what failed for now, and sends a delivery status notification (RFC 3464)
back when it gives up. Each recipient has its own state, the queue
survives a restart on bun:sqlite, PostgreSQL or Redis, and several
workers — on one machine, or on many with PostgreSQL or Redis — can
share one queue. No runtime dependency: only peers.
Install
bun add @bumail/queue @bumail/smtp @bumail/mime
bun add @bumail/dns # for direct delivery by MX: the resolver@bumail/smtp (its client) and @bumail/mime (for the DSNs) are required
peers. @bumail/dns is optional: any resolver of its shape will do, and a
queue that only uses a smarthost needs none.
Subpaths
| import | what it gives |
| --- | --- |
| @bumail/queue | createQueue, the QueueStore contract, QueueError and the types |
| @bumail/queue/memory | MemoryQueueStore: in memory, lost on a restart — for specs and trials |
| @bumail/queue/sqlite | SqliteQueueStore: on disk with bun:sqlite, shared by the processes of one machine |
| @bumail/queue/postgres | PostgresQueueStore: on PostgreSQL through Bun's own Bun.sql, shared by instances on several machines |
| @bumail/queue/redis | RedisQueueStore: on Redis through Bun's own Bun.redis, shared by instances on several machines |
Usage
import { nodeResolver } from '@bumail/dns';
import { createQueue } from '@bumail/queue';
import { SqliteQueueStore } from '@bumail/queue/sqlite';
const queue = createQueue({
store: SqliteQueueStore.open({ directory: '/var/lib/bumail/queue' }),
hostname: 'mail.example.net', // your public name: EHLO, and the DSN's Reporting-MTA
resolver: nodeResolver(),
});
queue.start();
// The message as it leaves: whole, headers first, CRLF, already DKIM-signed
// (with @bumail/auth's signDkim, say).
const signedMessage = await Bun.file('outgoing.eml').bytes();
const item = await queue.enqueue(signedMessage, {
from: '[email protected]',
to: ['[email protected]', '[email protected]'],
});
item.recipients; // [{ address: '[email protected]', status: 'pending' }, …]
process.on('SIGTERM', () => queue.stop()); // lets deliveries under way endThe message is a Uint8Array, a string or a ReadableStream<Uint8Array>,
with CRLF line ends: a bare CR or LF is refused, as sendMail would
refuse it. Every address is checked as one sendMail takes
(@bumail/smtp/client's isMailbox: an RFC 5321 Mailbox with no source
route, no control character (C0, DEL or C1), no >, no U+2028 or U+2029,
no Unicode format character (\p{Cf}), no lone surrogate and no IPv4
literal octet above 255), so a bad one is refused at enqueue, never at
delivery.
Delivery groups the recipients by domain: one session per domain per
attempt, with STARTTLS when offered.
Retries and bounces
A 4xx, a connection error or a timeout leaves the recipient deferred;
the queue tries it again after 30 minutes, then 1 hour, 2, 4, and every 4
hours, each plus up to 10% of jitter, and gives up after 5 days
(RFC 5321 §4.5.4.1). A 5xx fails the recipient at once.
import { nodeResolver } from '@bumail/dns';
import { createQueue } from '@bumail/queue';
import { MemoryQueueStore } from '@bumail/queue/memory';
const queue = createQueue({
store: new MemoryQueueStore(),
hostname: 'mail.example.net',
resolver: nodeResolver(),
retry: { first: 15 * 60_000, giveUpAfter: 3 * 24 * 3_600_000 },
dsn: { delayAfter: 4 * 3_600_000, returnContent: 'headers' },
});A failed recipient gets a DSN back to the sender, a multipart/report
with the original's header fields; a recipient still deferred after
delayAfter gets one "delayed" warning. A DSN is enqueued from the null
sender <>, and a message from <> never causes one.
Routing
import { createQueue } from '@bumail/queue';
import { MemoryQueueStore } from '@bumail/queue/memory';
const smtpPassword = (await Bun.file('/run/secrets/smtp').text()).trim(); // from your secrets
createQueue({
store: new MemoryQueueStore(),
hostname: 'mail.example.net',
// Port 25 blocked? Everything through a provider's submission port:
route: {
host: 'smtp.provider.example',
port: 587,
auth: { username: '[email protected]', password: smtpPassword },
},
// …or only some domains, the rest by MX (which needs the resolver):
routes: { 'partner.example': { host: 'relay.partner.example' } },
});route defaults to 'mx': each recipient domain's own mail hosts.
Credentials go only over TLS whose certificate checked out: createQueue
refuses auth with a tls other than 'required', and checks the ports,
the TLS modes and the timeouts as sendMail would. A route sendMail
still refuses (INVALID_OPTION) defers its recipients as 4.3.5 and
says so on the error event; it does not bounce them until
retry.giveUpAfter, when they fail as 4.4.7 with a DSN like any
deferred recipient.
Several workers
import { nodeResolver } from '@bumail/dns';
import { createQueue } from '@bumail/queue';
import { SqliteQueueStore } from '@bumail/queue/sqlite';
// In each process, on the same directory: a claim leases an item to one worker.
const queue = createQueue({
store: SqliteQueueStore.open({ directory: '/var/lib/bumail/queue' }),
hostname: 'mail.example.net',
resolver: nodeResolver(),
});
queue.start();A worker claims a due item with a lease, renews it while it delivers, and
lets go of it with the outcome; a renewal that fails is told on error
and tried again. If it crashes, another worker claims the item once the
lease expires (leaseMs, 10 minutes by default). A directory the store
makes is 0700; one that exists keeps its mode.
concurrency (20) bounds the items one worker delivers at once — each
item opens one session per recipient domain — and perDomain (2) its
sessions to one recipient domain.
PostgreSQL, for several machines
import { nodeResolver } from '@bumail/dns';
import { createQueue } from '@bumail/queue';
import { PostgresQueueStore } from '@bumail/queue/postgres';
// In each instance, on the same database. A Bun.SQL client of your own works too:
// PostgresQueueStore.open({ sql: new Bun.SQL({ url, max: 10 }) }).
const store = PostgresQueueStore.open({
sql: Bun.env['DATABASE_URL'] ?? 'postgres://bumail@localhost:5432/mail',
tablePrefix: 'bumail_queue_', // the default: bumail_queue_items, …_messages, …_schema
});
await store.migrate(); // makes the tables; the first call would anyway
const queue = createQueue({ store, hostname: 'mail.example.net', resolver: nodeResolver() });
queue.start();
process.on('SIGTERM', async () => {
await queue.stop();
await store.close(); // closes the client it opened for the URL, never yours
});A claim is one UPDATE … RETURNING whose item a SELECT … FOR UPDATE
SKIP LOCKED picks: two instances never take the same item, and none
waits on one another is taking. An instance that crashes loses its items
when their leases expire, as with bun:sqlite. Every write runs at
READ COMMITTED, so a client whose sessions default to repeatable
read or serializable serves the queue too. No driver to install:
Bun.sql is Bun's.
Redis, for several machines
import { nodeResolver } from '@bumail/dns';
import { createQueue } from '@bumail/queue';
import { RedisQueueStore } from '@bumail/queue/redis';
// In each instance, on the same Redis. A Bun.RedisClient of your own works too:
// RedisQueueStore.open({ client: new Bun.RedisClient(url) }).
const store = RedisQueueStore.open({
url: Bun.env['REDIS_URL'] ?? 'redis://localhost:6379',
keyPrefix: 'bumail:queue:', // the default: bumail:queue:items, …:item:<id>, …
});
const queue = createQueue({ store, hostname: 'mail.example.net', resolver: nodeResolver() });
queue.start();
process.on('SIGTERM', async () => {
await queue.stop();
await store.close(); // closes the client it opened for the URL, never yours
});Every operation that writes is one Lua script, which Redis runs whole:
two instances never take the same item, and a crashed instance's items
are claimed again once their leases expire. The message is kept byte for
byte. One Redis, or a primary with replicas — not Redis Cluster. Redis
acknowledges a write once it is in memory: with appendfsync everysec
a crash loses up to a second of writes, and a failover to a replica can
lose acknowledged ones, so keep appendonly yes, appendfsync always
and maxmemory-policy noeviction where a lost message matters (the
guide's Redis section). No driver to install: Bun.redis is Bun's.
Events and admin
// queue: the one created under Usage.
queue.on('delivered', ({ id, recipient, reply }) => console.info({ id, recipient, reply }));
queue.on('deferred', ({ recipient, reply, nextAttemptAt }) => {});
queue.on('failed', ({ recipient, reply }) => {}); // reply: { code?, status?, text, host? }
queue.on('dsn', ({ kind, of, to }) => {});
queue.on('error', ({ error, id }) => console.error(id, error)); // the store, a lost lease, a bad route
const [next] = await queue.list({ limit: 50 }); // the next due first
if (next) await queue.retryNow(next.id); // false while a worker delivers it
if (next) await queue.cancel(next.id); // no DSNTesting
import { createQueue, type Sender } from '@bumail/queue';
import { MemoryQueueStore } from '@bumail/queue/memory';
let now = Date.UTC(2026, 0, 1);
const send: Sender = async (_message, options) => ({
accepted: [options.to].flat().map((recipient) => ({ recipient, reply: { code: 250, text: 'OK' } })),
rejected: [],
reply: { code: 250, text: 'Queued' },
host: 'mx.test',
port: 25,
tls: false,
authenticated: false,
});
const queue = createQueue({
store: new MemoryQueueStore(),
hostname: 'mail.test',
route: { host: 'mx.test' },
clock: { now: () => now },
send,
});
await queue.enqueue('Subject: hi\r\n\r\nhello\r\n', { from: 'a@test', to: 'b@test' });
await queue.deliverDue(); // one pass, by hand: no timer, no network
now += 30 * 60_000;Limits
import { nodeResolver } from '@bumail/dns';
import { createQueue } from '@bumail/queue';
import { MemoryQueueStore } from '@bumail/queue/memory';
createQueue({
store: new MemoryQueueStore(),
hostname: 'mail.example.net',
resolver: nodeResolver(),
limits: { maxMessageSize: 10 * 1024 * 1024, maxRecipients: 50, maxItems: 100_000 },
});Everything is bounded, through limits: the message (25 MiB), the
recipients per message (100), the items in the store (none by default),
the reply text kept per recipient (512 characters, control characters
replaced by spaces), and the original a DSN returns (64 KiB, each line cut
at 998 bytes). Nothing a
remote server says reaches a DSN's header fields with a CR, an LF or a
control character in it.
Traps
- Bun only. The stores use
bun:sqlite,Bun.sqlandBun.redis, so the package runs on Bun 1.4.2 or later, not on Node. - It sends what you enqueue, to anyone. The queue is not a relay policy: enqueue only what an authenticated user submitted, or what your own server writes, never what an unauthenticated client handed you.
API
| export | |
| --- | --- |
| createQueue(options), Queue | the queue: enqueue, deliverDue, start, stop, on, list, get, retryNow, cancel, owner |
| QueueOptions, Route, Smarthost, RetrySchedule, DsnOptions, QueueLimits, Clock, Sender | what createQueue takes |
| MessageSource, QueueEnvelope | what enqueue takes |
| QueueEvents, QueueListener, RecipientEvent, DeferredEvent, DsnEvent, QueueErrorEvent | the events |
| QueueStore | the contract a store answers: add, get, list, count, readMessage, claim, renew, complete, reschedule, cancel |
| QueueItem, RecipientState, RecipientStatus, FinalStatus, Diagnostic, Lease | an item, and where each recipient stands |
| NewQueueItem, AddOptions, ClaimRequest, AttemptResult, RecipientUpdate, QueueListOptions | what a store takes |
| QueueError, QueueErrorCode | INVALID, MESSAGE_TOO_BIG, TOO_MANY_RECIPIENTS, QUEUE_FULL, CLOSED, LEASE_LOST, MESSAGE_UNREADABLE |
| MemoryQueueStore | from @bumail/queue/memory |
| SqliteQueueStore, SqliteQueueStoreOptions | from @bumail/queue/sqlite: SqliteQueueStore.open({ directory, busyTimeout? }), close() |
| PostgresQueueStore, PostgresQueueStoreOptions, PostgresClient, PostgresQueryable | from @bumail/queue/postgres: PostgresQueueStore.open({ sql, tablePrefix? }), migrate(), close(); sql a Bun.SQL client or a postgres:// URL |
| RedisQueueStore, RedisQueueStoreOptions, RedisQueueClientOptions, RedisQueueUrlOptions, RedisQueueClient | from @bumail/queue/redis: RedisQueueStore.open({ client \| url, keyPrefix? }), close(); client a Bun.RedisClient, url a redis:// URL |
Documentation
These pages ship in the package, under docs/.
- Index: the pages, and when to read each.
- Guide: enqueuing, delivery and routing, the retry schedule, DSNs, several workers and leases, events and admin, the stores (PostgreSQL and Redis included), testing, and writing a store of your own.
- Troubleshooting: every
QueueError, and what to do about it. - Roadmap: what is coming, and what is not planned.
