perf: per-room channel delivery for the socket.io Redis manager (#28818)

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.

<!--
🚨 DO NOT DELETE THE TEXT BELOW 🚨
Keep the "Contributor License Agreement" confirmation text intact.
Deleting it will trigger the CLA-Bot to INVALIDATE your PR.

Your PR will NOT be reviewed or merged until you check the box below confirming that you have read and agree to the terms of the CLA.
-->

- [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.
This commit is contained in:
Classic298
2026-09-22 14:48:56 -04:00
committed by GitHub
parent fe56ab24f3
commit 1a74a9f46c
3 changed files with 92 additions and 1 deletions
+6
View File
@@ -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:
+4 -1
View File
@@ -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',
@@ -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']