(
self,
queue_name: str,
exchange_name: str,
durable: bool = True,
exclusive: bool = False,
auto_delete: bool = False,
)
| 138 | await message.reject(requeue=False) |
| 139 | |
| 140 | async def consume( |
| 141 | self, |
| 142 | queue_name: str, |
| 143 | exchange_name: str, |
| 144 | durable: bool = True, |
| 145 | exclusive: bool = False, |
| 146 | auto_delete: bool = False, |
| 147 | ) -> None: |
| 148 | if not self._is_initialized or not self._conn: |
| 149 | raise RuntimeError("🛑 RabbitMQ consumer is not initialized") |
| 150 | try: |
| 151 | async with self._conn: |
| 152 | channel = await self._conn.channel() # type: ignore |
| 153 | await channel.set_qos(prefetch_count=1) |
| 154 | |
| 155 | exchange = await channel.declare_exchange( |
| 156 | name=exchange_name, |
| 157 | type=ExchangeType.DIRECT, |
| 158 | durable=durable, |
| 159 | auto_delete=auto_delete, |
| 160 | internal=False, |
| 161 | ) |
| 162 | |
| 163 | queue = await channel.declare_queue( |
| 164 | queue_name, |
| 165 | durable=durable, |
| 166 | exclusive=exclusive, |
| 167 | auto_delete=auto_delete, |
| 168 | ) |
| 169 | |
| 170 | await queue.bind(exchange, queue_name) |
| 171 | |
| 172 | await queue.consume(lambda message: self._process_message(message, queue_name), no_ack=False) |
| 173 | |
| 174 | self._running = True |
| 175 | self.logger.info(f"🚀 RabbitMQ consumer started consuming messages from queue: {queue_name}") |
| 176 | await self._shutdown_event.wait() |
| 177 | except Exception as e: |
| 178 | self.logger.error( |
| 179 | f"🛑 Failed to consume messages from queue {queue_name}: {e}", |
| 180 | error=traceback.extract_tb(e.__traceback__)[-1], |
| 181 | ) |
| 182 | raise |
no test coverage detected