From 1a74a9f46c31e7bcf05ad38610a8a5c0fc9aaf1e Mon Sep 17 00:00:00 2001 From: Classic298 <27028174+Classic298@users.noreply.github.com> Date: Tue, 22 Sep 2026 20:48:56 +0200 Subject: [PATCH] perf: per-room channel delivery for the socket.io Redis manager (#28818) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Reworked from the ground up after the feedback that the registry implementation did not land. The Redis room registry and its whole recovery protocol (heartbeats, liveness keys, pruning, distrust windows, cache invalidation) are gone; the change is now ~105 lines with no state kept outside the process. With WEBSOCKET_MANAGER=redis every emit is published on one shared channel and every instance JSON-decodes every message: a 16 instance fleet decodes each streamed token delta 16 times and 15 discard it. py-spy across a loaded fleet (16 instances, ~4000 users) puts ~31% of all active CPU samples in the pubsub listener parse chain, the largest bucket. Room-targeted emits are now published on a per-room channel instead; every instance keeps one static pattern subscription covering all room channels and drops messages for rooms without local members by channel name, paying a set lookup instead of a JSON parse. No state leaves the process, so recovery paths and loss windows are identical to the stock manager; acks and control messages stay on the shared channel and sio.call works across instances unchanged. This is the delivery scheme the official socket.io Redis adapter for Node.js ships by default. Enabled by default; WEBSOCKET_REDIS_ROOM_CHANNELS=false restores shared-channel-only delivery. All instances must run the same mode, so the switch rides the full-stop upgrade this release already requires for its migration; in a mixed fleet, room emits from updated instances would not reach not-yet-updated ones. Verified end to end with two instances on a real Redis: cross-instance token streams delivered with the shared channel completely silent. Ref #28173. - [x] By submitting this pull request, I confirm that I have read and fully agree to the [Contributor License Agreement (CLA)](https://github.com/open-webui/open-webui/blob/main/CONTRIBUTOR_LICENSE_AGREEMENT), and I am providing my contributions under its terms. > [!NOTE] > Deleting the CLA section will lead to immediate closure of your PR and it will not be merged in. --- backend/open_webui/env.py | 6 ++ backend/open_webui/socket/main.py | 5 +- .../open_webui/socket/redis_room_channels.py | 82 +++++++++++++++++++ 3 files changed, 92 insertions(+), 1 deletion(-) create mode 100644 backend/open_webui/socket/redis_room_channels.py diff --git a/backend/open_webui/env.py b/backend/open_webui/env.py index fe5f8efb21..dc160e1592 100644 --- a/backend/open_webui/env.py +++ b/backend/open_webui/env.py @@ -493,6 +493,12 @@ else: WEBSOCKET_REDIS_URL = os.getenv('WEBSOCKET_REDIS_URL', REDIS_URL) WEBSOCKET_REDIS_CLUSTER = os.getenv('WEBSOCKET_REDIS_CLUSTER', str(REDIS_CLUSTER)).lower() == 'true' +# publishes room-targeted emits on per-room redis channels so instances skip +# messages for rooms without local members; must be identical across the fleet +# (toggle with a full restart, not a rolling one), set false for the previous +# shared-channel-only delivery +WEBSOCKET_REDIS_ROOM_CHANNELS = os.getenv('WEBSOCKET_REDIS_ROOM_CHANNELS', 'True').lower() == 'true' + websocket_redis_lock_timeout = os.getenv('WEBSOCKET_REDIS_LOCK_TIMEOUT', '60') try: diff --git a/backend/open_webui/socket/main.py b/backend/open_webui/socket/main.py index 659dc3da49..796c08b9cf 100644 --- a/backend/open_webui/socket/main.py +++ b/backend/open_webui/socket/main.py @@ -24,6 +24,7 @@ from open_webui.env import ( WEBSOCKET_REDIS_CLUSTER, WEBSOCKET_REDIS_LOCK_TIMEOUT, WEBSOCKET_REDIS_OPTIONS, + WEBSOCKET_REDIS_ROOM_CHANNELS, WEBSOCKET_REDIS_URL, WEBSOCKET_SENTINEL_HOSTS, WEBSOCKET_SENTINEL_PORT, @@ -38,6 +39,7 @@ from open_webui.models.chats import Chats from open_webui.models.folders import Folders from open_webui.models.notes import Notes, NoteUpdateForm from open_webui.models.users import UserNameResponse, Users +from open_webui.socket.redis_room_channels import AsyncRedisRoomChannelManager from open_webui.socket.utils import CachedRedisDict, RedisDict, RedisLock, YdocManager from open_webui.tasks import ( REDIS_PUBSUB_MAX_RECONNECT_INTERVAL, @@ -93,7 +95,8 @@ if WEBSOCKET_MANAGER == 'redis': if sentinel_hosts else WEBSOCKET_REDIS_URL ) - redis_manager = socketio.AsyncRedisManager(ws_redis_url, redis_options=WEBSOCKET_REDIS_OPTIONS, json=SOCKETIO_JSON) + manager_class = AsyncRedisRoomChannelManager if WEBSOCKET_REDIS_ROOM_CHANNELS else socketio.AsyncRedisManager + redis_manager = manager_class(ws_redis_url, redis_options=WEBSOCKET_REDIS_OPTIONS, json=SOCKETIO_JSON) sio = socketio.AsyncServer( cors_allowed_origins=SOCKETIO_CORS_ORIGINS, async_mode='asgi', diff --git a/backend/open_webui/socket/redis_room_channels.py b/backend/open_webui/socket/redis_room_channels.py new file mode 100644 index 0000000000..715ccc1ea9 --- /dev/null +++ b/backend/open_webui/socket/redis_room_channels.py @@ -0,0 +1,82 @@ +"""Per-room redis channels let instances skip the decode and packet encode for rooms with no local members.""" + +import asyncio + +from socketio import AsyncRedisManager + + +class AsyncRedisRoomChannelManager(AsyncRedisManager): + name = 'aioredisroomchannel' + + def __init__(self, *args, **kwargs): + super().__init__(*args, **kwargs) + self._local_room_channels = set() + + # collision-free while namespaces contain no '#' (socket.io default '/'); rooms may contain '#' + def _room_channel(self, namespace, room): + return f'{self.channel}#{namespace}#{room}'.encode() + + def basic_enter_room(self, sid, namespace, room, eio_sid=None): + super().basic_enter_room(sid, namespace, room, eio_sid=eio_sid) + if room is not None: + self._local_room_channels.add(self._room_channel(namespace, room)) + + def basic_leave_room(self, sid, namespace, room): + super().basic_leave_room(sid, namespace, room) + if room is not None and room not in self.rooms.get(namespace, {}): + self._local_room_channels.discard(self._room_channel(namespace, room)) + + async def _publish(self, data): + if data.get('method') == 'emit' and isinstance(data.get('room'), str): + channel = self._room_channel(data['namespace'], data['room']) + else: + channel = self.channel + _, error = self._get_redis_module_and_error() + for retries_left in range(1, -1, -1): # 2 attempts + try: + if not self.connected: + self._redis_connect() + return await self.redis.publish(channel, self.json.dumps(data)) + except error as exc: + if retries_left > 0: + self._get_logger().error('Cannot publish to redis... retrying', extra={'redis_exception': str(exc)}) + self.connected = False + else: + self._get_logger().error( + 'Cannot publish to redis... giving up', extra={'redis_exception': str(exc)} + ) + break + + async def _redis_listen_with_retries(self): + _, error = self._get_redis_module_and_error() + retry_sleep = 1 + subscribed = False + while True: + try: + if not subscribed: + self._redis_connect() + await self.pubsub.subscribe(self.channel) + await self.pubsub.psubscribe(f'{self.channel}#*') + retry_sleep = 1 + async for message in self.pubsub.listen(): + yield message + except error as exc: + self._get_logger().error( + f'Cannot receive from redis... retrying in {retry_sleep} secs', + extra={'redis_exception': str(exc)}, + ) + subscribed = False + await asyncio.sleep(retry_sleep) + retry_sleep *= 2 + if retry_sleep > 60: + retry_sleep = 60 + + async def _listen(self): + main_channel = self.channel.encode() + async for message in self._redis_listen_with_retries(): + if 'data' not in message: + continue + if (message['type'] == 'message' and message['channel'] == main_channel) or ( + message['type'] == 'pmessage' and message['channel'] in self._local_room_channels + ): + yield message['data']