from typing import Callable, Dict, List, Optional import aio_pika from urllib.parse import urlsplit, unquote from config import settings import ssl ExchangeResolver = Callable[[str], str] # exchangeName -> exchangeType # rabbit_base.py import aio_pika def _normalize_exchange_type(val: str) -> aio_pika.ExchangeType: # 1) Try attribute by NAME (DIRECT/FANOUT/TOPIC/HEADERS) if isinstance(val, str): name = val.upper() if hasattr(aio_pika.ExchangeType, name): return getattr(aio_pika.ExchangeType, name) # 2) Try enum by VALUE ("direct"/"fanout"/"topic"/"headers") try: return aio_pika.ExchangeType(val.lower()) except Exception: pass # 3) Default return aio_pika.ExchangeType.TOPIC def _parse_amqp_url(url: str) -> dict: 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)) # Build an SSLContext that DISABLES verification ssl_ctx = None if conn_kwargs.get("ssl"): ssl_ctx = ssl.create_default_context() ssl_ctx.check_hostname = False ssl_ctx.verify_mode = ssl.CERT_NONE # Pass ssl_context explicitly – this is what aio-pika supports self._conn = await aio_pika.connect_robust( **conn_kwargs, ssl_context=ssl_ctx # <- key bit ) 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_str = self._exchange_type_resolver(name) # e.g. "direct" ex_type = _normalize_exchange_type(ex_type_str) 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