mirror of
https://github.com/open-webui/open-webui.git
synced 2026-10-06 02:48:04 +00:00
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.
This commit is contained in:
parent
811df250ab
commit
987711a124
2 changed files with 14 additions and 1 deletions
|
|
@ -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,
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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) {
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue