mirror of
https://github.com/open-webui/open-webui.git
synced 2026-10-05 02:41:34 +00:00
fix(stream): keep seq state across terminal events, reply on auth fail, broaden terminal TTL
- Stopped deleting resumeSeqByMessageId entries on done / cancel / error. The map now only clears on chat/navigation transitions (loadChat, initNewChat, component unmount). Fixes a continuation- reuses-message_id duplication hole: if a message completes and its id is immediately reused for a continued/extended response, the log's 30s grace window can still serve old envelopes; resetting the client's lastSeq to 0 on done made those old envelopes replay. Keeping the seq alive until the message is truly retired avoids re-applying them. - resume_stream now always replies when the payload has a message_id, even on auth failure (empty envelopes). Silent early-return branches could previously leave the client fence raised for the full 10s fallback timeout during reconnect/auth churn. Silent drops only remain for cases where the client couldn't have raised a fence. - Terminal TTL shortening now also triggers on chat:tasks:cancel and error-finalized completions, not just done:True. Cancelled/errored streams no longer leak log+seq keys for the full 1h TTL.
This commit is contained in:
parent
fd1cf4221e
commit
9af2b478bb
2 changed files with 34 additions and 39 deletions
|
|
@ -671,36 +671,30 @@ async def chat_events(sid, data):
|
|||
|
||||
@sio.on('resume-stream')
|
||||
async def resume_stream(sid, data):
|
||||
"""One `resume-stream:replay` batch emit (payload + completion signal).
|
||||
"""Batch replay emit; also serves as the client's fence-clear signal.
|
||||
|
||||
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.
|
||||
Reply whenever we have a message_id to target so the client fence
|
||||
clears deterministically — even on auth failure. Silent drops only
|
||||
happen when the client couldn't have raised a fence in the first
|
||||
place (malformed payload, no message_id).
|
||||
"""
|
||||
if not isinstance(data, dict):
|
||||
return
|
||||
|
||||
user = SESSION_POOL.get(sid)
|
||||
if not user:
|
||||
return
|
||||
|
||||
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):
|
||||
last_seq = 0
|
||||
|
||||
if not user_id or not message_id:
|
||||
if not message_id:
|
||||
return
|
||||
|
||||
request_id = data.get('request_id')
|
||||
envelopes = []
|
||||
if REDIS is not None:
|
||||
user = SESSION_POOL.get(sid)
|
||||
user_id = user.get('id') if user else None
|
||||
if user_id and REDIS is not None:
|
||||
try:
|
||||
last_seq = int(data.get('last_seq') or 0)
|
||||
except (TypeError, ValueError):
|
||||
last_seq = 0
|
||||
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',
|
||||
{
|
||||
|
|
@ -1033,19 +1027,27 @@ async def get_event_emitter(request_info, update_db=True):
|
|||
await _stream_log_append(user_id, message_id, envelope, seq)
|
||||
await sio.emit('events', envelope, room=f'user:{user_id}')
|
||||
|
||||
# On done, shorten TTL so log + seq counter self-evict together.
|
||||
# Any terminal event shortens TTL so log + seq self-evict together.
|
||||
# Covers normal completion (done:True), explicit cancel, and
|
||||
# errored completions — all end-of-stream flows that the client
|
||||
# will no longer replay against.
|
||||
outer_type = event_data.get('type') if isinstance(event_data, dict) else None
|
||||
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:
|
||||
try:
|
||||
pipe = REDIS.pipeline(transaction=False)
|
||||
pipe.expire(_stream_key(user_id, message_id), RESUME_STREAM_DONE_TTL_SEC)
|
||||
pipe.expire(_stream_seq_key(user_id, message_id), RESUME_STREAM_DONE_TTL_SEC)
|
||||
await asyncio.wait_for(
|
||||
pipe.execute(), timeout=RESUME_STREAM_REDIS_TIMEOUT_SEC
|
||||
)
|
||||
except Exception as e:
|
||||
log.warning(f'stream resume log done-TTL shorten failed for {message_id}: {e}')
|
||||
is_terminal = (
|
||||
(isinstance(inner, dict) and inner.get('done') is True)
|
||||
or outer_type == 'chat:tasks:cancel'
|
||||
or (isinstance(inner, dict) and inner.get('error'))
|
||||
)
|
||||
if is_terminal and REDIS is not None and user_id and message_id:
|
||||
try:
|
||||
pipe = REDIS.pipeline(transaction=False)
|
||||
pipe.expire(_stream_key(user_id, message_id), RESUME_STREAM_DONE_TTL_SEC)
|
||||
pipe.expire(_stream_seq_key(user_id, message_id), RESUME_STREAM_DONE_TTL_SEC)
|
||||
await asyncio.wait_for(
|
||||
pipe.execute(), timeout=RESUME_STREAM_REDIS_TIMEOUT_SEC
|
||||
)
|
||||
except Exception as e:
|
||||
log.warning(f'stream resume log terminal-TTL shorten failed for {message_id}: {e}')
|
||||
|
||||
if update_db and message_id and not request_info.get('chat_id', '').startswith('local:'):
|
||||
event_type = event_data.get('type')
|
||||
|
|
|
|||
|
|
@ -485,12 +485,10 @@
|
|||
// Set all response messages to done
|
||||
for (const messageId of history.messages[message.parentId].childrenIds) {
|
||||
history.messages[messageId].done = true;
|
||||
resumeSeqByMessageId.delete(messageId);
|
||||
}
|
||||
await processNextInQueue($chatId);
|
||||
} else {
|
||||
message.done = true;
|
||||
resumeSeqByMessageId.delete(message.id);
|
||||
}
|
||||
} else if (type === 'chat:message:delta' || type === 'message') {
|
||||
message.content += data.content;
|
||||
|
|
@ -1908,8 +1906,6 @@
|
|||
if (done) {
|
||||
message.done = true;
|
||||
|
||||
resumeSeqByMessageId.delete(message.id);
|
||||
|
||||
if ($settings.responseAutoCopy) {
|
||||
copyToClipboard(message.content);
|
||||
}
|
||||
|
|
@ -2541,7 +2537,6 @@
|
|||
};
|
||||
|
||||
responseMessage.done = true;
|
||||
resumeSeqByMessageId.delete(responseMessageId);
|
||||
|
||||
history.messages[responseMessageId] = responseMessage;
|
||||
history.currentId = responseMessageId;
|
||||
|
|
@ -2609,7 +2604,6 @@
|
|||
content: $i18n.t(`Uh-oh! There was an issue with the response.`) + '\n' + errorMessage
|
||||
};
|
||||
responseMessage.done = true;
|
||||
resumeSeqByMessageId.delete(responseMessage.id);
|
||||
|
||||
if (responseMessage.statusHistory) {
|
||||
responseMessage.statusHistory = responseMessage.statusHistory.filter(
|
||||
|
|
@ -2643,7 +2637,6 @@
|
|||
if (responseMessage.parentId && history.messages[responseMessage.parentId]) {
|
||||
for (const messageId of history.messages[responseMessage.parentId].childrenIds) {
|
||||
history.messages[messageId].done = true;
|
||||
resumeSeqByMessageId.delete(messageId);
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue