Skip to main content

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)

Source distribution for openframe-adapters-queue-rabbitmq 0.1.0
File Size Uploaded
openframe_adapters_queue_rabbitmq-0.1.0.tar.gz 22.1 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for openframe-adapters-queue-rabbitmq 0.1.0
File Interpreter ABI Platform
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

Release history Release notifications | RSS feed

This release

0.1.0 This release

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page