openframe-adapters-queue-rabbitmq
RabbitMQ (AMQP) queue adapter for the OpenFrame Microservice Development Suite.
Part of the openframe-adapters monorepo.
What it provides
| Symbol | Purpose |
|---|---|
RabbitmqSettings |
Pydantic-settings subclass — reads all config from env vars |
RabbitmqProducer[T] |
Generic async message producer — BaseProducer[T] |
RabbitmqConsumer[T] |
Generic async message consumer — BaseConsumer[T] |
RabbitmqPlugin |
BasePort (Identity + Lifecycle) — structured lifecycle via PluginRegistry/ApplicationBootstrap |
Installation
# Via meta-package (recommended)
pip install "openframe-adapters[rabbitmq]"
# Or directly
pip install openframe-adapters-queue-rabbitmq
Requires openframe-core>=3.3 — the ApplicationBootstrap.compose()
wiring pattern shown below did not exist before that release.
Quick start
from openframe.adapters.queue.rabbitmq import RabbitmqSettings, RabbitmqProducer, RabbitmqConsumer
settings = RabbitmqSettings(rabbitmq_url="amqp://guest:guest@localhost:5672/")
# Produce
producer = RabbitmqProducer(settings)
await producer.start()
await producer.publish({"event": "item.created", "id": "abc"})
await producer.publish_batch([{"event": "x"}, {"event": "y"}])
await producer.close()
# Consume
consumer = RabbitmqConsumer(settings)
async def handle(event: dict) -> None:
print(f"Received: {event}")
await consumer.subscribe(handle) # runs until consumer.close() called
Publishing routes through RabbitMQ's default (nameless) exchange, using the configured queue name as the routing key — the simplest AMQP delivery pattern, requiring no exchange or binding setup.
Configuration
| Env var | Default | Description |
|---|---|---|
RABBITMQ_URL |
required | amqp://user:password@host:port/vhost |
RABBITMQ_QUEUE |
"openframe" |
Default queue for producer and consumer |
RABBITMQ_DURABLE |
true |
Declare the queue as durable (survives broker restart) |
RABBITMQ_PREFETCH_COUNT |
10 |
Consumer QoS — unacked messages in flight |
RABBITMQ_RECONNECT_INTERVAL |
5.0 |
Seconds between connect_robust() reconnect attempts |
Typed domain objects
from openframe.adapters.queue.rabbitmq import RabbitmqProducer, RabbitmqConsumer, RabbitmqSettings
from dataclasses import dataclass, asdict
@dataclass
class OrderEvent:
order_id: str
event_type: str
class OrderProducer(RabbitmqProducer[OrderEvent]):
def _serialise(self, message: OrderEvent) -> bytes:
import json
return json.dumps(asdict(message)).encode("utf-8")
class OrderConsumer(RabbitmqConsumer[OrderEvent]):
def _deserialise(self, raw: bytes) -> OrderEvent:
import json
return OrderEvent(**json.loads(raw.decode("utf-8")))
Plugin lifecycle (optional)
RabbitmqPlugin is a BasePort — wire it up with ApplicationBootstrap.compose(),
the recommended zero-subclass entry point (requires openframe-core>=3.3):
from openframe.core.runtime import ApplicationBootstrap
from openframe.core.ports import Capability
from openframe.adapters.queue.rabbitmq import RabbitmqPlugin, RabbitmqSettings
rabbitmq = RabbitmqPlugin(RabbitmqSettings())
async with ApplicationBootstrap.compose(rabbitmq) as app:
plugin = app.get(Capability.QUEUE)
producer = plugin.get_producer()
await producer.publish({"event": "item.created"})
consumer = plugin.make_consumer()
await consumer.subscribe(handler)
RabbitmqPlugin doubles as both the producer and consumer port for the QUEUE
capability (get_producer() / make_consumer()), so a single instance is
usually enough. If a service registers a separate producer-only and
consumer-only RabbitmqPlugin (e.g. different queues/settings for each), pass
both to compose(): ApplicationBootstrap.compose(producer_plugin, consumer_plugin). Reach for a subclassed ApplicationBootstrap with
configure() only when you need per-port config=/init_timeout= or
conditional registration order, and use app.registry as an escape hatch
for anything neither tier covers.
Consumer acknowledgement semantics
RabbitMQ has first-class per-message ack/nack built into the AMQP protocol, so there is no manual offset-commit bookkeeping (unlike Kafka):
| Outcome | Behaviour |
|---|---|
| Handler returns | message.ack() called → broker marks message consumed |
| Handler raises | message.nack(requeue=True) called → broker redelivers the message |
consumer.close() |
iteration loop exits, channel/connection closed cleanly |
Resilience (optional — requires openframe-core>=3.4)
openframe-core's openframe.core.resilience ships CircuitBreakerProxy,
which composes around a TracingProxy-wrapped plugin from the outside — no
adapter code changes needed:
from openframe.core.resilience import CircuitBreakerProxy
from openframe.core.tracing import TracingProxy
producer = CircuitBreakerProxy(
TracingProxy(plugin.get_producer(), prefix="rabbitmq"),
failure_threshold=5,
reset_timeout=30.0,
)
Compose the circuit breaker around the traced producer, not the reverse — a short-circuited call should never produce a misleading adapter span for a call that never reached RabbitMQ.
License
MIT — © Furious Meteors Engineering
Metadata
Release files for openframe-adapters-queue-rabbitmq 0.1.0
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| openframe_adapters_queue_rabbitmq-0.1.0.tar.gz | 22.1 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| openframe_adapters_queue_rabbitmq-0.1.0-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 38.1 kB
Release files / openframe_adapters_queue_rabbitmq-0.1.0.tar.gz
| Download URL | openframe_adapters_queue_rabbitmq-0.1.0.tar.gz |
|---|---|
| Size | 22.1 kB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
249aeb99ee7a5daf494468874e18d00e58eae175496ce29f29523a5e1f8d19bc
|
|
BLAKE2b-256 checksum How to use checksums |
eaa159f3c03eb9ab68cafb3b97e53cba4945413bb459787d21305d0dca1b5c26
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/7.0.0 CPython/3.13.14
|
Release files / openframe_adapters_queue_rabbitmq-0.1.0-py3-none-any.whl
| Download URL | openframe_adapters_queue_rabbitmq-0.1.0-py3-none-any.whl |
|---|---|
| Size | 16.0 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
733aa4caf6d05ae922f449c492c11a536e741da22dc28b34f33a3ea91a23445a
|
|
BLAKE2b-256 checksum How to use checksums |
9823340a2240972bae4a3423d2b06b7cc70278a7501e90df5952cfeb988c7278
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/7.0.0 CPython/3.13.14
|