| 6 | |
| 7 | |
| 8 | class RMWorker: |
| 9 | def __init__(self, consumer: AsyncRabbitMQConsumer, queue: str, exchange: str, log: Log) -> None: |
| 10 | self.consumer = consumer |
| 11 | self.queue = queue |
| 12 | self.exchange = exchange |
| 13 | self.logger = log |
| 14 | |
| 15 | async def initialize(self, loop: asyncio.AbstractEventLoop, handlers: list[MessageHandler]) -> None: |
| 16 | await self.consumer.initialize(loop=loop) |
| 17 | for handler in handlers: |
| 18 | self.consumer.register_handler(handler) |
| 19 | self.logger.info("🚀 RabbitMQ worker initialized successfully") |
| 20 | |
| 21 | async def start(self) -> None: |
| 22 | try: |
| 23 | await self.consumer.consume( |
| 24 | queue_name=self.queue, exchange_name=self.exchange, durable=True, exclusive=False, auto_delete=False |
| 25 | ) |
| 26 | except Exception as e: |
| 27 | self.logger.error(f"🛑 Failed to start RabbitMQ worker: {e}") |
| 28 | await self.stop() |
| 29 | raise |
| 30 | |
| 31 | async def stop(self) -> None: |
| 32 | await self.consumer.close() |
| 33 | self.logger.info("🚦 RabbitMQ worker stopped successfully") |