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

1"""Shared aiohttp ClientSession pool. 

2 

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. 

7 

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 

12 

13Usage: 

14 from open_webui.utils.session_pool import get_session, cleanup_response 

15 

16 session = await get_session() 

17 r = await session.request(...) 

18 # When done with the *response* (not the session): 

19 await cleanup_response(r) 

20 

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""" 

25 

26import logging 

27from typing import Optional 

28 

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 

38 

39log = logging.getLogger(__name__) 

40 

41_session: Optional[aiohttp.ClientSession] = None 

42 

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) 

48 

49 

50def get_client_timeout(stream: bool = False) -> aiohttp.ClientTimeout: 

51 return _CLIENT_STREAM_TIMEOUT if stream else _CLIENT_TIMEOUT 

52 

53 

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 

84 

85 

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 

93 

94 

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. 

100 

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 

117 

118 

119async def stream_wrapper(response, session=None, passthrough=False): 

120 """Wrap a stream to ensure cleanup happens even if streaming is interrupted. 

121 

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``. 

124 

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)