---
title: "JavaScript"
description: "Install queen-mq, construct a client, and use the fluent grammar end to end: push, pop, consume, transactions, buffering and shutdown."
---

> Documentation Index
> Fetch the complete documentation index at: https://queenmq.com/llms.txt
> Use this file to discover all available pages before exploring further.

# JavaScript

The JavaScript client is the reference implementation: every other SDK is a transliteration
of it. It is ESM-only, requires Node 22 or newer (`engines.node: ">=22.0.0"`), and ships the
broker client and the streaming SDK in one package.

```sh
npm install queen-mq
pnpm add queen-mq
yarn add queen-mq
bun add queen-mq
```

```js
import { Queen } from 'queen-mq'
```

The package also exports `Admin`, the five defaults objects (`CLIENT_DEFAULTS`,
`QUEUE_DEFAULTS`, `CONSUME_DEFAULTS`, `POP_DEFAULTS`, `BUFFER_DEFAULTS`), and the streaming
entry points `Stream`, `tokenBucketGate` and `slidingWindowGate`. Named ESM exports mean a
bundler can drop the streaming code if you only import `Queen`.

## Construct a client

The constructor accepts three shapes: a URL string, an array of URLs, or a config object.

```js
const client = new Queen('http://localhost:6632')
```

```js
const client = new Queen(['http://broker-a:6632', 'http://broker-b:6632'])
```

```js
const client = new Queen({
  urls: ['http://broker-a:6632', 'http://broker-b:6632'],
  bearerToken: process.env.QUEEN_TOKEN,
  timeoutMillis: 30000,
  loadBalancingStrategy: 'affinity',
  handleSignals: false,
})
```

With one URL the client talks to it directly. With more than one it builds a load balancer
and enables failover, so a `5xx` or a network error marks that backend unhealthy and the
request is retried against the next one. A `4xx` is never retried across backends, and a
`429` is retried in place against the same backend (see [429 and 403](#429-and-403)).

Every URL is validated at construction. Passing neither `url` nor `urls` throws
`Must provide urls or url in configuration`.

### Client options

Defaults come from `CLIENT_DEFAULTS` in `utils/defaults.js`.

| Option | Default | What it does |
| --- | --- | --- |
| `url` / `urls` | n/a | Required. One backend or many. |
| `timeoutMillis` | `30000` | Per-request timeout. A long-poll pop adds 5 s of slack on top so the server-side timeout always fires first. |
| `retryAttempts` | `3` | **Total** tries for `5xx` and network failures, with exponential backoff from `retryDelayMillis`. |
| `retryDelayMillis` | `1000` | First retry delay; doubles per attempt. |
| `loadBalancingStrategy` | `'affinity'` | `'round-robin'`, `'session'` or `'affinity'`. Affinity hashes the pop's `queue:partition:consumerGroup` key onto a ring so the same lane keeps hitting the same broker. |
| `affinityHashRing` | `128` | Virtual nodes per backend on the affinity ring. |
| `enableFailover` | `true` | Try the next backend on `5xx` / network error. |
| `healthRetryAfterMillis` | `5000` | How long a backend stays marked unhealthy. |
| `bearerToken` | `null` | Sent as `Authorization: Bearer <token>` on every request. |
| `headers` | `{}` | Extra headers on every request. |
| `hostHeader` | `null` | Advertise a different `Host` than the address dialed. |
| `handleSignals` | `true` | Register `SIGINT` / `SIGTERM` handlers that flush and exit. |
| `logger` | `null` | Any object with `info`/`warn`/`error`. |
| `structuredLogs` | `false` | Call the logger pino/bunyan style: `info(fields, message)`. |
| `retry429` | n/a | `{ maxAttempts, baseMs, capMs }`. See below. |

> **Note**
>
> `headers: { Host: '...' }` cannot work: `Host` is on the WHATWG forbidden-header list and
> `fetch` drops it silently. The client detects it, maps it onto `hostHeader`, and warns
> loudly on `console.warn` (not gated by the log env var), because a wrong `Host` at a
> multi-tenant proxy either 421s or lands your traffic on another tenant's cluster.
> `hostHeader` must be a bare authority (`acme.example:6711`), never a URL, and it needs
> `undici`. In a browser the client refuses to send rather than advertise the wrong host.

## Push

`push()` takes one item or an array. Each item's body is `data`, or `payload`, or, if
neither key is present, the item itself. A `transactionId` you do not set is minted as a
UUIDv7 client-side. A `traceId` is forwarded only if it is a valid UUID.

The call returns a thenable builder, so `await` executes it. The resolved value is the
broker's array, one entry per item in request order:

```js
const res = await client.queue('orders').push([{ data: { id: 1 } }])
// res[0] === { index: 0, message_id: '...', transaction_id: '...', queueName: 'orders', status: 'queued' }
```

`status` is `queued`, `duplicate`, `buffered` or `failed`. `buffered` means the broker's
database write failed and the message went to the broker's disk spool for replay. The HTTP
status is still `201`, so checking the code alone hides it. The push callbacks below only
bucket `queued`, `duplicate` and `failed`, so a `buffered` item fires none of them: read
`status` off the returned array if you need to see it.

Supply your own `transactionId` to make a push idempotent inside the queue's dedup window:

Before awaiting, you can attach callbacks that split the response by status:

```js
await client
  .queue('orders')
  .push(items)
  .onSuccess(queued => metrics.add(queued.length))
  .onDuplicate(dups => log.warn({ n: dups.length }, 'already accepted'))
  .onError((failed, err) => log.error(err))
```

Without an `onError` callback, a `failed` item makes the await throw. With one, it does not.

## Pop

`pop()` resolves to an array of messages, or `[]`. Each message carries `id`,
`transactionId`, `traceId`, `data`, `producerSub`, `createdAt`, `partitionId`, `partition`,
`leaseId` and `consumerGroup`.

`wait` defaults to `true` on the builder, so a bare `pop()` is a 30-second long poll. Set
`.wait(false)` for an immediate return and `.timeoutMillis(n)` to change the poll window.

> **Caution**
>
> `pop()` catches everything and returns `[]`: a network failure, an exhausted `429` budget
> and a terminal `403` all look like an empty queue. The error is logged with its status and
> `code`, but the caller cannot see it. Use `consume()` when you need failures to surface.

Two knobs matter for throughput. `.batch(n)` caps how many messages come back;
`.partitions(n)` claims up to `n` partitions in one call, with `batch` acting as a global cap
across all of them and one `leaseId` covering the lot, so renewing once extends them all.
`.group('name')` selects a consumer group; without one the pop uses queue mode.

`.autoAck(true)` on a pop commits the cursor inside the pop transaction. That is at-most-once
delivery: a crash after the response is a lost message.

## Consume

`consume(handler)` starts `concurrency` workers, each in its own poll loop, and resolves when
every worker stops. It never returns on its own unless you set `limit`, `idleMillis` or an
`AbortSignal`.

| Method | Default | Effect |
| --- | --- | --- |
| `.group(name)` | queue mode | Consumer group. Its cursor is what advances on ack. |
| `.concurrency(n)` | `1` | Independent poll loops. |
| `.batch(n)` | `1` | Messages per pop. |
| `.partitions(n)` | `1` | Claim up to `n` partitions per pop; `batch` is the global cap. |
| `.each()` | off | Call the handler once per message instead of once per batch. |
| `.limit(n)` | none | Stop the worker after `n` messages. |
| `.idleMillis(n)` | none | Stop the worker after `n` ms with no messages. |
| `.autoAck(bool)` | `true` | Ack on return, nack on throw. Client-side, never sent to the server. |
| `.wait(bool)` | `true` | Long-poll the pop. |
| `.timeoutMillis(n)` | `30000` | Long-poll window. Short values let a shutdown flag be observed sooner. |
| `.renewLease(true, ms)` | off | Renew the batch's lease every `ms` while the handler runs. |
| `.subscriptionMode(m)` | none | `all` or `new`, applied only when the group's cursor is created. |
| `.subscriptionFrom(t)` | none | `now` or an ISO timestamp, same first-contact-only rule. |
| `consume(fn, { signal })` | n/a | An `AbortSignal` that stops every worker. |

Without `.each()` the handler receives the **array** of popped messages, even when `batch` is
`1`. That differs from Python, Go and PHP, which pass a single message in that case.

Each delivered message also gets a `trace()` method that posts a trace event and never
throws:

```js
await msg.trace({ traceName: ['tenant-acme'], eventType: 'step', data: { stage: 'charged' } })
```

### What autoAck actually does

`autoAck` in `consume()` is client-side. The consumer never sends `autoAck` to the server;
it calls `POST /api/v1/ack` itself after your handler returns, with status `completed` on
return and `failed` on throw. Because an ack is an offset commit, a nack clamps the group's
cursor at the failed message, so **everything after it in the popped batch will be
redelivered**. With `.each()` the consumer detects this and abandons the rest of the batch
rather than producing duplicates and rejected acks. Without `.each()` your handler owns the
whole batch, so it should either succeed for all of it or throw.

Attaching `.onSuccess()` or `.onError()` to a consume forces `autoAck` off, on the assumption
that the callback acks. If it does not, nothing is acked and the lease expires.

## Acknowledge manually

```js
const messages = await client.queue('orders').group('billing').batch(10).pop()
for (const msg of messages) { await handle(msg) }
await client.ack(messages, true, { group: 'billing' })
```

`ack(message, status, context)` takes one message or an array. `status` is `true`
(`completed`), `false` (`failed`), or one of the strings `completed`, `failed`, `retry`,
`dlq`. `context.group` names the consumer group; `context.error` records a failure reason
that lands on the DLQ row if this nack exhausts the retry budget.

`partitionId` is mandatory on every ack: the client throws (batch) or returns
`{ success: false }` (single) without it, because a `transactionId` alone is not unique across
partitions. Pass the message objects you got from the pop and it is already there.

A rejected ack still arrives as HTTP 200 with `success: false` on the item, so the per-item
flag is the only signal that the broker accepted it. The client checks the response length
against the number of acks sent and throws if they disagree, rather than misattributing a
result to the wrong message.

Per-message statuses in one batch call:

```js
messages[0]._status = 'dlq'
messages[0]._error = 'unparseable payload'
await client.ack(messages, true, { group: 'billing' })
```

## Transactions

`transaction()` bundles pushes and acks into one PostgreSQL transaction, all or nothing. The
builder collects `.queue(name).partition(p).push(items)` and `.ack(messages, status)`
operations and sends them on `.commit()`, together with the `leaseId`s of every acked message
as `requiredLeases`, so the commit fails if a lease expired in the meantime. Committing with
no operations throws.

This is the primitive for a pipeline stage: ack the input and push the output atomically, so
a crash cannot produce a stage that consumed without emitting.

> **Note**
>
> A transaction is not exactly-once end to end. If the response is lost and you retry, you
> duplicate the pushes unless their `transactionId`s are deterministic and still inside the
> dedup window.

## Lease renewal

```js
await client.renew(messages)
```

`renew()` accepts a lease ID, a message, or an array, and de-duplicates the lease IDs. A
multi-partition pop shares one lease across every claimed partition, so passing the whole
array issues one HTTP call, not N.

> **Caution**
>
> The returned `newExpiresAt` is always `null`. The broker answers
> `{"renewed":N,"expiresAt":"..."}` and the client reads `newExpiresAt` / `lease_expires_at`.
> The renewal happens; the reported expiry does not exist. Schedule the next renewal from your
> own interval, which is what `.renewLease(true, ms)` does.

## Client-side buffering

```js
await client
  .queue('events')
  .buffer({ messageCount: 500, timeMillis: 200 })
  .push([{ data: { seen: Date.now() } }])
```

With `buffer()`, `push()` returns `{ buffered: true, count: N }` immediately and the messages
are flushed when the buffer reaches `messageCount` (default `100`) or `timeMillis` (default
`1000`) elapses since the first message, whichever comes first. Buffers are keyed per
`queue/partition`.

> **Caution**
>
> A buffered push gives up two things. You never see the per-item `queued` / `duplicate`
> verdict, because the response arrives after your call returned. And the buffer is in process
> memory: an unflushed buffer is lost if the process dies without running `close()`.

Flush explicitly with `await client.queue('events').flushBuffer()` for one queue/partition or
`await client.flushAllBuffers()` for all of them. `client.getBufferStats()` returns
`{ activeBuffers, totalBufferedMessages, oldestBufferAge, flushesPerformed }`.

## Admin

```js
const overview = await client.admin.getOverview()
const queues = await client.admin.listQueues({ limit: 50 })
const lagging = await client.admin.getLaggingConsumers(120)
```

`client.admin` is a lazily created singleton over the observability and management routes:
resources (overview, namespaces, tasks, queues, partitions), messages, traces, status and
analytics, consumer groups (list, get, lagging, delete-for-queue, seek), health, metrics and
the two maintenance-mode switches. The full route list, with the access level each one
requires, is in [Reference](/reference).

> **Caution**
>
> `admin.clearQueue()` and `admin.moveMessageToDLQ()` call paths the broker does not register
> (`DELETE /api/v1/queues/:queue/clear` and
> `POST /api/v1/messages/:partitionId/:transactionId/dlq`). Both return 404.
> `admin.retryMessage()` and `admin.deleteMessage()` are real.

## Dead-letter queue

```js
const { messages } = await client.queue('orders').dlq('billing').limit(50).offset(0).get()
```

`dlq()` reads `GET /api/v1/dlq` and swallows errors to `{ messages: [], total: 0 }`. Each row
carries `queue`, `partition`, `consumerGroup`, `errorMessage`, `retryCount`, `data`,
`createdAt` and `failedAt`.

> **Note**
>
> `total` is always `undefined`: the broker answers
> `{"messages":[…],"pagination":{"limit":…,"offset":…}}` and has no `total` field. And
> `.from()`, `.to()` and the builder's partition are sent as query parameters the broker's DLQ
> handler never reads (it collects only `queue`, `consumerGroup`, `limit` and `offset`), so
> they are silently ignored. Page with `limit`/`offset` and filter the rest client-side.

To reprocess one dead letter, `admin.retryMessage(partitionId, transactionId)` calls the
broker's replay route, which re-pushes the stored payload and then drops the DLQ row. The
replayed message is a new message: it gets a fresh identity and lands at the tail of its
partition, so ordering relative to the original is not preserved. There is no bulk redrive:
to replay many rows, read them and push them yourself.

## 429 and 403

```js
const client = new Queen({
  url: 'https://cell.example',
  bearerToken: token,
  retry429: { maxAttempts: 20, baseMs: 500, capMs: 30000 },
})
```

A `429` is retried in place against the same backend. The delay honours `Retry-After`
(seconds) when present, otherwise `baseMs * 2^attempt` capped at `capMs`, and both get ±20 %
jitter to avoid a synchronised herd. `maxAttempts` defaults to `10` for push and admin calls;
a long-poll pop (`wait: true`) retries unboundedly unless you set `maxAttempts`, in which case
it applies to both. A `429` never marks the backend unhealthy.

A `403` is terminal. The consume worker logs it with its `code` and throws out of the worker
rather than retrying, because `cluster_suspended`, `storage_quota_exceeded`, `feature_gated`
and `forbidden` never resolve on their own. Errors carry `.status` and, when the body has
one, `.code`, so you can branch without matching on message text.

## Logging

Console logging is off unless `QUEEN_CLIENT_LOG=true` (Node) or
`window.QUEEN_CLIENT_LOG === true` (browser). Passing a `logger` overrides the env gate
entirely, so the logger is then always active. It must implement `info`, `warn` and `error`;
`debug` falls back to `info`.

```js
const client = new Queen({ url, logger: pino(), structuredLogs: true })
```

By default the logger receives one formatted string, `[operation] {json}`. With
`structuredLogs: true` the detail keys become top-level fields and the operation becomes the
message, pino/bunyan style. Use it only with a logger that accepts a leading merge object.

## Graceful shutdown

```js
await client.close()
```

`close()` flushes every buffer, tears down the buffer manager, removes the signal handlers
and destroys the HTTP agent. That last step matters: the client owns its own `undici` Agent
so `close()` can release the keep-alive sockets. Without it, Node's event loop stays pinned
by open sockets and the process never exits on its own.

Unless you pass `handleSignals: false`, the constructor registers `SIGINT` and `SIGTERM`
handlers that call `close()` and `process.exit(0)`. A second signal exits immediately with
code 1. When Queen is embedded in an application with its own lifecycle, either construct
with `handleSignals: false` or call `client.disableGracefulShutdown()`;
`client.enableGracefulShutdown()` puts them back.

## Streaming

```js
import { Stream } from 'queen-mq'

const stream = Stream.from(client.queue('events').group('agg'))
```

The windowed streaming SDK ships in the same package and builds on the same client. It has
its own concepts (windows, watermarks, gates, state) and is documented separately in
[Use Queen](/use).

Source: https://queenmq.com/use/clients/javascript/index.mdx
