From 951ccdc7c018751f165d93d975b11fa6b9bf3fec Mon Sep 17 00:00:00 2001 From: Yun Seongmin Date: Sat, 26 Sep 2026 16:47:27 +0900 Subject: [PATCH] Resume informer watches safely with bounded recovery --- kubernetes/informer/informer.py | 194 ++++--- kubernetes/test/test_informer.py | 904 ++++++++++++++++++++++++++----- kubernetes/watch/watch.py | 71 ++- kubernetes/watch/watch_test.py | 259 ++++++++- 4 files changed, 1190 insertions(+), 238 deletions(-) diff --git a/kubernetes/informer/informer.py b/kubernetes/informer/informer.py index f0e83000e3..90d107eff0 100644 --- a/kubernetes/informer/informer.py +++ b/kubernetes/informer/informer.py @@ -20,6 +20,8 @@ """ import logging +import math +import random import threading import time @@ -58,7 +60,11 @@ class SharedInformer: or all-namespace list functions. resync_period: How often (seconds) to perform a full re-list from the API server. - Defaults to 0 which disables periodic resyncs. + Defaults to 0 which disables periodic resyncs. The interval starts + after a successful list completes; reconnects do not reset it. + Failed lists and watches use separate exponential retry delays + (0.5 to 60 seconds, with jitter). Network calls can extend the + interval beyond the configured period. label_selector: Optional label selector string forwarded to the API server. field_selector: @@ -90,7 +96,9 @@ def __init__( self._watch = None self._thread = None self._stop_event = threading.Event() - self._resource_version = None # most recent RV seen; None forces a full re-list + # Last applied RV; None forces a full re-list. + self._resource_version = None + self._last_list_time = None # ---------------------------------------------------------------- # # Public API # @@ -236,107 +244,149 @@ def _initial_list(self): self._resource_version = rv or "0" def _run_loop(self): - """Background loop: list then watch, reconnect on errors. + """Own LIST/WATCH retries and resume only from applied events. - A full re-list is only performed when ``self._resource_version`` is - ``None`` (first start or after a 410 Gone response). On all other - reconnects the most recent ``resourceVersion`` is reused so that no - events are missed and the API server does not need to send a full - object snapshot. + LIST deadlines are measured from successful LIST completion. A failed + periodic LIST leaves its previous RV usable while its retry backs off; + initial and expired-RV LISTs must succeed before another WATCH starts. """ + next_resync = float("inf") + if self._resync_period > 0: + last_list = self._last_list_time + if last_list is None: + last_list = time.monotonic() + next_resync = last_list + self._resync_period + list_retry_at = watch_retry_at = 0.0 + list_backoff = watch_backoff = 1.0 while not self._stop_event.is_set(): - # Full re-list only when we have no resource version to resume from. - if self._resource_version is None: + now = time.monotonic() + list_required = self._resource_version is None + list_due = max(next_resync, list_retry_at) + if list_required: + list_due = list_retry_at + if now >= list_due: try: self._initial_list() except Exception as exc: - logger.exception("Error during initial list; retrying") + logger.exception("Error during list; retrying") self._fire(ERROR, exc) - self._stop_event.wait(timeout=5) - continue - - # Watch loop - last_resync = time.monotonic() - self._watch = Watch() + list_retry_at = time.monotonic() + random.uniform( + list_backoff / 2, list_backoff) + list_backoff = min(60.0, list_backoff * 2) + else: + list_backoff = 1.0 + list_retry_at = 0.0 + self._last_list_time = time.monotonic() + next_resync = ( + self._last_list_time + self._resync_period + if self._resync_period > 0 else float("inf") + ) + continue + + if list_required: + self._stop_event.wait(timeout=max(0, list_due - now)) + continue + if now < watch_retry_at: + self._stop_event.wait( + timeout=max(0, min(watch_retry_at, list_due) - now)) + continue + + # One HTTP request per stream: retries and expired RV recovery + # belong here, so Watch's internal EOF retries cannot bypass + # delays. + self._watch = Watch(retry=False) kw = self._build_kwargs() kw["resource_version"] = self._resource_version - # When a resync period is configured, set a matching server-side - # watch timeout so that the stream exits after resync_period seconds - # even if no events arrive. Without this, a quiet period longer - # than resync_period would never trigger a resync because the check - # below only runs when the generator yields an event. if self._resync_period > 0: - kw["timeout_seconds"] = max(1, int(self._resync_period)) + kw["timeout_seconds"] = max(1, math.ceil(list_due - now)) + watch_started = time.monotonic() + watch_failed = False + resync_requested = False + stream = None try: - for event in self._watch.stream(self._list_func, **kw): + stream = self._watch.stream(self._list_func, **kw) + for event in stream: if self._stop_event.is_set(): break + # A busy stream must also yield to a due periodic LIST. + if time.monotonic() >= list_due: + resync_requested = True + break + if event is None: + continue evt_type = event.get("type") obj = event.get("object") - # Sync the most recent resource version from the Watch - # instance (updated by unmarshal_event before yielding). - # Do this before firing handlers so consumers that wake on - # an event immediately see the advanced resource version. - if self._watch is not None and self._watch.resource_version: - self._resource_version = self._watch.resource_version - if evt_type == ADDED: - self._cache._put(obj) - self._fire(ADDED, obj) - elif evt_type == MODIFIED: + if evt_type in (ADDED, MODIFIED): self._cache._put(obj) - self._fire(MODIFIED, obj) elif evt_type == DELETED: self._cache._remove(obj) - self._fire(DELETED, obj) elif evt_type == BOOKMARK: - # BOOKMARK events carry an updated resource version but - # no object state change; the Watch instance already - # records the new resource_version internally. - self._fire(BOOKMARK, event.get("raw_object", obj)) - elif evt_type == ERROR: - self._fire(ERROR, obj) + obj = event.get("raw_object", obj) + + if evt_type in (ADDED, MODIFIED, DELETED, BOOKMARK): + # Acknowledge only events applied to the cache (or a + # BOOKMARK), before notifying handlers. Watch advances + # its own RV during parsing, before cache mutation can + # fail or a stop request can interrupt processing. + if isinstance(obj, dict): + metadata = obj.get("metadata") or {} + resource_version = metadata.get("resourceVersion") + else: + metadata = getattr(obj, "metadata", None) + resource_version = getattr( + metadata, "resource_version", None) + if resource_version: + self._resource_version = resource_version + + self._fire(evt_type, obj) except ApiException as exc: + watch_failed = True if exc.status == 410: - # The stored resource version is too old; force a full re-list. logger.warning( "Watch expired (410 Gone); will re-list from scratch" ) self._resource_version = None else: logger.warning( - "Watch stream ended with ApiException (status=%s); reconnecting", - exc.status, - ) + "Watch failed (status=%s); reconnecting", exc.status) self._fire(ERROR, exc) except Exception as exc: - logger.exception("Unexpected error in watch loop; reconnecting") + watch_failed = True + logger.exception( + "Unexpected error in watch loop; reconnecting") self._fire(ERROR, exc) finally: - # Capture the most recent resource version seen by the Watch - # (updated on every ADDED/MODIFIED/DELETED/BOOKMARK event) so - # that the next watch connection can resume without re-listing. - # Do not overwrite a None that was set by a 410 handler above. - if ( - self._resource_version is not None - and self._watch is not None - and self._watch.resource_version - ): - self._resource_version = self._watch.resource_version - self._watch = None - - # Periodic resync: after the watch stream exits (whether due to the - # server-side timeout_seconds, a stop request, or an error) check if - # a resync is due. This path is what actually fires the resync when - # the cluster is quiet and no events arrive for resync_period seconds. - if ( - not self._stop_event.is_set() - and self._resource_version is not None # 410 already schedules a re-list - and self._resync_period > 0 - and (time.monotonic() - last_resync) >= self._resync_period - ): - logger.debug("Informer resync triggered") + watch_finished = time.monotonic() + # Explicitly finalize generators when stop/resync breaks the + # loop, so their HTTP response is released before the next + # LIST. try: - self._initial_list() + if hasattr(stream, "close"): + stream.close() except Exception as exc: - logger.exception("Error during resync list; continuing") + watch_failed = True + logger.exception( + "Error closing watch stream; reconnecting") self._fire(ERROR, exc) + finally: + self._watch = None + + duration = watch_finished - watch_started + # An event or BOOKMARK alone does not prove a healthy connection. + # Only a clean server timeout or long-lived EOF resets history; + # a slow failure is still a failure. + healthy = not watch_failed and ( + duration >= min(60, kw.get("timeout_seconds", 60)) + ) + # A planned relist is neither a failed connection nor evidence + # of recovery. Preserve prior failures without adding a delay. + if resync_requested and not watch_failed: + watch_retry_at = watch_finished + continue + if healthy: + watch_backoff = 1.0 + watch_retry_at = watch_finished + else: + watch_retry_at = watch_finished + random.uniform( + watch_backoff / 2, watch_backoff) + watch_backoff = min(60.0, watch_backoff * 2) diff --git a/kubernetes/test/test_informer.py b/kubernetes/test/test_informer.py index 8ea5461e74..b57355ab05 100644 --- a/kubernetes/test/test_informer.py +++ b/kubernetes/test/test_informer.py @@ -14,20 +14,18 @@ """Unit tests for kubernetes.informer.""" +import json import threading import time import unittest +from types import SimpleNamespace from unittest.mock import MagicMock, patch +from kubernetes.client.exceptions import ApiException from kubernetes.informer.cache import ObjectCache, _meta_namespace_key -from kubernetes.informer.informer import ( - ADDED, - BOOKMARK, - DELETED, - ERROR, - MODIFIED, - SharedInformer, -) +from kubernetes.informer.informer import (ADDED, BOOKMARK, DELETED, ERROR, + MODIFIED, SharedInformer) +from kubernetes.watch import Watch def _make_pod(namespace, name): @@ -367,50 +365,6 @@ def fake_stream(func, **kw): self.assertEqual(len(cached), 1) self.assertIs(cached[0], pod) - def test_bookmark_advances_resource_version(self): - """A BOOKMARK event causes the informer's _resource_version to advance. - - PR #2505 added BOOKMARK-aware handling to Watch.unmarshal_event: it - extracts resourceVersion from the raw BOOKMARK dict and stores it on - self.resource_version *without* deserialising the object (because - BOOKMARK events may be incomplete). The informer must read that value - back so that the next watch reconnect starts from the BOOKMARK's RV - rather than the initial-list RV. - """ - bookmark_obj = {"metadata": {"resourceVersion": "100"}} - - list_func = MagicMock() - list_resp = MagicMock() - list_resp.items = [] - list_resp.metadata = MagicMock(resource_version="5") - list_func.return_value = list_resp - - informer = SharedInformer(list_func=list_func) - - with patch("kubernetes.informer.informer.Watch") as MockWatch: - mock_w = MagicMock() - # Start at the initial-list RV; fake_stream will advance it to the - # BOOKMARK's RV, mirroring how Watch.unmarshal_event updates - # self.resource_version before yielding a BOOKMARK event. - mock_w.resource_version = "5" - - def fake_stream(func, **kw): - # Simulate Watch.unmarshal_event setting resource_version from - # the BOOKMARK metadata before the event is yielded. - mock_w.resource_version = "100" - yield {"type": "BOOKMARK", "object": bookmark_obj, "raw_object": bookmark_obj} - informer._stop_event.set() - - mock_w.stream.side_effect = fake_stream - MockWatch.return_value = mock_w - - informer.start() - informer._thread.join(timeout=3) - - # The informer must have synced the RV from the BOOKMARK, not the - # stale initial-list RV ("5"). - self.assertEqual(informer._resource_version, "100") - def test_bookmark_handler_receives_raw_dict(self): """BOOKMARK handlers receive the raw dict, not a deserialized model. @@ -460,8 +414,6 @@ def test_multiple_bookmarks_advance_resource_version_to_latest(self): informer = SharedInformer(list_func=list_func) - rv_sequence = iter(["10", "20", "30"]) - with patch("kubernetes.informer.informer.Watch") as MockWatch: mock_w = MagicMock() @@ -502,14 +454,8 @@ def test_resync_period_triggers_full_list(self): with patch("kubernetes.informer.informer.Watch") as MockWatch, \ patch("kubernetes.informer.informer.time") as mock_time: - # Sequence of time.monotonic() calls inside _run_loop: - # 1. last_resync = time.monotonic() → 0.0 (watch-loop start) - # 2. post-stream: time.monotonic() → 61.0 (≥60 → resync fires) - # 3. last_resync = time.monotonic() → 61.0 (second watch-loop start) - # The stop_event is set during the second stream, so the - # post-stream check is short-circuited and no further calls occur. - mock_time.monotonic.side_effect = [0.0, 61.0, 61.0] - mock_time.sleep = time.sleep # keep real sleep/wait working + clock = [0.0] + mock_time.monotonic.side_effect = lambda: clock[0] mock_w = MagicMock() mock_w.resource_version = "5" @@ -519,6 +465,7 @@ def fake_stream(func, **kw): if stream_calls["n"] == 1: # Simulate the stream timing out (timeout_seconds expired) # with no events – the resync should fire after this returns. + clock[0] = 61.0 return iter([]) # Second iteration: stop the informer. informer._stop_event.set() @@ -783,41 +730,6 @@ def fake_stream(func, **kw): self.assertIsNone(informer.cache.get_by_key("default/pod-delete")) self.assertIsNotNone(informer.cache.get_by_key("default/pod-keep")) - def test_resource_version_stored_from_watch(self): - """After the watch stream ends the latest RV is preserved for reconnect.""" - pod = _make_pod("default", "rv-pod") - events = [{"type": "ADDED", "object": pod}] - - list_func = MagicMock() - list_resp = MagicMock() - list_resp.items = [] - list_resp.metadata = MagicMock(resource_version="10") - list_func.return_value = list_resp - - informer = SharedInformer(list_func=list_func) - - call_count = {"n": 0} - - with patch("kubernetes.informer.informer.Watch") as MockWatch: - mock_w = MagicMock() - mock_w.resource_version = "99" - - def fake_stream(func, **kw): - call_count["n"] += 1 - yield from events - informer._stop_event.set() - - mock_w.stream.side_effect = fake_stream - MockWatch.return_value = mock_w - - informer.start() - informer._thread.join(timeout=3) - - # The Watch reported RV "99"; the informer should have stored it. - self.assertEqual(informer._resource_version, "99") - # list_func should have been called once for the initial list only. - self.assertEqual(list_func.call_count, 1) - def test_reconnect_skips_relist_when_rv_known(self): """On reconnect without 410 the informer must NOT call the list function again.""" pod = _make_pod("default", "reconnect-pod") @@ -855,41 +767,6 @@ def fake_stream(func, **kw): self.assertEqual(list_func.call_count, 1) self.assertEqual(stream_calls["n"], 2) - def test_410_gone_triggers_relist(self): - """A 410 Gone ApiException must reset resource_version and trigger re-list.""" - from kubernetes.client.exceptions import ApiException - - list_func = MagicMock() - list_resp = MagicMock() - list_resp.items = [] - list_resp.metadata = MagicMock(resource_version="3") - list_func.return_value = list_resp - - informer = SharedInformer(list_func=list_func) - - stream_calls = {"n": 0} - - with patch("kubernetes.informer.informer.Watch") as MockWatch: - mock_w = MagicMock() - mock_w.resource_version = "3" - - def fake_stream(func, **kw): - stream_calls["n"] += 1 - if stream_calls["n"] == 1: - raise ApiException(status=410, reason="Gone") - # Second stream (after re-list): stop cleanly - informer._stop_event.set() - return iter([]) - - mock_w.stream.side_effect = fake_stream - MockWatch.return_value = mock_w - - informer.start() - informer._thread.join(timeout=3) - - # list_func called twice: initial list + re-list after 410. - self.assertEqual(list_func.call_count, 2) - # ------------------------------------------------------------------ # Tests analogous to client-go shared_informer_test.go scenarios. # ------------------------------------------------------------------ @@ -1131,6 +1008,767 @@ def fake_stream(func, **kw): self.assertIsNotNone(informer.cache.get_by_key("default/stable-pod")) +class TestSharedInformerWatchCheckpoint(unittest.TestCase): + """Exercise checkpoints with the real Watch parser and retry scheduler.""" + + def setUp(self): + # Timing is verified separately with a fake clock below. + jitter = patch( + "kubernetes.informer.informer.random.uniform", + return_value=0) + jitter.start() + self.addCleanup(jitter.stop) + + def _make_checkpoint_object(self, resource_version): + return { + "metadata": { + "namespace": "default", + "name": "pod", + "resourceVersion": resource_version, + } + } + + def _make_event(self, event_type, resource_version): + return { + "type": event_type, + "object": self._make_checkpoint_object(resource_version), + } + + def _make_watch_response(self, *events): + response = MagicMock() + response.status = 200 + response.stream.return_value = iter( + ( + (json.dumps(event) if isinstance(event, dict) else event) + + "\n" + ).encode() + for event in events + ) + return response + + def _make_list_source(self, responses, initial_items=()): + responses = iter(responses) + requests = [] + + def list_func(**kwargs): + requests.append(kwargs) + if not kwargs.get("watch"): + return SimpleNamespace( + items=list(initial_items), + metadata=SimpleNamespace(resource_version="10"), + ) + return next(responses) + + return list_func, requests + + def _stop_on_error(self, informer): + errors = [] + + def on_error(error): + errors.append(error) + informer.stop() + + informer.add_event_handler(ERROR, on_error) + return errors + + def test_watch_reconnect_uses_applied_event_and_bookmark(self): + for return_type in [None, "V1Pod"]: + with self.subTest(return_type=return_type): + list_func, requests = self._make_list_source( + [ + self._make_watch_response( + self._make_event(ADDED, "11") + ), + self._make_watch_response( + { + "type": BOOKMARK, + "object": { + "metadata": {"resourceVersion": "12"} + }, + } + ), + self._make_watch_response( + self._make_event(MODIFIED, "13") + ), + ] + ) + informer = SharedInformer(list_func) + errors = self._stop_on_error(informer) + checkpoints = [] + bookmark_cache = [] + informer.add_event_handler( + ADDED, + lambda obj: checkpoints.append(informer._resource_version), + ) + + def on_bookmark(obj): + checkpoints.append(informer._resource_version) + bookmark_cache.extend(informer.cache.list_keys()) + + informer.add_event_handler(BOOKMARK, on_bookmark) + informer.add_event_handler( + MODIFIED, lambda obj: informer.stop() + ) + with patch( + "kubernetes.informer.informer.Watch", + side_effect=lambda **kw: Watch(return_type, **kw), + ) as factory: + informer._run_loop() + + self.assertFalse(errors) + # Each reconnect passes through the informer scheduler. + self.assertEqual(factory.call_count, 3) + self.assertEqual( + [ + request["resource_version"] + for request in requests + if request.get("watch") + ], + ["10", "11", "12"], + ) + self.assertEqual(checkpoints, ["11", "12"]) + self.assertEqual(bookmark_cache, ["default/pod"]) + self.assertEqual(informer._resource_version, "13") + + def test_empty_and_invalid_lines_do_not_interrupt_watch(self): + for return_type in [None, "V1Pod"]: + with self.subTest(return_type=return_type): + response = self._make_watch_response( + "", "not json", self._make_event(ADDED, "11") + ) + list_func, requests = self._make_list_source([response]) + informer = SharedInformer(list_func) + errors = self._stop_on_error(informer) + informer.add_event_handler(ADDED, lambda obj: informer.stop()) + + with patch( + "kubernetes.informer.informer.Watch", + lambda **kw: Watch(return_type, **kw), + ): + informer._run_loop() + + self.assertFalse(errors) + # One LIST and one WATCH, with no reconnect. + self.assertEqual(len(requests), 2) + self.assertIsNotNone(informer.cache.get_by_key("default/pod")) + self.assertEqual(informer._resource_version, "11") + response.close.assert_called_once() + response.release_conn.assert_called_once() + + def test_unknown_event_does_not_advance_internal_watch_checkpoint(self): + list_func, requests = self._make_list_source( + [ + self._make_watch_response(self._make_event("UNKNOWN", "11")), + self._make_watch_response(self._make_event(ADDED, "12")), + ] + ) + informer = SharedInformer(list_func) + errors = self._stop_on_error(informer) + informer.add_event_handler(ADDED, lambda obj: informer.stop()) + with patch( + "kubernetes.informer.informer.Watch", + lambda **kw: Watch("V1Pod", **kw) + ): + informer._run_loop() + + self.assertFalse(errors) + self.assertEqual( + [ + request["resource_version"] + for request in requests + if request.get("watch") + ], + ["10", "10"], + ) + self.assertEqual(informer._resource_version, "12") + + def test_cache_failure_reconnects_from_last_applied_event(self): + for return_type in [None, "V1Pod"]: + for event_type in [ADDED, MODIFIED, DELETED]: + with self.subTest( + return_type=return_type, event_type=event_type + ): + list_func, requests = self._make_list_source( + [ + self._make_watch_response( + self._make_event(event_type, "11") + ), + self._make_watch_response( + self._make_event(event_type, "11") + ), + ], + initial_items=( + [self._make_checkpoint_object("10")] + if event_type != ADDED + else [] + ), + ) + failed = False + + def key_func(obj): + nonlocal failed + resource_version = ( + obj["metadata"]["resourceVersion"] + if isinstance(obj, dict) + else obj.metadata.resource_version + ) + if resource_version == "11" and not failed: + failed = True + raise ValueError("cache key failed") + return _meta_namespace_key(obj) + + informer = SharedInformer(list_func, key_func=key_func) + errors = [] + informer.add_event_handler(ERROR, errors.append) + checkpoints = [] + + def on_applied(obj): + resource_version = ( + obj["metadata"]["resourceVersion"] + if isinstance(obj, dict) + else obj.metadata.resource_version + ) + if resource_version == "11": + checkpoints.append(informer._resource_version) + informer.stop() + + informer.add_event_handler(event_type, on_applied) + # Stop on unexpected exhaustion instead of retrying + # an invalid test fixture. + informer.add_event_handler( + ERROR, + lambda error: ( + informer.stop() + if isinstance(error, StopIteration) + else None + ), + ) + with patch( + "kubernetes.informer.informer.Watch", + lambda **kw: Watch(return_type, **kw), + ): + informer._run_loop() + + self.assertEqual(len(errors), 1) + self.assertIsInstance(errors[0], ValueError) + self.assertEqual( + [ + request["resource_version"] + for request in requests + if request.get("watch") + ], + ["10", "10"], + ) + self.assertEqual(checkpoints, ["11"]) + self.assertEqual( + informer.cache.get_by_key("default/pod") is None, + event_type == DELETED, + ) + + def test_handler_failure_preserves_applied_checkpoint_and_other_handlers( + self, + ): + list_func, requests = self._make_list_source( + [ + self._make_watch_response(self._make_event(ADDED, "11")), + self._make_watch_response(self._make_event(MODIFIED, "12")), + ] + ) + informer = SharedInformer(list_func) + errors = self._stop_on_error(informer) + checkpoints = [] + + def failing_handler(obj): + raise ValueError("handler failed after cache application") + + informer.add_event_handler(ADDED, failing_handler) + informer.add_event_handler( + ADDED, lambda obj: checkpoints.append(informer._resource_version) + ) + informer.add_event_handler(MODIFIED, lambda obj: informer.stop()) + informer._run_loop() + + # Handler failures retain the existing log-and-continue policy. + self.assertFalse(errors) + self.assertEqual(checkpoints, ["11"]) + self.assertEqual( + [ + request["resource_version"] + for request in requests + if request.get("watch") + ], + ["10", "11"], + ) + self.assertEqual(informer._resource_version, "12") + + def test_stop_after_deserialization_replays_unapplied_event_on_restart( + self, + ): + list_func, requests = self._make_list_source( + [ + self._make_watch_response(self._make_event(ADDED, "11")), + self._make_watch_response(self._make_event(ADDED, "11")), + ] + ) + informer = SharedInformer(list_func) + errors = self._stop_on_error(informer) + first_watch = Watch("V1Pod") + unmarshal_event = first_watch.unmarshal_event + + def stop_after_unmarshal(data, return_type): + event = unmarshal_event(data, return_type) + # Simulate stop between parsing and cache application. + informer.stop() + return event + + first_watch.unmarshal_event = stop_after_unmarshal + with patch( + "kubernetes.informer.informer.Watch", return_value=first_watch + ): + informer._run_loop() + + # Received, but not applied. + self.assertEqual(first_watch.resource_version, "11") + self.assertEqual(informer.cache.list(), []) + checkpoint_after_stop = informer._resource_version + informer.add_event_handler(ADDED, lambda obj: informer.stop()) + informer._stop_event.clear() + with patch( + "kubernetes.informer.informer.Watch", + lambda **kw: Watch("V1Pod", **kw) + ): + informer._run_loop() + + self.assertFalse(errors) + self.assertEqual(checkpoint_after_stop, "10") + self.assertEqual( + [ + request["resource_version"] + for request in requests + if request.get("watch") + ], + ["10", "10"], + ) + self.assertIsNotNone(informer.cache.get_by_key("default/pod")) + self.assertEqual(informer._resource_version, "11") + + def test_repeated_410_relists_without_restoring_expired_checkpoint(self): + expired = { + "type": ERROR, + "object": { + "code": 410, + "reason": "Gone", + "message": "resource version expired", + }, + } + list_func, requests = self._make_list_source( + [ + self._make_watch_response( + self._make_event(ADDED, "11"), expired + ), + self._make_watch_response(expired), + self._make_watch_response(self._make_event(MODIFIED, "12")), + ] + ) + informer = SharedInformer(list_func) + errors = [] + informer.add_event_handler(ERROR, errors.append) + informer.add_event_handler( + ERROR, + lambda error: ( + informer.stop() + if not isinstance(error, ApiException) + else None + ), + ) + informer.add_event_handler(MODIFIED, lambda obj: informer.stop()) + with patch( + "kubernetes.informer.informer.Watch", + lambda **kw: Watch("V1Pod", **kw) + ): + informer._run_loop() + + self.assertEqual(len(errors), 2) + self.assertTrue(all(isinstance(error, ApiException) + for error in errors)) + self.assertTrue(all(error.status == 410 for error in errors)) + self.assertEqual( + [ + request["resource_version"] + for request in requests + if request.get("watch") + ], + ["10", "10", "10"], + ) + self.assertEqual( + len([request for request in requests if not request.get("watch")]), + 3, + ) + self.assertEqual(informer._resource_version, "12") + + +class TestSharedInformerRetryScheduling(unittest.TestCase): + """Exercise scheduler deadlines without sleeping or replacing Watch.""" + + def setUp(self): + self.now = 0.0 + self.requests = [] + self.waits = [] + self.list_action = None + self.watch_action = None + clock = patch("kubernetes.informer.informer.time.monotonic", + side_effect=lambda: self.now) + jitter = patch("kubernetes.informer.informer.random.uniform", + side_effect=lambda lower, upper: upper) + clock.start() + self.jitter = jitter.start() + self.addCleanup(clock.stop) + self.addCleanup(jitter.stop) + + def _make_informer(self, resync_period=0): + def source(**kwargs): + watching = kwargs.get("watch", False) + self.requests.append(("watch" if watching else "list", + self.now, kwargs)) + if watching: + return self.watch_action() + if self.list_action: + return self.list_action() + return SimpleNamespace( + items=[], metadata=SimpleNamespace(resource_version="10")) + + informer = SharedInformer(source, resync_period=resync_period) + + def wait(timeout): + self.assertGreater(timeout, 0) + self.waits.append(timeout) + self.assertLess(len(self.waits), 30, "unexpected retry loop") + self.now += timeout + return False + + informer._stop_event.wait = MagicMock(side_effect=wait) + return informer + + def _response(self, duration=0, events=(), error=None): + def chunks(): + self.now += duration + for event in events: + line = json.dumps(event) if isinstance(event, dict) else event + yield (line + "\n").encode() + if error: + raise error + + response = MagicMock() + response.status = 200 + response.stream.side_effect = lambda **kwargs: chunks() + return response + + def _event(self, event_type=BOOKMARK): + return {"type": event_type, + "object": {"metadata": {"name": "pod", + "resourceVersion": "11"}}} + + def test_short_streams_back_off_even_after_events_or_bookmarks(self): + for event_type in [None, ADDED, BOOKMARK]: + with self.subTest(event_type=event_type): + self.now = 0 + self.requests.clear() + self.waits.clear() + informer = self._make_informer() + count = [0] + + def watch(): + count[0] += 1 + if count[0] == 4: + informer._stop_event.set() + events = [self._event(event_type)] if event_type else [] + return self._response(events=events) + + self.watch_action = watch + informer._run_loop() + self.assertEqual(self.waits, [1, 2, 4]) + self.assertEqual( + [t for op, t, _ in self.requests if op == "watch"], + [0, 1, 3, 7]) + + def test_slow_watch_errors_keep_failure_history_and_cap_jitter(self): + informer = self._make_informer() + count = [0] + + def watch(): + count[0] += 1 + if count[0] == 9: + informer._stop_event.set() + return self._response() + return self._response( + duration=70, error=RuntimeError("read failed")) + + self.watch_action = watch + informer._run_loop() + self.assertEqual(self.waits, [1, 2, 4, 8, 16, 32, 60, 60]) + self.assertEqual( + [entry.args for entry in self.jitter.call_args_list[:8]], + [(0.5, 1), (1, 2), (2, 4), (4, 8), (8, 16), + (16, 32), (30, 60), (30, 60)]) + + def test_long_clean_watch_resets_backoff(self): + informer = self._make_informer() + count = [0] + + def watch(): + count[0] += 1 + if count[0] == 5: + informer._stop_event.set() + return self._response(duration=60 if count[0] == 3 else 0) + + self.watch_action = watch + informer._run_loop() + self.assertEqual(self.waits, [1, 2, 1]) + self.assertEqual([t for op, t, _ in self.requests if op == "watch"], + [0, 1, 3, 63, 64]) + + def test_backoff_wakes_for_periodic_list_before_next_watch(self): + informer = self._make_informer(resync_period=5) + count = [0] + + def watch(): + count[0] += 1 + if count[0] == 4: + informer._stop_event.set() + return self._response() + + self.watch_action = watch + informer._run_loop() + self.assertEqual([(op, t) for op, t, _ in self.requests], + [("list", 0), ("watch", 0), ("watch", 1), + ("watch", 3), ("list", 5), ("watch", 7)]) + self.assertEqual([kw["timeout_seconds"] for op, _, kw in self.requests + if op == "watch"], [5, 4, 2, 3]) + + def test_periodic_list_failures_back_off_while_watch_keeps_delivering( + self): + informer = self._make_informer(resync_period=2) + list_count = [0] + + def list_objects(): + list_count[0] += 1 + if list_count[0] == 1: + return SimpleNamespace( + items=[], metadata=SimpleNamespace(resource_version="10")) + if list_count[0] == 5: + informer._stop_event.set() + raise RuntimeError("list unavailable") + + self.list_action = list_objects + self.watch_action = lambda: self._response( + duration=1, events=[self._event()]) + informer._run_loop() + self.assertEqual([t for op, t, _ in self.requests if op == "list"], + [0, 2, 3, 5, 9]) + self.assertEqual(informer._resource_version, "11") + self.assertGreater( + len([op for op, _, _ in self.requests if op == "watch"]), 3) + + def test_initial_and_expired_list_failures_prevent_watch(self): + for expired in [False, True]: + with self.subTest(expired=expired): + self.now = 0 + self.requests.clear() + self.waits.clear() + informer = self._make_informer() + count = [0] + + def list_objects(): + count[0] += 1 + if expired and count[0] == 1: + return SimpleNamespace( + items=[], metadata=SimpleNamespace( + resource_version="10")) + if count[0] == (4 if expired else 3): + informer._stop_event.set() + raise RuntimeError("list unavailable") + + self.list_action = list_objects + self.watch_action = lambda: self._response(events=[{ + "type": ERROR, "object": {"code": 410, "reason": "Gone", + "message": "expired"}}]) + informer._run_loop() + self.assertEqual(len([op for op, _, _ in self.requests + if op == "watch"]), int(expired)) + self.assertEqual(self.waits, [1, 2]) + self.assertIsNone(informer._resource_version) + + def test_slow_relist_does_not_make_short_watch_healthy(self): + informer = self._make_informer(resync_period=2) + list_count = [0] + watch_count = [0] + + def list_objects(): + list_count[0] += 1 + if list_count[0] > 1: + self.now += 70 + return SimpleNamespace(items=[], metadata=SimpleNamespace( + resource_version="10")) + + def watch(): + watch_count[0] += 1 + if watch_count[0] == 4: + informer._stop_event.set() + return self._response() + + self.list_action = list_objects + self.watch_action = watch + informer._run_loop() + # Delays keep growing across the slow LIST, and resync starts anew + # from each LIST completion, not its start time. + self.assertEqual([t for op, t, _ in self.requests if op == "watch"], + [0, 1, 72, 144]) + self.assertEqual([t for op, t, _ in self.requests if op == "list"], + [0, 2, 74]) + self.assertEqual( + [entry.args for entry in self.jitter.call_args_list[:3]], + [(0.5, 1), (1, 2), (2, 4)]) + + def test_fractional_period_busy_watch_does_not_accumulate_backoff(self): + informer = self._make_informer(resync_period=0.5) + count = [0] + + def watch(): + count[0] += 1 + if count[0] == 4: + informer._stop_event.set() + return self._response(duration=0.5, events=[self._event()]) + + self.watch_action = watch + informer._run_loop() + self.assertEqual([t for op, t, _ in self.requests if op == "watch"], + [0, 0.5, 1, 1.5]) + self.assertEqual(self.waits, []) + + def test_planned_relist_preserves_previous_watch_failure_history(self): + informer = self._make_informer(resync_period=0.5) + count = [0] + + def watch(): + count[0] += 1 + if count[0] == 4: + informer._stop_event.set() + if count[0] == 2: + return self._response(duration=0.5, events=[self._event()]) + return self._response() + + self.watch_action = watch + informer._run_loop() + self.assertEqual([t for op, t, _ in self.requests if op == "watch"], + [0, 1, 1.5, 3.5]) + self.assertEqual( + [entry.args for entry in self.jitter.call_args_list[:2]], + [(0.5, 1), (1, 2)]) + + def test_normal_idle_timeout_relists_without_failure_delay(self): + informer = self._make_informer(resync_period=2) + count = [0] + + def watch(): + count[0] += 1 + if count[0] == 2: + informer._stop_event.set() + return self._response(duration=2) + + self.watch_action = watch + informer._run_loop() + self.assertEqual([(op, t) for op, t, _ in self.requests], [ + ("list", 0), ("watch", 0), ("list", 2), ("watch", 2)]) + self.assertEqual(self.waits, []) + + def test_cancel_during_backoff_prevents_another_request(self): + informer = self._make_informer() + self.watch_action = lambda: self._response() + informer._stop_event.wait.side_effect = ( + lambda timeout: informer._stop_event.set()) + informer._run_loop() + self.assertEqual([op for op, _, _ in self.requests], ["list", "watch"]) + + def test_invalid_lines_followed_by_410_relist_without_internal_retry(self): + informer = self._make_informer() + count = [0] + + def watch(): + count[0] += 1 + if count[0] == 3: + informer._stop_event.set() + return self._response() + return self._response(events=["", "not json", { + "type": ERROR, "object": {"code": 410, "reason": "Gone", + "message": "expired"}}]) + + self.watch_action = watch + informer._run_loop() + self.assertEqual([op for op, _, _ in self.requests], + ["list", "watch", "list", "watch", "list", "watch"]) + self.assertEqual(self.waits, [1, 2]) + + def test_restart_preserves_periodic_list_deadline(self): + informer = self._make_informer(resync_period=5) + count = [0] + + def watch(): + count[0] += 1 + if count[0] in (1, 3): + informer._stop_event.set() + return self._response() + return self._response(duration=4) + + self.watch_action = watch + informer._run_loop() + self.now = 1 + informer._stop_event.clear() + informer._run_loop() + self.assertEqual([(op, t) for op, t, _ in self.requests], + [("list", 0), ("watch", 0), ("watch", 1), + ("list", 5), ("watch", 5)]) + self.assertEqual([kw["timeout_seconds"] for op, _, kw in self.requests + if op == "watch"], [5, 4, 5]) + + def test_response_close_failure_at_resync_is_retried(self): + informer = self._make_informer(resync_period=1) + errors = [] + informer.add_event_handler(ERROR, errors.append) + response = self._response(duration=1, events=[self._event()]) + response.close.side_effect = RuntimeError("close failed") + count = [0] + + def watch(): + count[0] += 1 + if count[0] == 1: + return response + informer._stop_event.set() + return self._response() + + self.watch_action = watch + informer._run_loop() + self.assertEqual(len(errors), 1) + self.assertEqual(str(errors[0]), "close failed") + response.release_conn.assert_called_once() + self.assertEqual(count[0], 2) + self.assertIsNone(informer._watch) + + def test_default_resync_accepts_callable_without_timeout_parameter(self): + informer = None + requests = [] + + def source(watch=False, resource_version=None, _preload_content=True): + requests.append(watch) + if watch: + informer._stop_event.set() + return self._response() + return SimpleNamespace(items=[], metadata=SimpleNamespace( + resource_version="10")) + + informer = SharedInformer(source) + informer._run_loop() + self.assertEqual(requests, [False, True]) + + if __name__ == "__main__": unittest.main() - diff --git a/kubernetes/watch/watch.py b/kubernetes/watch/watch.py index ce7174a57f..46fa5bd851 100644 --- a/kubernetes/watch/watch.py +++ b/kubernetes/watch/watch.py @@ -52,6 +52,20 @@ def _find_return_type(func): return "" +def _event_resource_version(event): + """Return a valid checkpoint from a supported event's raw metadata.""" + if not isinstance(event, dict) or event.get('type') not in ( + 'ADDED', 'MODIFIED', 'DELETED', 'BOOKMARK'): + return None + obj = event.get('raw_object', event.get('object')) + metadata = obj.get('metadata') if isinstance(obj, dict) else None + if isinstance(metadata, dict): + resource_version = metadata.get('resourceVersion') + if isinstance(resource_version, str) and resource_version: + return resource_version + return None + + def iter_resp_lines(resp): buffer = bytearray() for segment in resp.stream(amt=None, decode_content=False): @@ -86,7 +100,14 @@ def iter_resp_lines(resp): class Watch: - def __init__(self, return_type=None): + def __init__(self, return_type=None, retry=True): + """Create a watch, optionally disabling automatic reconnection. + + With ``retry=False``, each ``stream()`` makes at most one API call: + EOF ends the iterator and an ERROR event raises ``ApiException``. + This does not configure HTTP transport retries or request timeouts. + """ + self._retry = retry self._raw_return_type = return_type self._stop = False self._api_client = client.ApiClient() @@ -143,27 +164,19 @@ def unmarshal_event(self, data, return_type): if not return_type: return js - if js['type'] == 'BOOKMARK': - # Extract and store resource_version from BOOKMARK event for - # efficiency. No deserialization as event can be incomplete. - if isinstance(js['object'], dict) and 'metadata' in js['object']: - metadata = js['object']['metadata'] - if isinstance(metadata, dict) and 'resourceVersion' in metadata: - self.resource_version = metadata['resourceVersion'] - elif js['type'] != 'ERROR': + # BOOKMARK objects can be incomplete; preserve their raw form. + if js['type'] not in ('ERROR', 'BOOKMARK'): js['object'] = self._api_client.deserialize( json.dumps(js['raw_object']), return_type, 'application/json', ) - if hasattr(js['object'], 'metadata'): - self.resource_version = js['object'].metadata.resource_version - # For custom objects that we don't have model defined, json - # deserialization results in dictionary - elif (isinstance(js['object'], dict) and 'metadata' in js['object'] - and 'resourceVersion' in js['object']['metadata']): - self.resource_version = js['object']['metadata'][ - 'resourceVersion'] + + # Keep the previous checkpoint when raw metadata is missing or + # invalid, even if a model decoder would coerce its value. + resource_version = _event_resource_version(js) + if resource_version is not None: + self.resource_version = resource_version return js except json.JSONDecodeError: return None @@ -177,7 +190,8 @@ def stream(self, func, *args, **kwargs): ``code`` 410. In that case you have to recover yourself, probably by listing the API resource to obtain the latest state and then watching from that state on by setting ``resource_version`` to - one returned from listing. + one returned from listing. Events without a valid resource version + preserve the last checkpoint and do not reset the 410 retry limit. :param func: The API function pointer. Any parameter to the function can be passed after this parameter. @@ -212,7 +226,7 @@ def stream(self, func, *args, **kwargs): # Do not attempt retries if user specifies a timeout. # We want to ensure we are returning within that timeout. - disable_retries = ('timeout_seconds' in kwargs) + disable_retries = not self._retry or 'timeout_seconds' in kwargs retry_after_410 = False deserialize = kwargs.pop('deserialize', True) while True: @@ -246,19 +260,24 @@ def stream(self, func, *args, **kwargs): raise client.rest.ApiException( status=obj['code'], reason=reason) else: - retry_after_410 = False + if _event_resource_version(event) is not None: + retry_after_410 = False yield event else: - if line: + if line: yield line # Normal non-empty line - else: - yield '' # Only yield one empty line + else: + yield '' # Only yield one empty line if self._stop: break finally: - resp.close() - resp.release_conn() - self._resp = None + try: + resp.close() + finally: + try: + resp.release_conn() + finally: + self._resp = None if self.resource_version is not None: kwargs['resource_version'] = self.resource_version else: diff --git a/kubernetes/watch/watch_test.py b/kubernetes/watch/watch_test.py index 592678bd93..909c29dd75 100644 --- a/kubernetes/watch/watch_test.py +++ b/kubernetes/watch/watch_test.py @@ -15,9 +15,9 @@ import json import os import time +import unittest from types import SimpleNamespace from typing import Any, Optional -import unittest from unittest.mock import Mock, call from kubernetes import client, config @@ -126,8 +126,8 @@ def test_watch_with_interspersed_newlines(self): count = 0 # Consume all test events from the mock service, stopping when no more data is available. - # Note that "timeout_seconds" below is not a timeout; rather, it disables retries and is - # the only way to do so. Without that, the stream will re-read the test data forever. + # "timeout_seconds" also disables retries. Without it or + # retry=False, the stream will re-read the test data forever. for e in w.stream(fake_api.get_namespaces, timeout_seconds=1): # Here added a statement for exception for empty lines. if e is None: @@ -163,8 +163,8 @@ def test_watch_with_multibyte_utf8(self): count = 0 # Consume all test events from the mock service, stopping when no more data is available. - # Note that "timeout_seconds" below is not a timeout; rather, it disables retries and is - # the only way to do so. Without that, the stream will re-read the test data forever. + # "timeout_seconds" also disables retries. Without it or + # retry=False, the stream will re-read the test data forever. for event in w.stream(fake_api.get_configmaps, timeout_seconds=1): count += 1 self.assertEqual("MODIFIED", event['type']) @@ -209,8 +209,8 @@ def test_watch_with_invalid_utf8(self): count = 0 # Consume all test events from the mock service, stopping when no more data is available. - # Note that "timeout_seconds" below is not a timeout; rather, it disables retries and is - # the only way to do so. Without that, the stream will re-read the test data forever. + # "timeout_seconds" also disables retries. Without it or + # retry=False, the stream will re-read the test data forever. for event in w.stream(fake_api.get_configmaps, timeout_seconds=1): count += 1 self.assertEqual("MODIFIED", event['type']) @@ -483,6 +483,251 @@ def test_watch_with_error_event(self): fake_resp.close.assert_called_once() fake_resp.release_conn.assert_called_once() + def test_watch_retry_disabled_stops_at_eof_without_forwarding_option(self): + response = Mock() + response.stream.return_value = [ + '{"type":"ADDED","object":{"metadata":' + '{"name":"test","resourceVersion":"2"}}}\n' + ] + requests = [] + + def list_namespaces(*, resource_version, watch, _preload_content): + """:rtype: V1NamespaceList""" + requests.append((resource_version, watch, _preload_content)) + if len(requests) > 1: + self.fail("Watch retried after EOF with retry=False") + return response + + watch = Watch(retry=False) + events = list(watch.stream(list_namespaces, resource_version="1")) + + self.assertEqual([("1", True, False)], requests) + self.assertEqual(["ADDED"], [event["type"] for event in events]) + self.assertEqual("2", watch.resource_version) + response.close.assert_called_once() + response.release_conn.assert_called_once() + + def test_watch_retry_disabled_raises_first_expired_event(self): + response = Mock() + response.stream.return_value = [ + '{"type":"ERROR","object":{"code":410,"reason":"Gone",' + '"message":"expired"}}\n' + ] + list_namespaces = Mock(side_effect=[ + response, AssertionError("Watch retried with retry=False") + ]) + + with self.assertRaises(ApiException) as caught: + list(Watch(retry=False).stream( + list_namespaces, resource_version="1")) + + self.assertEqual(410, caught.exception.status) + list_namespaces.assert_called_once_with( + resource_version="1", watch=True, _preload_content=False) + response.close.assert_called_once() + response.release_conn.assert_called_once() + + def test_watch_invalid_events_do_not_reset_expired_retry(self): + expired = ( + '{"type":"ERROR","object":{"code":410,"reason":"Gone",' + '"message":"expired"}}\n' + ) + for ignored_line in ( + '\n', 'not-json\n', + '{"type":"UNKNOWN","object":{"metadata":' + '{"resourceVersion":"1"}}}\n'): + with self.subTest(line=ignored_line): + responses = [Mock(), Mock()] + for response in responses: + response.stream.return_value = [ignored_line, expired] + list_namespaces = Mock(side_effect=responses + [ + AssertionError("Invalid input reset the 410 retry limit") + ]) + list_namespaces.__doc__ = ':rtype: V1NamespaceList' + + with self.assertRaises(ApiException) as caught: + list(Watch().stream( + list_namespaces, resource_version="1")) + + self.assertEqual(410, caught.exception.status) + self.assertEqual(2, list_namespaces.call_count) + for response in responses: + response.close.assert_called_once() + response.release_conn.assert_called_once() + + def test_watch_valid_events_reset_expired_retry(self): + expired = ( + '{"type":"ERROR","object":{"code":410,"reason":"Gone",' + '"message":"expired"}}\n' + ) + for event_type in ('ADDED', 'MODIFIED', 'DELETED', 'BOOKMARK'): + for resource_version in ('1', '2'): + with self.subTest(event_type=event_type, rv=resource_version): + responses = [Mock(), Mock(), Mock()] + responses[0].stream.return_value = [expired] + responses[1].stream.return_value = [ + json.dumps({ + "type": event_type, + "object": {"metadata": { + "name": "test", + "resourceVersion": resource_version}} + }) + '\n', + expired + ] + responses[2].stream.return_value = [expired] + list_namespaces = Mock(side_effect=responses + [ + AssertionError("Watch retried consecutive 410 errors") + ]) + list_namespaces.__doc__ = ':rtype: V1NamespaceList' + received = [] + + with self.assertRaises(ApiException) as caught: + for event in Watch().stream( + list_namespaces, resource_version="1"): + received.append(event["type"]) + + self.assertEqual(410, caught.exception.status) + self.assertEqual([event_type], received) + self.assertEqual([ + call(resource_version="1", watch=True, + _preload_content=False), + call(resource_version="1", watch=True, + _preload_content=False), + call(resource_version=resource_version, watch=True, + _preload_content=False) + ], list_namespaces.call_args_list) + for response in responses: + response.close.assert_called_once() + response.release_conn.assert_called_once() + + def test_unmarshal_preserves_checkpoint_for_unsupported_or_invalid_rv( + self): + invalid_metadata = [ + None, {}, {"resourceVersion": None}, {"resourceVersion": ""}, + {"resourceVersion": 0}, {"resourceVersion": 42}, + ] + for return_type in ("V1Namespace", "object"): + for event_type in ("ADDED", "MODIFIED", "DELETED", "BOOKMARK"): + for metadata in invalid_metadata: + with self.subTest( + return_type=return_type, event_type=event_type, + metadata=metadata): + watch = Watch(return_type) + watch.resource_version = "1" + data = json.dumps({ + "type": event_type, + "object": {"metadata": metadata}, + }) + if (return_type == "V1Namespace" + and event_type != "BOOKMARK" + and metadata in ({"resourceVersion": 0}, + {"resourceVersion": 42})): + # Keep the model decoder's strict validation. + with self.assertRaises(ValueError): + watch.unmarshal_event(data, return_type) + self.assertEqual(watch.resource_version, "1") + continue + event = watch.unmarshal_event(data, return_type) + self.assertEqual(watch.resource_version, "1") + self.assertEqual(event["type"], event_type) + with self.subTest(return_type=return_type, event_type="UNKNOWN"): + watch = Watch(return_type) + watch.resource_version = "1" + event = watch.unmarshal_event(json.dumps({ + "type": "UNKNOWN", + "object": {"metadata": {"resourceVersion": "2"}}, + }), return_type) + self.assertEqual(watch.resource_version, "1") + if return_type == "V1Namespace": + self.assertIsInstance(event["object"], client.V1Namespace) + else: + self.assertIsInstance(event["object"], dict) + + def test_unmarshal_valid_checkpoints_and_raw_return_type_contract(self): + for return_type in ("V1Namespace", "object", None): + for event_type in ("ADDED", "MODIFIED", "DELETED", "BOOKMARK"): + with self.subTest(return_type=return_type, + event_type=event_type): + watch = Watch(return_type) + watch.resource_version = "1" + obj = {"metadata": {"resourceVersion": "2"}} + event = watch.unmarshal_event(json.dumps({ + "type": event_type, "object": obj, + }), return_type) + self.assertEqual(watch.resource_version, + "2" if return_type else "1") + self.assertEqual(event["raw_object"], obj) + if (return_type == "V1Namespace" + and event_type != "BOOKMARK"): + self.assertIsInstance( + event["object"], client.V1Namespace) + else: + self.assertEqual(event["object"], obj) + + def test_reconnect_preserves_checkpoint_after_invalid_event_rv(self): + cases = [ + ("UNKNOWN", {"resourceVersion": "2"}), + ("ADDED", {}), + ("BOOKMARK", {"resourceVersion": None}), + ] + for return_type in ("V1Namespace", "object"): + for event_type, metadata in cases: + with self.subTest(return_type=return_type, + event_type=event_type): + responses = [Mock(), Mock()] + responses[0].stream.return_value = [json.dumps({ + "type": event_type, + "object": {"metadata": metadata}, + }) + "\n"] + responses[1].stream.return_value = [json.dumps({ + "type": "ADDED", + "object": {"metadata": {"resourceVersion": "3"}}, + }) + "\n"] + source = Mock(side_effect=responses + [ + AssertionError("unexpected third request")]) + watch = Watch(return_type) + events = [] + for event in watch.stream(source, resource_version="1"): + events.append(event) + if len(events) == 2: + watch.stop() + + self.assertEqual([event_type, "ADDED"], + [event["type"] for event in events]) + self.assertEqual(watch.resource_version, "3") + self.assertEqual(source.call_args_list, [ + call(resource_version="1", watch=True, + _preload_content=False)] * 2) + for response in responses: + response.close.assert_called_once() + response.release_conn.assert_called_once() + + def test_invalid_rv_does_not_reset_consecutive_expired_retry(self): + expired = {"type": "ERROR", "object": { + "code": 410, "reason": "Gone", "message": "expired"}} + cases = (("ADDED", {}), ("BOOKMARK", {"resourceVersion": ""})) + for return_type in ("V1Namespace", "object"): + for event_type, metadata in cases: + with self.subTest(return_type=return_type, + event_type=event_type): + responses = [Mock(), Mock()] + for response in responses: + response.stream.return_value = [ + json.dumps({"type": event_type, "object": { + "metadata": metadata}}) + "\n", + json.dumps(expired) + "\n", + ] + source = Mock(side_effect=responses + [ + AssertionError("Invalid RV reset 410 retry")]) + watch = Watch(return_type) + with self.assertRaises(ApiException) as caught: + list(watch.stream(source, resource_version="1")) + self.assertEqual(caught.exception.status, 410) + self.assertEqual(watch.resource_version, "1") + self.assertEqual(source.call_args_list, [ + call(resource_version="1", watch=True, + _preload_content=False)] * 2) + def test_watch_retries_on_error_event(self): fake_resp = Mock() fake_resp.close = Mock()