---
title: "Rust client"
description: "Complete surface of the queen-mq crate: Config and its defaults, every builder method, the Admin calls that reach a route that exists, the error type, and the streams DSL."
---

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

# Rust client

Crate `queen-mq`, edition 2021, Rust 1.75 or newer. Async on tokio. TLS is rustls with the
`ring` provider, matching the broker, so the build needs no OpenSSL and no cmake.

The wire types come from `queen-protocol`, the crate the broker also depends on in its tests.
Types that appear in this client's signatures are re-exported, so a caller does not need a
direct dependency on it.

## Constructor

```rust
Queen::connect(config: Config) -> Result<Queen>
Queen::connect_to(url: impl Into<String>) -> Result<Queen>
```

`connect` validates and builds the HTTP clients; it opens no connection. `Queen` is `Clone`,
and clones share one pool, one load balancer and one set of buffers.

### `Config`

| Field | Type | Default | Meaning |
| --- | --- | --- | --- |
| `urls` | `Vec<String>` | required | Broker URLs. More than one enables load balancing and failover. |
| `timeout` | `Duration` | 30 s | Per-request timeout. A long-poll pop adds 5 s of slack on top of its poll window. |
| `retry_attempts` | `u32` | 3 | Attempts against `5xx` and transport faults. Must be at least 1. |
| `retry_delay` | `Duration` | 1 s | First retry delay; doubles per attempt. |
| `strategy` | `Strategy` | `Affinity` | `RoundRobin`, `Session` or `Affinity`. |
| `affinity_hash_ring` | `usize` | 128 | Virtual nodes per backend on the consistent-hash ring. |
| `failover` | `bool` | `true` | Try another backend when one fails. |
| `health_retry_after` | `Duration` | 5 s | How long a failed backend stays out of the pool. |
| `bearer_token` | `Option<String>` | none | Sent as `Authorization: Bearer <token>`. |
| `headers` | `Vec<(String, String)>` | empty | Sent on every request. A `Host` entry is rejected. |
| `host_header` | `Option<HostHeader>` | none | Host to advertise, independent of the address dialled. |
| `retry_429` | `Retry429` | see below | Backoff policy for HTTP 429. |

Builder methods: `timeout`, `retry_attempts`, `retry_delay` is a field, `strategy`, `failover`,
`bearer_token`, `header`, `host_header` (returns `Result`), `retry_429`.

### `Retry429`

| Field | Type | Default | Meaning |
| --- | --- | --- | --- |
| `max_attempts` | `Option<u32>` | `Some(10)` | `None` means unbounded. A long-poll pop is unbounded unless this is set explicitly. |
| `base` | `Duration` | 500 ms | Exponential backoff base. |
| `cap` | `Duration` | 30 s | Backoff ceiling. |

`Retry-After` wins when the server sends one. Both paths get ±20 % jitter, so a fleet of
consumers does not come back in lockstep. A 429 is retried against the same backend and never
triggers failover: it is a quota signal, not a health signal.

### `HostHeader`

`HostHeader::parse` accepts a bare authority: `acme`, `acme.eu1.queenmq.cloud`,
`acme.local:6711`, `[::1]:6711`. Anything URL-shaped is rejected. The connection still goes to
the configured address; only the request authority and the TLS SNI are rewritten. With several
backends the client builds one connection-pinned HTTP client per backend, so failover keeps
working under a single virtual Host.

## `Queen`

| Method | Returns | Notes |
| --- | --- | --- |
| `queue(name)` | `QueueBuilder` | |
| `queue_opt(Option<String>)` | `QueueBuilder` | For discovery pops with no queue name. |
| `admin()` | `Admin` | |
| `transaction()` | `TransactionBuilder` | |
| `ack(&Message)` | `Result<AckResult>` | Completed. |
| `nack(&Message, reason)` | `Result<AckResult>` | Failed, with a DLQ reason. |
| `ack_with(&Message, AckStatus, Option<String>)` | `Result<AckResult>` | |
| `ack_all(&[Message])` | `Result<Vec<AckResult>>` | One request. All messages must share a group. |
| `nack_all(&[Message], reason)` | `Result<Vec<AckResult>>` | |
| `renew(&Message, Option<i32>)` | `Result<bool>` | Errors when the message holds no lease. |
| `flush_all_buffers()` | `Result<Vec<PushResult>>` | |
| `buffer_stats()` | `BufferStats` | |
| `close()` | `Result<()>` | Flushes buffers. |
| `shutdown_on_signal()` | `Result<()>` | Feature `signals`. Resolves on `SIGINT` or `SIGTERM`, then flushes. |

The consumer group and lease for an ack are read from the `Message`, not passed separately.

## `QueueBuilder`

Cloneable and consuming: every configuration method takes `self` and returns it.

### Addressing

`partition(name)`, `namespace(name)`, `task(name)`, `group(name)`, `name() -> Option<&str>`.

The partition defaults to `Default`, which the broker treats as the absence of a partition, so
`partition("Default")` and no partition at all address the same lane.

### Consumer options

| Method | Default | Meaning |
| --- | --- | --- |
| `concurrency(n)` | 1 | Parallel workers, each with its own poll loop. |
| `batch(n)` | 1 | Messages per poll. |
| `partitions(n)` | 1 | Partitions claimed per poll, sharing the batch budget and one lease. |
| `limit(n)` | none | Stop after this many messages, **across all workers**. |
| `idle(Duration)` | none | Stop after this long without a message. |
| `auto_ack(bool)` | `true` | Ack on handler success, nack on handler error. |
| `wait(bool)` | `true` | Long-poll instead of returning empty. |
| `poll_timeout(Duration)` | 30 s | Long-poll window. |
| `renew_lease(Duration)` | off | Extend the lease every interval while a handler runs. |
| `lease_seconds(i32)` | queue default | Per-request lease override. |
| `subscription_mode(SubscriptionMode)` | none | `New` or `All`, applied only when seeding a new cursor. |
| `subscription_from(String)` | none | `now` or an ISO-8601 timestamp. |
| `cancel(Cancel)` | none | Cooperative shutdown. |
| `buffer(BufferOptions)` | off | Client-side push batching. |

### Producing

| Method | Returns |
| --- | --- |
| `push(payload: impl Serialize)` | `Result<Vec<PushResult>>` |
| `push_many(impl IntoIterator<Item = impl Serialize>)` | `Result<Vec<PushResult>>` |
| `push_items(Vec<PushItem>)` | `Result<Vec<PushResult>>` |
| `flush_buffer()` | `Result<Vec<PushResult>>` |

With a buffer configured, `push` returns an empty vector: the messages are queued locally and
the per-item verdict arrives with the flush.

### Popping

| Method | Returns |
| --- | --- |
| `pop()` | `Result<Vec<Message>>` |
| `pop_auto_ack()` | `Result<Vec<Message>>` |

`pop` returns `Err` on failure rather than an empty vector. An empty claim and a claim refused
because pop maintenance is on both return `Ok` with nothing in it.

### The consume loop

```rust
consume<F, Fut, E>(handler: F) -> Result<ConsumeSummary>
    where F: Fn(Message) -> Fut, Fut: Future<Output = Result<(), E>>, E: Display

consume_batch<F, Fut, E>(handler: F) -> Result<ConsumeSummary>
    where F: Fn(Vec<Message>) -> Fut, ...
```

`ConsumeSummary` carries `processed`, `acked`, `nacked` and `stopped_by`, which is
`StopReason::Limit`, `Idle`, `Cancelled` or `Ended`.

### Lifecycle and DLQ

| Method | Returns |
| --- | --- |
| `configure(QueueOptions)` | `Result<Value>` |
| `create()` | `Result<Value>` |
| `delete()` | `Result<Value>` |
| `dlq(limit: Option<i32>, offset: Option<i32>)` | `Result<DlqResponse>` |

`configure` reports a failure that the broker returns inside a 200 body as an `Err`.

The DLQ builder takes only `limit` and `offset` because those, plus `queue` and
`consumerGroup`, are the only filters `GET /api/v1/dlq` reads. The other SDKs also send
`partition`, `from` and `to`, which the broker drops.

## `TransactionBuilder`

| Method | Notes |
| --- | --- |
| `ack(&Message)` | Completed. Collects the lease into `requiredLeases`. |
| `ack_with(&Message, AckStatus)` | `Retry` and `Dlq` survive to SQL rather than collapsing to a boolean. |
| `ack_all(impl IntoIterator<Item = &Message>)` | |
| `push(queue, payload)` | Returns `Result<Self>`; the payload is serialized here. |
| `push_to(queue, partition, payload)` | |
| `push_item(TxnPushItem)` | The only place `traceId` is honoured. A non-UUID is rejected. |
| `len()`, `is_empty()` | Staged operation count. |
| `commit()` | `Result<TransactionResponse>` |

Consecutive pushes to the same queue and partition merge into one operation, matching how the
broker groups frames. A rolled-back transaction is HTTP 200 with `success: false`, which
`commit` surfaces as an `Err`.

## `Admin`

Most methods return `serde_json::Value` verbatim; the endpoints whose shape is part of the
contract are typed.

### Resources

`overview()`, `namespaces()`, `tasks()`, `list_queues(params)`, `queue(name)`,
`partitions(params)`.

### Messages, DLQ and traces

`list_messages(params)`, `message(pid, txn)`, `delete_message(pid, txn)`,
`retry_message(pid, txn)`, `dlq(DlqParams)`, `record_trace(&TraceRequest)`,
`trace_names(params)`, `traces_by_name(name, params)`, `traces_for_message(pid, txn)`.

`dlq` returns `DlqResponse` and `record_trace` returns `TraceResponse`.

`retry_message` is a dead-letter replay: it re-pushes the DLQ snapshot and drops the DLQ row.
It errors on a message that is not in the DLQ.

### Consumer groups

`list_consumer_groups()`, `consumer_group(name)`, `lagging_consumers(min_lag_seconds)`,
`delete_consumer_group(group, delete_metadata)`,
`delete_consumer_group_for_queue(group, queue, delete_metadata)`,
`seek_consumer_group(group, queue, &SeekRequest)`,
`set_subscription_timestamp(group, timestamp)`, `refresh_consumer_stats()`.

Seeking is how a replay is done. `subscription_from` on a pop only seeds a cursor that does not
exist yet; it never moves one that does.

### Leases, status and system

`renew_lease(lease_id, seconds)`, `status(params)`,
`queue_stats(params)`, `queue_detail(name, params)`, `analytics(params)`,
`system_metrics(params)`, `worker_metrics(params)`, `postgres_stats()`, `health()`,
`metrics() -> String`, `maintenance()`, `set_maintenance(bool)`, `pop_maintenance()`,
`set_pop_maintenance(bool)`.

`RenewLeaseResponse::expires_at()` reads the expiry from whichever of `newExpiresAt`,
`expiresAt` or `lease_expires_at` is present. The broker writes the same value under all three,
one per SDK that looks for a different key.

`MaintenanceResponse::push_paused()` and `pop_paused()` return `Option<bool>`, because which of
the two keys a reply carries depends on the route: the `GET` on `/system/maintenance` reports
both, and each `POST` reports only its own. An absent key is `None` rather than `false`.

### Not present

`clear_queue` and `move_message_to_dlq` are absent. Their routes are not registered by the
broker and answer 404: `DELETE /api/v1/queues/:queue/clear` and
`POST /api/v1/messages/:partitionId/:transactionId/dlq`. Drop a queue with `QueueBuilder::delete`, and
dead-letter a message with `AckStatus::Dlq`.

## Errors

`queen_mq::Error` is an enum: `Config`, `Http { status, message, code, retry_after_seconds }`,
`Timeout`, `Network`, `Decode`, `AllBackendsFailed { attempted, last }`, `Invalid`.

| Method | Meaning |
| --- | --- |
| `status()` | The HTTP status, when there was one. |
| `code()` | The proxy's `ErrorCode`, when it sent one. |
| `is_rate_limited()` | 429. |
| `is_terminal_refusal()` | 403: suspended, gated, out of storage quota. |
| `is_retryable()` | 429, `5xx`, timeout or transport fault. Any other 4xx is not. |

`ErrorCode` covers `rate_limited`, `quota_exceeded`, `cluster_suspended`,
`storage_quota_exceeded`, `feature_gated` and `forbidden`; an unrecognized code parses into
`Other` and is treated as not retryable, which is the safe direction.

## `Message`

| Field | Type |
| --- | --- |
| `id` | `String` |
| `transaction_id` | `String` |
| `trace_id` | `Option<String>` |
| `data` | `serde_json::Value` |
| `producer_sub` | `Option<String>` |
| `created_at` | `String` |
| `partition_id` | `String` |
| `partition` | `String` |
| `lease_id` | `String` |
| `consumer_group` | `String` |

Plus `is_leased()`.

`trace_id` is always `None` for a message pushed through `/api/v1/push`: that path cannot store
one. A message pushed inside a transaction can carry it.

`lease_id` is empty on an `autoAck` delivery, which is what `is_leased()` reports.

## Streams

`queen_mq::streams::Stream`, built from a `QueueBuilder`.

| Stage | Methods |
| --- | --- |
| Stateless | `map`, `filter`, `flat_map` |
| Keying | `key_by` |
| Windows | `window_tumbling(seconds)`, `window_sliding(size, slide)`, `window_session(gap)`, `window_cron(Every)` |
| Window options | `grace_seconds`, `idle_flush_ms`, `event_time`, `allowed_lateness`, `on_late`, `window_options` |
| Reducing | `reduce(initial, fold)`, `aggregate_count`, `aggregate_sum`, `aggregate_min`, `aggregate_max`, `aggregate_avg` |
| Gating | `gate` |
| Terminal | `to`, `to_partitioned`, `foreach` |
| Run | `config_hash()`, `run(&Queen, RunOptions)` |

Chain rules are checked at `run`: a terminal must be last, `reduce` needs a window in front of
it, at most one window, reducer, `key_by` and `gate` each, and `gate` cannot share a stream
with windowing.

### `RunOptions`

| Field | Default | Meaning |
| --- | --- | --- |
| `query_id` | required | Durable identity. Two processes sharing it share state and cursor. |
| `batch_size` | 200 | Messages per cycle. |
| `max_partitions` | 4 | Partitions claimed per poll. |
| `max_wait` | 1000 ms | Long-poll wait on the source. |
| `subscription_mode`, `subscription_from` | none | Seed a new cursor. |
| `reset` | `false` | Wipe state when the chain's shape changed. |
| `consumer_group` | `streams.{query_id}` | |
| `cancel` | none | |

`run` returns a `StreamHandle` with `stop()`, `metrics() -> StreamMetrics` and `query_id()`.
`StreamMetrics` carries `cycles`, `flush_cycles`, `messages`, `push_items`, `state_ops`,
`late_events`, `gate_allowed`, `gate_denied` and `errors`.

Registering a chain whose `config_hash` differs from the stored one returns an `Err` naming
`reset`, because the existing state was computed under the old shape.

## Same surface, other languages

- [JavaScript](/reference/sdk/javascript)
- [Python](/reference/sdk/python)
- [Go](/reference/sdk/go)
- [PHP](/reference/sdk/php)
- [C++](/reference/sdk/cpp)

Source: https://queenmq.com/reference/sdk/rust/index.mdx
