Skip to content
Merged
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
11 changes: 11 additions & 0 deletions src/lightkube/core/generic_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -66,8 +66,13 @@ def get_request(self, timeout):
br = self._br
if self._version is not None:
br.params["resourceVersion"] = self._version
else:
br.params.pop("resourceVersion", None)
return self._build_request(br.method, br.url, params=br.params, timeout=timeout)

def reset_version(self):
self._version = None

def process_one_line(self, line):
event = decoders.decode_watch_event(line)
if event.type == "ERROR":
Expand Down Expand Up @@ -327,6 +332,9 @@ def watch(self, br: BasicRequest, on_error: OnErrorHandler = on_error_raise):
raise
if handle_error.action is OnErrorAction.STOP:
break
if isinstance(e, ApiError) and e.status.code == 410:
# 410 Gone means the resourceVersion is too old, we need to restart the watch from scratch
wd.reset_version()
if handle_error.sleep > 0:
time.sleep(handle_error.sleep)
continue
Expand Down Expand Up @@ -411,6 +419,9 @@ async def watch(self, br: BasicRequest, on_error: OnErrorHandler = on_error_rais
raise
if handle_error.action is OnErrorAction.STOP:
break
if isinstance(e, ApiError) and e.status.code == 410:
# 410 Gone means the resourceVersion is too old, we need to restart the watch from scratch
wd.reset_version()
if handle_error.sleep > 0:
await asyncio.sleep(handle_error.sleep)
continue
Expand Down
7 changes: 6 additions & 1 deletion src/lightkube/types.py
Original file line number Diff line number Diff line change
Expand Up @@ -65,7 +65,12 @@ def on_error_stop(e: Exception, count: int) -> OnErrorResult:


def on_error_retry(e: Exception, count: int) -> OnErrorResult:
"""Retry to perform the API call again from the last version"""
"""Immediately retries the API call using the last known resource version.
If the retry fails with an unrecoverable error (such as
"too old resource version"), it restarts without specifying a resource
version. The watch then delivers the current state and continues streaming
updates from that point onward.
"""
return OnErrorResult(OnErrorAction.RETRY)


Expand Down
17 changes: 17 additions & 0 deletions tests/test_async_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -229,6 +229,23 @@ async def test_watch_error(client: lightkube.AsyncClient, httpx2_mock: respx.Rou
await client.close()


@pytest.mark.asyncio
async def test_watch_error_restarts_from_scratch(client: lightkube.AsyncClient, httpx2_mock: respx.Router):
stale_request = httpx2_mock.get("https://localhost:9443/api/v1/nodes?watch=true&resourceVersion=1").respond(
content=make_watch_error()
)
initial_request = httpx2_mock.get("https://localhost:9443/api/v1/nodes?watch=true").respond(content=make_watch_list(1))

watch = client.watch(Node, on_error=types.on_error_retry)
await anext(watch)
await anext(watch)

assert len(initial_request.calls) == 2
assert len(stale_request.calls) == 1
assert "resourceVersion" not in initial_request.calls[1][0].url.params
await client.close()


@pytest.mark.asyncio
async def test_watch_version(client: lightkube.AsyncClient, httpx2_mock: respx.Router):
httpx2_mock.get("https://localhost:9443/api/v1/nodes?resourceVersion=1&watch=true").respond(status_code=404)
Expand Down
15 changes: 15 additions & 0 deletions tests/test_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -324,6 +324,21 @@ def test_watch_error(client: lightkube.Client, httpx2_mock: respx.Router) -> Non
assert exc.value.status.code == 410


def test_watch_error_restarts_from_scratch(client: lightkube.Client, httpx2_mock: respx.Router) -> None:
stale_request = httpx2_mock.get("https://localhost:9443/api/v1/nodes?watch=true&resourceVersion=1").respond(
content=make_watch_error()
)
initial_request = httpx2_mock.get("https://localhost:9443/api/v1/nodes?watch=true").respond(content=make_watch_list(1))

watch = client.watch(Node, on_error=types.on_error_retry)
next(watch)
next(watch)

assert len(initial_request.calls) == 2
assert len(stale_request.calls) == 1
assert "resourceVersion" not in initial_request.calls[1][0].url.params


def test_watch_version(client: lightkube.Client, httpx2_mock: respx.Router) -> None:
httpx2_mock.get("https://localhost:9443/api/v1/nodes?resourceVersion=1&watch=true").respond(status_code=404)
httpx2_mock.get("https://localhost:9443/api/v1/nodes?resourceVersion=2&watch=true").respond(content=make_watch_list())
Expand Down
Loading