diff --git a/backend/open_webui/routers/ollama.py b/backend/open_webui/routers/ollama.py index f9b1468458..b02205f709 100644 --- a/backend/open_webui/routers/ollama.py +++ b/backend/open_webui/routers/ollama.py @@ -100,6 +100,8 @@ async def send_request( key: str | None = None, user: UserModel = None, stream: bool = False, + # passthrough must stay False for /api/chat: middleware parses it per line + passthrough: bool = False, content_type: str | None = None, metadata: dict | None = None, api_config: dict | None = None, @@ -171,7 +173,7 @@ async def send_request( streaming = True return StreamingResponse( - stream_wrapper(r), + stream_wrapper(r, passthrough=passthrough), status_code=r.status, headers=response_headers, ) @@ -658,6 +660,7 @@ async def pull_model( key=get_api_key(url_idx, url, (await Config.get('ollama.api_configs', {}))), user=user, stream=True, + passthrough=True, ) @@ -697,6 +700,7 @@ async def push_model( key=get_api_key(url_idx, url, (await Config.get('ollama.api_configs', {}))), user=user, stream=True, + passthrough=True, ) @@ -729,6 +733,7 @@ async def create_model( key=get_api_key(url_idx, url, (await Config.get('ollama.api_configs', {}))), user=user, stream=True, + passthrough=True, ) @@ -1011,6 +1016,7 @@ async def generate_completion( key=get_api_key(url_idx, url, api_configs), user=user, stream=True, + passthrough=True, ) @@ -1242,6 +1248,7 @@ async def generate_openai_completion( key=get_api_key(url_idx, url, api_configs), user=user, stream=payload.get('stream', False), + passthrough=True, metadata=metadata, api_config=api_config, request=request, @@ -1350,6 +1357,7 @@ async def generate_openai_chat_completion( key=get_api_key(url_idx, url, api_configs), user=user, stream=payload.get('stream', False), + passthrough=True, metadata=metadata, api_config=api_config, request=request, @@ -1401,6 +1409,7 @@ async def generate_anthropic_messages( key=get_api_key(url_idx, url, api_configs), user=user, stream=payload.get('stream', False), + passthrough=True, content_type='text/event-stream' if payload.get('stream', False) else None, api_config=api_config, request=request, @@ -1458,6 +1467,7 @@ async def generate_responses( key=get_api_key(url_idx, url, api_configs), user=user, stream=payload.get('stream', False), + passthrough=True, content_type='text/event-stream' if payload.get('stream', False) else None, api_config=api_config, request=request, diff --git a/backend/open_webui/routers/openai.py b/backend/open_webui/routers/openai.py index b5aac9e385..73a5f7ecc3 100644 --- a/backend/open_webui/routers/openai.py +++ b/backend/open_webui/routers/openai.py @@ -1520,7 +1520,7 @@ async def embeddings(request: Request, form_data: dict, user): if 'text/event-stream' in r.headers.get('Content-Type', ''): streaming = True return StreamingResponse( - stream_wrapper(r), + stream_wrapper(r, passthrough=True), status_code=r.status, headers=_clean_proxy_headers(r.headers), ) @@ -1647,7 +1647,7 @@ async def responses( if 'text/event-stream' in r.headers.get('Content-Type', ''): streaming = True return StreamingResponse( - stream_wrapper(r), + stream_wrapper(r, passthrough=True), status_code=r.status, headers=_clean_proxy_headers(r.headers), ) @@ -1769,7 +1769,7 @@ async def proxy(path: str, request: Request, user=Depends(get_verified_user)): if 'text/event-stream' in r.headers.get('Content-Type', ''): streaming = True return StreamingResponse( - stream_wrapper(r), + stream_wrapper(r, passthrough=True), status_code=r.status, headers=_clean_proxy_headers(r.headers), ) diff --git a/backend/open_webui/utils/session_pool.py b/backend/open_webui/utils/session_pool.py index 414c45ed48..980435e033 100644 --- a/backend/open_webui/utils/session_pool.py +++ b/backend/open_webui/utils/session_pool.py @@ -112,14 +112,23 @@ async def cleanup_response( await result -async def stream_wrapper(response, session=None, content_handler=None): +async def stream_wrapper(response, session=None, content_handler=None, passthrough=False): """Wrap a stream to ensure cleanup happens even if streaming is interrupted. This is more reliable than BackgroundTask which may not run if the client disconnects. When using the shared pool, ``session`` should be ``None``. + + ``passthrough=True`` yields raw network chunks (iter_any) instead of + lines: byte-identical output without a buffer scan, slice and copy per + line. Only for streams no internal consumer parses line-by-line. """ try: - stream = content_handler(response.content) if content_handler else response.content + if content_handler: + stream = content_handler(response.content) + elif passthrough: + stream = response.content.iter_any() + else: + stream = response.content async for chunk in stream: yield chunk finally: