This commit is contained in:
Classic298 2026-09-30 21:07:03 +04:00 • committed by GitHub
commit 23a6cc2482
No known key found for this signature in database
GPG key ID: B5690EEEBB952194

View file

@ -2,6 +2,7 @@
import asyncio
from redis.exceptions import NoPermissionError
from socketio import AsyncRedisManager
@ -11,6 +12,7 @@ class AsyncRedisRoomChannelManager(AsyncRedisManager):
def __init__(self, *args, **kwargs):
super().__init__(*args, **kwargs)
self._local_room_channels = set()
self._use_room_channels: bool = True
# collision-free while namespaces contain no '#' (socket.io default '/'); rooms may contain '#'
def _room_channel(self, namespace, room):
@ -26,8 +28,16 @@ class AsyncRedisRoomChannelManager(AsyncRedisManager):
if room is not None and room not in self.rooms.get(namespace, {}):
self._local_room_channels.discard(self._room_channel(namespace, room))
def _fall_back_to_shared_channel(self) -> None:
self._use_room_channels = False
self._get_logger().error(
'Cannot use redis room channels... falling back to the shared channel '
'(add &%s#* to the redis ACL or set WEBSOCKET_REDIS_ROOM_CHANNELS=False)',
self.channel,
)
async def _publish(self, data):
if data.get('method') == 'emit' and isinstance(data.get('room'), str):
if self._use_room_channels and data.get('method') == 'emit' and isinstance(data.get('room'), str):
channel = self._room_channel(data['namespace'], data['room'])
else:
channel = self.channel
@ -38,6 +48,9 @@ class AsyncRedisRoomChannelManager(AsyncRedisManager):
self._redis_connect()
return await self.redis.publish(channel, self.json.dumps(data))
except error as exc:
if isinstance(exc, NoPermissionError) and channel != self.channel:
self._fall_back_to_shared_channel()
return await self._publish(data)
if retries_left > 0:
self._get_logger().error('Cannot publish to redis... retrying', extra={'redis_exception': str(exc)})
self.connected = False
@ -56,11 +69,15 @@ class AsyncRedisRoomChannelManager(AsyncRedisManager):
if not subscribed:
self._redis_connect()
await self.pubsub.subscribe(self.channel)
await self.pubsub.psubscribe(f'{self.channel}#*')
if self._use_room_channels:
await self.pubsub.psubscribe(f'{self.channel}#*')
retry_sleep = 1
async for message in self.pubsub.listen():
yield message
except error as exc:
if isinstance(exc, NoPermissionError) and self._use_room_channels:
self._fall_back_to_shared_channel()
continue
self._get_logger().error(
f'Cannot receive from redis... retrying in {retry_sleep} secs',
extra={'redis_exception': str(exc)},