mirror of
https://github.com/open-webui/open-webui.git
synced 2026-10-05 02:41:34 +00:00
fix: replies stop reaching other instances when the Redis user is limited by an ACL
Since WEBSOCKET_REDIS_ROOM_CHANNELS became the default, a Redis user whose ACL only allows the socketio and open-webui:* channels is refused the new socketio#* channels. Instances then only log a repeating "Cannot receive from redis" retry and stop hearing each other, so a reply streamed on one instance never reaches a tab connected to another. Now the instance logs an error naming the Redis ACL rule to add (&socketio#*) and falls back to the single socketio channel, as with WEBSOCKET_REDIS_ROOM_CHANNELS=False. Replies reach every tab again, and adding the rule and restarting brings back the per-room channels and the CPU they save.
This commit is contained in:
parent
3819dc2793
commit
813d6efda5
1 changed files with 19 additions and 2 deletions
|
|
@ -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)},
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue