From 29cceeeba71c98efb82fbe05519e22f756deea20 Mon Sep 17 00:00:00 2001 From: Tommaso Barbugli Date: Thu, 1 Oct 2026 13:41:54 +0200 Subject: [PATCH] perf: raise the default connection pool size from 5 to 100 max_conns_per_host is applied as httpx's hard max_connections, so the old default of 5 let a client run only 5 requests at a time; the rest queued for up to the 30s pool timeout. Raise it to 100 (httpx's own default) and keep the keep-alive pool the same size so busy apps reuse connections instead of re-handshaking. Idle connections still expire after 55s. Also define the pool defaults once in getstream.base, add tests against a real local server for concurrency, the cap, and keep-alive reuse, and add a README section on reusing the client and tuning the pool. --- README.md | 19 +++++++ getstream/base.py | 7 ++- getstream/stream.py | 17 +++--- tests/test_http_client.py | 107 ++++++++++++++++++++++++++++++++++---- 4 files changed, 129 insertions(+), 21 deletions(-) diff --git a/README.md b/README.md index 71f8765e..40944d1e 100644 --- a/README.md +++ b/README.md @@ -99,6 +99,25 @@ response: StreamResponse[StartClosedCaptionsResponse] = call.start_closed_captio response.data # Gives the StartClosedCaptionsResponse model ``` +### Reusing the client and connection pooling + +Create one `Stream` (or `AsyncStream`) client when your app starts and reuse it everywhere, for example as a module-level object or on your web framework's app state. The client keeps a pool of keep-alive HTTP connections, so most requests skip the TCP and TLS handshake. Creating a new client per request throws that pool away and adds latency to every call. The sub-clients (`client.video`, `client.chat`, `client.moderation`, `client.feeds`) all share the same pool. Call `client.close()` (or `await client.aclose()`) when your app shuts down. + +By default the pool allows up to 100 concurrent connections, closes connections that have been idle for 55 seconds, and uses a 10 second connect timeout and a 30 second request timeout. When every connection is busy, further requests wait for one to free up. If your app sends more concurrent requests than that, raise the limit: + +```python +client = Stream( + api_key="your_api_key", + api_secret="your_api_secret", + max_conns_per_host=200, # max concurrent connections, also kept alive for reuse + idle_timeout=55.0, # seconds an idle connection stays open + connect_timeout=10.0, # seconds for the TCP + TLS handshake + request_timeout=30.0, # seconds per request +) +``` + +The same settings can be set with the `STREAM_MAX_CONNS_PER_HOST`, `STREAM_IDLE_TIMEOUT`, `STREAM_CONNECT_TIMEOUT` and `STREAM_REQUEST_TIMEOUT` environment variables. For full control, pass your own `httpx.Client` (or `httpx.AsyncClient`) as `http_client=`; the settings above are then ignored and your client's configuration is used as-is. + ### Logging The SDK emits structured log events (`client.initialized`, `http.request.sent`, `http.response.received`, `http.request.failed`) through the stdlib `logging` module. By default nothing is printed: pass a `logging.Logger` to see them. diff --git a/getstream/base.py b/getstream/base.py index 4dbd48e1..71389734 100644 --- a/getstream/base.py +++ b/getstream/base.py @@ -38,9 +38,12 @@ # ── Connection pool defaults (CHA-2956) ────────────────────────────── -# Kept in sync with getstream.stream constants; duplicated here so BaseClient/AsyncBaseClient can be instantiated standalone (e.g. by sub-clients constructed directly without going through Stream/AsyncStream). -DEFAULT_MAX_CONNS_PER_HOST = 5 +# Defined here (not in getstream.stream) so BaseClient/AsyncBaseClient can be instantiated standalone (e.g. by sub-clients constructed directly without going through Stream/AsyncStream). +# DEFAULT_MAX_CONNS_PER_HOST is a hard cap on in-flight requests per client (excess requests wait up to the pool timeout), and also the keep-alive pool size so connections are reused instead of re-handshaked under sustained load. Matches httpx's own default max_connections. +DEFAULT_MAX_CONNS_PER_HOST = 100 +# DEFAULT_IDLE_TIMEOUT sits below the typical 60s LB idle timeout with a 5s safety margin. DEFAULT_IDLE_TIMEOUT = 55.0 +# DEFAULT_CONNECT_TIMEOUT caps TCP + TLS handshake duration. DEFAULT_CONNECT_TIMEOUT = 10.0 diff --git a/getstream/stream.py b/getstream/stream.py index d25848ff..d0bdea34 100644 --- a/getstream/stream.py +++ b/getstream/stream.py @@ -11,7 +11,13 @@ import jwt from pydantic_settings import BaseSettings, SettingsConfigDict -from getstream.base import _log_client_initialized, _resolve_logger +from getstream.base import ( + DEFAULT_CONNECT_TIMEOUT, + DEFAULT_IDLE_TIMEOUT, + DEFAULT_MAX_CONNS_PER_HOST, + _log_client_initialized, + _resolve_logger, +) from getstream.common import telemetry from getstream.config import RetryConfig from getstream.chat.client import ChatClient @@ -31,14 +37,9 @@ BASE_URL = "https://chat.stream-io-api.com/" # ── Connection pool defaults (CHA-2956) ────────────────────────────── +# Pool defaults (DEFAULT_MAX_CONNS_PER_HOST etc.) live in getstream.base. # DEFAULT_REQUEST_TIMEOUT is the default per-request timeout (was 6.0 prior to 3.5.0). DEFAULT_REQUEST_TIMEOUT = 30.0 -# DEFAULT_MAX_CONNS_PER_HOST caps concurrent TCP connections per host. -DEFAULT_MAX_CONNS_PER_HOST = 5 -# DEFAULT_IDLE_TIMEOUT sits below the typical 60s LB idle timeout with a 5s safety margin. -DEFAULT_IDLE_TIMEOUT = 55.0 -# DEFAULT_CONNECT_TIMEOUT caps TCP + TLS handshake duration. -DEFAULT_CONNECT_TIMEOUT = 10.0 class Settings(BaseSettings): @@ -116,7 +117,7 @@ def __init__( http_client: Optional pre-built ``httpx`` client. Mutually exclusive with ``transport``. When provided, sub-clients (video/chat/moderation) reuse it instead of opening their own. token: Pre-minted user JWT. Mutually exclusive with ``api_secret``. request_timeout: Default per-request timeout in seconds. Default 30.0. Replaces the older ``timeout`` kwarg; ``timeout`` is kept as an alias for backward compatibility. - max_conns_per_host: Max concurrent TCP connections per host. Default 5. Ignored when ``http_client`` is set. + max_conns_per_host: Max concurrent TCP connections per host, i.e. max in-flight requests for this client; also the number of idle keep-alive connections kept for reuse. Default 100. Ignored when ``http_client`` is set. idle_timeout: Idle connection lifetime in seconds. Default 55.0 (sits 5s under the typical 60s LB idle timeout). Ignored when ``http_client`` is set. connect_timeout: TCP + TLS handshake timeout in seconds. Default 10.0. Ignored when ``http_client`` is set. logger: Optional stdlib ``logging.Logger`` for the SDK's structured log events (``client.initialized``, ``http.request.sent``, ``http.response.received``, ``http.request.failed``). Defaults to ``logging.getLogger("getstream")``, which is a no-op until the caller attaches a handler. diff --git a/tests/test_http_client.py b/tests/test_http_client.py index fbe7d8e0..b318deae 100644 --- a/tests/test_http_client.py +++ b/tests/test_http_client.py @@ -1,5 +1,9 @@ +import asyncio import logging import os +import threading +import time +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer import httpx import pytest @@ -33,8 +37,8 @@ class TestSyncPoolDefaults: def test_default_limits_applied(self): client = Stream(api_key="k", api_secret="s", base_url="http://test") pool = client.client._transport._pool - assert pool._max_connections == 5 - assert pool._max_keepalive_connections == 5 + assert pool._max_connections == 100 + assert pool._max_keepalive_connections == 100 assert pool._keepalive_expiry == 55.0 def test_default_timeout_is_30s(self): @@ -51,8 +55,8 @@ class TestAsyncPoolDefaults: async def test_default_limits_applied(self): client = AsyncStream(api_key="k", api_secret="s", base_url="http://test") pool = client.client._transport._pool - assert pool._max_connections == 5 - assert pool._max_keepalive_connections == 5 + assert pool._max_connections == 100 + assert pool._max_keepalive_connections == 100 assert pool._keepalive_expiry == 55.0 await client.aclose() @@ -176,19 +180,100 @@ async def test_async_sub_client_pools_match_configured_knobs(self): def test_sync_sub_client_pools_match_defaults(self): # Even with no explicit knobs, sub-clients must carry the SDK defaults - # (5/55/10/30), not whatever a freshly-built sub-client would default to. + # (100/55/10/30), not whatever a freshly-built sub-client would default to. client = Stream(api_key="k", api_secret="s", base_url="http://test") for name in ("video", "chat", "moderation", "feeds"): sub = getattr(client, name) assert sub.client is client.client pool = sub.client._transport._pool - assert pool._max_connections == 5 + assert pool._max_connections == 100 assert pool._keepalive_expiry == 55.0 assert sub.client.timeout.connect == 10.0 assert sub.client.timeout.read == 30.0 client.close() +# ── pool behavior against a real local server ──────────────────────── + + +class _TrackingHandler(BaseHTTPRequestHandler): + protocol_version = "HTTP/1.1" + + def do_GET(self): + server = self.server + with server.lock: + server.in_flight += 1 + server.peak = max(server.peak, server.in_flight) + server.client_ports.add(self.client_address[1]) + time.sleep(server.delay) + with server.lock: + server.in_flight -= 1 + body = b"{}" + self.send_response(200) + self.send_header("Content-Type", "application/json") + self.send_header("Content-Length", str(len(body))) + self.end_headers() + self.wfile.write(body) + + def log_message(self, *args): + pass + + +class _TrackingServer(ThreadingHTTPServer): + daemon_threads = True + request_queue_size = 128 + + def __init__(self, delay): + super().__init__(("127.0.0.1", 0), _TrackingHandler) + self.delay = delay + self.lock = threading.Lock() + self.in_flight = 0 + self.peak = 0 + self.client_ports = set() + + +@pytest.fixture +def tracking_server(): + server = _TrackingServer(delay=0.2) + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + yield server + server.shutdown() + server.server_close() + + +@pytest.mark.asyncio +class TestPoolBehavior: + def _make(self, server, **kw): + return AsyncStream( + api_key="k", + api_secret="s", + base_url=f"http://127.0.0.1:{server.server_port}", + **kw, + ) + + async def test_default_pool_allows_high_concurrency(self, tracking_server): + client = self._make(tracking_server) + await asyncio.gather(*(client.get("/app") for _ in range(30))) + assert tracking_server.peak == 30 + await client.aclose() + + async def test_max_conns_per_host_caps_concurrency(self, tracking_server): + client = self._make(tracking_server, max_conns_per_host=3) + await asyncio.gather(*(client.get("/app") for _ in range(12))) + assert tracking_server.peak == 3 + assert len(tracking_server.client_ports) == 3 + await client.aclose() + + async def test_sequential_requests_reuse_connection(self, tracking_server): + tracking_server.delay = 0 + client = self._make(tracking_server) + for _ in range(5): + await client.get("/app") + assert len(tracking_server.client_ports) == 1 + await client.aclose() + + # ── transport (primary API) ────────────────────────────────────────── @@ -423,7 +508,7 @@ def test_info_log_emitted_with_defaults(self, caplog): assert len(infos) == 1 r = infos[0] assert r.getMessage() == "client.initialized" - assert getattr(r, "stream.client.max_conns_per_host") == 5 + assert getattr(r, "stream.client.max_conns_per_host") == 100 assert getattr(r, "stream.client.idle_timeout_seconds") == 55.0 assert getattr(r, "stream.client.connect_timeout_seconds") == 10.0 assert getattr(r, "stream.client.request_timeout_seconds") == 30.0 @@ -481,7 +566,7 @@ async def test_info_log_emitted_with_defaults(self, caplog): await client.aclose() infos = [r for r in caplog.records if r.name == "getstream"] assert len(infos) == 1 - assert getattr(infos[0], "stream.client.max_conns_per_host") == 5 + assert getattr(infos[0], "stream.client.max_conns_per_host") == 100 assert getattr(infos[0], "stream.client.user_http_client") is False @@ -678,12 +763,12 @@ def _clear_stream_env(monkeypatch): def test_sync_constructs_with_spec_defaults(self, monkeypatch): self._clear_stream_env(monkeypatch) client = Stream(api_key="k", api_secret="s", base_url="http://test") - assert client.max_conns_per_host == 5 + assert client.max_conns_per_host == 100 assert client.idle_timeout == 55.0 assert client.connect_timeout == 10.0 assert client.request_timeout == 30.0 pool = client.client._transport._pool - assert pool._max_connections == 5 + assert pool._max_connections == 100 assert pool._keepalive_expiry == 55.0 client.close() @@ -691,7 +776,7 @@ def test_sync_constructs_with_spec_defaults(self, monkeypatch): async def test_async_constructs_with_spec_defaults(self, monkeypatch): self._clear_stream_env(monkeypatch) client = AsyncStream(api_key="k", api_secret="s", base_url="http://test") - assert client.max_conns_per_host == 5 + assert client.max_conns_per_host == 100 assert client.idle_timeout == 55.0 assert client.connect_timeout == 10.0 assert client.request_timeout == 30.0