A Python client for STOMP asynchronous messaging protocol that is:
- asynchronous,
- not abandoned,
- has typed, modern, comprehensible API.
Before you start using stompman, make sure you have it installed. If you optionally want to use
stompman over a websocket, you can install with stompman[ws] instead of stompman:
uv add stompman
poetry add stompmanInitialize a client:
async with stompman.Client(
servers=[
stompman.ConnectionParameters(host="171.0.0.1", port=61616, login="user1", passcode="passcode1"),
stompman.ConnectionParameters(host="172.0.0.1", port=61616, login="user2", passcode="passcode2"),
],
# SSL — can be either `None` (default), `True`, or `ssl.SSLContext'
ssl=None,
# Error frame handler:
on_error_frame=lambda error_frame: print(error_frame.body),
# Optional parameters with sensible defaults:
heartbeat=stompman.Heartbeat(will_send_interval_ms=1000, want_to_receive_interval_ms=1000),
connect_retry_attempts=3,
connect_retry_interval=1,
connect_timeout=2,
connection_confirmation_timeout=2,
disconnect_confirmation_timeout=2,
write_retry_attempts=3,
check_server_alive_interval_factor=3,
no_message_restart_interval=datetime.timedelta(hours=1), # None to disable
keep_alive_on_connection_failure=False,
) as client:
...Initialize a client with a custom connection class, for example, connecting to a stomp producer over websocket:
# uv/poetry add stompman[ws] to get WebScoketConnection support
from stompman.connection_ws import WebSocketConnection
async with stompman.Client(
servers=[
stompman.ConnectionParameters(host="171.0.0.1", port=8080, login="", passcode="", ws_uri_path="/ws/path"),
],
connection_class=WebSocketConnection,
...
) as client:
...To send a message, use the following code:
await client.send(b"hi there!", destination="DLQ", headers={"persistent": "true"})Or, to send messages in a transaction:
async with client.begin() as transaction:
for _ in range(10):
await transaction.send(body=b"hi there!", destination="DLQ", headers={"persistent": "true"})
await asyncio.sleep(0.1)By default client.send() returns once the frame has been written to the socket. That says nothing about whether the
broker accepted it. Pass receipt_timeout to wait for the broker's
STOMP receipt instead:
await client.send(b"hi there!", destination="DLQ", headers={"persistent": "true"}, receipt_timeout=3.0)The default receipt_timeout=None preserves the existing write-only behavior. A timeout must be finite and positive.
It covers writing SEND and waiting for its receipt, after a connection is available. The client generates the
receipt header itself, so passing your own receipt header with receipt_timeout is an error.
A failure raises stompman.SendReceiptError, with reason equal to rejected, timeout, or connection_lost.
Raw broker error frames are available through error.frame, but are excluded from the exception's representation.
Only rejected is a definitive answer. On timeout and connection_lost the outcome is ambiguous: the frame is
fully buffered before the socket is drained, so the broker may well have accepted the message even though the
confirmation was lost. A confirmed send still waits for a connection through connect_retry_attempts, but the write
itself is attempted exactly once: write_retry_attempts does not apply, because a connection can fail after the frame
reached the broker and no connection type can tell the two apart. Nothing is silently replayed on a new connection —
retrying is your decision. Give each application event a stable ID and deduplicate in the broker or downstream
consumer; JMSCorrelationID alone does not deduplicate anything.
A receipt means the broker took ownership of the message. It does not mean any consumer has received or processed it.
Concurrent confirmed sends never consume each other's receipts, and a receipt is bound to the connection it was issued
on, so a late receipt from a previous connection cannot confirm a newer send. An ERROR without receipt-id fails all
pending confirmations on that connection, the same policy as subscriptions. Broker errors still reach the
on_error_frame callback.
Confirmation is not available for transaction.send(): a receipt for a transactional SEND says nothing about whether
the COMMIT succeeded.
Now, let's subscribe to a destination and listen for messages:
async def handle_message_from_dlq(message_frame: stompman.MessageFrame) -> None:
print(message_frame.body)
await client.subscribe("DLQ", handle_message_from_dlq, on_suppressed_exception=print)Entered stompman.Client will block forever waiting for messages if there are any active subscriptions.
Sometimes it's useful to avoid that:
dlq_subscription = await client.subscribe("DLQ", handle_message_from_dlq, on_suppressed_exception=print)
await dlq_subscription.unsubscribe()By default, subscription have ACK mode "client-individual". If handler successfully processes the message, an ACK frame will be sent. If handler raises an exception, a NACK frame will be sent. You can catch (and log) exceptions using on_suppressed_exception parameter:
await client.subscribe(
"DLQ",
handle_message_from_dlq,
on_suppressed_exception=lambda exception, message_frame: print(exception, message_frame),
)You can change the ack mode used by specifying the ack parameter:
# Server will assume that all messages sent to the subscription before the ACK'ed message are received and processed:
await client.subscribe("DLQ", handle_message_from_dlq, ack="client", on_suppressed_exception=print)
# Server will assume that messages are received as soon as it send them to client:
await client.subscribe("DLQ", handle_message_from_dlq, ack="auto", on_suppressed_exception=print)You can pass custom headers to client.subscribe():
await client.subscribe(
"DLQ",
handle_message_from_dlq,
ack="client",
headers={"selector": "location = 'Europe'"},
on_suppressed_exception=print,
)If you want to send ACK and NACK frames yourself, you can use client.subscribe_with_manual_ack():
async def handle_message_from_dlq(message_frame: stompman.AckableMessageFrame) -> None:
print(message_frame.body)
await message_frame.ack()
await client.subscribe_with_manual_ack("DLQ", handle_message_from_dlq, ack="client")Note that this way exceptions won't be suppressed automatically.
Pass receipt_timeout to either subscription method to wait for the broker to
accept the subscription before returning. This uses standard
STOMP receipts,
not broker-specific error messages.
subscription = await client.subscribe_with_manual_ack(
"DLQ",
handle_message_from_dlq,
receipt_timeout=3.0,
on_subscription_error=lambda error: print(error.reason),
)
# It is now safe to publish a request that requires this response subscription.The default receipt_timeout=None preserves the existing write-only behavior.
A timeout must be finite and positive. It covers writing SUBSCRIBE and waiting
for its receipt, after a connection is available. The client generates its own
receipt header in this mode.
An initial failure raises SubscriptionError, with reason equal to
rejected, timeout, connection_lost, or unsubscribed. The optional,
synchronous on_subscription_error callback also reports failures during
automatic resubscription, when there is no caller awaiting subscribe().
The rejected subscription is removed before the callback runs. Callbacks should
not block; their exceptions are logged without terminating the frame reader.
Raw broker error frames are available through error.frame, but are excluded
from the exception's representation.
Confirmed subscriptions are restored after reconnect with fresh receipt IDs.
Unconfirmed or rejected subscriptions are not blindly replayed. Timeouts and
cancellation remove local state and attempt bounded cleanup on the same
connection. An ERROR without receipt-id fails all pending confirmations on
that connection; it does not remove previously confirmed subscriptions.
Neither subscription confirmation nor a publish receipt proves downstream
business processing.
The handler concurrency limit remains in effect while confirmations are pending. The reader temporarily buffers message handlers waiting for capacity so it can reach interleaved receipts/errors, then resumes normal backpressure. Use broker prefetch/consumer-window settings to bound deliveries on the wire.
stompman takes care of cleaning up resources automatically. When you leave the context of async context managers stompman.Client(), or client.begin(), the necessary frames will be sent to the server.
-
If multiple servers were provided, stompman will attempt to connect to each one simultaneously and will use the first that succeeds. If all servers fail to connect, an
stompman.FailedAllConnectAttemptsErrorwill be raised. In normal situation it doesn't need to be handled: tune retry and timeout parameters instompman.Client()to your needs. -
When connection is lost, stompman will attempt to handle it automatically.
stompman.FailedAllConnectAttemptsErrorwill be raised if all connection attempts fail.stompman.FailedAllWriteAttemptsErrorwill be raised if connection succeeds but sending a frame or heartbeat lead to losing connection. -
Set
keep_alive_on_connection_failure=Trueto keep background heartbeat and read recovery running after a retry cycle is exhausted. The default remainsFalse, and errors fromClient.send()still followconnect_retry_attemptsandwrite_retry_attempts. -
Connections that succeed and immediately fail are spaced by
connect_retry_intervalas well, preventing a tight reconnect loop. -
If no messages are received for
no_message_restart_interval(defaults to 1 hour), stompman will force a reconnect. Set toNoneto disable. -
To implement health checks, use
stompman.Client.is_alive()— it will returnTrueif everything is OK andFalseif server is not responding. -
stompmanwill write log warnings when connection is lost, after successful reconnection or invalid state during ack/nack.
- stompman supports Python 3.11 and newer.
- It implements STOMP 1.2 — the latest version of the protocol.
- Heartbeats are required, and sent automatically in background (defaults to 1 second).
Also, I want to pointed out that:
- Protocol parsing is inspired by aiostomp (meaning: consumed by me and refactored from).
- stompman is tested and used with ActiveMQ Artemis and ActiveMQ Classic.
- Caveat: a message sent by a Stomp client is converted into a JMS
TextMessage/BytesMessagebased on thecontent-lengthheader (see the docs here). In order to send aTextMessage,Client.sendneeds to be invoked withadd_content_lengthheader set toFalse
- Caveat: a message sent by a Stomp client is converted into a JMS
- Specification says that headers in CONNECT and CONNECTED frames shouldn't be escaped for backwards compatibility. stompman escapes headers in CONNECT frame (outcoming), but does not unescape headers in CONNECTED (outcoming).
An implementation of STOMP broker for FastStream.
See examples in examples/.