Metadata-Version: 2.4
Name: pyrmqsod
Version: 0.1.3
Summary: Shared aio-pika RabbitMQ client: connection, topology, publisher, consumer, router, health.
Requires-Python: >=3.10
Requires-Dist: aio-pika>=9.4.3
Description-Content-Type: text/markdown

# pyrmqsod

Shared aio-pika RabbitMQ client for use in FastAPI and Django: connection, topology, publisher, consumer, type-based router, health check.

## Install

From git (editable):

```bash
pip install pyrmqsod'
```


## Config

Build a `RabbitMQConfig` from your app settings (e.g. pydantic-settings or Django):

```python
from pyrmqsod import RabbitMQConfig
from pyrmqsod.config import ExtraQueueConfig

config = RabbitMQConfig(
    url=settings.RABBITMQ_URL,
    in_queue=settings.RABBITMQ_IN_QUEUE,
    out_queue=settings.RABBITMQ_OUT_QUEUE,
    dlx=settings.RABBITMQ_DLX,       # optional
    dlq=settings.RABBITMQ_DLQ,        # optional
    # Declare any number of extra queues (e.g. upload/download completion queues)
    extra_queues=[
        ExtraQueueConfig(
            name=settings.RABBITMQ_UPLOAD_COMPLETED_QUEUE,
        )
    ]
    if settings.RABBITMQ_UPLOAD_COMPLETED_QUEUE
    else None,
)
```

## Usage

- **Connection**: `connect_robust(config)`, `with_channel(connection, callback)`, `declare_topology(channel, config)`.
- **Publish**: `publish_message(channel, routing_key, payload, message_id)`.
- **Consumer**: `start_consumer(channel, config, outbox_writer, handler)`.
  `outbox_writer` is `async (message_id, routing_key, payload) -> None`; call it when publish fails to store for retry.
- **Router**: `make_router(registry)` returns a `MessageHandler` that dispatches by `payload["type"]`; `UnknownMessageTypeError` for unknown types.
- **Health**: `check_rabbitmq_health(config)` returns `True`/`False`.

## Django example

```python
# settings
RABBITMQ_URL = os.environ.get("RABBITMQ_URL", "")
RABBITMQ_IN_QUEUE = os.environ.get("RABBITMQ_IN_QUEUE", "")
RABBITMQ_OUT_QUEUE = os.environ.get("RABBITMQ_OUT_QUEUE", "")

# rabbitmq.py
from django.conf import settings
from pyrmqsod import (
    RabbitMQConfig,
    connect_robust,
    with_channel,
    declare_topology,
    start_consumer,
    make_router,
)

def get_config():
    return RabbitMQConfig(
        url=settings.RABBITMQ_URL,
        in_queue=settings.RABBITMQ_IN_QUEUE,
        out_queue=settings.RABBITMQ_OUT_QUEUE,
    )

async def run_consumer(handler):
    config = get_config()
    conn = await connect_robust(config)
    async def run(channel):
        await declare_topology(channel, config)
        async def outbox_writer(message_id, routing_key, payload):
            pass  # or your outbox implementation
        await start_consumer(channel, config, outbox_writer, handler)
    await with_channel(conn, run)
```

## Requirements

- Python >= 3.10
- aio-pika >= 9.4.3

## Development & tests

From the project root:

```bash
python -m venv .venv
source .venv/bin/activate  # on Windows: .venv\Scripts\activate

python -m pip install -e ".[dev]"
python -m pytest
```

## Type checking

With dev dependencies installed (`python -m pip install -e ".[dev]"`), you can run `ty` from the project root:

```bash
ty check src/pyrmqsod tests --ignore unresolved-import
```

Or via pre-commit:

```bash
pre-commit run ty --all-files
```
