--- title: "Broker bindings" description: "How each broker carries the canonical envelope natively — the Redis, RabbitMQ, Amazon SQS, Azure Service Bus, Pulsar, Kafka and Artemis bindings. The body is identical everywhere; a binding maps the contract onto each broker's native metadata so consumers route and trace without decoding the body." source: https://babelqueue.com/docs/spec/1.x/broker-bindings/ updated: 2026-06-15T00:00:00.000Z --- # Broker bindings The [envelope](/docs/spec/1.x/envelope/) is **broker-agnostic** — it is just the message body. A **binding** says how a given broker carries that body natively, so a consumer can route by URN and continue a trace **without decoding the body**, and so a message one SDK produces is consumed byte-for-byte by another over the same broker. Two rules hold for every binding: - The **body is always the canonical envelope** — compact UTF-8 JSON, byte-identical across SDKs. Native metadata is a *redundant, routable projection* of the body, never a replacement for it. - The envelope is **frozen at `schema_version: 1`**. A new broker is purely additive — it never changes the wire format. ## Redis Reliable-queue pattern: `RPUSH` to produce, a blocking move to a `:processing` list to reserve, then remove on ack. Redis lists carry no per-message metadata, so routing and tracing read the envelope body directly (`job`/`urn`, `trace_id`). At-least-once via the processing list. ## RabbitMQ (AMQP 0-9-1) The body is the envelope; the contract is also projected onto AMQP properties and headers so a consumer can route on the broker's metadata: | AMQP field | Carries | | :--- | :--- | | `type` | the URN (`job`) | | `correlation_id` | `trace_id` | | `message_id` | `meta.id` | | header `x-attempts` | `attempts` | | header `x-schema-version` | `meta.schema_version` | | header `x-source-lang` | `meta.lang` | Messages are `application/json`, persistent delivery, on durable queues; consumed with `basic.get` + manual ack (at-least-once). ## Amazon SQS The canonical envelope is the **`MessageBody`**. On produce, the transport projects the envelope onto native **`MessageAttributes`** — a routable view of the body. Ids and strings are `DataType` **String**; counters are `DataType` **Number**: | MessageAttribute | DataType | Value | | :--- | :--- | :--- | | `bq-job` | String | the URN (`job`) | | `bq-trace-id` | String | `trace_id` | | `bq-message-id` | String | `meta.id` | | `bq-schema-version` | Number | `meta.schema_version` (`1`) | | `bq-source-lang` | String | `meta.lang` | | `bq-created-at` | Number | `meta.created_at` as epoch milliseconds | **Consuming** uses SQS's native delivery semantics: - A receive **reserves** the message for the visibility timeout; a successfully handled message is removed with `DeleteMessage`. A failing handler simply does **not** delete it — SQS redelivers after the visibility timeout (at-least-once). - **`attempts` is reconciled** to `max(body.attempts, ApproximateReceiveCount − 1)`. The first delivery reads `0`; a runtime-incremented count is **never lowered**; an absent or non-numeric `ApproximateReceiveCount` is ignored. A drop-in driver that instead surfaces the broker's native count (e.g. Laravel's `SqsJob::attempts()`) documents that divergence. **FIFO queues** (`.fifo`) set `MessageGroupId` (the configured group, else the queue name) and `MessageDeduplicationId` = `meta.id` — unless the queue uses content-based deduplication. Delayed delivery uses `DelaySeconds`, capped at SQS's 900-second maximum. This binding is implemented identically across every SDK that ships an SQS transport (Go, Python, Node, Java, PHP, .NET) and is locked by the [conformance suite](https://github.com/BabelQueue/conformance), so the projected attributes and the reconciled `attempts` are guaranteed to match. ## Azure Service Bus The canonical envelope is the **`Body`**. Azure Service Bus has first-class native slots for almost every envelope concept, so the binding maps onto native message fields and needs only two custom application properties: | Field | Value | | :--- | :--- | | `Subject` (a.k.a. Label) | the URN (`job`) — **route on `Subject` without reading the body** | | `CorrelationId` | `trace_id` | | `MessageId` | `meta.id` (enables ASB duplicate detection) | | `ContentType` | `application/json` | | `DeliveryCount` | broker-maintained, **1-based** — the native attempts source | | `ApplicationProperties["bq-schema-version"]` | `meta.schema_version` (`1`) | | `ApplicationProperties["bq-source-lang"]` | `meta.lang` | | `ApplicationProperties["bq-created-at"]` | `meta.created_at` (epoch ms, convenience mirror) | Application properties are **native AMQP-typed values** (numbers stay numbers), not the `DataType`-wrapped strings SQS uses. **Consuming** uses ASB's PeekLock model: - A receive **reserves** the message for the lock duration; a handled message is removed with `Complete`. A failing handler `Abandon`s it — the broker redelivers and **increments `DeliveryCount`** (at-least-once). At `MaxDeliveryCount` ASB auto-moves it to the native `$DeadLetterQueue` sub-queue. - **`attempts` is reconciled** to `max(body.attempts, DeliveryCount − 1)`: `DeliveryCount` (1-based) is the native floor (first delivery reads `0`), and a runtime-incremented body count is never lowered. The rule is identical for the native-consumer SDKs (.NET, Java, Node) and the runtime-transport SDKs (Python, Go). **Delayed delivery** is native — set `ScheduledEnqueueTime` to `now + delay`. **Auth** is a connection string or Azure AD (`DefaultAzureCredential`); transport is AMQP 1.0 over TLS (or WebSockets). An optional `SessionId` (FIFO) is opt-in and not a contract field. This binding is implemented identically across every SDK that ships an ASB transport (.NET, Java, Python, Node, Go) and is locked by the [conformance suite](https://github.com/BabelQueue/conformance). PHP is deferred (no modern official client). ## Apache Pulsar The canonical envelope is the **message payload**. Pulsar message properties are **string→string**, so the binding projects a redundant, routable view of the body onto `bq-` properties (every value stringified) and keeps the body authoritative: | Property | Value | | :--- | :--- | | `bq-job` | the URN (`job`) — **route on `bq-job` without reading the body** | | `bq-trace-id` | `trace_id` | | `bq-message-id` | `meta.id` | | `bq-schema-version` | `meta.schema_version` (`"1"`) | | `bq-source-lang` | `meta.lang` | | `bq-attempts` | `attempts` — the **authoritative** count, carried in the body | | `publishTime` | mirrors `meta.created_at` (broker-set; body authoritative) | Properties are strings (Pulsar has no typed properties), so numbers are stringified — unlike ASB's native AMQP-typed values. **Consuming** receives one message at a time: - A handled message is `acknowledge`d. A failing handler `negativeAcknowledge`s it — the broker redelivers it (at-least-once) and **increments `RedeliveryCount`**. With a native `DeadLetterPolicy` it eventually moves to the cross-language `.dlq` topic. - **`attempts` is reconciled** to `max(body.attempts, RedeliveryCount)`. `RedeliveryCount` is **0-based** (0 on first delivery), so it maps directly with **no −1** — and a runtime-incremented body count is never lowered. The rule is identical for the native-consumer SDKs (.NET, Java, Node) and the runtime-transport SDKs (Python, Go). **Delayed delivery** is native — `deliverAfter` (relative) or `deliverAt` (absolute), also mirrored on a `bq-delay` property. The default subscription is **`Shared`**, named `babelqueue`; topics default to `persistent://public/default/`. **Auth** is via the service URL (`pulsar://` or `pulsar+ssl://`) plus any client-configured TLS/token. This binding is implemented identically across every SDK that ships a Pulsar transport (.NET, Java, Python, Node, Go) and is locked by the [conformance suite](https://github.com/BabelQueue/conformance). **PHP** reaches Pulsar over Pulsar's native **WebSocket API** (`PulsarTransport` to produce, `PulsarConsumer` to consume) — the same envelope and `bq-` properties, just a different access path; it round-trips live against the native SDKs. ## Apache Kafka Kafka is a partitioned, append-only **log** with consumer-group **offset commits** — not a queue with per-message ack. It has **no native** per-message ack, delayed delivery, dead-letter queue, or delivery counter, so this binding absorbs all four in the transport layer (the envelope stays `schema_version: 1`). The canonical envelope is the record **value**; the contract fields move to `bq-` **record headers** (UTF-8 byte strings, so integers are stringified): | Field | Value | | :--- | :--- | | `value` | the canonical envelope JSON | | `timestamp` | mirrors `meta.created_at` (Unix ms) | | header `bq-job` | the URN (`job`) — **route on `bq-job` without reading the body** | | header `bq-trace-id` | `trace_id` | | header `bq-message-id` | `meta.id` | | header `bq-schema-version` | `meta.schema_version` (`"1"`) | | header `bq-source-lang` | `meta.lang` | | header `bq-attempts` | `attempts` — the **authoritative** retry counter (Kafka has no native one) | **Consuming** is **process-then-commit** (manual commit, `enable.auto.commit = false`): a record is reserved by being polled, the handler runs, and only then is the offset committed (`Commit`). A crash before commit redelivers it on the next poll — at-least-once; handlers MUST be idempotent (dedupe on `meta.id`). - **`attempts` is reconciled** to the **`bq-attempts` header when present** (authoritative), else the body's own `attempts` (the fallback for a non-BabelQueue producer). This is **not a max** — the header overrides the body. - **Retry / delay** is SDK-owned, since Kafka has neither: a failing handler republishes the envelope to a tiered **retry topic** `.retry.` with `bq-attempts + 1` (and `bq-delay` / `bq-original-topic`), then commits; a retry-topic worker re-injects it after the tier delay. A delay or release with **no retry topics configured raises** rather than silently dropping. - **Terminal failure → DLQ:** at max-tries the envelope goes to `.dlq` with the additive `dead_letter` block (opt-in; if disabled, terminal failures degrade to commit-and-drop). **Connection** is `bootstrap.servers` + a `group.id`; **auth** is `SASL_SSL` / `SSL` via the native client. This binding is implemented identically across Java, Go, Node, Python and .NET and is locked by the [conformance suite](https://github.com/BabelQueue/conformance). **PHP** ships both halves over **`ext-rdkafka`** (opt-in — the one binding that relaxes the zero-extension rule): `KafkaTransport` to produce, `KafkaConsumer` (process-then-commit) to consume, plus the SDK-owned retry-topic machinery (`KafkaRetryRouter` + `KafkaRetryConsumer`); proven live Java→PHP. ## Apache ActiveMQ Artemis Artemis speaks **AMQP 1.0** (not RabbitMQ's 0-9-1), and gives the binding native primitives — per-message settlement, scheduled delivery, a delivery counter and a dead-letter address — so it maps onto them rather than re-implementing. The canonical envelope is the **message body**; the contract fields ride the slots a JMS peer reads, plus the string-valued `bq_` application properties: | Field | Value | | :--- | :--- | | body | the canonical envelope JSON | | annotation `x-opt-jms-type` | the URN (`job`) → JMSType — **route on it without reading the body** | | `correlation-id` | `trace_id` → JMSCorrelationID | | `creation-time` | `meta.created_at` (Unix ms) → JMSTimestamp | | property `bq_schema_version` | `meta.schema_version` (`"1"`) | | property `bq_source_lang` | `meta.lang` | | property `bq_attempts` | `attempts` — a 0-based mirror (the body stays authoritative) | | property `bq_app_id` | `"babelqueue"` | These property names use **underscores**, not the hyphens of the Kafka/Pulsar bindings: a JMS property name must be a valid Java identifier (no `-`), and the Java binding produces/consumes over JMS, so every Artemis SDK uses the same JMS-legal `bq_` form for cross-protocol parity. Unlike Kafka/Pulsar, the URN / `trace_id` / message-id are **not** `bq_` properties — they ride the JMS-native slots (`x-opt-jms-type`, `correlation-id`, the broker-set message-id), so a plain JMS consumer routes and correlates with no BabelQueue awareness. **Consuming** reserves one message at a time and settles it per message: - A handled message is **`accept`ed** (acknowledged). A failing handler **`release`s** it — the broker redelivers it (at-least-once) and **increments the AMQP `delivery-count`**. At max-tries the envelope goes to the cross-language `.dlq` with the additive `dead_letter` block (alongside Artemis's own dead-letter address). - **`attempts` is reconciled** to `max(body.attempts, delivery-count)`. The AMQP `delivery-count` is **0-based** (0 on first delivery), so it maps directly with **no −1** — and a runtime-incremented body count is never lowered. The Java binding consumes over **JMS** and reads the 1-based `JMSXDeliveryCount`, subtracting 1 to arrive at the **same** 0-based `attempts`. The rule is identical for the native-consumer SDKs (.NET, Java, Node) and the runtime-transport SDKs (Python, Go). **Delayed delivery** is native — JMS 2.0 `setDeliveryDelay` (Java) or the `x-opt-delivery-time` annotation (AMQP), also mirrored on a `bq_delay` property. **Connection** is the broker's AMQP acceptor (`amqp://` / `amqps://`, default port 5672); Java connects over JMS via the Artemis client. This binding is implemented identically across Java (JMS), .NET, Python, Node and Go (all AMQP 1.0) and is locked by the [conformance suite](https://github.com/BabelQueue/conformance). **PHP** reaches Artemis over **STOMP** (`StompTransport` to produce; Laravel ships the `babelqueue-artemis` drop-in consume driver) — Artemis bridges STOMP ↔ AMQP 1.0 ↔ JMS on the same address, so a STOMP-produced message is consumed natively by the JMS and AMQP-1.0 SDKs, and vice versa. ## Other brokers With Amazon SQS, Azure Service Bus, Apache Pulsar, Apache Kafka and Apache ActiveMQ Artemis all shipped, the binding catalogue covers every broker on the expansion roadmap. New brokers follow the same shape — body identical, contract fields projected onto native metadata — and ship as additive MINOR releases. The envelope stays `schema_version: 1`.