GradLLM / rabbit_base.py
johnbridges's picture
update rabbitmq code to make it more robust
66c4f69
# rabbit_base.py
from typing import Callable, Dict, List, Optional
from urllib.parse import urlsplit, unquote
import ssl
import json
import logging
import aio_pika
from config import settings
ExchangeResolver = Callable[[str], str]
logger = logging.getLogger(__name__)
def _normalize_exchange_type(val: str) -> aio_pika.ExchangeType:
if isinstance(val, str):
name = val.upper()
if hasattr(aio_pika.ExchangeType, name):
return getattr(aio_pika.ExchangeType, name)
try:
return aio_pika.ExchangeType(val.lower())
except Exception:
pass
return aio_pika.ExchangeType.TOPIC
def _parse_amqp_url(url: str) -> dict:
parts = urlsplit(url)
return {
"scheme": parts.scheme or "amqp",
"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
)
def is_connected(self) -> bool:
return bool(
self._conn and not self._conn.is_closed and
self._chan and not self._chan.is_closed
)
async def close(self) -> None:
try:
if self._chan and not self._chan.is_closed:
logger.info("Closing AMQP channel")
await self._chan.close()
finally:
self._chan = None
try:
if self._conn and not self._conn.is_closed:
logger.info("Closing AMQP connection")
await self._conn.close()
finally:
self._conn = None
logger.info("AMQP connection closed")
async def connect(self) -> None:
if self._conn and not self._conn.is_closed and self._chan and not self._chan.is_closed:
return
conn_kwargs = _parse_amqp_url(str(settings.AMQP_URL))
safe_target = {
"scheme": conn_kwargs["scheme"],
"host": conn_kwargs["host"],
"port": conn_kwargs["port"],
"virtualhost": conn_kwargs["virtualhost"],
"ssl": conn_kwargs["ssl"],
"login": conn_kwargs["login"],
}
logger.info("AMQP connect -> %s", json.dumps(safe_target))
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
logger.warning("AMQP TLS verification is DISABLED (CERT_NONE)")
try:
self._conn = await aio_pika.connect_robust(
host=conn_kwargs["host"],
port=conn_kwargs["port"],
login=conn_kwargs["login"],
password=conn_kwargs["password"],
virtualhost=conn_kwargs["virtualhost"],
ssl=conn_kwargs["ssl"],
ssl_context=ssl_ctx,
heartbeat=60, # keepalive during long CPU work
timeout=30,
client_properties={"connection_name": "hf_backend_publisher"},
)
logger.info("AMQP connection established")
self._chan = await self._conn.channel()
logger.info("AMQP channel created")
await self._chan.set_qos(prefetch_count=settings.RABBIT_PREFETCH)
logger.info("AMQP QoS set (prefetch=%s)", settings.RABBIT_PREFETCH)
except Exception:
logger.exception("AMQP connection/channel setup failed")
raise
async def ensure_exchange(self, name: str) -> aio_pika.Exchange:
await self.connect()
if name in self._exchanges:
ex = self._exchanges[name]
if ex.channel and not ex.channel.is_closed:
return ex
# drop stale cache and recreate
self._exchanges.pop(name, None)
ex_type_str = self._exchange_type_resolver(name)
ex_type = _normalize_exchange_type(ex_type_str)
try:
ex = await self._chan.declare_exchange(name, ex_type, durable=True)
self._exchanges[name] = ex
logger.info("Exchange declared: name=%s type=%s durable=true", name, ex_type.value)
return ex
except Exception:
logger.exception("Failed declaring exchange: %s (%s)", name, ex_type_str)
raise
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
try:
q = await self._chan.declare_queue(
queue_name,
durable=True,
exclusive=False,
auto_delete=True,
arguments=args,
)
logger.info(
"Queue declared: name=%s durable=true auto_delete=true args=%s",
queue_name, args or {}
)
for rk in routing_keys or [""]:
await q.bind(ex, rk)
logger.info("Queue bound: queue=%s exchange=%s rk='%s'", queue_name, exchange, rk)
return q
except Exception:
logger.exception(
"Failed declare/bind queue: queue=%s exchange=%s rks=%s args=%s",
queue_name, exchange, routing_keys, args or {}
)
raise