GradLLM / rabbit_base.py
johnbridges's picture
fixed paths
ba71442
raw
history blame
1.76 kB
import asyncio, json, uuid
import aio_pika
from typing import Callable, Dict, List, Optional
from config import settings
ExchangeResolver = Callable[[str], str] # exchangeName -> exchangeType
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