fix: run API outlet filters when streaming clients disconnect on [DONE]

Direct API clients that close the connection as soon as they receive the SSE `[DONE]` frame (hermes-agent does this on every request) never had their outlet filters run. On disconnect the server cancels the stream task, and the outlet was only awaited after the stream had been read to the end, so it was cancelled before it started or while it ran.

The outlet now starts as its own task when the `[DONE]` frame passes through, before it is sent to the client, and the stream awaits that task at the end. A disconnect still cancels the stream and still closes the upstream connection right away as before, but the outlet task keeps running to completion. Clients that read to the end see no change: the response bytes are identical and the outlet has finished before the stream closes. Streams without a `[DONE]` frame keep the existing behaviour.

Only shielding the existing outlet call does not fix it: after `[DONE]` the stream still waits on the upstream EOF and connection cleanup, and a disconnect landing there skips the outlet entirely.

Fixes #29869
This commit is contained in:
Classic298 2026-09-24 12:13:22 +02:00
parent bbfa876afd
commit 77ad55a581

View file

@ -6650,6 +6650,7 @@ async def streaming_chat_response_handler(response, ctx):
try:
assistant_message = {}
outlet_task = None
filter_context = FilterContext()
has_api_outlet_filters = ENABLE_API_OUTLET_FILTERS and bool(filter_functions)
if ENABLE_API_OUTLET_FILTERS and not has_api_outlet_filters:
@ -6700,9 +6701,21 @@ async def streaming_chat_response_handler(response, ctx):
if data:
if has_api_outlet_filters:
update_assistant_message_from_stream(assistant_message, data)
# Clients may disconnect right after [DONE]; a task outlives that cancellation
if has_api_outlet_filters and assistant_message and outlet_task is None:
line = data.decode('utf-8', 'replace') if isinstance(data, bytes) else data
if isinstance(line, str) and any(
part.removeprefix('data:').strip() == '[DONE]' for part in line.splitlines()
):
ctx['assistant_message'] = assistant_message
outlet_task = asyncio.create_task(outlet_filter_handler(ctx))
yield data
if has_api_outlet_filters and assistant_message:
if outlet_task is not None:
await asyncio.shield(outlet_task)
elif has_api_outlet_filters and assistant_message:
ctx['assistant_message'] = assistant_message
await outlet_filter_handler(ctx)
except Exception as e: