Metadata-Version: 2.5
Name: openframe-adapters-queue-rabbitmq
Version: 0.1.0
Summary: OpenFrame Microservice Suite - RabbitMQ queue adapter.
Project-URL: Homepage, https://github.com/Furious-Meteors/openframe-adapters
Project-URL: Documentation, https://furious-meteors.github.io/openframe-adapters/
Project-URL: Repository, https://github.com/Furious-Meteors/openframe-adapters
Project-URL: Changelog, https://github.com/Furious-Meteors/openframe-adapters/blob/production/.github/CHANGELOG.md
Project-URL: Bug Tracker, https://github.com/Furious-Meteors/openframe-adapters/issues
Author-email: Furious Meteors Engineering <engineering@furiousmeteors.dev>
Maintainer-email: Furious Meteors Engineering <engineering@furiousmeteors.dev>
License: MIT
Keywords: amqp,hexagonal,messaging,microservice,openframe,queue,rabbitmq
Classifier: Development Status :: 4 - Beta
Classifier: Intended Audience :: Developers
Classifier: License :: OSI Approved :: MIT License
Classifier: Operating System :: OS Independent
Classifier: Programming Language :: Python :: 3
Classifier: Programming Language :: Python :: 3.11
Classifier: Programming Language :: Python :: 3.12
Classifier: Programming Language :: Python :: 3.13
Classifier: Topic :: Software Development :: Libraries :: Python Modules
Classifier: Topic :: System :: Distributed Computing
Classifier: Typing :: Typed
Requires-Python: >=3.11
Requires-Dist: aio-pika>=9.4
Requires-Dist: openframe-core<4,>=3.3
Provides-Extra: dev
Requires-Dist: pytest-asyncio>=0.23; extra == 'dev'
Requires-Dist: pytest-mock>=3.14; extra == 'dev'
Requires-Dist: pytest>=8.0; extra == 'dev'
Description-Content-Type: text/markdown

# openframe-adapters-queue-rabbitmq

RabbitMQ (AMQP) queue adapter for the **OpenFrame Microservice Development Suite**.

Part of the [`openframe-adapters`](https://github.com/Furious-Meteors/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

```bash
# 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

```python
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

```python
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`**):

```python
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:

```python
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
