Skip to content

Streaming responses never return their connection to the pool (closed before the chunked terminator is read) #4040

Description

@mrembalski

Confirm this is an issue with the Python library and not an underlying OpenAI API

  • This is an issue with the Python library

Describe the bug

Follow-up to #3440.

Summary

Stream.__stream__ / AsyncStream.__stream__ stop iterating as soon as they see data: [DONE] and then close the response. At that point the HTTP/1.1 chunked-encoding terminator (0\r\n\r\n) is still unread, so the transport cannot mark the message complete and closes the TCP connection instead of returning it to the pool.

Result: every streaming chat completion opens a new TCP connection (plus a TLS handshake over https), regardless of keepalive_expiry, pool limits, or server keep-alive settings. Non-streaming requests pool normally, and a raw httpx streaming request read to EOF pools normally, so this is specific to the SDK's stream handling.

Regression since 6132922 ("fix(client): close streams without requiring full consumption", first released in v2.7.0).

Root cause

# AsyncStream.__stream__
try:
    async for sse in iterator:
        if sse.data.startswith("[DONE]"):
            break
        ...
finally:
    await response.aclose()

Stream.__stream__ does the same synchronously.

On close, httpcore decides whether the connection can be pooled from the h11 state machine (httpcore/_async/http11.py, _response_closed): the connection goes back to IDLE only if our_state is h11.DONE and their_state is h11.DONE; otherwise it is closed.

their_state only reaches DONE after h11 has parsed the 0\r\n\r\n terminator, i.e. after the body generator yields EndOfMessage. Because __stream__ stops at the [DONE] Data event, that never happens. httpcore2 keeps the same rule.

Proposed fix

Drain only in the [DONE] branch, never in finally, so abandoned or errored streams keep today's immediate-close behavior:

if sse.data.startswith("[DONE]"):
    # The server has signalled completion; only the body terminator remains.
    # Read it so the connection can return to the pool.
    async for _ in iterator:
        pass
    break

and the same for _ in iterator: pass in Stream.__stream__.

On the "adds waiting" concern from #3440: this reads bytes the server has already committed to sending, and each read is bounded by the request's existing read timeout exactly like every other chunk of the stream.

If an explicit bound is needed, two options that keep the common path deterministic:

  1. Time bound: wrap the drain in anyio.move_on_after(1.0) for the async stream (anyio is already a dependency); fall through to the existing close when it expires.
  2. Byte bound: stop draining once more than a few hundred bytes arrive after [DONE] (a healthy server sends less), which bounds work without a timer and applies to sync and async alike.

Impact

I measured 0.1-0.25 s of extra time-to-first-token per request that disappears once the connection is reused. Keep-alive and pool tuning on either side cannot help until the client reads the terminator.

To Reproduce

import asyncio
import json
import re
import sys
import threading

import openai
from openai import AsyncOpenAI, OpenAI

CHUNKS = [
    {"id": "c", "object": "chat.completion.chunk", "created": 0, "model": "m",
     "choices": [{"index": 0, "delta": {"content": "hi"}, "finish_reason": None}]},
    {"id": "c", "object": "chat.completion.chunk", "created": 0, "model": "m",
     "choices": [{"index": 0, "delta": {}, "finish_reason": "stop"}]},
]
COMPLETION = {"id": "c", "object": "chat.completion", "created": 0, "model": "m",
              "choices": [{"index": 0, "message": {"role": "assistant", "content": "hi"},
                           "finish_reason": "stop"}]}
MSGS = [{"role": "user", "content": "hi"}]
connections = 0


async def serve(reader: asyncio.StreamReader, writer: asyncio.StreamWriter) -> None:
    global connections
    connections += 1
    try:
        while True:
            head = await reader.readuntil(b"\r\n\r\n")
            length = int(re.search(rb"content-length: (\d+)", head, re.I).group(1))
            body = json.loads(await reader.readexactly(length))
            if body.get("stream"):
                writer.write(b"HTTP/1.1 200 OK\r\ncontent-type: text/event-stream\r\n"
                             b"transfer-encoding: chunked\r\n\r\n")
                events = [f"data: {json.dumps(c)}\n\n".encode() for c in CHUNKS]
                events.append(b"data: [DONE]\n\n")
                for ev in events:
                    writer.write(f"{len(ev):x}\r\n".encode() + ev + b"\r\n")
                    await writer.drain()
                    await asyncio.sleep(0.005)
                writer.write(b"0\r\n\r\n")
            else:
                payload = json.dumps(COMPLETION).encode()
                writer.write(b"HTTP/1.1 200 OK\r\ncontent-type: application/json\r\n"
                             + f"content-length: {len(payload)}\r\n\r\n".encode() + payload)
            await writer.drain()
    except (asyncio.IncompleteReadError, ConnectionError):
        pass
    finally:
        writer.close()


def start_server() -> tuple[int, callable]:
    ready = threading.Event()
    state: dict = {}

    def run() -> None:
        async def main() -> None:
            server = await asyncio.start_server(serve, "127.0.0.1", 0)
            state["port"] = server.sockets[0].getsockname()[1]
            state["loop"] = asyncio.get_running_loop()
            state["stop"] = asyncio.Event()
            ready.set()
            await state["stop"].wait()
            server.close()
            server.close_clients()
            await server.wait_closed()

        asyncio.run(main())

    thread = threading.Thread(target=run, daemon=True)
    thread.start()
    ready.wait()

    def stop() -> None:
        state["loop"].call_soon_threadsafe(state["stop"].set)
        thread.join(timeout=5)

    return state["port"], stop


def report(label: str) -> None:
    global connections
    print(f"{label:<44} -> {connections} TCP connection(s)")
    connections = 0


async def async_cases(base_url: str) -> None:
    client = AsyncOpenAI(api_key="x", base_url=base_url)
    for _ in range(5):
        async for _ in await client.chat.completions.create(model="m", messages=MSGS, stream=True):
            pass
    report("AsyncOpenAI streaming x5, fully consumed")

    client = AsyncOpenAI(api_key="x", base_url=base_url)
    for _ in range(5):
        async with await client.chat.completions.create(model="m", messages=MSGS, stream=True) as s:
            async for _ in s:
                pass
    report("AsyncOpenAI streaming x5, async with")

    client = AsyncOpenAI(api_key="x", base_url=base_url)
    for _ in range(5):
        await client.chat.completions.create(model="m", messages=MSGS)
    report("AsyncOpenAI non-streaming x5")

    try:
        import httpx
    except ImportError:
        return
    async with httpx.AsyncClient() as raw:
        for _ in range(5):
            async with raw.stream("POST", f"{base_url}/chat/completions",
                                  json={"model": "m", "messages": MSGS, "stream": True}) as r:
                async for _ in r.aiter_lines():
                    pass
    report(f"raw httpx {httpx.__version__} stream x5, read to EOF")


def sync_cases(base_url: str) -> None:
    client = OpenAI(api_key="x", base_url=base_url)
    for _ in range(5):
        for _ in client.chat.completions.create(model="m", messages=MSGS, stream=True):
            pass
    report("OpenAI (sync) streaming x5, fully consumed")

    client = OpenAI(api_key="x", base_url=base_url)
    for _ in range(5):
        client.chat.completions.create(model="m", messages=MSGS)
    report("OpenAI (sync) non-streaming x5")


if __name__ == "__main__":
    port, stop = start_server()
    base_url = f"http://127.0.0.1:{port}/v1"
    print(f"openai {openai.__version__}, python {sys.version.split()[0]}")
    asyncio.run(async_cases(base_url))
    sync_cases(base_url)
    stop()

Output on Python 3.13.12, expected 1 connection in every row:

case (5 requests each) 2.6.1 2.7.0 2.45.0 3.26.1
AsyncOpenAI streaming, fully consumed 1 5 5 5
AsyncOpenAI streaming, async with + consumed 1 5 5 5
AsyncOpenAI non-streaming 1 1 1 1
raw httpx 0.28.1 stream, read to EOF 1 1 1 n/a (httpx2)
OpenAI (sync) streaming, fully consumed 1 5 5 5
OpenAI (sync) non-streaming 1 1 1 1

2.6.1 is the last release that pools streaming connections.

Code snippets

OS

macOS

Python version

Python v3.13.12

Library version

3.26.1 (httpx2 2.13.1, httpcore2 2.13.1)

Activity

  1. nightcityblade commented on Oct 8, 2026

    @nightcityblade
    Contributor

    Hi, I'd like to work on this. I'll submit a PR shortly.

  2. nightcityblade commented on Oct 8, 2026

    @nightcityblade
    Contributor

    I need to step back: the repository's PR template confirms pull requests are limited to repository collaborators, so I can't submit the promised upstream PR. This issue is unclaimed from my side.

  3. marcuswood-oai commented on Oct 8, 2026

    @marcuswood-oai
    Contributor

    Thanks for the detailed reproduction! We confirmed the connection-reuse issue. The earlier drain fixes risked adding waits or errors after a completed stream.

    Enabling HTTP/2 restored connection reuse in our tests, without SDK changes. For your reported version:

    pip install 'httpx2[http2]==2.13.1'
    from openai import AsyncOpenAI, DefaultAsyncHttpx2Client
    
    client = AsyncOpenAI(
        http_client=DefaultAsyncHttpx2Client(http2=True),
    )

    For sync usage, use OpenAI and DefaultHttpx2Client. Keep reusing the client and preserve your existing settings.

    Could you try this, confirm stream.response.http_version reports HTTP/2, and let us know whether the extra latency disappears? If it reports HTTP/1.1, the connection fell back and this workaround won’t help.

    Related: #3476 — Consider enabling HTTP2 by default.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    bugSomething isn't working

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions