diff --git a/backend/open_webui/socket/main.py b/backend/open_webui/socket/main.py index 95d7ca962d..cc2c57d1ee 100644 --- a/backend/open_webui/socket/main.py +++ b/backend/open_webui/socket/main.py @@ -282,7 +282,11 @@ async def _stream_seq_allocate(user_id: str, message_id: str): seq = await asyncio.wait_for( REDIS.incr(key), timeout=RESUME_STREAM_REDIS_TIMEOUT_SEC ) - _breaker_record_success() + # Don't record breaker success here — only append success counts + # as "Redis is healthy end-to-end." If INCR kept succeeding but + # XADD kept failing, counting INCR successes would reset the + # breaker on every frame and it would never trip for the exact + # failure mode we're trying to short-circuit. if seq == 1: try: await asyncio.wait_for( @@ -806,6 +810,7 @@ async def resume_stream(sid, data): capped.append(env) total_bytes += size capped.reverse() + truncated = len(capped) < len(envelopes) await sio.emit( 'resume-stream:replay', @@ -813,6 +818,7 @@ async def resume_stream(sid, data): 'message_id': message_id, 'request_id': request_id, 'envelopes': capped, + 'truncated': truncated, }, to=sid, ) diff --git a/src/lib/components/chat/Chat.svelte b/src/lib/components/chat/Chat.svelte index 68053ccd4e..5258af6507 100644 --- a/src/lib/components/chat/Chat.svelte +++ b/src/lib/components/chat/Chat.svelte @@ -738,6 +738,13 @@ // clears the fence if we're still the active request; if a newer // request has taken over, let its own lifecycle handle cleanup. const ownRequestId = payload?.request_id; + if (payload?.truncated) { + // Replay omitted older entries that wouldn't fit in the Socket.IO + // buffer. The final done checkpoint reconciles most content; + // non-DB-backed side-channel events (sources/embeds) may be + // missing until a full chat reload. + console.warn('resume-stream replay truncated for', messageId); + } const envelopes = Array.isArray(payload?.envelopes) ? payload.envelopes : []; try { for (const envelope of envelopes) {