From 987711a124ff86c90763f8ee3c4ce3f27c7e0e7b Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 15 Apr 2026 08:42:03 +0000 Subject: [PATCH] fix(stream): breaker only counts append success, surface replay truncation Warning: breaker was being reset by _stream_seq_allocate's INCR success before each append attempt. Under "INCR always succeeds, XADD always times out" the failure count never accumulated to the threshold, so the breaker never tripped for the exact failure mode it was supposed to short-circuit. Removed the success call from seq allocation; only end-to-end append success now counts as "Redis is healthy," which means a sustained append-only failure pattern will now correctly trip the breaker after 3 frames. Warning: replay byte-cap truncation was silent. Added \`truncated\` flag to the resume-stream:replay payload so the client can log / react to the case. Frontend currently logs a console warning; a future improvement could trigger an automatic chat reload to recover non-DB-backed side-channel events (sources, embeds) that the final done checkpoint doesn't include. Deferred: prune resumeSeqByMessageId on per-message terminal. Fourth round of flip-flop on this one; the current design (no per-message prune) was chosen to avoid continuation-reuses-message_id replay duplication. Memory is int-per-message bounded by chat size. --- backend/open_webui/socket/main.py | 8 +++++++- src/lib/components/chat/Chat.svelte | 7 +++++++ 2 files changed, 14 insertions(+), 1 deletion(-) 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) {