Queues & jobs
@basaltkit/queue runs work in the background through a small, driver-agnostic core. You define typed jobs, dispatch them from anywhere, and workers process them — on Redis (BullMQ) in production, inline in dev/test, or on RabbitMQ, Kafka, or Amazon SQS through a driver package. The backend is one line to swap; your jobs never change.
Install
# the core, plus the backend you picked and its client
pnpm add @basaltkit/queue @basaltkit/queue-bullmq bullmq@basaltkit/queue is always required. It is not one of the backends — it is the contract they all implement, and the package your job code imports:
| You get from | What |
|---|---|
@basaltkit/queue — always | defineJob, dispatch, the QUEUE token, QueueManager, workers, context propagation, and the sync driver |
| a backend package — one of them | the driver and its plugin: where those jobs actually run |
The backend package depends on the core; it does not replace it. Adding one is choosing an execution target, not swapping libraries.
| Backend | Package | Plugin | Client to install |
|---|---|---|---|
| BullMQ (Redis) | @basaltkit/queue-bullmq | bullmqQueuePlugin | bullmq |
| RabbitMQ | @basaltkit/queue-rabbitmq | rabbitmqQueuePlugin | amqplib |
| Amazon SQS | @basaltkit/queue-sqs | sqsQueuePlugin | @aws-sdk/client-sqs |
| Kafka | @basaltkit/queue-kafka | kafkaQueuePlugin | kafkajs |
| none — inline, dev/tests | (the core already has it) | queuePlugin | — nothing |
The core knows no broker, so an app on SQS never installs — or loads — BullMQ and its ioredis weight. And with no backend package at all you still have a working queue: the sync driver, which is all dev and tests need.
The payoff is that your job code never mentions the backend. Swapping RabbitMQ for Redis is one changed import in app.ts; every defineJob and every dispatch stays exactly as written.
Define a job
import { defineJob } from '@basaltkit/queue'
import { z } from 'zod'
export const SendWelcome = defineJob({
name: 'send-welcome',
queue: 'welcome', // which queue/worker handles it (default 'default')
schema: z.object({ userId: z.string() }),
attempts: 3, // retry up to 3 times
backoff: { type: 'exponential', delay: '30s' },
async handle({ userId }) {
// ... do the work
},
})The schema makes the payload type-safe end to end — handle's argument and dispatch's payload are both inferred from it, and the payload is validated on dispatch.
attempts must be a positive integer (defineJob throws on 0, a negative or a fraction). The job name is the routing key between producer and worker, so it must be unique per manager: registering a different definition under a taken name throws DuplicateJobError (registering the same one twice is a no-op).
Register it
Your backend's plugin registers a QueueManager under the QUEUE token, starts the declared workers on boot, and closes everything on shutdown:
import { createApp } from '@basaltkit/core'
import { bullmqQueuePlugin } from '@basaltkit/queue-bullmq'
import { SendWelcome } from './jobs/send-welcome.js'
const app = await createApp({
plugins: [
bullmqQueuePlugin({
connection: process.env.REDIS_URL!, // Redis URL or ioredis options
jobs: [SendWelcome], // jobs this process produces and/or runs
workers: [{ queue: 'welcome', concurrency: 5 }], // start a worker for this queue
}),
],
}).boot()Swapping backend is swapping that one import: rabbitmqQueuePlugin, sqsQueuePlugin and kafkaQueuePlugin accept the same jobs/workers keys alongside their own connection options. Your jobs never change.
Use the core's queuePlugin directly when you want the sync driver, or a driver of your own:
import { queuePlugin } from '@basaltkit/queue'
queuePlugin({ jobs: [SendWelcome] }) // sync driver — dev and tests
queuePlugin({ driver: myDriver, jobs, workers }) // a driver you wroteWith no driver, the plugin uses the sync driver: dispatch runs handle inline in the same process — no Redis, ideal for dev and tests. Know its semantics before relying on it: it is at-most-once (a job that exhausts its inline retries is lost), handler errors reject the dispatch() call (your request fails instead of a background retry), and it is not meant for production — a production deploy (any NODE_ENV but an explicit development/test, unset included) that falls back to it without a driver logs a warning at boot (pass driver: new SyncQueueDriver() to opt in deliberately). A worker's queue must match a job's queue, or the job lands in the backend but nothing consumes it.
Producer and worker in separate processes
In production the API process usually only produces (calls dispatch) while a separate process consumes. Both must register the same jobs — the worker needs each job's handle, and a job that reaches a worker that hasn't registered it throws UnknownJobError. Only the consumer declares workers:
// API process — produces only (no `workers`)
bullmqQueuePlugin({ jobs: [SendWelcome, GenerateInvoice], connection: process.env.REDIS_URL! })
// worker process — consumes
bullmqQueuePlugin({
jobs: [SendWelcome, GenerateInvoice],
connection: process.env.REDIS_URL!,
workers: [
{ queue: 'welcome', concurrency: 5 },
{ queue: 'billing', concurrency: 2 },
],
})Dispatch
import { ctx } from '@basaltkit/core'
import { QUEUE } from '@basaltkit/queue'
await ctx().container.get(QUEUE).dispatch(SendWelcome, { userId: 'u-1' })
// or straight off the job (once it's registered):
await SendWelcome.dispatch({ userId: 'u-1' }, { delay: '5m', priority: 5 })dispatch returns as soon as the job is enqueued. Request context (requestId, tenantId, …) is captured and restored inside the worker.
Delivery is at-least-once on every broker driver: a crash between "handler finished" and "broker acknowledged" runs the job again, and dispatch has no idempotency key. Write handlers so a second run is harmless (upserts, "already sent?" checks, a unique constraint on the effect).
Context in the worker — and the trust boundary
The worker does not restore the envelope's context wholesale. It rebuilds it from an allowlist, validating each field:
| Field | Restored as | If malformed |
|---|---|---|
requestId, correlationId, traceId | the same key (short, printable strings) | dropped — they only label logs |
tenant / tenantId | tenant: { id } and tenantId — the id must match the tenancy grammar (or your validateTenantId), and the two must agree | the job is rejected (JobContextError) — dropping the tenant would run the job in the central scope |
userId (from user.id at dispatch) | userId and a minimal actor user: { id } | the job is rejected (JobContextError) |
Everything else in the message is ignored. The actor is the id only: audit records it as actorId, and gate.actor() re-reads that user's roles from the permission store in the job's tenant — roles are never taken from the message. So a job dispatched from a request runs with the dispatching user's current permissions, and a job dispatched outside a request has no actor.
The broker is trusted by default. Without a signing key, anyone who can write to the queue backend can enqueue a job, pick its payload, and name any valid tenant and user id. Close that with a signing key shared by producers and workers:
bullmqQueuePlugin({
connection: process.env.REDIS_URL!,
jobs,
signingKey: process.env.QUEUE_SIGNING_KEY!, // ≥ 32 bytes
// rotation: [newKey, oldKey] — the first signs, every key verifies
})Each envelope then carries an HMAC-SHA256 (sig) over the job name, payload and context, and the worker rejects a job whose signature is missing or wrong (JobSignatureError) before the handler runs — a tampered tenant, a foreign producer and a signed payload replayed under another job's name all fail. The signature does not stop a replay of an identical, genuine message by someone who can read the broker; idempotent handlers cover that.
Rolling it out: deploy the key to producers and workers together; jobs already queued unsigned are rejected by a worker that has a key, so drain the queues first (or retry them from the dead-letter/failed set after the rollout). An app that set a custom tenant-id grammar in tenancyPlugin passes the same function as validateTenantId.
Drivers
The backend is chosen by the plugin you register — each one builds its driver for you. Reach for queuePlugin({ driver }) only for a driver of your own.
| Driver | Package | delayed | priority | retries | backoff |
|---|---|---|---|---|---|
| BullMQ (Redis) | @basaltkit/queue-bullmq | ✅ | ✅ | ✅ | ✅ |
| RabbitMQ | @basaltkit/queue-rabbitmq | ✅ | ✅ | ✅ | ✅ |
| Amazon SQS | @basaltkit/queue-sqs | ✅ (≤15 min) | ❌ | ✅ | ✅ |
| Kafka | @basaltkit/queue-kafka | ❌ | ❌ | ✅ | ❌ |
| Sync (dev/test) | @basaltkit/queue | ❌ | ❌ | ✅ | ❌ |
Every bundled driver follows the same observability norm: infrastructure faults — broker down, a worker's connect failing at boot, a retry re-publish failing — surface through an onError-family option with a contextual console.error default (prefix [basalt:queue]). They are never unhandled rejections that kill the process, and never silent.
BullMQ (Redis)
bullmqQueuePlugin needs bullmq installed (see Install). Every driver option sits on the plugin beside jobs/workers, so observability costs you nothing extra:
import { bullmqQueuePlugin } from '@basaltkit/queue-bullmq'
bullmqQueuePlugin({
connection: process.env.REDIS_URL!,
jobs,
workers,
onError: (error, { queue, source }) => log.error({ queue, source, error }, 'queue infra error'),
onJobFailed: ({ queue, job, jobId, error }) => alertDeadJob(queue, job, jobId, error),
})The driver class is exported too, for the rare case where you build it yourself (sharing one driver between two plugins, or wrapping it):
import { BullmqQueueDriver } from '@basaltkit/queue-bullmq'
import { queuePlugin } from '@basaltkit/queue'
queuePlugin({ driver: new BullmqQueueDriver({ connection: process.env.REDIS_URL! }), jobs, workers })| Option | Type | Default | Why |
|---|---|---|---|
connection | string | ConnectionOptions | — (required) | Redis URL (redis:///rediss://, TLS inferred) or ioredis options. Percent-encoded credentials (p%40ss for p@ss) are decoded before they reach Redis. |
onError | (error, { queue, source: 'worker' | 'queue' }) => void | console.error with context | BullMQ emits infra errors (Redis down) as EventEmitter 'error' events — unhandled, they crash the process. The driver always attaches a listener; this option routes it to your logger/alerting. |
onJobFailed | ({ queue, job, jobId?, error }) => void | console.error with context | Fires once, when a job exhausts its retries (or throws BullMQ's UnrecoverableError). BullMQ emits 'failed' after every attempt; the driver skips the ones it is about to retry. Without it, permanently failed jobs were only visible by polling queue:stats. |
RabbitMQ
import { rabbitmqQueuePlugin } from '@basaltkit/queue-rabbitmq'
rabbitmqQueuePlugin({ url: process.env.AMQP_URL!, jobs, workers })Retries and backoff use a per-queue delay queue (<queue>.delay) whose messages TTL-expire back into the main queue; exhausted jobs land in <queue>.dead. Priority uses x-max-priority. Delivery safety: the driver prefers a publisher-confirm channel and only acks a message after the broker has confirmed any retry/dead-letter re-publish — acking before the publish is confirmed would be a silent job-loss window. close() drains in-flight handlers first; anything unfinished stays unacked, so the broker redelivers it.
If the channel or connection closes under the driver (broker restart, network cut, a channel-level protocol error), it reconnects: the next add() opens a fresh channel, and workers are re-subscribed with exponential backoff (reconnectDelayMs, doubling up to 30 s). Messages unacked on the dead channel are redelivered by the broker. The retry headers a message carries are untrusted: the attempt number is clamped to 1..50 and a negative backoff to 0.
| Option | Type | Default | Why |
|---|---|---|---|
url | string | — (required) | AMQP URL, e.g. amqp://user:pass@host:5672. |
onError | (error, { source: 'connection' | 'channel' }) => void | console.error with context | amqplib surfaces broker faults as EventEmitter 'error' events — unhandled, they crash the process. Also receives a worker's connect/consume failure at boot (otherwise the app would report healthy with zero workers) and a failed re-publish/ack after a job failure (the durable copy stays on the broker and is redelivered). |
maxPriority | number | 10 | x-max-priority for the priority queues. |
drainTimeoutMs | number | 10_000 | How long close() waits for in-flight handlers so their acks land on a live channel. Past it, unfinished jobs stay unacked and are redelivered — bounded shutdown without job loss. |
reconnectDelayMs | number | 1000 | First wait before re-subscribing workers after the channel/connection closed; doubles per consecutive failure, capped at 30 s. |
connect | AmqpConnect | amqplib | Injectable connector — tests run without a broker. |
Mixed delays at scale
The delay queue relies on per-message TTL, which only releases a message at the queue head (head-of-line blocking). For many different delays on one queue, prefer RabbitMQ's delayed-message-exchange plugin.
Amazon SQS
import { sqsQueuePlugin } from '@basaltkit/queue-sqs'
sqsQueuePlugin({ region: 'eu-west-1', queueUrl: (q) => QUEUE_URLS[q], jobs, workers })SQS has native per-message delay (≤ 15 minutes) but no priority. Retries and backoff are handled at the app level for parity with the other drivers: a failed message is re-sent with an incremented attempt and a DelaySeconds backoff (clamped to 15 min); an exhausted job goes to the dead-letter queue (<queue><deadSuffix>). The queueUrl resolver must map every queue name — including the DLQ names — to its SQS URL.
| Option | Type | Default | Why |
|---|---|---|---|
queueUrl | (queue: string) => string | — (required) | Resolves queue names (and <queue>-dead) to SQS URLs. |
region | string | SDK default | AWS region for the default client. |
deadSuffix | string | '-dead' | Suffix of the dead-letter queue's name. |
waitTimeSeconds | number | 20 | Long-poll wait per receive. |
visibilityTimeout | number | 30 | How long a received message stays hidden while processed. |
onError | (error, { queue, stage? }) => void | console.error with context | An SQS call failed; stage is 'receive' (network, credentials, queue deleted — without it the poller used to retry immediately and silently, a hot spin), 'delete' (the job succeeded but its message could not be deleted — SQS redelivers it after the visibility timeout; it is not re-sent as a failure) or 'reroute' (a retry/dead-letter re-send failed — the original is kept, so it is retried after the visibility timeout). Never fatal: the poller keeps running. |
errorPauseMs | number | 1000 | Pause between consecutive failed receives — bounds the retry rate against a broken endpoint. |
api | SqsApi | AWS SDK | Injectable API — tests run without AWS. |
A user-requested delay over 15 minutes throws SqsDelayTooLongError at dispatch (a backoff delay is clamped instead, so retries never throw).
Kafka
import { kafkaQueuePlugin } from '@basaltkit/queue-kafka'
kafkaQueuePlugin({ brokers: ['localhost:9092'], jobs, workers })Kafka is a log, not a task queue, and the driver is deliberately honest about that: no delayed, no priority, no backoff (Kafka cannot defer a message). Retries publish to a retry topic (<queue>.retry) the worker also consumes; exhausted jobs go to <queue>.dead. Worker concurrency maps to partitionsConsumedConcurrently, so effective parallelism is bounded by the topic's partition count.
| Option | Type | Default | Why |
|---|---|---|---|
brokers | string[] | — (required) | Kafka bootstrap brokers. |
clientId | string | 'basalt' | kafkajs client id. |
groupId | string | 'basalt-queue' | Consumer group workers join. |
retrySuffix | string | '.retry' | Suffix of the retry topic. |
deadSuffix | string | '.dead' | Suffix of the dead-letter topic. |
onError | (error, { source: 'consumer' | 'producer', queue? }) => void | console.error with context | source: 'consumer': the worker's connect/subscribe/run failed at boot — without this the app reports healthy with zero workers and the floating rejection is process-fatal. source: 'producer': a retry/dead-letter re-publish failed inside the consume callback (see below). |
client | KafkaClient | kafkajs | Injectable client — tests run without a broker. |
When a failed job's re-publish to the retry/dead topic itself fails (producer or broker outage), the driver reports it through onError and then rethrows so the message's offset is not committed — Kafka redelivers the message (at-least-once) instead of the job silently vanishing. Expect redelivery during a producer outage, never loss. (RabbitMQ keeps the message unacked for the same reason; not committing the offset is Kafka's equivalent.)
Sync (dev/test)
The inline driver's semantics are covered above: at-most-once, handler errors reject dispatch(), and an implicit production fallback warns at boot. For test assertions it records every execution in driver.executed ({ queue, jobName, attempts }), capped at the 1000 most recent entries (oldest evicted) so a long-lived process on this driver cannot leak memory.
Capability checks
Backends differ — Kafka has no message priority, SQS caps delays at 15 minutes, the sync driver runs inline. Rather than silently drop an option a backend can't honor, each driver declares its capabilities and the queue checks every dispatch against them.
queuePlugin({
driver: new KafkaQueueDriver({ brokers }),
onUnsupported: 'throw', // 'warn' (default) · 'throw' · 'ignore'
})
// a delayed job on Kafka:
await Job.dispatch(payload, { delay: '5m' })
// onUnsupported: 'warn' → logs once, runs immediately
// onUnsupported: 'throw' → throws UnsupportedJobOptionError
// onUnsupported: 'ignore'→ silently proceeds (legacy)Use 'throw' in production for a hard guarantee; the default 'warn' never breaks a dev run but never hides a dropped option either.
Job retention in Redis
With the BullMQ driver, finished jobs stay in Redis so you can inspect and retry them. By default completed jobs keep the last 1000, and failed jobs are kept forever — which means the failed set can grow unbounded. Control it with removeOnComplete / removeOnFail, globally on the plugin or per job:
// Global default for every job
bullmqQueuePlugin({
connection: process.env.REDIS_URL!,
jobs: [SendWelcome],
removeOnComplete: { age: '7d' }, // keep completed for 7 days
removeOnFail: { age: '14d' }, // failed no longer grow forever
})
// Per job — overrides the global default
defineJob({
name: 'email.welcome',
removeOnComplete: true, // remove as soon as it finishes
removeOnFail: { count: 500 }, // keep the last 500 failures
handle: () => {},
})Each option accepts true (remove on finish), false (keep all), a number (keep that many most-recent), or { age, count } where age is a duration like '14d'. Left unset, the defaults above apply. The sync driver ignores retention — it stores nothing. (The queue's own bull:<queue>:* structure keys always exist once the queue is created; that's BullMQ, not leftover jobs.)
Inspecting a queue — "did my job actually run?"
Two supported commands, and they answer different questions — counts and which jobs:
basalt queue:stats --queue orders
# → { waiting, active, completed, failed, delayed }
basalt queue:jobs --queue orders
# → id / name / state / attempts / age, newest firstRead the numbers with the lifecycle in mind, because this is where most debugging goes wrong:
dispatch ──▶ waiting ──▶ active ──▶ completed ← finished jobs live HERE
│ └────▶ failed (after exhausting attempts)
└── delayed (when dispatched with `delay`)A worker drains waiting in milliseconds, so a healthy queue shows waiting: 0, active: 0 almost always. That is success, not silence — the work is in completed. Inspecting only waiting/active is the classic false alarm: you conclude "nothing ran" when everything ran.
Reading it right
waiting climbing and completed flat → no worker is consuming (check workers: [{ queue }] matches defineJob({ queue })). completed climbing → the worker is working. failed climbing → jobs are exhausting their attempts; wire onJobFailed.
Listing the individual jobs
queue:stats gives you numbers; queue:jobs gives you the jobs themselves:
basalt queue:jobs --queue orders --states failed --limit 10id name state attempts age reason
1041 order.reconcile failed 3 2m Timeout after 30000ms| Flag | Default | What it does |
|---|---|---|
--queue | default | Which queue to inspect. |
--states | completed,failed,waiting,active | Comma-separated: waiting, active, completed, failed, delayed. An unknown state is rejected with the valid list. |
--limit | 20 | Maximum rows in total (newest first), capped at 1000. Must be a positive integer — 0, a negative or a non-number is refused (for queue:retry too, where --limit 0 used to re-enqueue every failed job). |
--payload | off | Also print each job's payload. Off by default — see the warning below. |
Why completed and failed are in the default states. Per the lifecycle above, a healthy queue has waiting: 0, active: 0. Defaulting to only those would print "no jobs" on a queue that is working perfectly — the same false alarm as reading only those counters. delayed is deliberately not in the default: it is a separate question, so ask for it (--states delayed).
In code, the same thing through the QUEUE manager:
import { QUEUE } from '@basaltkit/queue'
const jobs = await ctx().container.get(QUEUE).list('orders', {
states: ['failed'],
limit: 10,
})
if (!jobs) {
// The active driver cannot list — an honest "unsupported", NOT an empty queue.
} else {
for (const job of jobs) {
console.log(job.id, job.name, job.state, job.attemptsMade, job.payload)
}
}Each entry is a driver-neutral JobSummary — { id, name, state, attemptsMade, timestamp, payload, context?, failedReason? }. Two things it does for you:
payloadis your data, already unwrapped.dispatchwraps what you pass in an envelope so the request context survives the hop to the worker:jsonc{ "payload": { /* exactly what you passed to dispatch() */ }, "context": { /* requestId, tenantId, … — restored around handle() */ } }list()opens that for you:job.payloadis your object andjob.contextis the captured context. (Read the broker directly and you get the raw envelope, so your data sits atjob.data.payload.)It is the same shape on every driver. No BullMQ
Jobleaks through, so the code above doesn't break when you change brokers.
Job payloads are data
A payload can hold whatever you dispatched — including personal data. That's why basalt queue:jobs hides payloads unless you pass --payload, and why list() should be treated like the records it came from: don't log the result wholesale. Any endpoint you build on top of it must be authenticated (meta: { auth: true }) and authorized per tenant, and should return counts by default with raw payloads behind an explicit flag.
Which drivers can do this
stats() / retryFailed() / list() are optional driver capabilities. A driver that can't do one omits it, the manager returns undefined, and the CLI prints "Not supported" — an honest gap instead of a guess.
| Driver | Can list jobs? | Why |
|---|---|---|
bullmq | ✅ | Redis keeps the jobs; reading them changes nothing. |
sync | ❌ | Runs inline and stores nothing. |
rabbitmq | ❌ | AMQP has no non-destructive read — basic.get/consume hide the message from real workers and mark it redelivered. |
sqs | ❌ | ReceiveMessage starts the visibility timeout and bumps ApproximateReceiveCount; peeking could redrive jobs into the DLQ. |
kafka | ❌ | Reading is harmless, but a log has no per-message state — any waiting/completed would be invented. |
The rule the framework follows: looking at a queue must never change it. A list() built on a destructive read would be a debugging command that perturbs production, so those drivers omit it and point you at their own tooling and dead-letter destinations instead.
Going under the hood (and why you shouldn't need to)
Before queue:jobs existed, the only way to see individual jobs was the broker's own client — roughly 20 lines of ioredis + new Queue() + getJobs() + teardown. That works, and it also couples your app to one broker and hands you the raw envelope. Prefer list(). If you do go direct, the two classic traps are asking only for ['waiting','active'] (empty on a healthy queue) and forgetting that your data is at job.data.payload. And point the client at the same connection your app uses: a bare new Queue('orders') silently defaults to localhost:6379 and, wherever Redis lives elsewhere, reads a different queue and reports it as empty.
Run domain events on the queue
queuedOn bridges @basaltkit/events → queue: emit just enqueues a job, and the handler runs on the worker with retries and restored context. It returns the unsubscribe function; the created job is named listener:<event>. Job names are unique per manager, so a second queued listener on the same event needs its own name ({ name: 'order.created:crm' }) — without it, queuedOn throws DuplicateJobError instead of silently replacing the first listener's handler.
import { EventBus, defineEvent } from '@basaltkit/events'
import { QUEUE, queuedOn } from '@basaltkit/queue'
import { ctx } from '@basaltkit/core'
import { z } from 'zod'
const bus = new EventBus()
const manager = ctx().container.get(QUEUE)
const OrderCreated = defineEvent('order.created', z.object({ orderId: z.string() }))
const unsubscribe = queuedOn(bus, manager, OrderCreated, async ({ orderId }) => {
// runs on the worker, with the driver's retry/backoff
}, { queue: 'orders', attempts: 3 })
await bus.emit(OrderCreated, { orderId: 'o-1' })Errors
| Class | Code | When |
|---|---|---|
JobValidationError | JOB_INVALID | Payload fails the job's schema (thrown at dispatch; has .job and .issues) |
JobNotRegisteredError | QUEUE_JOB_NOT_REGISTERED | dispatch before the job was registered with a manager (add it to jobs) |
UnknownJobError | QUEUE_UNKNOWN_JOB | A job reached a worker that hasn't registered it (producer/worker job lists differ) |
UnsupportedJobOptionError | — | A dispatch requested an option the driver can't honor, under onUnsupported: 'throw' |
DuplicateJobError | QUEUE_DUPLICATE_JOB | A different job was registered under a name already taken in this manager (thrown at register/dispatch/queuedOn) |
JobSignatureError | QUEUE_BAD_SIGNATURE | With a signingKey: a job arrived unsigned or with a signature that does not verify (thrown in the worker; the handler never runs) |
JobContextError | QUEUE_INVALID_CONTEXT | A job's context carried a malformed tenant/tenantId/userId (thrown in the worker; the handler never runs) |
SqsDelayTooLongError | — | A delay over SQS's 15-minute maximum (@basaltkit/queue-sqs, thrown at dispatch) |
Failure modes & troubleshooting
| If you see | It means | Do |
|---|---|---|
Boot warning [basalt:queue] No 'connection' (Redis) configured… in production | The plugin silently fell back to the inline sync driver: at-most-once, no background retries, handler errors fail the dispatching request | Configure a Redis connection, or pass driver: new SyncQueueDriver() to opt in deliberately |
dispatch() rejects with your handler's error | Sync-driver semantics: errors propagate to the dispatcher by design (a broker driver would return immediately and retry in background) | Expected in dev/test; use a broker driver where you need background retries |
| A job is enqueued but never runs | No worker declared for the job's queue, or the worker queue name doesn't match the job's | Align defineJob({ queue }) with workers: [{ queue }] |
UnknownJobError in worker logs | The job reached a worker that hasn't registered it — producer and worker jobs lists differ | Register the same jobs array in both processes |
JobSignatureError in worker logs | The worker has a signingKey and the job was unsigned (queued before the rollout, or by a producer without the key) or altered on the broker | Give every producer the same key (or add the old key to the rotation list); investigate unexpected writers to the broker |
JobContextError in worker logs | The job's tenant id fails the worker's grammar (e.g. a custom tenancyPlugin grammar) or its tenant/user fields are malformed | Pass the same validateTenantId to the queue plugin as to tenancyPlugin; investigate unexpected writers |
Repeated [basalt:queue] bullmq worker error (queue "…") | Redis infra fault (connectivity, failover); BullMQ reconnects on its own | Route onError to alerting; check Redis |
[basalt:queue] job "…" on queue "…" failed permanently | The job exhausted its attempts; it stays in the failed set (default retention keeps all) | Inspect, fix the cause, basalt queue:retry --queue <q>; route onJobFailed to alerting |
UnsupportedJobOptionError at dispatch | The driver can't honor a requested option (e.g. delay on Kafka) under onUnsupported: 'throw' | Drop the option or switch drivers |
SqsDelayTooLongError at dispatch | A delay beyond SQS's 15-minute cap | Cap the delay, or use BullMQ/RabbitMQ for long delays |
[basalt:queue] kafka consumer error (queue "…") at boot | The worker's broker connect/subscribe failed — the process stays up but consumes nothing | Fix brokers/network and restart the worker; alert on this log line |
The same Kafka message is redelivered repeatedly, with [basalt:queue] kafka producer error alongside | A failed job's retry/dead-letter re-publish is failing, so the driver refuses to commit the offset — redelivery instead of silent loss | Restore the producer/broker; the backlog drains itself |
[basalt:queue] rabbitmq channel error after a job failure | The retry re-publish wasn't confirmed or the ack failed; nothing was acked, so the broker redelivers the durable copy | Check broker health; no action needed for the job itself |
Writing a driver
A driver is any object implementing the QueueDriver seam — four methods and an optional capability declaration:
import type { QueueDriver, DriverCapabilities, JobExecutor, AddJobOptions } from '@basaltkit/queue'
export class MyQueueDriver implements QueueDriver {
readonly name = 'my-backend'
// Declare what the backend honors. Omit it and the driver is assumed fully
// capable (back-compat) — but then nothing is checked, so prefer declaring it.
readonly capabilities: DriverCapabilities = { delayed: false, priority: false, retries: true, backoff: false }
private executor: JobExecutor | undefined
// The QueueManager calls this once, handing you how to run a received job.
setExecutor(executor: JobExecutor): void {
this.executor = executor
}
// Enqueue. `options` carries attempts/backoff/delayMs/priority — honor what
// your `capabilities` claim; the QueueManager has already applied its
// onUnsupported policy for the rest.
async add(queue: string, jobName: string, data: unknown, options: AddJobOptions): Promise<void> {
// publish { jobName, data, options } to your backend
}
// Start consuming `queue`. For each received job call
// `this.executor(jobName, data)`; on success remove it, on failure retry or
// dead-letter per your backend's model.
startWorker(queue: string, options?: { concurrency?: number }): void {
// consume → await this.executor?.(jobName, data)
}
async close(): Promise<void> {
// disconnect producers/consumers
}
}Then plug it in:
queuePlugin({ driver: new MyQueueDriver(), jobs, workers })Guidance for a faithful driver:
- Be honest in
capabilities. If the backend can't defer a message, setdelayed: false— the compatibility check turns a silent drop into a loud one. The bundled drivers are a reference:@basaltkit/queue-rabbitmq(delay + retries via a dead-letter queue),@basaltkit/queue-sqs(native delay, no priority),@basaltkit/queue-kafka(a log, so no delay/priority; retries via a retry topic). - Carry retry state in the message.
attempts/backoffcome fromadd; stamp the current attempt into message metadata so the worker knows when to retry versus dead-letter. Treat what you read back as untrusted: clamp the attempt to an integer in1..MAX_JOB_ATTEMPTS(a negative attempt would otherwise buy unlimited retries) and a backoff to a non-negative bound. - Carry
dataopaquely. It is the dispatch envelope (payload, context and, with asigningKey, its signature); round-trip it as JSON, unchanged. - Make the client injectable. Each bundled driver takes an injectable connector (
connect/client/api), so its retry and dead-letter logic is unit-tested without a running broker. Do the same and your driver is testable in CI.
See also
- Scheduled tasks — dispatch jobs on a cron schedule with
schedule.job(...). - Notes SaaS cookbook — queues wired into a real app (BullMQ + Redis, mailer off the request).