diff --git a/src/lightkube/core/generic_client.py b/src/lightkube/core/generic_client.py index f591154..685644e 100644 --- a/src/lightkube/core/generic_client.py +++ b/src/lightkube/core/generic_client.py @@ -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": @@ -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 @@ -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 diff --git a/src/lightkube/types.py b/src/lightkube/types.py index 4870867..ae2f7b4 100644 --- a/src/lightkube/types.py +++ b/src/lightkube/types.py @@ -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) diff --git a/tests/test_async_client.py b/tests/test_async_client.py index c5a49c1..f5ca73a 100644 --- a/tests/test_async_client.py +++ b/tests/test_async_client.py @@ -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) diff --git a/tests/test_client.py b/tests/test_client.py index f387ce6..293afe7 100644 --- a/tests/test_client.py +++ b/tests/test_client.py @@ -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())