From 6155105971501a2dc6de9864a9e82553765bbf18 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 14 Apr 2026 21:46:12 +0000 Subject: [PATCH] style(stream): tighten comments on delta-emit change --- backend/open_webui/utils/middleware.py | 73 ++++---------------------- 1 file changed, 10 insertions(+), 63 deletions(-) diff --git a/backend/open_webui/utils/middleware.py b/backend/open_webui/utils/middleware.py index 54aece9489..02d70669e4 100644 --- a/backend/open_webui/utils/middleware.py +++ b/backend/open_webui/utils/middleware.py @@ -3647,15 +3647,7 @@ async def streaming_chat_response_handler(response, ctx): int(metadata.get('params', {}).get('stream_delta_chunk_size') or 1), ) last_delta_data = None - # Plain-text tokens accumulate here instead of being - # re-serialized with the full chat history on every SSE - # event. Emitted as `chat:message:delta` (which the - # frontend already handles via string append), dropping - # per-token WS payload from O(size_so_far) to O(new_chars). - # Structural events (reasoning, tool calls, images, ...) - # still take the legacy full-serialize path below and - # clear this buffer when they do, since the full-content - # checkpoint already contains the accumulated text. + # Plain text chars buffered for chat:message:delta emit. pending_text_delta = '' async def flush_pending_delta_data(threshold: int = 0): @@ -3664,10 +3656,7 @@ async def streaming_chat_response_handler(response, ctx): nonlocal pending_text_delta if delta_count >= threshold: - # Emit the full-content checkpoint first (reasoning - # updates, tool-call state, etc.). This overwrites - # the frontend `message.content` with the canonical - # backend state at that point. + # Full-content checkpoint first, then appended text delta. if last_delta_data is not None: await event_emitter( { @@ -3677,10 +3666,6 @@ async def streaming_chat_response_handler(response, ctx): ) last_delta_data = None - # Then emit plain-text tokens accumulated AFTER - # the last checkpoint. The frontend appends these - # to `message.content`, matching what the next - # serialize_output checkpoint would produce. if pending_text_delta: await event_emitter( { @@ -3720,9 +3705,6 @@ async def streaming_chat_response_handler(response, ctx): if data: if 'event' in data and not getattr(request.state, 'direct', False): - # Flush any accumulated text delta so this - # out-of-band event stays correctly - # ordered relative to streamed content. await flush_pending_delta_data() await event_emitter(data.get('event', {})) @@ -3767,11 +3749,7 @@ async def streaming_chat_response_handler(response, ctx): processed_data.update(response_metadata) processed_data.pop('done', None) - # Responses API emits a full-content - # checkpoint; flush pending text delta - # first and drop the pending buffer, - # since `processed_data.content` already - # includes everything accumulated. + # processed_data.content subsumes any pending text. await flush_pending_delta_data() pending_text_delta = '' await event_emitter( @@ -3839,11 +3817,6 @@ async def streaming_chat_response_handler(response, ctx): url = url_citation.get('url', '') title = url_citation.get('title', url) - # Citation references text - # that was streamed before it — - # flush pending text delta so - # the citation arrives after - # the content it cites. await flush_pending_delta_data() await event_emitter( { @@ -3939,10 +3912,6 @@ async def streaming_chat_response_handler(response, ctx): if message_files is None: message_files = image_file_list - # Flush pending text delta so the new - # file attachment lands in the - # correct position relative to the - # streamed text around it. await flush_pending_delta_data() await event_emitter( { @@ -3959,12 +3928,7 @@ async def streaming_chat_response_handler(response, ctx): or delta.get('thinking') ) - # Snapshot output structure before this - # event mutates it. Used at the emit - # decision below to detect whether this - # event was a pure text append to an - # existing message block (the hot path - # that gets the delta optimization). + # Snapshot pre-event output for the delta classifier below. prev_output_len_before_value = len(output) prev_last_type_before_value = ( output[-1].get('type') if output else None @@ -4156,19 +4120,10 @@ async def streaming_chat_response_handler(response, ctx): }, ) else: - # Plain-text delta classifier. - # Streaming-hot-path fix: a single - # SSE event that only appended text - # characters to the currently - # active `message` block can be - # emitted as a tiny `chat:message:delta` - # instead of re-serializing the - # whole chat. Any other case - # (reasoning block update, tag - # block, block boundary change, - # tool-call, etc.) falls through - # to the legacy full-serialize - # path so semantics stay identical. + # Pure text-append-to-active-message-block events + # take the delta fast path. Everything else (reasoning, + # tag blocks, structural changes) falls through to + # the legacy full-serialize emit. is_plain_text_delta = ( not reasoning_content and not inside_tag_block @@ -4188,23 +4143,15 @@ async def streaming_chat_response_handler(response, ctx): if delta: delta_count += 1 if value and is_plain_text_delta: - # Text accumulated into pending_text_delta; - # do NOT set last_delta_data (no full-content - # checkpoint for this event). + # Already buffered in pending_text_delta. pass else: + # data subsumes any pending text. last_delta_data = data - # The full-content checkpoint in `data` - # already includes any previously - # queued pending text, so drop it to - # avoid double-emitting those chars. pending_text_delta = '' if delta_count >= delta_chunk_size: await flush_pending_delta_data(delta_chunk_size) else: - # Non-delta event (no choices[].delta): - # flush any pending streaming state first - # so this emission lands in order. await flush_pending_delta_data() await event_emitter( {