KubeMQ
ConnectorsRabbitMQ (AMQP 0-9-1)How-to guides

Authentication

How a RabbitMQ (AMQP 0-9-1) client authenticates to KubeMQ — PLAIN/AMQPLAIN with a JWT or a connector user, certificate login, token expiry, and authorization.

The KubeMQ RabbitMQ connector offers SASL PLAIN and AMQPLAIN always, and EXTERNAL (certificate login) once it is switched on. What the password means depends on one thing checked first: whether connector-local users are configured. With none (the default), the password carries a KubeMQ JWT on a secured broker and the username is cosmetic; with users configured, both are checked against that store. When authentication is disabled and no users are configured (the dev default), any credentials are accepted, so the examples clone-and-run with guest:guest.

On a stock dev broker, authentication is off and any credentials are accepted — guest:guest completes a full round-trip. SASL PLAIN with a KubeMQ JWT in the password is the one credentialed form that also runs on a stock broker. To encrypt the JWT on the wire, use the TLS listener (amqps://:5671) — see TLS and mTLS.

The four login paths

connection.start advertises PLAIN AMQPLAIN always; EXTERNAL is added only when Connectors.Amqp.SslCertLogin = true and the TLS handshake presented a certificate the listener verified. ANONYMOUS is never offered, in any configuration — RabbitMQ 4.3.4 advertises it by default, so this is a recorded deviation; a client that selects it is refused connection.close(503) naming the mechanism.

The credential store is consulted first — Authentication.Enable is not the discriminator:

PathSelected whenWhat the credential meansWhat becomes the identity
Connector-local credential storeConnectors.Amqp.Credentials has at least one entry — whatever Authentication.Enable isBoth the username and the password are checked against the store; the platform token is not consultedamqp-_user- + the username, encoded — see Users and permissions
Platform tokenno store and Authentication.Enable = trueThe password is a KubeMQ JWT; the username is recorded for audit and user-id checks but never authenticatedThe token's ClientID claim, verbatim and unprefixed. A token with an empty ClientID is refused
Anonymousno store and Authentication.Enable = falseAny credentials are accepted; nothing is checkedamqp-{uuid8}, minted by the server fresh on every connection — the client's connection_name never reaches it
SASL EXTERNALSslCertLogin = true and a verified client certificate; independent of Authentication.Enable, but consults the credential store when one is configuredNo password; the verified certificate is the credentialamqp-_cert- + the certificate field named by SslCertLoginFrom, encoded

AMQPLAIN is the same two credential rules as PLAIN — it exists for older client libraries that never adopted PLAIN's SASL encoding; its initial response is an AMQP field table with LOGIN and PASSWORD fields rather than PLAIN's null-separated triple.

Failures are indistinguishable on the wire. A wrong password, an unknown username, and no matching entry at all close the connection with the identical connection.close(403) and ACCESS_REFUSED - authentication failed (sent before the TCP close, since authentication_failure_close is advertised), plus an auth.failure audit event. Credential comparison is constant-time. This matches RabbitMQ 4.3.4, measured directly.

Password = KubeMQ JWT, username = cosmetic

With no connector-local users configured, the username slot is informational and the password slot carries the credential. Most AMQP 0-9-1 clients accept a (username, password) pair in the connection URL — put the JWT in the password:

amqp://<username-ignored>:<KubeMQ-JWT>@host:5672/<vhost>
        └─ cosmetic ─┘  └─ authenticated ─┘
FieldRole
UsernameCosmetic. Recorded for audit and user-id checks; ignored for authorization.
PasswordThe KubeMQ JWT — passed to the auth service.
IdentityClientID from the JWT's claims. This single ClientID covers both connector-level and channel-level authorization.

Token expiry ends the connection. A connection whose JWT expires is closed with connection.close(320) — reply text CONNECTION_FORCED - credential expired — at the instant the token's exp passes, matching RabbitMQ 4.3.4's OAuth 2 behavior. There is no grace period and no way to switch it off. A token carrying no exp claim sets no deadline. Refresh with connection.update-secret before the old token expires — see Refreshing a credential.

The JWT travels in the SASL password in cleartext at the AMQP layer. Production deployments that use authentication MUST use the TLS listener (amqps://:5671), otherwise the JWT is exposed on the wire. See TLS and mTLS and Auth & security.

Put the JWT in the password slot regardless of client library:

// amqp091-go — username is cosmetic, password is the KubeMQ JWT.
conn, err := amqp.Dial(fmt.Sprintf(
    "amqp://audit-user:%s@broker:5672/", os.Getenv("KUBEMQ_AMQP_JWT")))
if err != nil {
    log.Fatalf("dial (bad/expired JWT? auth-disabled broker?): %v", err)
}
defer conn.Close()
# pika — credentials: username audit-only, password is the KubeMQ JWT.
creds = pika.PlainCredentials("audit-user", os.environ["KUBEMQ_AMQP_JWT"])
params = pika.ConnectionParameters(host="broker", port=5672, credentials=creds)
conn = pika.BlockingConnection(params)
// amqp-client — setUsername is cosmetic; setPassword carries the KubeMQ JWT.
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("broker");
factory.setPort(5672);
factory.setUsername("audit-user");
factory.setPassword(System.getenv("KUBEMQ_AMQP_JWT"));
Connection connection = factory.newConnection();
// amqplib — username audit-only, password is the KubeMQ JWT.
const connection = await amqp.connect(
  `amqp://audit-user:${process.env.KUBEMQ_AMQP_JWT}@broker:5672/`);
// RabbitMQ.Client — UserName is cosmetic; Password carries the KubeMQ JWT.
var factory = new ConnectionFactory
{
    HostName = "broker",
    Port = 5672,
    UserName = "audit-user",
    Password = Environment.GetEnvironmentVariable("KUBEMQ_AMQP_JWT"),
};
using var connection = factory.CreateConnection();
# bunny — username audit-only, password is the KubeMQ JWT.
conn = Bunny.new(
  host: "broker", port: 5672,
  user: "audit-user", password: ENV.fetch("KUBEMQ_AMQP_JWT"))
conn.start
// lapin — username audit-only, password is the KubeMQ JWT.
let uri = format!("amqp://audit-user:{}@broker:5672/",
    std::env::var("KUBEMQ_AMQP_JWT")?);
let conn = Connection::connect(&uri, ConnectionProperties::default()).await?;

Refreshing a credential

connection.update-secret re-authenticates an established connection with a new token, so a long-lived connection can outlive a short-lived token without reconnecting. On success the expiry deadline moves to the new token's exp — it is not cleared.

OutcomeResponse
The connection authenticated with a platform token and the new secret authenticates as the same identityconnection.update-secret-ok
The new secret does not authenticate, or authenticates with an empty ClientIDconnection.close(530), NOT_ALLOWED - Secret update failed
The new secret is valid but resolves to a different identityconnection.close(530), NOT_ALLOWED - New secret was refused by one of the backends — reconnect instead; RabbitMQ 4.3.4 refuses the same case with the same code
The connection holds no token to replace — it authenticated with a client certificate, against the connector-local credential store, or with Authentication.Enable = falseconnection.close(530), NOT_ALLOWED - Secret update failed, whatever the new secret says

Certificate login (SASL EXTERNAL)

EXTERNAL authenticates by client TLS certificate instead of a password. Three things must all hold, and the first is the one that catches people:

  1. Connectors.Amqp.SslCertLogin = true. It defaults to false, and it is deliberately not inferred from the server's TLS settings — enabling mutual TLS for another transport must not silently add a second way to authenticate to AMQP. It requires TlsPort to be set.
  2. Mutual TLS on the server: Security.Cert, Security.Key and Security.Ca all set. Without the CA the listener never asks for a client certificate; without the certificate and key there is no TLS listener at all. Either way EXTERNAL is never offered, the server still starts, and it logs one WARN naming the security mode it found.
  3. The handshake presented a certificate the listener verified. An unverified certificate never grants identity.

The identity comes from Connectors.Amqp.SslCertLoginFrom — RabbitMQ's own ssl_cert_login_from setting name and both of its values, so a line copied from a migrating RabbitMQ configuration means the same thing here:

  • distinguished_name (the default, and RabbitMQ's): the certificate's full RFC-2253 distinguished name, e.g. CN=svc,OU=eng,O=Acme.
  • common_name: the subject's common name alone. The subject must carry exactly one non-empty common name; a certificate with none, an empty one, or several is refused rather than logged in under one of them. If your certificates can carry several, use distinguished_name, which keeps them all.
  • subject_alternative_name is not supported — any other value is a configuration error at startup, not a silent fallback.

When a credential store is configured, EXTERNAL consults it: a certificate whose derived identity names no configured user is refused. The AMQP 1.0 connector on the same ports differs on both mechanisms — it offers ANONYMOUS whenever authentication is disabled, and EXTERNAL on every verified certificate with no opt-in.

Auth disabled — the dev default

When authentication is disabled and no users are configured (the default on a dev broker), any credentials are accepted — guest:guest completes a full round-trip. The ClientID is then amqp-{uuid8}, minted by the server per connection; the client's connection_name property is a display and log field only and never reaches the identity. Because every connection gets a fresh random subject, no authorization policy can single out one client in this mode — to give AMQP clients different permissions, configure Connectors.Amqp.Credentials or enable authentication.

The user-id property

When authentication is enabled and the user-id property is set on basic.publish, it MUST equal the connection's identity; a mismatch is rejected with 406 precondition-failed. When auth is disabled, user-id passes through unvalidated.

Authorization

Two layers can apply, in this order:

  1. Connector-level permissions — RabbitMQ's own per-vhost configure/write/read triple, enforced only for a connector-local user that carries Permissions entries. See Users and permissions.
  2. Platform (Casbin) authorization, when enabled, checked per channel against the connection's identity:
OperationPlatform permissionOn denial
PublishWrite on the exchange, checked before routing and before the exchange's existence is resolved403 access-refused closes the channel, naming the exchange (audited amqp.publish.denied). The publish is never acknowledged and never returned as 312.
Publish — each destination queueWrite on the queue's channel, by the shared queue layer every KubeMQ protocol goes throughThe whole publish fails as one outcome: basic.nack in confirm mode, 541 channel close without confirms. Never a partial delivery.
Consume / basic.getRead on the queue403 access-refused.
queue.declare / queue.purge / queue.deleteWrite on the queue403 access-refused (audited amqp.topology.denied).
queue.bind / queue.unbindWrite on the queue and Read on the exchange403 access-refused.
exchange.declare / exchange.deleteWrite on the exchange403 access-refused.
Passive declare (queue or exchange)any single permission — Read or Write403, so a denied caller cannot enumerate names (RabbitMQ 4.3.1+ behavior).

When authorization is disabled, all operations are allowed. Exchanges are an authorization resource of their own, so a policy that grants publish by naming only a queue is missing the exchange grant — add a Write rule naming the exchange and keep the queue-scoped rules; the queue layer still relies on them. This is stricter than RabbitMQ, which authorizes only the exchange.

The subject a policy must name depends on the login path. A platform-token login is its client id, verbatim (^orders-producer$); a connector-local user is the encoded username under amqp-_user- (^amqp-_user-orders__producer$ for orders_producer); a certificate login is the encoded certificate field under amqp-_cert-. A policy written for one path grants nothing under another, and ^amqp-.*$ grants every configured user everything. See Which subject a policy must name.

Example policy

{"Queues":true,"ClientID":"^orders-producer$","Channel":"^amqp\\.default\\..*$","Read":true,"Write":true}

A platform-token client with this policy can declare, publish to, and consume under amqp.default.*. A client without Read on a channel gets 403 on consume; a client without Write on the exchange gets 403 on publish, and one without Write on a destination queue gets its publish nacked (or a 541 close outside confirm mode).

Quick decision guide

You want…Do this
Clone-and-run on a stock dev brokerConnect with any credentials (guest:guest)
Authenticate with a KubeMQ identitySASL PLAIN or AMQPLAIN, JWT in the password slot; username is audit-only
RabbitMQ-style users with passwords and per-vhost permissionsConfigure Connectors.Amqp.Credentials — see Users and permissions
Log in with a client certificateSslCertLogin = true, mutual TLS, and SASL EXTERNAL
Keep a connection alive across token rotationconnection.update-secret with the new token before the old one expires
Encrypt the credential on the wireUse amqps://:5671 — see TLS and mTLS
Diagnose a permission failure403 on consume means no Read; 403 on publish means no Write on the exchange; a nack or 541 on publish means no Write on a destination queue

Was this page helpful?

On this page