mirror of
https://github.com/open-webui/open-webui.git
synced 2026-09-16 23:43:03 +00:00
style(stream): tighten comments on delta-emit change
This commit is contained in:
parent
d7b36be54d
commit
6155105971
1 changed files with 10 additions and 63 deletions
|
|
@ -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(
|
||||
{
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue