From fd1cf4221ec1b6ac13ee163add4c23341db8c9e3 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 14 Apr 2026 22:49:12 +0000 Subject: [PATCH] fix(stream): env-tunable timeouts, request_id correlation, doc concurrent-emitter limit MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Timeouts split and env-configurable. RESUME_STREAM_REDIS_TIMEOUT_SEC (hot path, default 100ms) stays tight to protect live streaming. RESUME_STREAM_READ_TIMEOUT_SEC (replay XRANGE, default 1s) is looser since a resume is user-blocking and silently timing out is worse than a brief extra wait. Cross-region / non-colocated Redis setups can override either. - resume-stream now carries a request_id; the server echoes it in resume-stream:replay and the client ignores replies whose id does not match the active request for that message. Prevents reconnect churn where a stale reply could clear a fence belonging to a newer in-flight request. - Documented the concurrent-emitter live ordering caveat in get_event_emitter. Single-emitter flows are safe (streaming loop awaits sequentially). Replay is safe (sorted by seq on read). Overlapping emitters for the same message_id can produce live frames out of seq order — accepting as a known limitation; proper fix needs either a distributed per-message lock or a client reorder buffer, neither justified by the rarity of the scenario. --- backend/open_webui/socket/main.py | 36 +++++++++++++++++++++++++---- src/lib/components/chat/Chat.svelte | 17 ++++++++++++++ 2 files changed, 48 insertions(+), 5 deletions(-) diff --git a/backend/open_webui/socket/main.py b/backend/open_webui/socket/main.py index 28d0c7558d..1703ec95e8 100644 --- a/backend/open_webui/socket/main.py +++ b/backend/open_webui/socket/main.py @@ -1,5 +1,6 @@ import asyncio import json +import os import random import socketio @@ -179,9 +180,16 @@ RESUME_STREAM_MAXLEN = 2000 RESUME_STREAM_TTL_SEC = 3600 RESUME_STREAM_DONE_TTL_SEC = 30 RESUME_STREAM_TTL_REFRESH_EVERY = 64 -# Tight upper bound on any Redis call in the streaming hot path so a -# slow Redis degrades resume but can't stall live token delivery. -RESUME_STREAM_REDIS_TIMEOUT_SEC = 0.1 +# Hot-path timeout: tight so a slow Redis can't stall live tokens. +# Replay read timeout: looser since a resume is user-blocking anyway +# and silent timeout here is worse than a brief extra wait. Both +# configurable for infra where Redis isn't colocated. +RESUME_STREAM_REDIS_TIMEOUT_SEC = float( + os.environ.get('RESUME_STREAM_REDIS_TIMEOUT_SEC', '0.1') +) +RESUME_STREAM_READ_TIMEOUT_SEC = float( + os.environ.get('RESUME_STREAM_READ_TIMEOUT_SEC', '1.0') +) def _stream_key(user_id: str, message_id: str) -> str: @@ -268,7 +276,7 @@ async def _stream_log_read(user_id: str, message_id: str, after_seq: int): try: entries = await asyncio.wait_for( REDIS.xrange(_stream_key(user_id, message_id), min='-', max='+'), - timeout=RESUME_STREAM_REDIS_TIMEOUT_SEC, + timeout=RESUME_STREAM_READ_TIMEOUT_SEC, ) except asyncio.TimeoutError: log.warning(f'stream resume log read timed out for {message_id}') @@ -678,6 +686,7 @@ async def resume_stream(sid, data): user_id = user.get('id') message_id = data.get('message_id') + request_id = data.get('request_id') try: last_seq = int(data.get('last_seq') or 0) except (TypeError, ValueError): @@ -690,9 +699,15 @@ async def resume_stream(sid, data): if REDIS is not None: envelopes = await _stream_log_read(user_id, message_id, last_seq) + # Echo request_id so the client can ignore stale replies that + # correspond to a superseded request (reconnect churn). await sio.emit( 'resume-stream:replay', - {'message_id': message_id, 'envelopes': envelopes}, + { + 'message_id': message_id, + 'request_id': request_id, + 'envelopes': envelopes, + }, to=sid, ) @@ -987,6 +1002,17 @@ async def disconnect(sid): async def get_event_emitter(request_info, update_db=True): + # A single emitter's calls are serial (the streaming loop awaits each + # before the next), so seq ordering is guaranteed within one emitter. + # CONCURRENT emitters for the same (user_id, message_id) — e.g. a + # duplicate request leaking past frontend dedup, or a retry path that + # overlaps with the original — can interleave INCR/XADD/emit across + # tasks and cause live frames to arrive out of seq order. Replay + # reads already sort by seq so resume is safe, but live streaming in + # that edge case can drop a frame via the client dedupe guard. Fixing + # it properly needs either a distributed per-message lock or a + # client-side reorder buffer; neither is worth the complexity for a + # configuration OWUI doesn't normally produce. async def __event_emitter__(event_data): user_id = request_info['user_id'] chat_id = request_info['chat_id'] diff --git a/src/lib/components/chat/Chat.svelte b/src/lib/components/chat/Chat.svelte index 6863074fcd..ceabf2da95 100644 --- a/src/lib/components/chat/Chat.svelte +++ b/src/lib/components/chat/Chat.svelte @@ -634,6 +634,9 @@ // Fallback timer in case the replay ack never arrives. const RESUME_FENCE_TIMEOUT_MS = 10000; const resumeFenceTimerByMessageId = new Map(); + // Active request_id per message — lets us ignore stale replay + // responses from superseded requests (reconnect churn). + const resumeActiveRequestIdByMessageId = new Map(); // Drop fence + timer without flushing (lifecycle transitions). const dropResumeFence = (messageId) => { @@ -643,6 +646,7 @@ resumeFenceTimerByMessageId.delete(messageId); } resumeQueueByMessageId.delete(messageId); + resumeActiveRequestIdByMessageId.delete(messageId); }; // Drop fence AND flush buffered events (happy path + timeout). @@ -683,8 +687,14 @@ clearResumeFence(message.id); }, RESUME_FENCE_TIMEOUT_MS) ); + const requestId = + typeof crypto !== 'undefined' && crypto.randomUUID + ? crypto.randomUUID() + : `${message.id}-${Date.now()}-${Math.random()}`; + resumeActiveRequestIdByMessageId.set(message.id, requestId); $socket.emit('resume-stream', { message_id: message.id, + request_id: requestId, last_seq: resumeSeqByMessageId.get(message.id) ?? 0 }); }; @@ -692,6 +702,13 @@ const onResumeStreamReplay = async (payload) => { const messageId = payload?.message_id; if (!messageId) return; + // Drop stale replies from superseded requests so they can't + // clear a fence that belongs to a newer in-flight request. + const expected = resumeActiveRequestIdByMessageId.get(messageId); + if (expected && payload?.request_id && payload.request_id !== expected) { + return; + } + resumeActiveRequestIdByMessageId.delete(messageId); const envelopes = Array.isArray(payload?.envelopes) ? payload.envelopes : []; try { for (const envelope of envelopes) {