|
|
|
|
|
from typing import Callable, Dict, List, Optional |
|
|
import aio_pika |
|
|
from urllib.parse import urlsplit, unquote |
|
|
from config import settings |
|
|
|
|
|
ExchangeResolver = Callable[[str], str] |
|
|
|
|
|
|
|
|
def _parse_amqp_url(url: str) -> dict: |
|
|
""" |
|
|
Convert an AMQP URL into kwargs for aio_pika.connect_robust, |
|
|
so we don't pass the raw URL (and leak the password in logs). |
|
|
""" |
|
|
parts = urlsplit(url) |
|
|
return { |
|
|
"host": parts.hostname or "localhost", |
|
|
"port": parts.port or (5671 if parts.scheme == "amqps" else 5672), |
|
|
"login": parts.username or "guest", |
|
|
"password": parts.password or "guest", |
|
|
"virtualhost": unquote(parts.path[1:] or "/"), |
|
|
"ssl": parts.scheme == "amqps", |
|
|
} |
|
|
|
|
|
|
|
|
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) -> None: |
|
|
if self._conn and not self._conn.is_closed: |
|
|
return |
|
|
conn_kwargs = _parse_amqp_url(str(settings.AMQP_URL)) |
|
|
self._conn = await aio_pika.connect_robust(**conn_kwargs) |
|
|
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: Dict[str, int] = {} |
|
|
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 |
|
|
|