fix(stream): env-tunable timeouts, request_id correlation, doc concurrent-emitter limit

- Timeouts split and env-configurable. RESUME_STREAM_REDIS_TIMEOUT_SEC
  (hot path, default 100ms) stays tight to protect live streaming.
  RESUME_STREAM_READ_TIMEOUT_SEC (replay XRANGE, default 1s) is looser
  since a resume is user-blocking and silently timing out is worse
  than a brief extra wait. Cross-region / non-colocated Redis setups
  can override either.

- resume-stream now carries a request_id; the server echoes it in
  resume-stream:replay and the client ignores replies whose id does
  not match the active request for that message. Prevents reconnect
  churn where a stale reply could clear a fence belonging to a
  newer in-flight request.

- Documented the concurrent-emitter live ordering caveat in
  get_event_emitter. Single-emitter flows are safe (streaming loop
  awaits sequentially). Replay is safe (sorted by seq on read).
  Overlapping emitters for the same message_id can produce live frames
  out of seq order — accepting as a known limitation; proper fix
  needs either a distributed per-message lock or a client reorder
  buffer, neither justified by the rarity of the scenario.
This commit is contained in:
Claude 2026-04-14 22:49:12 +00:00
parent 7b30714941
commit fd1cf4221e
No known key found for this signature in database
2 changed files with 48 additions and 5 deletions

View file

@ -1,5 +1,6 @@
import asyncio
import json
import os
import random
import socketio
@ -179,9 +180,16 @@ RESUME_STREAM_MAXLEN = 2000
RESUME_STREAM_TTL_SEC = 3600
RESUME_STREAM_DONE_TTL_SEC = 30
RESUME_STREAM_TTL_REFRESH_EVERY = 64
# 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.1
# Hot-path timeout: tight so a slow Redis can't stall live tokens.
# Replay read timeout: looser since a resume is user-blocking anyway
# and silent timeout here is worse than a brief extra wait. Both
# configurable for infra where Redis isn't colocated.
RESUME_STREAM_REDIS_TIMEOUT_SEC = float(
os.environ.get('RESUME_STREAM_REDIS_TIMEOUT_SEC', '0.1')
)
RESUME_STREAM_READ_TIMEOUT_SEC = float(
os.environ.get('RESUME_STREAM_READ_TIMEOUT_SEC', '1.0')
)
def _stream_key(user_id: str, message_id: str) -> str:
@ -268,7 +276,7 @@ async def _stream_log_read(user_id: str, message_id: str, after_seq: int):
try:
entries = await asyncio.wait_for(
REDIS.xrange(_stream_key(user_id, message_id), min='-', max='+'),
timeout=RESUME_STREAM_REDIS_TIMEOUT_SEC,
timeout=RESUME_STREAM_READ_TIMEOUT_SEC,
)
except asyncio.TimeoutError:
log.warning(f'stream resume log read timed out for {message_id}')
@ -678,6 +686,7 @@ async def resume_stream(sid, data):
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):
@ -690,9 +699,15 @@ async def resume_stream(sid, data):
if REDIS is not None:
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',
{'message_id': message_id, 'envelopes': envelopes},
{
'message_id': message_id,
'request_id': request_id,
'envelopes': envelopes,
},
to=sid,
)
@ -987,6 +1002,17 @@ async def disconnect(sid):
async def get_event_emitter(request_info, update_db=True):
# A single emitter's calls are serial (the streaming loop awaits each
# before the next), so seq ordering is guaranteed within one emitter.
# CONCURRENT emitters for the same (user_id, message_id) — e.g. a
# duplicate request leaking past frontend dedup, or a retry path that
# overlaps with the original — can interleave INCR/XADD/emit across
# tasks and cause live frames to arrive out of seq order. Replay
# reads already sort by seq so resume is safe, but live streaming in
# that edge case can drop a frame via the client dedupe guard. Fixing
# it properly needs either a distributed per-message lock or a
# client-side reorder buffer; neither is worth the complexity for a
# configuration OWUI doesn't normally produce.
async def __event_emitter__(event_data):
user_id = request_info['user_id']
chat_id = request_info['chat_id']

View file

@ -634,6 +634,9 @@
// Fallback timer in case the replay ack never arrives.
const RESUME_FENCE_TIMEOUT_MS = 10000;
const resumeFenceTimerByMessageId = new Map();
// Active request_id per message — lets us ignore stale replay
// responses from superseded requests (reconnect churn).
const resumeActiveRequestIdByMessageId = new Map();
// Drop fence + timer without flushing (lifecycle transitions).
const dropResumeFence = (messageId) => {
@ -643,6 +646,7 @@
resumeFenceTimerByMessageId.delete(messageId);
}
resumeQueueByMessageId.delete(messageId);
resumeActiveRequestIdByMessageId.delete(messageId);
};
// Drop fence AND flush buffered events (happy path + timeout).
@ -683,8 +687,14 @@
clearResumeFence(message.id);
}, RESUME_FENCE_TIMEOUT_MS)
);
const requestId =
typeof crypto !== 'undefined' && crypto.randomUUID
? crypto.randomUUID()
: `${message.id}-${Date.now()}-${Math.random()}`;
resumeActiveRequestIdByMessageId.set(message.id, requestId);
$socket.emit('resume-stream', {
message_id: message.id,
request_id: requestId,
last_seq: resumeSeqByMessageId.get(message.id) ?? 0
});
};
@ -692,6 +702,13 @@
const onResumeStreamReplay = async (payload) => {
const messageId = payload?.message_id;
if (!messageId) return;
// Drop stale replies from superseded requests so they can't
// clear a fence that belongs to a newer in-flight request.
const expected = resumeActiveRequestIdByMessageId.get(messageId);
if (expected && payload?.request_id && payload.request_id !== expected) {
return;
}
resumeActiveRequestIdByMessageId.delete(messageId);
const envelopes = Array.isArray(payload?.envelopes) ? payload.envelopes : [];
try {
for (const envelope of envelopes) {