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
« 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.
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"""
8from __future__ import annotations
10from contextlib import asynccontextmanager
11from typing import TYPE_CHECKING, Any, Self, cast
13import aiohttp
15from lexigram.contracts.web.http_models import HttpResponse, HttpStatusError
16from lexigram.logging import (
17 get_logger,
18)
20if TYPE_CHECKING:
21 from collections.abc import AsyncIterator
23 from lexigram.contracts.infra.resilience.protocols import (
24 CircuitBreakerProtocol,
25 RetryPolicyProtocol,
26 )
28logger = get_logger(__name__)
31class _StreamContext:
32 """Line-oriented streaming response context for SSE/NDJSON streams.
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.
38 Args:
39 resp: The raw aiohttp ClientResponse, owned by the enclosing
40 context manager for the duration of the stream.
41 """
43 __slots__ = ("_resp",)
45 def __init__(self, resp: aiohttp.ClientResponse) -> None:
46 self._resp = resp
48 @property
49 def status(self) -> int:
50 """HTTP status code of the response."""
51 return self._resp.status
53 def raise_for_status(self) -> None:
54 """Raise :class:`~lexigram.contracts.web.models.HttpStatusError` for 4xx/5xx.
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 )
70 async def aiter_lines(self) -> AsyncIterator[str]:
71 """Yield decoded text lines from the streaming response body.
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")
80class ResilientHTTPClient:
81 """Thin async HTTP client for LLM provider communication.
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.
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 """
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
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
123 async def __aenter__(self) -> Self:
124 return self
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()
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}"
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
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 )
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`.
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 )
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 )
216 return await self._raw_request(method, url, **kwargs)
218 async def get(self, url: str, **kwargs: Any) -> HttpResponse:
219 """Perform a GET request."""
220 return await self.request("GET", url, **kwargs)
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)
230 async def put(self, url: str, **kwargs: Any) -> HttpResponse:
231 """Perform a PUT request."""
232 return await self.request("PUT", url, **kwargs)
234 async def delete(self, url: str, **kwargs: Any) -> HttpResponse:
235 """Perform a DELETE request."""
236 return await self.request("DELETE", url, **kwargs)
238 async def patch(self, url: str, **kwargs: Any) -> HttpResponse:
239 """Perform a PATCH request."""
240 return await self.request("PATCH", url, **kwargs)
242 async def head(self, url: str, **kwargs: Any) -> HttpResponse:
243 """Perform a HEAD request."""
244 return await self.request("HEAD", url, **kwargs)
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.
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.
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=``, …).
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)
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
283 async def stop(self) -> None:
284 """Alias for :meth:`close`."""
285 await self.close()