|
21 | 21 | import aiohttp |
22 | 22 | import aiohttp_retry |
23 | 23 |
|
| 24 | +from kubernetes.aio.client._retry import ( |
| 25 | + is_retry_after_response, |
| 26 | + on_retry_after_error, |
| 27 | + retry_after_backoff, |
| 28 | +) |
24 | 29 | from kubernetes.aio.client.exceptions import ApiException, ApiValueError |
25 | 30 |
|
26 | 31 | RESTResponseType = aiohttp.ClientResponse |
@@ -284,14 +289,55 @@ async def request( |
284 | 289 | self.pool_manager = self._create_pool_manager() |
285 | 290 | pool_manager = self.pool_manager |
286 | 291 |
|
287 | | - if self._effective_retry_options is not None and method in ALLOW_RETRY_METHODS: |
| 292 | + client_go_read_retries = ( |
| 293 | + method in ['GET', 'HEAD'] |
| 294 | + and getattr(self.configuration, 'client_go_retries', False) |
| 295 | + ) |
| 296 | + |
| 297 | + if ( |
| 298 | + self._effective_retry_options is not None |
| 299 | + and method in ALLOW_RETRY_METHODS |
| 300 | + and not client_go_read_retries |
| 301 | + ): |
288 | 302 | if self.retry_client is None: |
289 | 303 | self.retry_client = aiohttp_retry.RetryClient( |
290 | 304 | client_session=self.pool_manager, |
291 | 305 | retry_options=self._effective_retry_options |
292 | 306 | ) |
293 | 307 | pool_manager = self.retry_client |
294 | 308 |
|
295 | | - r = await pool_manager.request(**args) |
| 309 | + async def read_request(check_retry_status=False): |
| 310 | + response = await self.pool_manager.request(**args) |
| 311 | + if check_retry_status: |
| 312 | + self._raise_retry_after_response(response) |
| 313 | + return response |
| 314 | + |
| 315 | + if client_go_read_retries: |
| 316 | + backoff = retry_after_backoff( |
| 317 | + getattr(self.configuration, 'retries', None), |
| 318 | + getattr(self.configuration, 'client_go_retry_backoff', None), |
| 319 | + ) |
| 320 | + r = await on_retry_after_error( |
| 321 | + backoff, self._is_read_retryable, lambda: read_request(True)) |
| 322 | + else: |
| 323 | + r = await pool_manager.request(**args) |
296 | 324 |
|
297 | 325 | return RESTResponse(r) |
| 326 | + |
| 327 | + @classmethod |
| 328 | + def _is_read_retryable(cls, error): |
| 329 | + return is_retry_after_response(error) |
| 330 | + |
| 331 | + @staticmethod |
| 332 | + def _retry_after_error(response): |
| 333 | + error = ApiException(status=response.status, reason=response.reason) |
| 334 | + error.headers = response.headers |
| 335 | + return error |
| 336 | + |
| 337 | + @classmethod |
| 338 | + def _raise_retry_after_response(cls, response): |
| 339 | + error = cls._retry_after_error(response) |
| 340 | + if not is_retry_after_response(error): |
| 341 | + return |
| 342 | + response.release() |
| 343 | + raise error |
0 commit comments