Skip to content

Commit 50f6582

Browse files
committed
rest: apply client-go retry generation patches
Signed-off-by: Dr. Stefan Schimanski <[email protected]>
1 parent 9bc5eac commit 50f6582

10 files changed

Lines changed: 510 additions & 16 deletions

kubernetes/aio/client/configuration.py

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -355,6 +355,22 @@ def __init__(
355355
self.retries = retries
356356
"""Retry configuration
357357
"""
358+
self.client_go_retries = False
359+
"""Enable Kubernetes client-go-compatible retry semantics.
360+
361+
When enabled, GET and HEAD requests retry Retry-After responses.
362+
The retry ceiling is read from ``retries`` when set; otherwise it
363+
follows the client-go default of at most 10 retries.
364+
"""
365+
self.client_go_retry_backoff = None
366+
"""Backoff for Kubernetes client-go-compatible retries.
367+
368+
If unset, client-go-compatible GET and HEAD retries use the
369+
client-go default retry ceiling with no additional client-side
370+
delay beyond Retry-After. When set, ``retries`` still overrides
371+
the retry ceiling if it is not None.
372+
"""
373+
358374
self.trace_configs = trace_configs
359375
"""aiohttp.TraceConfig list forwarded to ClientSession for tracing.
360376
"""

kubernetes/aio/client/rest.py

Lines changed: 41 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,11 @@
2222
import aiohttp_retry
2323

2424
from kubernetes.aio.client.exceptions import ApiException, ApiValueError
25+
from kubernetes.aio.utils.retry import (
26+
is_retry_after_response,
27+
on_retry_after_error,
28+
retry_after_backoff,
29+
)
2530

2631
RESTResponseType = aiohttp.ClientResponse
2732

@@ -292,6 +297,41 @@ async def request(
292297
)
293298
pool_manager = self.retry_client
294299

295-
r = await pool_manager.request(**args)
300+
async def read_request(check_retry_status=False):
301+
response = await self.pool_manager.request(**args)
302+
if check_retry_status:
303+
self._raise_retry_after_response(response)
304+
return response
305+
306+
if (
307+
method in ['GET', 'HEAD']
308+
and getattr(self.configuration, 'client_go_retries', False)
309+
):
310+
backoff = retry_after_backoff(
311+
getattr(self.configuration, 'retries', None),
312+
getattr(self.configuration, 'client_go_retry_backoff', None),
313+
)
314+
r = await on_retry_after_error(
315+
backoff, self._is_read_retryable, lambda: read_request(True))
316+
else:
317+
r = await pool_manager.request(**args)
296318

297319
return RESTResponse(r)
320+
321+
@classmethod
322+
def _is_read_retryable(cls, error):
323+
return is_retry_after_response(error)
324+
325+
@staticmethod
326+
def _retry_after_error(response):
327+
error = ApiException(status=response.status, reason=response.reason)
328+
error.headers = response.headers
329+
return error
330+
331+
@classmethod
332+
def _raise_retry_after_response(cls, response):
333+
error = cls._retry_after_error(response)
334+
if not is_retry_after_response(error):
335+
return
336+
response.release()
337+
raise error

kubernetes/aio/test/test_generated_api.py

Lines changed: 26 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -46,6 +46,8 @@ async def asyncSetUp(self):
4646
'items': [],
4747
}
4848
self.response_status = 200
49+
self.response_headers = {}
50+
self.responses = None
4951
app = web.Application()
5052
app.router.add_route('*', '/{path:.*}', self._handle_request)
5153
self.runner = web.AppRunner(app)
@@ -91,7 +93,17 @@ async def _handle_request(self, request):
9193
}).encode() + b'\n')
9294
await response.write_eof()
9395
return response
94-
return web.json_response(self.response, status=self.response_status)
96+
if self.responses is not None:
97+
index = len(self.requests) - 1
98+
response, status, headers = self.responses[
99+
min(index, len(self.responses) - 1)
100+
]
101+
return web.json_response(response, status=status, headers=headers)
102+
return web.json_response(
103+
self.response,
104+
status=self.response_status,
105+
headers=self.response_headers,
106+
)
95107

96108
async def test_bearer_alias_supports_synchronous_token_refresh(self):
97109
self.configuration.api_key['authorization'] = 'expired-token'
@@ -123,6 +135,19 @@ async def refresh(configuration):
123135
self.requests[-1][0].headers['Authorization'],
124136
)
125137

138+
async def test_client_go_retry_retries_get_retry_after_response(self):
139+
self.configuration.client_go_retries = True
140+
self.configuration.retries = 1
141+
self.responses = [
142+
({'message': 'retry later'}, 429, {'Retry-After': '0'}),
143+
(self.response, 200, {}),
144+
]
145+
146+
namespaces = await CoreV1Api(self.api_client).list_namespace()
147+
148+
self.assertEqual([], namespaces.items)
149+
self.assertEqual(2, len(self.requests))
150+
126151
async def test_delete_job_accepts_job_and_status_responses(self):
127152
responses = (
128153
(

kubernetes/client/configuration.py

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -357,6 +357,22 @@ def __init__(
357357
self.retries = retries
358358
"""Retry configuration
359359
"""
360+
self.client_go_retries = False
361+
"""Enable Kubernetes client-go-compatible retry semantics.
362+
363+
When enabled, GET and HEAD requests retry Retry-After responses.
364+
The retry ceiling is read from ``retries`` when set; otherwise it
365+
follows the client-go default of at most 10 retries.
366+
"""
367+
self.client_go_retry_backoff = None
368+
"""Backoff for Kubernetes client-go-compatible retries.
369+
370+
If unset, client-go-compatible GET and HEAD retries use the
371+
client-go default retry ceiling with no additional client-side
372+
delay beyond Retry-After. When set, ``retries`` still overrides
373+
the retry ceiling if it is not None.
374+
"""
375+
360376
# Enable client side validation
361377
self.client_side_validation = client_side_validation
362378

kubernetes/client/rest.py

Lines changed: 86 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -20,8 +20,14 @@
2020
from urllib.parse import urlparse
2121

2222
import urllib3
23+
from urllib3.util.retry import Retry
2324

2425
from kubernetes.client.exceptions import ApiException, ApiValueError
26+
from kubernetes.utils.retry import (
27+
is_retry_after_response,
28+
on_retry_after_error,
29+
retry_after_backoff,
30+
)
2531

2632
SUPPORTED_SOCKS_PROXIES = {"socks5", "socks5h", "socks4", "socks4a"}
2733
RESTResponseType = urllib3.HTTPResponse
@@ -105,6 +111,8 @@ def getheader(self, name, default=None):
105111
class RESTClientObject:
106112

107113
def __init__(self, configuration) -> None:
114+
self.configuration = configuration
115+
108116
# urllib3.PoolManager will pass all kw parameters to connectionpool
109117
# https://github.com/shazow/urllib3/blob/f9409436f83aeb79fbaf090181cd81b784f1b8ce/urllib3/poolmanager.py#L75 # noqa: E501
110118
# https://github.com/shazow/urllib3/blob/f9409436f83aeb79fbaf090181cd81b784f1b8ce/urllib3/connectionpool.py#L680 # noqa: E501
@@ -217,6 +225,30 @@ def request(
217225
read=_request_timeout[1]
218226
)
219227

228+
client_go_retries = getattr(self.configuration, 'client_go_retries', False)
229+
request_retries = None
230+
if client_go_retries:
231+
request_retries = self._urllib3_retries_without_status(
232+
getattr(self.configuration, 'retries', None))
233+
234+
def pool_request(*args, **kwargs):
235+
if request_retries is not None:
236+
kwargs['retries'] = request_retries
237+
return self.pool_manager.request(*args, **kwargs)
238+
239+
def read_request(check_retry_status=False):
240+
response = pool_request(
241+
method,
242+
url,
243+
fields={},
244+
timeout=timeout,
245+
headers=headers,
246+
preload_content=False
247+
)
248+
if check_retry_status:
249+
self._raise_retry_after_response(response)
250+
return response
251+
220252
try:
221253
# For `POST`, `PUT`, `PATCH`, `OPTIONS`, `DELETE`
222254
if method in ['POST', 'PUT', 'PATCH', 'OPTIONS', 'DELETE']:
@@ -247,7 +279,7 @@ def request(
247279
request_body = None
248280
if body is not None:
249281
request_body = json.dumps(body)
250-
r = self.pool_manager.request(
282+
r = pool_request(
251283
method,
252284
url,
253285
body=request_body,
@@ -256,7 +288,7 @@ def request(
256288
preload_content=False
257289
)
258290
elif content_type == 'application/x-www-form-urlencoded':
259-
r = self.pool_manager.request(
291+
r = pool_request(
260292
method,
261293
url,
262294
fields=post_params,
@@ -272,7 +304,7 @@ def request(
272304
del headers['Content-Type']
273305
# Ensures that dict objects are serialized
274306
post_params = [(a, json.dumps(b)) if isinstance(b, dict) else (a,b) for a, b in post_params]
275-
r = self.pool_manager.request(
307+
r = pool_request(
276308
method,
277309
url,
278310
fields=post_params,
@@ -285,7 +317,7 @@ def request(
285317
# other content types than JSON when `body` argument is
286318
# provided in serialized form.
287319
elif isinstance(body, str) or isinstance(body, bytes):
288-
r = self.pool_manager.request(
320+
r = pool_request(
289321
method,
290322
url,
291323
body=body,
@@ -295,7 +327,7 @@ def request(
295327
)
296328
elif headers['Content-Type'].startswith('text/') and isinstance(body, bool):
297329
request_body = "true" if body else "false"
298-
r = self.pool_manager.request(
330+
r = pool_request(
299331
method,
300332
url,
301333
body=request_body,
@@ -310,16 +342,57 @@ def request(
310342
raise ApiException(status=0, reason=msg)
311343
# For `GET`, `HEAD`
312344
else:
313-
r = self.pool_manager.request(
314-
method,
315-
url,
316-
fields={},
317-
timeout=timeout,
318-
headers=headers,
319-
preload_content=False
320-
)
345+
if client_go_retries:
346+
backoff = retry_after_backoff(
347+
getattr(self.configuration, 'retries', None),
348+
getattr(self.configuration, 'client_go_retry_backoff', None),
349+
)
350+
r = on_retry_after_error(
351+
backoff, self._is_read_retryable,
352+
lambda: read_request(True))
353+
else:
354+
r = read_request()
321355
except urllib3.exceptions.SSLError as e:
322356
msg = "\n".join([type(e).__name__, str(e)])
323357
raise ApiException(status=0, reason=msg)
324358

325359
return RESTResponse(r)
360+
361+
@classmethod
362+
def _is_read_retryable(cls, error):
363+
return is_retry_after_response(error)
364+
365+
@staticmethod
366+
def _retry_after_error(response):
367+
error = ApiException(status=response.status, reason=response.reason)
368+
error.headers = response.getheaders()
369+
return error
370+
371+
@classmethod
372+
def _raise_retry_after_response(cls, response):
373+
error = cls._retry_after_error(response)
374+
if not is_retry_after_response(error):
375+
return
376+
try:
377+
response.drain_conn()
378+
except Exception:
379+
response.release_conn()
380+
raise error
381+
382+
@staticmethod
383+
def _urllib3_retries_without_status(retries):
384+
if retries is False:
385+
return False
386+
if retries is None:
387+
retries = Retry.DEFAULT
388+
elif retries is True:
389+
retries = Retry.DEFAULT
390+
elif isinstance(retries, int):
391+
retries = Retry.from_int(retries)
392+
if isinstance(retries, Retry):
393+
return retries.new(
394+
status=0,
395+
status_forcelist=(),
396+
respect_retry_after_header=False,
397+
)
398+
return retries

kubernetes/test/test_generated_api.py

Lines changed: 36 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -80,10 +80,25 @@ def do_DELETE(self):
8080
self._respond()
8181

8282
def _respond(self):
83+
request_index = self.server.request_count
84+
self.server.request_count += 1
8385
response = self.server.response_body
84-
self.send_response(self.server.response_status)
86+
response_bodies = getattr(self.server, 'response_bodies', None)
87+
if response_bodies is not None:
88+
response = response_bodies[min(request_index, len(response_bodies) - 1)]
89+
status = self.server.response_status
90+
response_statuses = getattr(self.server, 'response_statuses', None)
91+
if response_statuses is not None:
92+
status = response_statuses[min(request_index, len(response_statuses) - 1)]
93+
response_headers = self.server.response_headers
94+
response_header_sets = getattr(self.server, 'response_header_sets', None)
95+
if response_header_sets is not None:
96+
response_headers = response_header_sets[min(request_index, len(response_header_sets) - 1)]
97+
self.send_response(status)
8598
self.send_header('Content-Type', 'application/json')
8699
self.send_header('Content-Length', str(len(response)))
100+
for key, value in response_headers.items():
101+
self.send_header(key, value)
87102
self.end_headers()
88103
self.wfile.write(response)
89104

@@ -386,6 +401,8 @@ def setUp(self):
386401
('127.0.0.1', 0), _RecordingHandler)
387402
self.server.response_body = b'{}'
388403
self.server.response_status = 200
404+
self.server.response_headers = {}
405+
self.server.request_count = 0
389406
self.thread = threading.Thread(target=self.server.serve_forever)
390407
self.thread.start()
391408
self.api_client = ApiClient(Configuration(
@@ -411,6 +428,24 @@ def test_continue_parameter_keeps_public_name(self):
411428
self.server.request_path,
412429
)
413430

431+
def test_client_go_retry_retries_get_retry_after_response(self):
432+
self.api_client.configuration.client_go_retries = True
433+
self.api_client.configuration.retries = 1
434+
self.server.response_statuses = [429, 200]
435+
self.server.response_bodies = [
436+
b'{"message": "retry later"}',
437+
b'{"items": []}',
438+
]
439+
self.server.response_header_sets = [
440+
{'Retry-After': '0'},
441+
{},
442+
]
443+
444+
namespaces = CoreV1Api(self.api_client).list_namespace()
445+
446+
self.assertEqual([], namespaces.items)
447+
self.assertEqual(2, self.server.request_count)
448+
414449
def test_deferred_api_validation_rejects_invalid_first_call(self):
415450
with self.assertRaises(ValidationError) as raised:
416451
CoreV1Api(self.api_client).list_namespace(limit='invalid')

0 commit comments

Comments
 (0)