|
1 | 1 | import functools |
2 | 2 | import pika |
| 3 | +from pika import exceptions as pika_exceptions |
3 | 4 | import time |
4 | 5 |
|
5 | 6 | from varys.utils import varys_message |
@@ -71,12 +72,44 @@ def run(self): |
71 | 72 |
|
72 | 73 | self._connection = pika.BlockingConnection(self._parameters) |
73 | 74 | self._channel = self._connection.channel() |
74 | | - self._channel.exchange_declare( |
75 | | - exchange=self._exchange, |
76 | | - exchange_type=self._exchange_type, |
77 | | - durable=True, |
78 | | - ) |
79 | | - self._channel.queue_declare(queue=self._queue, durable=True) |
| 75 | + try: |
| 76 | + self._channel.exchange_declare( |
| 77 | + exchange=self._exchange, |
| 78 | + exchange_type=self._exchange_type, |
| 79 | + durable=True, |
| 80 | + passive=True, |
| 81 | + ) |
| 82 | + except pika_exceptions.ChannelClosed as e: |
| 83 | + if e.reply_code != 404: |
| 84 | + raise |
| 85 | + |
| 86 | + self._log.info( |
| 87 | + f"Exchange {self._exchange} does not exist, creating it..." |
| 88 | + ) |
| 89 | + self._channel = self._connection.channel() |
| 90 | + self._channel.exchange_declare( |
| 91 | + exchange=self._exchange, |
| 92 | + exchange_type=self._exchange_type, |
| 93 | + durable=True, |
| 94 | + ) |
| 95 | + try: |
| 96 | + self._channel.queue_declare(queue=self._queue, durable=True) |
| 97 | + except pika_exceptions.ChannelClosed as e: |
| 98 | + if e.reply_code != 404: |
| 99 | + raise |
| 100 | + |
| 101 | + self._log.info( |
| 102 | + f"Queue {self._queue} does not exist, creating it..." |
| 103 | + ) |
| 104 | + self._channel = self._connection.channel() |
| 105 | + self._channel.exchange_declare( |
| 106 | + exchange=self._exchange, |
| 107 | + exchange_type=self._exchange_type, |
| 108 | + durable=True, |
| 109 | + passive=True, |
| 110 | + ) |
| 111 | + self._channel.queue_declare(queue=self._queue, durable=True) |
| 112 | + |
80 | 113 | self._channel.queue_bind( |
81 | 114 | queue=self._queue, |
82 | 115 | exchange=self._exchange, |
|
0 commit comments