mirror of
https://github.com/open-webui/open-webui.git
synced 2026-10-04 02:33:43 +00:00
style+fix(stream): trim comments, address review round
Tone down verbose comments added over the last few review rounds (-79 lines net). Logic unchanged by the tone-down itself. Also addresses three new findings: - Hot-path Redis timeout dropped 500ms → 200ms. Halves worst-case stall under degraded Redis; healthy Redis (<1ms) is unaffected. - onDestroy now calls dropAllResumeFences so pending fence timers can't fire post-unmount against stale component state. - Fence flush uses insertion order instead of partition+sort. Live frames were emitted sequentially by one emitter and preserved by socket.io transport order, so push-order already == seq-order; sorting was mixing seq and seq-less frames incorrectly in the Redis-degraded path.
This commit is contained in:
parent
06f4656f5e
commit
8c377fd774
2 changed files with 38 additions and 117 deletions
|
|
@ -172,47 +172,33 @@ YDOC_MANAGER = YdocManager(
|
|||
)
|
||||
|
||||
|
||||
# Bounded Redis stream log keyed by message_id. Clients that reconnect
|
||||
# mid-stream can replay missed events from here. No-op without Redis.
|
||||
# Bounded Redis stream log keyed by message_id. Clients replay from it
|
||||
# on reconnect. No-op without Redis. Also no-op in REALTIME_CHAT_SAVE
|
||||
# mode — DB is authoritative there and replay would double-apply.
|
||||
RESUME_STREAM_MAXLEN = 2000
|
||||
RESUME_STREAM_TTL_SEC = 3600
|
||||
RESUME_STREAM_DONE_TTL_SEC = 30
|
||||
# Refresh TTL every N appends to keep the hot path at one Redis RTT.
|
||||
RESUME_STREAM_TTL_REFRESH_EVERY = 64
|
||||
# Upper bound on any single Redis call in the streaming hot path. If
|
||||
# Redis goes unhealthy we'd rather degrade resume gracefully than stall
|
||||
# live token delivery behind a slow RTT.
|
||||
RESUME_STREAM_REDIS_TIMEOUT_SEC = 0.5
|
||||
# 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.2
|
||||
|
||||
|
||||
def _stream_key(user_id: str, message_id: str) -> str:
|
||||
# user_id-scoped: only the owning user's session can construct the
|
||||
# key, so resume doesn't need a DB chat/message auth check.
|
||||
# user_id-scoped key — auth is implicit from the session's user.
|
||||
return f'{REDIS_KEY_PREFIX}:stream:{user_id}:{message_id}'
|
||||
|
||||
|
||||
def _stream_seq_key(user_id: str, message_id: str) -> str:
|
||||
# Separate key for the atomic seq counter so INCR is cheap and
|
||||
# doesn't touch the stream. Lives and dies alongside the stream log.
|
||||
return f'{REDIS_KEY_PREFIX}:stream:{user_id}:{message_id}:seq'
|
||||
|
||||
|
||||
async def _stream_seq_allocate(user_id: str, message_id: str):
|
||||
"""Atomically allocate the next seq for this message_id.
|
||||
"""Allocate the next seq via atomic INCR, or None when resume is off.
|
||||
|
||||
Global INCR (not a per-emitter counter) so overlapping emitters for
|
||||
the same message_id can't both emit the same seq and cause client
|
||||
dedupe to silently drop one of them as a duplicate.
|
||||
|
||||
Returns the allocated seq, or None if Redis is unavailable or the
|
||||
call times out. Callers treat None as "emit without seq" — live
|
||||
streaming continues but that frame isn't eligible for resume dedupe.
|
||||
|
||||
Also returns None when ENABLE_REALTIME_CHAT_SAVE is on: in that
|
||||
mode the DB is already the authoritative per-token store, so the
|
||||
frontend loads the up-to-date content on refresh. Logging the
|
||||
resume stream in parallel would cause duplicate content on refresh
|
||||
(resume replays from seq=0 on top of DB-loaded content).
|
||||
Returns None when Redis is unavailable, times out, or
|
||||
ENABLE_REALTIME_CHAT_SAVE is set. Callers emit without seq in that
|
||||
case — live streaming continues, no dedupe/resume for that frame.
|
||||
"""
|
||||
if ENABLE_REALTIME_CHAT_SAVE:
|
||||
return None
|
||||
|
|
@ -223,8 +209,6 @@ 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
|
||||
)
|
||||
# Set TTL on first INCR; EXPIRE is NX-like in effect (idempotent
|
||||
# enough here — occasional reset is fine).
|
||||
if seq == 1:
|
||||
try:
|
||||
await asyncio.wait_for(
|
||||
|
|
@ -243,13 +227,7 @@ async def _stream_seq_allocate(user_id: str, message_id: str):
|
|||
|
||||
|
||||
async def _stream_log_append(user_id: str, message_id: str, envelope: dict, seq: int) -> None:
|
||||
"""Append an envelope to the resume log.
|
||||
|
||||
Uses Redis-generated stream IDs so overlapping emitters can't collide.
|
||||
Pipelined with the periodic EXPIRE. Wrapped in a short timeout so a
|
||||
slow Redis can't block live token delivery for the user — on timeout
|
||||
we drop the log entry and let the emit proceed.
|
||||
"""
|
||||
"""Append envelope to resume log. Timeout drops the entry, not the emit."""
|
||||
if REDIS is None or not user_id or not message_id:
|
||||
return
|
||||
try:
|
||||
|
|
@ -319,13 +297,9 @@ async def _stream_log_read(user_id: str, message_id: str, after_seq: int):
|
|||
except Exception:
|
||||
continue
|
||||
out.append((entry_seq, envelope))
|
||||
# Sort by seq before returning envelopes. Stream append order can
|
||||
# diverge from seq order because INCR + XADD aren't atomic: a
|
||||
# concurrent emitter can win the INCR race for a later seq yet lose
|
||||
# the XADD race and land in the stream ahead of an earlier seq. If
|
||||
# we handed envelopes to the client in append order, the client's
|
||||
# `incomingSeq <= lastSeq` guard would permanently drop the later-
|
||||
# arriving-but-earlier-numbered frame.
|
||||
# Stream append order can diverge from seq order under concurrent
|
||||
# emitters (INCR/XADD aren't atomic); sort so the client gets frames
|
||||
# in seq order.
|
||||
out.sort(key=lambda pair: pair[0])
|
||||
return [envelope for _seq, envelope in out]
|
||||
|
||||
|
|
@ -684,21 +658,11 @@ async def chat_events(sid, data):
|
|||
|
||||
@sio.on('resume-stream')
|
||||
async def resume_stream(sid, data):
|
||||
"""Replay missed log entries in a single batch.
|
||||
"""One `resume-stream:replay` batch emit (payload + completion signal).
|
||||
|
||||
One `resume-stream:replay` emit carries all envelopes with seq >
|
||||
last_seq (possibly zero) and serves as the completion signal. Sent
|
||||
whenever we have a message_id to target — including the Redis-down
|
||||
and no-log cases — so the client's fence clears reliably instead of
|
||||
waiting on its fallback timeout.
|
||||
|
||||
Auth-rejection branches (non-dict payload, missing session, missing
|
||||
message_id) intentionally return without a reply: the client either
|
||||
never raised a fence for this call (because it had no message_id)
|
||||
or isn't a legitimate session, so a silent drop is correct there.
|
||||
|
||||
Stream auth is implicit: the key is scoped by user_id, so an
|
||||
authenticated session can only ever read its own logs. No DB lookup.
|
||||
Sent whenever we have a valid (user, message_id); auth-rejection
|
||||
branches drop silently since the client never raised a fence for
|
||||
those. Key scoping by user_id means no DB auth check is needed.
|
||||
"""
|
||||
if not isinstance(data, dict):
|
||||
return
|
||||
|
|
@ -1023,12 +987,6 @@ async def get_event_emitter(request_info, update_db=True):
|
|||
chat_id = request_info['chat_id']
|
||||
message_id = request_info['message_id']
|
||||
|
||||
# Atomically allocate this frame's seq. Global INCR (not an
|
||||
# in-process counter) so overlapping emitters for the same
|
||||
# message_id can't both produce the same seq. If Redis is
|
||||
# unavailable or slow, seq is None and we emit without it — live
|
||||
# streaming continues, but this frame is excluded from resume
|
||||
# dedupe and replay (graceful degradation).
|
||||
seq = await _stream_seq_allocate(user_id, message_id)
|
||||
|
||||
envelope = {
|
||||
|
|
@ -1036,17 +994,15 @@ async def get_event_emitter(request_info, update_db=True):
|
|||
'message_id': message_id,
|
||||
'data': event_data,
|
||||
}
|
||||
# Log before emit so a reconnecting client can't resume-read past
|
||||
# a frame that hasn't been persisted yet. Client seq guard drops
|
||||
# duplicates from the inverted race.
|
||||
if seq is not None:
|
||||
envelope['seq'] = seq
|
||||
# Log before emit: an inverted order could let a reconnecting
|
||||
# client resume-read before the append lands and permanently
|
||||
# miss the frame. Duplicates are dropped by the client seq
|
||||
# guard.
|
||||
await _stream_log_append(user_id, message_id, envelope, seq)
|
||||
await sio.emit('events', envelope, room=f'user:{user_id}')
|
||||
|
||||
# On done, shorten TTL on both the stream log and the seq counter
|
||||
# so they self-evict together after the grace window.
|
||||
# On done, shorten TTL so log + seq counter self-evict together.
|
||||
inner = event_data.get('data') if isinstance(event_data, dict) else None
|
||||
if isinstance(inner, dict) and inner.get('done') is True:
|
||||
if REDIS is not None and user_id and message_id:
|
||||
|
|
|
|||
|
|
@ -169,15 +169,9 @@
|
|||
|
||||
let taskIds = null;
|
||||
|
||||
// Last-seen WS seq per message. Off-message so it never hits the
|
||||
// persisted `history`.
|
||||
const resumeSeqByMessageId = new Map();
|
||||
// While a resume replay is in flight for a message, live frames for
|
||||
// that message are buffered here and flushed after the server's
|
||||
// single-batch `resume-stream:replay` arrives. Prevents a racing
|
||||
// live frame from advancing seq past unreplayed frames and causing
|
||||
// them to be dropped by the dedupe guard.
|
||||
const resumeQueueByMessageId = new Map();
|
||||
// Resume-protocol state. Off-message so it doesn't leak into `history`.
|
||||
const resumeSeqByMessageId = new Map(); // highest seq applied
|
||||
const resumeQueueByMessageId = new Map(); // live frames buffered during replay
|
||||
|
||||
// Chat Input
|
||||
let prompt = '';
|
||||
|
|
@ -457,18 +451,14 @@
|
|||
let message = history.messages[event.message_id];
|
||||
|
||||
if (message) {
|
||||
// If a replay is in flight for this message, buffer any
|
||||
// live (non-replayed) frame until the server's complete
|
||||
// ack arrives. Otherwise a live frame with seq > missed
|
||||
// replay frames would advance lastSeq and cause the
|
||||
// replay frames to be dropped by the dedupe guard.
|
||||
// Buffer live frames during replay; _replayed frames skip the fence.
|
||||
const queue = resumeQueueByMessageId.get(event.message_id);
|
||||
if (queue && !event?._replayed) {
|
||||
queue.push(event);
|
||||
return;
|
||||
}
|
||||
|
||||
// Track highest seq and drop replays we've already applied.
|
||||
// Dedupe.
|
||||
const incomingSeq = typeof event?.seq === 'number' ? event.seq : null;
|
||||
if (incomingSeq !== null) {
|
||||
const lastSeq = resumeSeqByMessageId.get(event.message_id) ?? 0;
|
||||
|
|
@ -641,18 +631,11 @@
|
|||
}
|
||||
};
|
||||
|
||||
// Ask the server to replay any frames we missed for a message.
|
||||
// Idempotent — the seq guard in chatEventHandler drops duplicates.
|
||||
// Safety net: if a resume-stream:replay ack never arrives (server
|
||||
// crash mid-handler, network loss between emit and ack, handler not
|
||||
// registered yet, etc.), clear the fence and flush anyway after this
|
||||
// timeout so live UI updates don't freeze indefinitely.
|
||||
// Fallback timer in case the replay ack never arrives.
|
||||
const RESUME_FENCE_TIMEOUT_MS = 10000;
|
||||
const resumeFenceTimerByMessageId = new Map();
|
||||
|
||||
// Drop the fence + timer without applying any buffered events.
|
||||
// Used on disconnect / chat-change / init transitions where the
|
||||
// buffered events are about state that is about to become stale.
|
||||
// Drop fence + timer without flushing (lifecycle transitions).
|
||||
const dropResumeFence = (messageId) => {
|
||||
const timer = resumeFenceTimerByMessageId.get(messageId);
|
||||
if (timer) {
|
||||
|
|
@ -662,8 +645,10 @@
|
|||
resumeQueueByMessageId.delete(messageId);
|
||||
};
|
||||
|
||||
// Drop fence AND flush buffered events through chatEventHandler.
|
||||
// Used on the happy path (replay ack arrived) and on timeout fallback.
|
||||
// Drop fence AND flush buffered events (happy path + timeout).
|
||||
// Flush in insertion order: live frames were emitted sequentially by
|
||||
// one emitter and preserved by socket.io, so push-order == seq-order
|
||||
// and no sort is needed (mixing with seq-less frames would reorder).
|
||||
const clearResumeFence = async (messageId) => {
|
||||
const timer = resumeFenceTimerByMessageId.get(messageId);
|
||||
if (timer) {
|
||||
|
|
@ -673,19 +658,7 @@
|
|||
const queue = resumeQueueByMessageId.get(messageId);
|
||||
if (!queue) return;
|
||||
resumeQueueByMessageId.delete(messageId);
|
||||
// Partition into seq-bearing and seq-less events so the sort
|
||||
// comparator is a strict ordering on its domain (avoids the
|
||||
// engine-dependent behavior of returning 0 for mixed pairs).
|
||||
// Apply seq-ordered events first, then seq-less (graceful
|
||||
// degradation frames) in insertion order.
|
||||
const withSeq = [];
|
||||
const withoutSeq = [];
|
||||
for (const event of queue) {
|
||||
if (typeof event?.seq === 'number') withSeq.push(event);
|
||||
else withoutSeq.push(event);
|
||||
}
|
||||
withSeq.sort((a, b) => a.seq - b.seq);
|
||||
for (const event of [...withSeq, ...withoutSeq]) {
|
||||
try {
|
||||
await chatEventHandler(event);
|
||||
} catch (e) {
|
||||
|
|
@ -697,9 +670,7 @@
|
|||
const requestResumeForMessage = (message) => {
|
||||
if (!message || !message.id || message.done) return;
|
||||
if (!$socket || !$socket.connected) return;
|
||||
// Raise the fence BEFORE emitting so any live frame arriving
|
||||
// while the server reads XRANGE is buffered, not raced past the
|
||||
// replay frames.
|
||||
// Fence up before emit so racing live frames get buffered.
|
||||
if (!resumeQueueByMessageId.has(message.id)) {
|
||||
resumeQueueByMessageId.set(message.id, []);
|
||||
}
|
||||
|
|
@ -723,16 +694,12 @@
|
|||
if (!messageId) return;
|
||||
const envelopes = Array.isArray(payload?.envelopes) ? payload.envelopes : [];
|
||||
try {
|
||||
// Replay frames bypass the fence (they're what the fence is
|
||||
// waiting for) but still pass through chatEventHandler's
|
||||
// dedupe guard so anything already applied is a no-op.
|
||||
for (const envelope of envelopes) {
|
||||
envelope._replayed = true;
|
||||
await chatEventHandler(envelope);
|
||||
}
|
||||
} finally {
|
||||
// Always clear the fence, even if replay application threw
|
||||
// partway through, so live updates aren't frozen forever.
|
||||
// Finally-clear so a mid-replay throw can't freeze live updates.
|
||||
await clearResumeFence(messageId);
|
||||
}
|
||||
};
|
||||
|
|
@ -848,10 +815,7 @@
|
|||
// Resume any in-flight streams on reconnect.
|
||||
$socket?.on('connect', requestResumeForAllInProgress);
|
||||
$socket?.on('resume-stream:replay', onResumeStreamReplay);
|
||||
// Drop any stale fences on disconnect so the reconnect path
|
||||
// starts from a clean slate instead of inheriting a timer that
|
||||
// could fire after the new resume request has already raised a
|
||||
// fresh fence.
|
||||
// Clean slate on reconnect so stale timers don't fire across runs.
|
||||
$socket?.on('disconnect', dropAllResumeFences);
|
||||
|
||||
$audioQueue?.destroy();
|
||||
|
|
@ -973,6 +937,7 @@
|
|||
$socket?.off('connect', requestResumeForAllInProgress);
|
||||
$socket?.off('resume-stream:replay', onResumeStreamReplay);
|
||||
$socket?.off('disconnect', dropAllResumeFences);
|
||||
dropAllResumeFences();
|
||||
audioQueueInstance?.destroy();
|
||||
audioQueue.set(null);
|
||||
} catch (e) {
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue