Coverage for /home/admin/Documents/AI/applications/lexigram-dev/lexigram/experimental/ai/lexigram-ai-llm/src/lexigram/ai/llm/http/client.py: 33%

99 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-08-25 07:19 +0800

1"""Thin async HTTP client for AI model API communication. 

2 

3Wraps ``aiohttp.ClientSession`` with base-URL resolution and default-header 

4merging. Retry and circuit-breaker concerns belong to the caller or a 

5dedicated resilience layer (``lexigram-resilience``). 

6""" 

7 

8from __future__ import annotations 

9 

10from contextlib import asynccontextmanager 

11from typing import TYPE_CHECKING, Any, Self, cast 

12 

13import aiohttp 

14 

15from lexigram.contracts.web.http_models import HttpResponse, HttpStatusError 

16from lexigram.logging import ( 

17 get_logger, 

18) 

19 

20if TYPE_CHECKING: 

21 from collections.abc import AsyncIterator 

22 

23 from lexigram.contracts.infra.resilience.protocols import ( 

24 CircuitBreakerProtocol, 

25 RetryPolicyProtocol, 

26 ) 

27 

28logger = get_logger(__name__) 

29 

30 

31class _StreamContext: 

32 """Line-oriented streaming response context for SSE/NDJSON streams. 

33 

34 Wraps an :class:`aiohttp.ClientResponse` and exposes :attr:`status`, 

35 :meth:`raise_for_status`, and :meth:`aiter_lines` without leaking the 

36 underlying aiohttp type to callers. 

37 

38 Args: 

39 resp: The raw aiohttp ClientResponse, owned by the enclosing 

40 context manager for the duration of the stream. 

41 """ 

42 

43 __slots__ = ("_resp",) 

44 

45 def __init__(self, resp: aiohttp.ClientResponse) -> None: 

46 self._resp = resp 

47 

48 @property 

49 def status(self) -> int: 

50 """HTTP status code of the response.""" 

51 return self._resp.status 

52 

53 def raise_for_status(self) -> None: 

54 """Raise :class:`~lexigram.contracts.web.models.HttpStatusError` for 4xx/5xx. 

55 

56 Raises: 

57 HttpStatusError: When the status code is 400 or higher. 

58 """ 

59 if self._resp.status >= 400: 

60 raise HttpStatusError( 

61 f"HTTP {self._resp.status}", 

62 status=self._resp.status, 

63 response=HttpResponse( 

64 status=self._resp.status, 

65 url=str(self._resp.url), 

66 method=self._resp.method.upper(), 

67 ), 

68 ) 

69 

70 async def aiter_lines(self) -> AsyncIterator[str]: 

71 """Yield decoded text lines from the streaming response body. 

72 

73 Yields: 

74 Each line decoded from UTF-8 with trailing CR/LF stripped. 

75 """ 

76 async for raw in self._resp.content: 

77 yield raw.decode("utf-8").rstrip("\r\n") 

78 

79 

80class ResilientHTTPClient: 

81 """Thin async HTTP client for LLM provider communication. 

82 

83 Wraps ``aiohttp.ClientSession`` with base-URL resolution and 

84 default-header merging. Optional *retry* and *circuit_breaker* instances 

85 from ``lexigram-resilience`` can be injected to add retry and circuit-breaker 

86 protection without coupling this class to a concrete implementation. 

87 

88 Args: 

89 base_url: Base URL prepended to every relative path. 

90 headers: Default request headers merged with per-request headers. 

91 timeout: Total request timeout in seconds. 

92 name: Logical name used in log messages. 

93 retry: Optional retry policy applied to every request. 

94 circuit_breaker: Optional circuit breaker wrapping every request. 

95 When *retry* is also set, the retry policy wraps the circuit breaker. 

96 """ 

97 

98 def __init__( 

99 self, 

100 base_url: str = "", 

101 headers: dict[str, str] | None = None, 

102 timeout: float = 300.0, 

103 name: str = "ai-http-client", 

104 retry: RetryPolicyProtocol | None = None, 

105 circuit_breaker: CircuitBreakerProtocol | None = None, 

106 ) -> None: 

107 self.base_url = base_url.rstrip("/") if base_url else "" 

108 self.headers = headers or {} 

109 self.timeout = timeout 

110 self.name = name 

111 self._retry = retry 

112 self._circuit_breaker = circuit_breaker 

113 self._session: aiohttp.ClientSession | None = None 

114 

115 async def _get_session(self) -> aiohttp.ClientSession: 

116 """Return the active session, creating it lazily on first use.""" 

117 if self._session is None: 

118 self._session = aiohttp.ClientSession( 

119 timeout=aiohttp.ClientTimeout(total=self.timeout), 

120 ) 

121 return self._session 

122 

123 async def __aenter__(self) -> Self: 

124 return self 

125 

126 async def __aexit__( 

127 self, 

128 exc_type: type[BaseException] | None, 

129 exc_val: BaseException | None, 

130 exc_tb: object, 

131 ) -> None: 

132 await self.close() 

133 

134 def _build_url(self, path: str) -> str: 

135 """Resolve *path* against :attr:`base_url`.""" 

136 if not path: 

137 return self.base_url 

138 if path.startswith(("http://", "https://")): 

139 return path 

140 if not path.startswith("/"): 

141 path = f"/{path}" 

142 return f"{self.base_url}{path}" 

143 

144 def _merge_headers(self, extra: dict[str, str] | None = None) -> dict[str, str]: 

145 """Merge *extra* headers with the client default headers.""" 

146 merged = self.headers.copy() 

147 if extra: 

148 merged.update(extra) 

149 return merged 

150 

151 async def _raw_request( 

152 self, 

153 method: str, 

154 url: str, 

155 **kwargs: Any, 

156 ) -> HttpResponse: 

157 """Execute a raw HTTP request without any resilience wrapping.""" 

158 full_url = self._build_url(url) 

159 headers = self._merge_headers(kwargs.pop("headers", None)) 

160 session = await self._get_session() 

161 async with session.request(method, full_url, headers=headers, **kwargs) as resp: 

162 try: 

163 json_data: Any = await resp.json(content_type=None) 

164 except (ValueError, UnicodeDecodeError): 

165 json_data = None 

166 body = await resp.read() 

167 resp_headers = {str(k): str(v) for k, v in resp.headers.items()} 

168 return HttpResponse( 

169 status=resp.status, 

170 headers=resp_headers, 

171 body=body, 

172 text=body.decode("utf-8", errors="replace"), 

173 json=json_data, 

174 url=str(resp.url), 

175 method=method.upper(), 

176 ) 

177 

178 async def request( 

179 self, 

180 method: str, 

181 url: str, 

182 **kwargs: Any, 

183 ) -> HttpResponse: 

184 """Execute an HTTP request and return an :class:`~lexigram.contracts.web.models.HttpResponse`. 

185 

186 When a *retry* policy is configured it wraps the innermost call. 

187 When a *circuit_breaker* is configured it is called inside the retry 

188 boundary. Both, either, or neither may be active. 

189 """ 

190 if self._retry is not None: 

191 # Retry wraps the circuit-breaker-protected call (or raw call) 

192 if self._circuit_breaker is not None: 

193 return cast( 

194 "HttpResponse", 

195 await self._retry.execute( 

196 self._circuit_breaker.call, 

197 self._raw_request, 

198 method, 

199 url, 

200 **kwargs, 

201 ), 

202 ) 

203 return cast( 

204 "HttpResponse", 

205 await self._retry.execute(self._raw_request, method, url, **kwargs), 

206 ) 

207 

208 if self._circuit_breaker is not None: 

209 return cast( 

210 "HttpResponse", 

211 await self._circuit_breaker.call( 

212 self._raw_request, method, url, **kwargs 

213 ), 

214 ) 

215 

216 return await self._raw_request(method, url, **kwargs) 

217 

218 async def get(self, url: str, **kwargs: Any) -> HttpResponse: 

219 """Perform a GET request.""" 

220 return await self.request("GET", url, **kwargs) 

221 

222 async def post(self, url: str, data: Any = None, **kwargs: Any) -> HttpResponse: 

223 """Perform a POST request.""" 

224 if isinstance(data, dict): 

225 kwargs["json"] = data 

226 elif data is not None: 

227 kwargs["data"] = data 

228 return await self.request("POST", url, **kwargs) 

229 

230 async def put(self, url: str, **kwargs: Any) -> HttpResponse: 

231 """Perform a PUT request.""" 

232 return await self.request("PUT", url, **kwargs) 

233 

234 async def delete(self, url: str, **kwargs: Any) -> HttpResponse: 

235 """Perform a DELETE request.""" 

236 return await self.request("DELETE", url, **kwargs) 

237 

238 async def patch(self, url: str, **kwargs: Any) -> HttpResponse: 

239 """Perform a PATCH request.""" 

240 return await self.request("PATCH", url, **kwargs) 

241 

242 async def head(self, url: str, **kwargs: Any) -> HttpResponse: 

243 """Perform a HEAD request.""" 

244 return await self.request("HEAD", url, **kwargs) 

245 

246 @asynccontextmanager 

247 async def stream( 

248 self, 

249 method: str, 

250 url: str, 

251 **kwargs: Any, 

252 ) -> AsyncIterator[_StreamContext]: 

253 """Async context manager for line-oriented streaming responses. 

254 

255 Opens a persistent connection and wraps the response in a 

256 :class:`_StreamContext` that exposes :meth:`~_StreamContext.aiter_lines` 

257 for consuming SSE / NDJSON streams — the format used by LLM provider 

258 streaming APIs. 

259 

260 Args: 

261 method: HTTP method (typically ``"POST"`` for LLM completions). 

262 url: Relative path or absolute URL. 

263 **kwargs: Forwarded to aiohttp (``json=``, ``data=``, 

264 ``headers=``, …). 

265 

266 Yields: 

267 :class:`_StreamContext` with :attr:`~_StreamContext.status`, 

268 :meth:`~_StreamContext.raise_for_status`, and 

269 :meth:`~_StreamContext.aiter_lines`. 

270 """ 

271 full_url = self._build_url(url) 

272 headers = self._merge_headers(kwargs.pop("headers", None)) 

273 session = await self._get_session() 

274 async with session.request(method, full_url, headers=headers, **kwargs) as resp: 

275 yield _StreamContext(resp) 

276 

277 async def close(self) -> None: 

278 """Close the underlying ``aiohttp`` session.""" 

279 if self._session is not None: 

280 await self._session.close() 

281 self._session = None 

282 

283 async def stop(self) -> None: 

284 """Alias for :meth:`close`.""" 

285 await self.close()