Skip to content

Commit 80bb5c9

Browse files
Catch all exceptionsin redis and rabbitmq managers (Fixes #1581)
1 parent 3f1d509 commit 80bb5c9

4 files changed

Lines changed: 54 additions & 57 deletions

File tree

src/socketio/async_aiopika_manager.py

Lines changed: 16 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -97,18 +97,20 @@ async def _publish(self, data):
9797
), routing_key='*',
9898
)
9999
break
100-
except aio_pika.AMQPException:
100+
except aio_pika.exceptions.ChannelInvalidStateError:
101+
# aio_pika raises this exception when the task is cancelled
102+
raise asyncio.CancelledError()
103+
except Exception as exc:
101104
if retry:
102-
self._get_logger().error('Cannot publish to rabbitmq... '
103-
'retrying')
105+
self._get_logger().error(
106+
'Cannot publish to rabbitmq... retrying',
107+
extra={"rabbitmq_exception": str(exc)})
104108
retry = False
105109
else:
106110
self._get_logger().error(
107-
'Cannot publish to rabbitmq... giving up')
111+
'Cannot publish to rabbitmq... giving up',
112+
extra={"rabbitmq_exception": str(exc)})
108113
break
109-
except aio_pika.exceptions.ChannelInvalidStateError:
110-
# aio_pika raises this exception when the task is cancelled
111-
raise asyncio.CancelledError()
112114

113115
async def _listen(self):
114116
retry_sleep = 1
@@ -125,12 +127,13 @@ async def _listen(self):
125127
async with message.process():
126128
yield message.body
127129
retry_sleep = 1
128-
except aio_pika.AMQPException:
129-
self._get_logger().error(
130-
'Cannot receive from rabbitmq... '
131-
'retrying in {} secs'.format(retry_sleep))
132-
await asyncio.sleep(retry_sleep)
133-
retry_sleep = min(retry_sleep * 2, 60)
134130
except aio_pika.exceptions.ChannelInvalidStateError:
135131
# aio_pika raises this exception when the task is cancelled
136132
raise asyncio.CancelledError()
133+
except Exception as exc:
134+
self._get_logger().error(
135+
'Cannot receive from rabbotmq... retrying in '
136+
f'{retry_sleep} secs',
137+
extra={"rabbitmq_exception": str(exc)})
138+
await asyncio.sleep(retry_sleep)
139+
retry_sleep = min(retry_sleep * 2, 60)

src/socketio/async_redis_manager.py

Lines changed: 14 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -75,21 +75,21 @@ def __init__(self, url='redis://localhost:6379/0', channel='socketio',
7575
self.redis = None
7676
self.pubsub = None
7777

78-
def _get_redis_module_and_error(self):
78+
def _get_redis_module(self):
7979
parsed_url = urlparse(self.redis_url)
8080
scheme = parsed_url.scheme.split('+', 1)[0].lower()
8181
if scheme in ['redis', 'rediss']:
8282
if aioredis is None or RedisError is None:
8383
raise RuntimeError('Redis package is not installed '
8484
'(Run "pip install redis" '
8585
'in your virtualenv).')
86-
return aioredis, RedisError
86+
return aioredis
8787
if scheme in ['valkey', 'valkeys']:
8888
if aiovalkey is None or ValkeyError is None:
8989
raise RuntimeError('Valkey package is not installed '
9090
'(Run "pip install valkey" '
9191
'in your virtualenv).')
92-
return aiovalkey, ValkeyError
92+
return aiovalkey
9393
if scheme == 'unix':
9494
if aioredis is None or RedisError is None:
9595
if aiovalkey is None or ValkeyError is None:
@@ -98,14 +98,14 @@ def _get_redis_module_and_error(self):
9898
'or "pip install valkey" '
9999
'in your virtualenv).')
100100
else:
101-
return aiovalkey, ValkeyError
101+
return aiovalkey
102102
else:
103-
return aioredis, RedisError
103+
return aioredis
104104
error_msg = f'Unsupported Redis URL scheme: {scheme}'
105105
raise ValueError(error_msg)
106106

107107
def _redis_connect(self):
108-
module, _ = self._get_redis_module_and_error()
108+
module = self._get_redis_module()
109109
parsed_url = urlparse(self.redis_url)
110110
if parsed_url.scheme in {"redis+sentinel", "valkey+sentinel"}:
111111
sentinels, service_name, connection_kwargs = \
@@ -121,30 +121,25 @@ def _redis_connect(self):
121121
self.connected = True
122122

123123
async def _publish(self, data): # pragma: no cover
124-
_, error = self._get_redis_module_and_error()
125124
for retries_left in range(1, -1, -1): # 2 attempts
126125
try:
127126
if not self.connected:
128127
self._redis_connect()
129128
return await self.redis.publish(
130129
self.channel, self.json.dumps(data))
131-
except error as exc:
130+
except Exception as exc:
132131
if retries_left > 0:
133132
self._get_logger().error(
134-
'Cannot publish to redis... '
135-
'retrying',
133+
'Cannot publish to redis... retrying',
136134
extra={"redis_exception": str(exc)})
137135
self.connected = False
138136
else:
139137
self._get_logger().error(
140-
'Cannot publish to redis... '
141-
'giving up',
138+
'Cannot publish to redis... giving up',
142139
extra={"redis_exception": str(exc)})
143-
144140
break
145141

146142
async def _redis_listen_with_retries(self): # pragma: no cover
147-
_, error = self._get_redis_module_and_error()
148143
retry_sleep = 1
149144
subscribed = False
150145
while True:
@@ -155,11 +150,11 @@ async def _redis_listen_with_retries(self): # pragma: no cover
155150
retry_sleep = 1
156151
async for message in self.pubsub.listen():
157152
yield message
158-
except error as exc:
159-
self._get_logger().error('Cannot receive from redis... '
160-
'retrying in '
161-
f'{retry_sleep} secs',
162-
extra={"redis_exception": str(exc)})
153+
except Exception as exc:
154+
self._get_logger().error(
155+
'Cannot receive from redis... retrying in '
156+
f'{retry_sleep} secs',
157+
extra={"redis_exception": str(exc)})
163158
subscribed = False
164159
await asyncio.sleep(retry_sleep)
165160
retry_sleep *= 2

src/socketio/kombu_manager.py

Lines changed: 10 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -114,14 +114,16 @@ def _publish(self, data):
114114
self.publisher_connection)
115115
producer_publish(self.json.dumps(data))
116116
break
117-
except (OSError, kombu.exceptions.KombuError):
117+
except Exception as exc:
118118
if retry:
119-
self._get_logger().error('Cannot publish to rabbitmq... '
120-
'retrying')
119+
self._get_logger().error(
120+
'Cannot publish to rabbitmq... retrying',
121+
extra={"rabbitmq_exception": str(exc)})
121122
retry = False
122123
else:
123124
self._get_logger().error(
124-
'Cannot publish to rabbitmq... giving up')
125+
'Cannot publish to rabbitmq... giving up',
126+
extra={"rabbitmq_exception": str(exc)})
125127
break
126128

127129
def _listen(self):
@@ -136,9 +138,10 @@ def _listen(self):
136138
message.ack()
137139
yield message.payload
138140
retry_sleep = 1
139-
except (OSError, kombu.exceptions.KombuError):
141+
except Exception as exc:
140142
self._get_logger().error(
141-
'Cannot receive from rabbitmq... '
142-
'retrying in {} secs'.format(retry_sleep))
143+
'Cannot receive from rabbotmq... retrying in '
144+
f'{retry_sleep} secs',
145+
extra={"rabbitmq_exception": str(exc)})
143146
time.sleep(retry_sleep)
144147
retry_sleep = min(retry_sleep * 2, 60)

src/socketio/redis_manager.py

Lines changed: 14 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,3 @@
1-
import logging
21
import time
32
from urllib.parse import urlparse
43

@@ -18,8 +17,6 @@
1817

1918
from .pubsub_manager import PubSubManager
2019

21-
logger = logging.getLogger('socketio')
22-
2320

2421
def parse_redis_sentinel_url(url):
2522
"""Parse a Redis Sentinel URL with the format:
@@ -112,21 +109,21 @@ def initialize(self): # pragma: no cover
112109
'Redis requires a monkey patched socket library to work '
113110
'with ' + self.server.async_mode)
114111

115-
def _get_redis_module_and_error(self):
112+
def _get_redis_module(self):
116113
parsed_url = urlparse(self.redis_url)
117114
scheme = parsed_url.scheme.split('+', 1)[0].lower()
118115
if scheme in ['redis', 'rediss']:
119116
if redis is None or RedisError is None:
120117
raise RuntimeError('Redis package is not installed '
121118
'(Run "pip install redis" '
122119
'in your virtualenv).')
123-
return redis, RedisError
120+
return redis
124121
if scheme in ['valkey', 'valkeys']:
125122
if valkey is None or ValkeyError is None:
126123
raise RuntimeError('Valkey package is not installed '
127124
'(Run "pip install valkey" '
128125
'in your virtualenv).')
129-
return valkey, ValkeyError
126+
return valkey
130127
if scheme == 'unix':
131128
if redis is None or RedisError is None:
132129
if valkey is None or ValkeyError is None:
@@ -135,14 +132,14 @@ def _get_redis_module_and_error(self):
135132
'or "pip install valkey" '
136133
'in your virtualenv).')
137134
else:
138-
return valkey, ValkeyError
135+
return valkey
139136
else:
140-
return redis, RedisError
137+
return redis
141138
error_msg = f'Unsupported Redis URL scheme: {scheme}'
142139
raise ValueError(error_msg)
143140

144141
def _redis_connect(self):
145-
module, _ = self._get_redis_module_and_error()
142+
module = self._get_redis_module()
146143
parsed_url = urlparse(self.redis_url)
147144
if parsed_url.scheme in {"redis+sentinel", "valkey+sentinel"}:
148145
sentinels, service_name, connection_kwargs = \
@@ -158,28 +155,26 @@ def _redis_connect(self):
158155
self.connected = True
159156

160157
def _publish(self, data): # pragma: no cover
161-
_, error = self._get_redis_module_and_error()
162158
for retries_left in range(1, -1, -1): # 2 attempts
163159
try:
164160
if not self.connected:
165161
self._redis_connect()
166162
return self.redis.publish(self.channel, self.json.dumps(data))
167-
except error as exc:
163+
except Exception as exc:
168164
if retries_left > 0:
169-
logger.error(
165+
self._get_logger().error(
170166
'Cannot publish to redis... retrying',
171167
extra={"redis_exception": str(exc)}
172168
)
173169
self.connected = False
174170
else:
175-
logger.error(
171+
self._get_logger().error(
176172
'Cannot publish to redis... giving up',
177173
extra={"redis_exception": str(exc)}
178174
)
179175
break
180176

181177
def _redis_listen_with_retries(self): # pragma: no cover
182-
_, error = self._get_redis_module_and_error()
183178
retry_sleep = 1
184179
subscribed = False
185180
while True:
@@ -189,10 +184,11 @@ def _redis_listen_with_retries(self): # pragma: no cover
189184
self.pubsub.subscribe(self.channel)
190185
retry_sleep = 1
191186
yield from self.pubsub.listen()
192-
except error as exc:
193-
logger.error('Cannot receive from redis... '
194-
f'retrying in {retry_sleep} secs',
195-
extra={"redis_exception": str(exc)})
187+
except Exception as exc:
188+
self._get_logger().error(
189+
'Cannot receive from redis... '
190+
f'retrying in {retry_sleep} secs',
191+
extra={"redis_exception": str(exc)})
196192
subscribed = False
197193
time.sleep(retry_sleep)
198194
retry_sleep *= 2

0 commit comments

Comments
 (0)