Edit on GitHub

agent_search_gateway.providers.http

Shared HTTP execution boundary for provider adapters.

  1"""Shared HTTP execution boundary for provider adapters."""
  2
  3import asyncio
  4import logging
  5import time
  6from collections.abc import Awaitable, Callable, Mapping
  7from typing import Any
  8
  9import httpx
 10
 11from ..errors import ErrorCode, ExecutionFailure, ProtocolFailure
 12from ..models import RetryPolicy
 13from ..observability import elapsed_ms, http_endpoint_for_log, log_event
 14from ..retry import retry_async
 15
 16
 17class _RetryableStatus(Exception):
 18    def __init__(self, status_code: int) -> None:
 19        super().__init__(str(status_code))
 20        self.status_code = status_code
 21
 22
 23class HttpStatusFailure(ExecutionFailure):
 24    """Internal HTTP failure with a machine-readable terminal status code."""
 25
 26    def __init__(self, provider_name: str, stage: str, status_code: int) -> None:
 27        super().__init__(
 28            ErrorCode.ALL_PROVIDERS_FAILED,
 29            f"{provider_name}/{stage}: HTTP status {status_code}",
 30        )
 31        self.status_code = status_code
 32
 33
 34class HttpJsonExecutor:
 35    """Execute JSON or text HTTP requests with one retry and logging policy."""
 36
 37    def __init__(
 38        self,
 39        client: httpx.AsyncClient,
 40        retry_policy: RetryPolicy,
 41        *,
 42        provider_name: str,
 43        logger: logging.Logger | None = None,
 44        sleep: Callable[[float], Awaitable[None]] = asyncio.sleep,
 45        monotonic: Callable[[], float] = time.monotonic,
 46    ) -> None:
 47        self._client = client
 48        self._retry_policy = retry_policy
 49        self._provider_name = provider_name
 50        self._logger = logger or logging.getLogger(__name__)
 51        self._sleep = sleep
 52        self._monotonic = monotonic
 53
 54    async def request_json(
 55        self,
 56        method: str,
 57        url: str,
 58        *,
 59        stage: str,
 60        headers: Mapping[str, str] | None = None,
 61        params: Mapping[str, Any] | None = None,
 62        json_body: object | None = None,
 63    ) -> object:
 64        response, attempt, attempt_started = await self._request_response(
 65            method,
 66            url,
 67            stage=stage,
 68            headers=headers,
 69            params=params,
 70            json_body=json_body,
 71        )
 72        try:
 73            return response.json()
 74        except ValueError as exc:
 75            self._log_failed(
 76                stage,
 77                http_endpoint_for_log(url),
 78                attempt,
 79                attempt_started,
 80                "decode",
 81            )
 82            raise ProtocolFailure(
 83                ErrorCode.PROTOCOL_ERROR,
 84                f"{self._provider_name}/{stage}: response was not valid JSON",
 85            ) from exc
 86
 87    async def request_text(
 88        self,
 89        method: str,
 90        url: str,
 91        *,
 92        stage: str,
 93        headers: Mapping[str, str] | None = None,
 94        params: Mapping[str, Any] | None = None,
 95        json_body: object | None = None,
 96    ) -> str:
 97        response, _, _ = await self._request_response(
 98            method,
 99            url,
100            stage=stage,
101            headers=headers,
102            params=params,
103            json_body=json_body,
104        )
105        return response.text
106
107    async def _request_response(
108        self,
109        method: str,
110        url: str,
111        *,
112        stage: str,
113        headers: Mapping[str, str] | None,
114        params: Mapping[str, Any] | None,
115        json_body: object | None,
116    ) -> tuple[httpx.Response, int, float]:
117        attempt = 0
118        attempt_started = self._monotonic()
119        log_endpoint = http_endpoint_for_log(url)
120
121        def before_attempt(current_attempt: int) -> None:
122            nonlocal attempt, attempt_started
123            attempt = current_attempt
124            attempt_started = self._monotonic()
125            log_event(
126                self._logger,
127                logging.DEBUG,
128                "http_attempt_started",
129                provider=self._provider_name,
130                stage=stage,
131                endpoint=log_endpoint,
132                attempt=attempt,
133            )
134
135        def on_retry(current_attempt: int, exc: BaseException, delay: float) -> None:
136            delay_ms = max(0, int(delay * 1000))
137            attempt_elapsed_ms = elapsed_ms(self._monotonic, attempt_started)
138            if isinstance(exc, _RetryableStatus):
139                log_event(
140                    self._logger,
141                    logging.WARNING,
142                    "http_retrying",
143                    provider=self._provider_name,
144                    stage=stage,
145                    endpoint=log_endpoint,
146                    attempt=current_attempt,
147                    delay_ms=delay_ms,
148                    elapsed_ms=attempt_elapsed_ms,
149                    category="status",
150                    status=exc.status_code,
151                )
152                return
153            log_event(
154                self._logger,
155                logging.WARNING,
156                "http_retrying",
157                provider=self._provider_name,
158                stage=stage,
159                endpoint=log_endpoint,
160                attempt=current_attempt,
161                delay_ms=delay_ms,
162                elapsed_ms=attempt_elapsed_ms,
163                category="transport",
164            )
165
166        async def operation() -> httpx.Response:
167            response = await self._client.request(
168                method,
169                url,
170                headers=headers,
171                params=params,
172                json=json_body,
173                timeout=self._retry_policy.request_timeout_seconds,
174            )
175            log_event(
176                self._logger,
177                logging.DEBUG,
178                "http_attempt_completed",
179                provider=self._provider_name,
180                stage=stage,
181                endpoint=log_endpoint,
182                attempt=attempt,
183                status=response.status_code,
184                elapsed_ms=elapsed_ms(self._monotonic, attempt_started),
185            )
186            if response.status_code in {408, 429} or response.status_code >= 500:
187                raise _RetryableStatus(response.status_code)
188            return response
189
190        try:
191            response = await retry_async(
192                self._retry_policy,
193                operation,
194                is_retryable=self._is_retryable,
195                sleep=self._sleep,
196                before_attempt=before_attempt,
197                on_retry=on_retry,
198            )
199        except _RetryableStatus as exc:
200            self._log_failed(
201                stage,
202                log_endpoint,
203                attempt,
204                attempt_started,
205                "status",
206                status=exc.status_code,
207            )
208            raise self._status_failure(stage, exc.status_code) from exc
209        except (httpx.TimeoutException, httpx.TransportError) as exc:
210            self._log_failed(stage, log_endpoint, attempt, attempt_started, "transport")
211            raise self._execution_failure(stage, "HTTP transport failure") from exc
212
213        if response.status_code < 200 or response.status_code >= 300:
214            self._log_failed(
215                stage,
216                log_endpoint,
217                attempt,
218                attempt_started,
219                "status",
220                status=response.status_code,
221            )
222            raise self._status_failure(stage, response.status_code)
223        return response, attempt, attempt_started
224
225    def _log_failed(
226        self,
227        stage: str,
228        endpoint: str,
229        attempt: int,
230        started: float,
231        category: str,
232        *,
233        status: int | None = None,
234    ) -> None:
235        attempt_elapsed_ms = elapsed_ms(self._monotonic, started)
236        if status is None:
237            log_event(
238                self._logger,
239                logging.DEBUG,
240                "http_failed",
241                provider=self._provider_name,
242                stage=stage,
243                endpoint=endpoint,
244                attempt=attempt,
245                category=category,
246                elapsed_ms=attempt_elapsed_ms,
247            )
248            return
249        log_event(
250            self._logger,
251            logging.DEBUG,
252            "http_failed",
253            provider=self._provider_name,
254            stage=stage,
255            endpoint=endpoint,
256            attempt=attempt,
257            category=category,
258            elapsed_ms=attempt_elapsed_ms,
259            status=status,
260        )
261
262    async def aclose(self) -> None:
263        await self._client.aclose()
264
265    @staticmethod
266    def _is_retryable(exc: BaseException) -> bool:
267        return isinstance(
268            exc,
269            _RetryableStatus | httpx.TimeoutException | httpx.TransportError,
270        )
271
272    def _status_failure(self, stage: str, status_code: int) -> HttpStatusFailure:
273        return HttpStatusFailure(self._provider_name, stage, status_code)
274
275    def _execution_failure(self, stage: str, reason: str) -> ExecutionFailure:
276        return ExecutionFailure(
277            ErrorCode.ALL_PROVIDERS_FAILED,
278            f"{self._provider_name}/{stage}: {reason}",
279        )
class HttpStatusFailure(agent_search_gateway.errors.ExecutionFailure):
24class HttpStatusFailure(ExecutionFailure):
25    """Internal HTTP failure with a machine-readable terminal status code."""
26
27    def __init__(self, provider_name: str, stage: str, status_code: int) -> None:
28        super().__init__(
29            ErrorCode.ALL_PROVIDERS_FAILED,
30            f"{provider_name}/{stage}: HTTP status {status_code}",
31        )
32        self.status_code = status_code

Internal HTTP failure with a machine-readable terminal status code.

HttpStatusFailure(provider_name: str, stage: str, status_code: int)
27    def __init__(self, provider_name: str, stage: str, status_code: int) -> None:
28        super().__init__(
29            ErrorCode.ALL_PROVIDERS_FAILED,
30            f"{provider_name}/{stage}: HTTP status {status_code}",
31        )
32        self.status_code = status_code
status_code
class HttpJsonExecutor:
 35class HttpJsonExecutor:
 36    """Execute JSON or text HTTP requests with one retry and logging policy."""
 37
 38    def __init__(
 39        self,
 40        client: httpx.AsyncClient,
 41        retry_policy: RetryPolicy,
 42        *,
 43        provider_name: str,
 44        logger: logging.Logger | None = None,
 45        sleep: Callable[[float], Awaitable[None]] = asyncio.sleep,
 46        monotonic: Callable[[], float] = time.monotonic,
 47    ) -> None:
 48        self._client = client
 49        self._retry_policy = retry_policy
 50        self._provider_name = provider_name
 51        self._logger = logger or logging.getLogger(__name__)
 52        self._sleep = sleep
 53        self._monotonic = monotonic
 54
 55    async def request_json(
 56        self,
 57        method: str,
 58        url: str,
 59        *,
 60        stage: str,
 61        headers: Mapping[str, str] | None = None,
 62        params: Mapping[str, Any] | None = None,
 63        json_body: object | None = None,
 64    ) -> object:
 65        response, attempt, attempt_started = await self._request_response(
 66            method,
 67            url,
 68            stage=stage,
 69            headers=headers,
 70            params=params,
 71            json_body=json_body,
 72        )
 73        try:
 74            return response.json()
 75        except ValueError as exc:
 76            self._log_failed(
 77                stage,
 78                http_endpoint_for_log(url),
 79                attempt,
 80                attempt_started,
 81                "decode",
 82            )
 83            raise ProtocolFailure(
 84                ErrorCode.PROTOCOL_ERROR,
 85                f"{self._provider_name}/{stage}: response was not valid JSON",
 86            ) from exc
 87
 88    async def request_text(
 89        self,
 90        method: str,
 91        url: str,
 92        *,
 93        stage: str,
 94        headers: Mapping[str, str] | None = None,
 95        params: Mapping[str, Any] | None = None,
 96        json_body: object | None = None,
 97    ) -> str:
 98        response, _, _ = await self._request_response(
 99            method,
100            url,
101            stage=stage,
102            headers=headers,
103            params=params,
104            json_body=json_body,
105        )
106        return response.text
107
108    async def _request_response(
109        self,
110        method: str,
111        url: str,
112        *,
113        stage: str,
114        headers: Mapping[str, str] | None,
115        params: Mapping[str, Any] | None,
116        json_body: object | None,
117    ) -> tuple[httpx.Response, int, float]:
118        attempt = 0
119        attempt_started = self._monotonic()
120        log_endpoint = http_endpoint_for_log(url)
121
122        def before_attempt(current_attempt: int) -> None:
123            nonlocal attempt, attempt_started
124            attempt = current_attempt
125            attempt_started = self._monotonic()
126            log_event(
127                self._logger,
128                logging.DEBUG,
129                "http_attempt_started",
130                provider=self._provider_name,
131                stage=stage,
132                endpoint=log_endpoint,
133                attempt=attempt,
134            )
135
136        def on_retry(current_attempt: int, exc: BaseException, delay: float) -> None:
137            delay_ms = max(0, int(delay * 1000))
138            attempt_elapsed_ms = elapsed_ms(self._monotonic, attempt_started)
139            if isinstance(exc, _RetryableStatus):
140                log_event(
141                    self._logger,
142                    logging.WARNING,
143                    "http_retrying",
144                    provider=self._provider_name,
145                    stage=stage,
146                    endpoint=log_endpoint,
147                    attempt=current_attempt,
148                    delay_ms=delay_ms,
149                    elapsed_ms=attempt_elapsed_ms,
150                    category="status",
151                    status=exc.status_code,
152                )
153                return
154            log_event(
155                self._logger,
156                logging.WARNING,
157                "http_retrying",
158                provider=self._provider_name,
159                stage=stage,
160                endpoint=log_endpoint,
161                attempt=current_attempt,
162                delay_ms=delay_ms,
163                elapsed_ms=attempt_elapsed_ms,
164                category="transport",
165            )
166
167        async def operation() -> httpx.Response:
168            response = await self._client.request(
169                method,
170                url,
171                headers=headers,
172                params=params,
173                json=json_body,
174                timeout=self._retry_policy.request_timeout_seconds,
175            )
176            log_event(
177                self._logger,
178                logging.DEBUG,
179                "http_attempt_completed",
180                provider=self._provider_name,
181                stage=stage,
182                endpoint=log_endpoint,
183                attempt=attempt,
184                status=response.status_code,
185                elapsed_ms=elapsed_ms(self._monotonic, attempt_started),
186            )
187            if response.status_code in {408, 429} or response.status_code >= 500:
188                raise _RetryableStatus(response.status_code)
189            return response
190
191        try:
192            response = await retry_async(
193                self._retry_policy,
194                operation,
195                is_retryable=self._is_retryable,
196                sleep=self._sleep,
197                before_attempt=before_attempt,
198                on_retry=on_retry,
199            )
200        except _RetryableStatus as exc:
201            self._log_failed(
202                stage,
203                log_endpoint,
204                attempt,
205                attempt_started,
206                "status",
207                status=exc.status_code,
208            )
209            raise self._status_failure(stage, exc.status_code) from exc
210        except (httpx.TimeoutException, httpx.TransportError) as exc:
211            self._log_failed(stage, log_endpoint, attempt, attempt_started, "transport")
212            raise self._execution_failure(stage, "HTTP transport failure") from exc
213
214        if response.status_code < 200 or response.status_code >= 300:
215            self._log_failed(
216                stage,
217                log_endpoint,
218                attempt,
219                attempt_started,
220                "status",
221                status=response.status_code,
222            )
223            raise self._status_failure(stage, response.status_code)
224        return response, attempt, attempt_started
225
226    def _log_failed(
227        self,
228        stage: str,
229        endpoint: str,
230        attempt: int,
231        started: float,
232        category: str,
233        *,
234        status: int | None = None,
235    ) -> None:
236        attempt_elapsed_ms = elapsed_ms(self._monotonic, started)
237        if status is None:
238            log_event(
239                self._logger,
240                logging.DEBUG,
241                "http_failed",
242                provider=self._provider_name,
243                stage=stage,
244                endpoint=endpoint,
245                attempt=attempt,
246                category=category,
247                elapsed_ms=attempt_elapsed_ms,
248            )
249            return
250        log_event(
251            self._logger,
252            logging.DEBUG,
253            "http_failed",
254            provider=self._provider_name,
255            stage=stage,
256            endpoint=endpoint,
257            attempt=attempt,
258            category=category,
259            elapsed_ms=attempt_elapsed_ms,
260            status=status,
261        )
262
263    async def aclose(self) -> None:
264        await self._client.aclose()
265
266    @staticmethod
267    def _is_retryable(exc: BaseException) -> bool:
268        return isinstance(
269            exc,
270            _RetryableStatus | httpx.TimeoutException | httpx.TransportError,
271        )
272
273    def _status_failure(self, stage: str, status_code: int) -> HttpStatusFailure:
274        return HttpStatusFailure(self._provider_name, stage, status_code)
275
276    def _execution_failure(self, stage: str, reason: str) -> ExecutionFailure:
277        return ExecutionFailure(
278            ErrorCode.ALL_PROVIDERS_FAILED,
279            f"{self._provider_name}/{stage}: {reason}",
280        )

Execute JSON or text HTTP requests with one retry and logging policy.

HttpJsonExecutor( client: httpx.AsyncClient, retry_policy: agent_search_gateway.models.RetryPolicy, *, provider_name: str, logger: logging.Logger | None = None, sleep: Callable[[float], Awaitable[None]] = <function sleep>, monotonic: Callable[[], float] = <built-in function monotonic>)
38    def __init__(
39        self,
40        client: httpx.AsyncClient,
41        retry_policy: RetryPolicy,
42        *,
43        provider_name: str,
44        logger: logging.Logger | None = None,
45        sleep: Callable[[float], Awaitable[None]] = asyncio.sleep,
46        monotonic: Callable[[], float] = time.monotonic,
47    ) -> None:
48        self._client = client
49        self._retry_policy = retry_policy
50        self._provider_name = provider_name
51        self._logger = logger or logging.getLogger(__name__)
52        self._sleep = sleep
53        self._monotonic = monotonic
async def request_json( self, method: str, url: str, *, stage: str, headers: Mapping[str, str] | None = None, params: Mapping[str, typing.Any] | None = None, json_body: object | None = None) -> object:
55    async def request_json(
56        self,
57        method: str,
58        url: str,
59        *,
60        stage: str,
61        headers: Mapping[str, str] | None = None,
62        params: Mapping[str, Any] | None = None,
63        json_body: object | None = None,
64    ) -> object:
65        response, attempt, attempt_started = await self._request_response(
66            method,
67            url,
68            stage=stage,
69            headers=headers,
70            params=params,
71            json_body=json_body,
72        )
73        try:
74            return response.json()
75        except ValueError as exc:
76            self._log_failed(
77                stage,
78                http_endpoint_for_log(url),
79                attempt,
80                attempt_started,
81                "decode",
82            )
83            raise ProtocolFailure(
84                ErrorCode.PROTOCOL_ERROR,
85                f"{self._provider_name}/{stage}: response was not valid JSON",
86            ) from exc
async def request_text( self, method: str, url: str, *, stage: str, headers: Mapping[str, str] | None = None, params: Mapping[str, typing.Any] | None = None, json_body: object | None = None) -> str:
 88    async def request_text(
 89        self,
 90        method: str,
 91        url: str,
 92        *,
 93        stage: str,
 94        headers: Mapping[str, str] | None = None,
 95        params: Mapping[str, Any] | None = None,
 96        json_body: object | None = None,
 97    ) -> str:
 98        response, _, _ = await self._request_response(
 99            method,
100            url,
101            stage=stage,
102            headers=headers,
103            params=params,
104            json_body=json_body,
105        )
106        return response.text
async def aclose(self) -> None:
263    async def aclose(self) -> None:
264        await self._client.aclose()