| import asyncio, json, uuid |
| import aio_pika |
| from typing import Callable, Dict, List, Optional |
| from .config import settings |
|
|
| ExchangeResolver = Callable[[str], str] |
|
|
| class RabbitBase: |
| def __init__(self, exchange_type_resolver: Optional[ExchangeResolver] = None): |
| self._conn: Optional[aio_pika.RobustConnection] = None |
| self._chan: Optional[aio_pika.RobustChannel] = None |
| self._exchanges: Dict[str, aio_pika.Exchange] = {} |
| self._exchange_type_resolver = exchange_type_resolver or (lambda _: settings.RABBIT_EXCHANGE_TYPE) |
|
|
| async def connect(self): |
| if self._conn and not self._conn.is_closed: |
| return |
| self._conn = await aio_pika.connect_robust(str(settings.AMQP_URL)) |
| self._chan = await self._conn.channel() |
| await self._chan.set_qos(prefetch_count=settings.RABBIT_PREFETCH) |
|
|
| async def ensure_exchange(self, name: str) -> aio_pika.Exchange: |
| await self.connect() |
| if name in self._exchanges: |
| return self._exchanges[name] |
| ex_type = self._exchange_type_resolver(name) |
| ex = await self._chan.declare_exchange(name, getattr(aio_pika.ExchangeType, ex_type), durable=True) |
| self._exchanges[name] = ex |
| return ex |
|
|
| async def declare_queue_bind(self, exchange: str, queue_name: str, routing_keys: List[str], ttl_ms: Optional[int]): |
| await self.connect() |
| ex = await self.ensure_exchange(exchange) |
| args = {} |
| if ttl_ms: |
| args["x-message-ttl"] = ttl_ms |
| q = await self._chan.declare_queue(queue_name, durable=True, exclusive=False, auto_delete=True, arguments=args) |
| for rk in routing_keys or [""]: |
| await q.bind(ex, rk) |
| return q |
|
|