Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 16 additions & 0 deletions kubernetes/aio/client/configuration.py
Original file line number Diff line number Diff line change
Expand Up @@ -355,6 +355,22 @@ def __init__(
self.retries = retries
"""Retry configuration
"""
self.client_go_retries = False
"""Enable Kubernetes client-go-compatible retry semantics.

When enabled, GET and HEAD requests retry Retry-After responses.
The retry ceiling is read from ``retries`` when set; otherwise it
follows the client-go default of at most 10 retries.
"""
self.client_go_retry_backoff = None
"""Backoff for Kubernetes client-go-compatible retries.

If unset, client-go-compatible GET and HEAD retries use the
client-go default retry ceiling with no additional client-side
delay beyond Retry-After. When set, ``retries`` still overrides
the retry ceiling if it is not None.
"""

self.trace_configs = trace_configs
"""aiohttp.TraceConfig list forwarded to ClientSession for tracing.
"""
Expand Down
50 changes: 48 additions & 2 deletions kubernetes/aio/client/rest.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,11 @@
import aiohttp
import aiohttp_retry

from kubernetes.aio.utils.retry import (
is_retry_after_response,
on_retry_after_error,
retry_after_backoff,
)
from kubernetes.aio.client.exceptions import ApiException, ApiValueError

RESTResponseType = aiohttp.ClientResponse
Expand Down Expand Up @@ -284,14 +289,55 @@ async def request(
self.pool_manager = self._create_pool_manager()
pool_manager = self.pool_manager

if self._effective_retry_options is not None and method in ALLOW_RETRY_METHODS:
client_go_read_retries = (
method in ['GET', 'HEAD']
and getattr(self.configuration, 'client_go_retries', False)
)

if (
self._effective_retry_options is not None
and method in ALLOW_RETRY_METHODS
and not client_go_read_retries
):
if self.retry_client is None:
self.retry_client = aiohttp_retry.RetryClient(
client_session=self.pool_manager,
retry_options=self._effective_retry_options
)
pool_manager = self.retry_client

r = await pool_manager.request(**args)
async def read_request(check_retry_status=False):
response = await self.pool_manager.request(**args)
if check_retry_status:
self._raise_retry_after_response(response)
return response

if client_go_read_retries:
backoff = retry_after_backoff(
getattr(self.configuration, 'retries', None),
getattr(self.configuration, 'client_go_retry_backoff', None),
)
r = await on_retry_after_error(
backoff, self._is_read_retryable, lambda: read_request(True))
else:
r = await pool_manager.request(**args)

return RESTResponse(r)

@classmethod
def _is_read_retryable(cls, error):
return is_retry_after_response(error)

@staticmethod
def _retry_after_error(response):
error = ApiException(status=response.status, reason=response.reason)
error.headers = response.headers
return error

@classmethod
def _raise_retry_after_response(cls, response):
error = cls._retry_after_error(response)
if not is_retry_after_response(error):
return
response.release()
raise error
28 changes: 27 additions & 1 deletion kubernetes/aio/test/test_generated_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,8 @@ async def asyncSetUp(self):
'items': [],
}
self.response_status = 200
self.response_headers = {}
self.responses = None
app = web.Application()
app.router.add_route('*', '/{path:.*}', self._handle_request)
self.runner = web.AppRunner(app)
Expand Down Expand Up @@ -91,7 +93,17 @@ async def _handle_request(self, request):
}).encode() + b'\n')
await response.write_eof()
return response
return web.json_response(self.response, status=self.response_status)
if self.responses is not None:
index = len(self.requests) - 1
response, status, headers = self.responses[
min(index, len(self.responses) - 1)
]
return web.json_response(response, status=status, headers=headers)
return web.json_response(
self.response,
status=self.response_status,
headers=self.response_headers,
)

async def test_bearer_alias_supports_synchronous_token_refresh(self):
self.configuration.api_key['authorization'] = 'expired-token'
Expand Down Expand Up @@ -123,6 +135,20 @@ async def refresh(configuration):
self.requests[-1][0].headers['Authorization'],
)

async def test_client_go_retry_retries_get_retry_after_response(self):
self.configuration.client_go_retries = True
self.configuration.retries = 1
self.responses = [
({'message': 'retry later'}, 429, {'Retry-After': '0'}),
(self.response, 200, {}),
]

namespaces = await CoreV1Api(self.api_client).list_namespace()

self.assertEqual([], namespaces.items)
self.assertEqual(2, len(self.requests))
self.assertIsNone(self.api_client.rest_client.retry_client)

async def test_delete_job_accepts_job_and_status_responses(self):
responses = (
(
Expand Down
6 changes: 4 additions & 2 deletions kubernetes/aio/utils/create_from_yaml.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,8 +18,6 @@

import yaml

from kubernetes.aio import client


async def create_from_yaml(
k8s_client,
Expand Down Expand Up @@ -115,6 +113,8 @@ async def create_from_dict(
processing of the request.
Valid values are: - All: all dry run stages will be processed
"""
from kubernetes.aio import client

api_exceptions = []
k8s_objects = []

Expand Down Expand Up @@ -156,6 +156,8 @@ async def create_from_yaml_single_item(
verbose=False,
namespace="default",
**kwargs):
from kubernetes.aio import client

group, _, version = yml_object["apiVersion"].partition("/")
if version == "":
version = group
Expand Down
11 changes: 11 additions & 0 deletions kubernetes/aio/utils/retry_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
is_too_many_requests,
on_error,
on_retry_after_error,
retry_after_max_retries,
retry_after_seconds,
retry_on_conflict,
)
Expand All @@ -33,6 +34,12 @@ def __init__(self, status, headers=None):
self.headers = headers or {}


class RetryOptions:

def __init__(self, attempts):
self.attempts = attempts


class AioRetryTest(unittest.IsolatedAsyncioTestCase):

def test_default_retry_matches_client_go(self):
Expand All @@ -41,6 +48,10 @@ def test_default_retry_matches_client_go(self):
Backoff(steps=5, duration=0.01, factor=1.0, jitter=0.1),
)

def test_retry_after_backoff_uses_aio_retry_attempts(self):
self.assertEqual(retry_after_max_retries(RetryOptions(attempts=3)), 2)
self.assertEqual(retry_after_max_retries(RetryOptions(attempts=1)), 0)

def test_retry_after_seconds_parses_delay_seconds(self):
error = FakeError(429, {"Retry-After": "7"})

Expand Down
22 changes: 16 additions & 6 deletions kubernetes/base/retry.py
Original file line number Diff line number Diff line change
Expand Up @@ -77,7 +77,8 @@ def retry_after_backoff(
``None``, it is interpreted as the retry ceiling and overrides
``backoff.steps``. ``False`` and ``0`` disable retries by returning a
single-attempt backoff. Integer values and urllib3-style Retry objects with
a ``total`` value are treated as the retry ceiling.
a ``total`` value are treated as retry ceilings. aiohttp-retry-style
options with an ``attempts`` value are treated as request attempt ceilings.
"""

if backoff is None:
Expand Down Expand Up @@ -106,13 +107,22 @@ def retry_after_max_retries(retries: Any = None) -> int:
if isinstance(retries, int):
return max(0, retries)

total = getattr(retries, "total", None)
if total is False:
if hasattr(retries, "total"):
total = getattr(retries, "total", None)
if total is False:
return 0
if total is True or total is None:
return 10
if isinstance(total, int):
return max(0, total)

attempts = getattr(retries, "attempts", None)
if attempts is False:
return 0
if total is True or total is None:
if attempts is True:
return 10
if isinstance(total, int):
return max(0, total)
if isinstance(attempts, int):
return max(0, attempts - 1)
return 10


Expand Down
16 changes: 16 additions & 0 deletions kubernetes/client/configuration.py
Original file line number Diff line number Diff line change
Expand Up @@ -357,6 +357,22 @@ def __init__(
self.retries = retries
"""Retry configuration
"""
self.client_go_retries = False
"""Enable Kubernetes client-go-compatible retry semantics.

When enabled, GET and HEAD requests retry Retry-After responses.
The retry ceiling is read from ``retries`` when set; otherwise it
follows the client-go default of at most 10 retries.
"""
self.client_go_retry_backoff = None
"""Backoff for Kubernetes client-go-compatible retries.

If unset, client-go-compatible GET and HEAD retries use the
client-go default retry ceiling with no additional client-side
delay beyond Retry-After. When set, ``retries`` still overrides
the retry ceiling if it is not None.
"""

# Enable client side validation
self.client_side_validation = client_side_validation

Expand Down
101 changes: 93 additions & 8 deletions kubernetes/client/rest.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,13 @@
from urllib.parse import urlparse

import urllib3
from urllib3.util.retry import Retry

from kubernetes.utils.retry import (
is_retry_after_response,
on_retry_after_error,
retry_after_backoff,
)
from kubernetes.client.exceptions import ApiException, ApiValueError

SUPPORTED_SOCKS_PROXIES = {"socks5", "socks5h", "socks4", "socks4a"}
Expand Down Expand Up @@ -105,6 +111,8 @@ def getheader(self, name, default=None):
class RESTClientObject:

def __init__(self, configuration) -> None:
self.configuration = configuration

# urllib3.PoolManager will pass all kw parameters to connectionpool
# https://github.com/shazow/urllib3/blob/f9409436f83aeb79fbaf090181cd81b784f1b8ce/urllib3/poolmanager.py#L75 # noqa: E501
# https://github.com/shazow/urllib3/blob/f9409436f83aeb79fbaf090181cd81b784f1b8ce/urllib3/connectionpool.py#L680 # noqa: E501
Expand Down Expand Up @@ -217,6 +225,32 @@ def request(
read=_request_timeout[1]
)

client_go_read_retries = (
method in ['GET', 'HEAD']
and getattr(self.configuration, 'client_go_retries', False)
)
read_retries = None
if client_go_read_retries:
read_retries = self._urllib3_retries_without_status(
getattr(self.configuration, 'retries', None))

def read_request(check_retry_status=False):
kwargs = {}
if read_retries is not None:
kwargs['retries'] = read_retries
response = self.pool_manager.request(
method,
url,
fields={},
timeout=timeout,
headers=headers,
preload_content=False,
**kwargs
)
if check_retry_status:
self._raise_retry_after_response(response)
return response

try:
# For `POST`, `PUT`, `PATCH`, `OPTIONS`, `DELETE`
if method in ['POST', 'PUT', 'PATCH', 'OPTIONS', 'DELETE']:
Expand Down Expand Up @@ -310,16 +344,67 @@ def request(
raise ApiException(status=0, reason=msg)
# For `GET`, `HEAD`
else:
r = self.pool_manager.request(
method,
url,
fields={},
timeout=timeout,
headers=headers,
preload_content=False
)
if client_go_read_retries:
backoff = retry_after_backoff(
getattr(self.configuration, 'retries', None),
getattr(self.configuration, 'client_go_retry_backoff', None),
)
r = on_retry_after_error(
backoff, self._is_read_retryable,
lambda: read_request(True))
else:
r = read_request()
except urllib3.exceptions.SSLError as e:
msg = "\n".join([type(e).__name__, str(e)])
raise ApiException(status=0, reason=msg)

return RESTResponse(r)

@classmethod
def _is_read_retryable(cls, error):
return is_retry_after_response(error)

@staticmethod
def _retry_after_error(response):
error = ApiException(status=response.status, reason=response.reason)
error.headers = response.getheaders()
return error

@classmethod
def _raise_retry_after_response(cls, response):
error = cls._retry_after_error(response)
if not is_retry_after_response(error):
return
cls._read_and_close_retry_response(response)
raise error

@staticmethod
def _read_and_close_retry_response(response):
try:
try:
content_length = int(response.getheaders().get(
'Content-Length', '-1'))
except (TypeError, ValueError):
content_length = -1
if content_length <= 2 << 10:
response.read(2 << 10)
finally:
response.close()

@staticmethod
def _urllib3_retries_without_status(retries):
if retries is False:
return False
if retries is None:
retries = Retry.DEFAULT
elif retries is True:
retries = Retry.DEFAULT
elif isinstance(retries, int):
retries = Retry.from_int(retries)
if isinstance(retries, Retry):
return retries.new(
status=0,
status_forcelist=(),
respect_retry_after_header=False,
)
return retries
Loading