Coverage for open_webui/utils/session_pool.py: 53%
49 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 05:07 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 05:07 +0000
1"""Shared aiohttp ClientSession pool.
3Instead of creating a new ClientSession (and TCPConnector) per request,
4callers acquire a long-lived session from this module. The pool manages
5a single TCPConnector with configurable limits, enabling TCP/SSL connection
6reuse, shared DNS cache, and bounded concurrency.
8All pool parameters are configurable via environment variables:
9 - AIOHTTP_POOL_CONNECTIONS (default 100) — max total connections
10 - AIOHTTP_POOL_CONNECTIONS_PER_HOST (default 30) — per-host limit
11 - AIOHTTP_POOL_DNS_TTL (default 300) — DNS cache TTL in seconds
13Usage:
14 from open_webui.utils.session_pool import get_session, cleanup_response
16 session = await get_session()
17 r = await session.request(...)
18 # When done with the *response* (not the session):
19 await cleanup_response(r)
21IMPORTANT: Callers must NOT close the shared session. Only the response
22needs cleanup. The session is closed once during application shutdown
23via ``close_session()``.
24"""
26import logging
27from typing import Optional
29import aiohttp
30from open_webui.env import (
31 AIOHTTP_CLIENT_STREAM_IDLE_TIMEOUT,
32 AIOHTTP_CLIENT_TIMEOUT,
33 AIOHTTP_POOL_CONNECTIONS,
34 AIOHTTP_POOL_CONNECTIONS_PER_HOST,
35 AIOHTTP_POOL_DNS_TTL,
36)
37from open_webui.utils.misc import stream_chunks_handler
39log = logging.getLogger(__name__)
41_session: Optional[aiohttp.ClientSession] = None
43_CLIENT_TIMEOUT = aiohttp.ClientTimeout(total=AIOHTTP_CLIENT_TIMEOUT)
44_CLIENT_STREAM_TIMEOUT = aiohttp.ClientTimeout(
45 total=AIOHTTP_CLIENT_TIMEOUT,
46 sock_read=AIOHTTP_CLIENT_STREAM_IDLE_TIMEOUT,
47)
50def get_client_timeout(stream: bool = False) -> aiohttp.ClientTimeout:
51 return _CLIENT_STREAM_TIMEOUT if stream else _CLIENT_TIMEOUT
54async def get_session() -> aiohttp.ClientSession:
55 """Return the shared aiohttp ClientSession, creating it lazily."""
56 global _session
57 if _session is None or _session.closed:
58 connector_kwargs = {
59 'ttl_dns_cache': AIOHTTP_POOL_DNS_TTL,
60 'enable_cleanup_closed': True,
61 }
62 if AIOHTTP_POOL_CONNECTIONS is not None: 62 ↛ 63line 62 didn't jump to line 63 because the condition on line 62 was never true
63 connector_kwargs['limit'] = AIOHTTP_POOL_CONNECTIONS
64 else:
65 connector_kwargs['limit'] = 0 # aiohttp: 0 = unlimited
66 if AIOHTTP_POOL_CONNECTIONS_PER_HOST is not None: 66 ↛ 67line 66 didn't jump to line 67 because the condition on line 66 was never true
67 connector_kwargs['limit_per_host'] = AIOHTTP_POOL_CONNECTIONS_PER_HOST
68 else:
69 connector_kwargs['limit_per_host'] = 0 # aiohttp: 0 = unlimited
70 connector = aiohttp.TCPConnector(**connector_kwargs)
71 timeout = get_client_timeout()
72 _session = aiohttp.ClientSession(
73 connector=connector,
74 timeout=timeout,
75 trust_env=True,
76 )
77 log.info(
78 'Created shared aiohttp session pool (limit=%s, per_host=%s, dns_ttl=%d)',
79 AIOHTTP_POOL_CONNECTIONS or 'unlimited',
80 AIOHTTP_POOL_CONNECTIONS_PER_HOST or 'unlimited',
81 AIOHTTP_POOL_DNS_TTL,
82 )
83 return _session
86async def close_session():
87 """Close the shared session. Called during application shutdown."""
88 global _session
89 if _session and not _session.closed: 89 ↛ exitline 89 didn't return from function 'close_session' because the condition on line 89 was always true
90 await _session.close()
91 log.info('Closed shared aiohttp session pool')
92 _session = None
95async def cleanup_response(
96 response: Optional[aiohttp.ClientResponse],
97 session: Optional[aiohttp.ClientSession] = None,
98):
99 """Release and close an aiohttp response, optionally closing the session.
101 When using the shared pool, ``session`` should be ``None`` (the pool
102 session is never closed per-request). When a caller creates its own
103 one-off session, pass it here to close it after the response.
104 """
105 if response: 105 ↛ 106line 105 didn't jump to line 106 because the condition on line 105 was never true
106 if not response.closed:
107 # aiohttp 3.9+ made ClientResponse.close() synchronous (returns None).
108 # Older versions returned a coroutine. Handle both gracefully.
109 result = response.close()
110 if result is not None:
111 await result
112 if session: 112 ↛ 113line 112 didn't jump to line 113 because the condition on line 112 was never true
113 if not session.closed:
114 result = session.close()
115 if result is not None:
116 await result
119async def stream_wrapper(response, session=None, passthrough=False):
120 """Wrap a stream to ensure cleanup happens even if streaming is interrupted.
122 This is more reliable than BackgroundTask which may not run if the client
123 disconnects. When using the shared pool, ``session`` should be ``None``.
125 ``passthrough=True`` yields raw network chunks (iter_any) instead of
126 lines: byte-identical output without a buffer scan, slice and copy per
127 line. Only for streams no internal consumer parses line-by-line.
128 """
129 try:
130 if passthrough:
131 stream = response.content.iter_any()
132 else:
133 stream = stream_chunks_handler(response.content)
134 async for chunk in stream:
135 yield chunk
136 finally:
137 await cleanup_response(response, session)