1111"""
1212
1313import asyncio
14+ import threading
15+ from concurrent .futures import Future , ThreadPoolExecutor
1416from typing import Any
1517
1618from aiolimiter import AsyncLimiter
@@ -54,7 +56,6 @@ def __init__(
5456 watch_timeout_seconds: Watch timeout before reconnect (default: 60)
5557 watch_reconnect_delay_seconds: Delay after watch failure (default: 5)
5658 """
57- self ._api_client = api_client
5859 self ._group = group
5960 self ._version = version
6061 self ._plural = plural
@@ -73,6 +74,11 @@ def __init__(
7374 self ._cache_lock = asyncio .Lock ()
7475 self ._watch_task = None
7576 self ._initialized = False
77+ self ._event_queue : asyncio .Queue | None = None
78+ self ._watch_executor = ThreadPoolExecutor (max_workers = 1 , thread_name_prefix = "k8s-watch-" )
79+ self ._stop_event = threading .Event ()
80+ self ._resource_version : str | None = None
81+ self ._resource_version_lock = threading .Lock ()
7682
7783 async def start (self ):
7884 """Start the API client and initialize cache watch."""
@@ -105,74 +111,167 @@ async def _list_and_sync_cache(self) -> str:
105111 name = item .get ("metadata" , {}).get ("name" )
106112 if name :
107113 self ._cache [name ] = item
108- return resource_version
114+ self ._advance_resource_version (resource_version )
115+ return self ._resource_version
109116
110- async def _watch_resources (self ):
111- """Background task to watch resources and maintain cache .
117+ def _watch_in_thread (self , resource_version : str | None , event_queue : asyncio . Queue , loop : asyncio . AbstractEventLoop ):
118+ """Watch resources in a background thread and put events to queue .
112119
113- Implements Kubernetes Informer pattern:
114- 1. Initial list-and-sync to populate cache
115- 2. Continuous watch for ADDED/MODIFIED/DELETED events
116- 3. Auto-reconnect on watch timeout or network failures
117- 4. Re-sync on reconnect to avoid event loss
120+ Args:
121+ resource_version: The resource version to watch from
122+ event_queue: Queue to put events for async processing
123+ loop: The asyncio event loop to use for queue operations
118124 """
119- resource_version = None
120125 try :
121- resource_version = await self ._list_and_sync_cache ()
122- logger .info (
123- f"Initial cache populated with { len (self ._cache )} resources, resourceVersion={ resource_version } "
124- )
126+ w = watch .Watch ()
127+ for event in w .stream (
128+ self ._custom_api .list_namespaced_custom_object ,
129+ group = self ._group ,
130+ version = self ._version ,
131+ namespace = self ._namespace ,
132+ plural = self ._plural ,
133+ resource_version = resource_version ,
134+ timeout_seconds = self ._watch_timeout_seconds ,
135+ ):
136+ if self ._stop_event .is_set ():
137+ break
138+ # Put event to queue for async processing
139+ asyncio .run_coroutine_threadsafe (event_queue .put (("event" , event )), loop )
125140 except Exception as e :
126- logger .error (f"Failed to populate initial cache: { e } " )
141+ logger .warning (f"Watch thread error: { e } " )
142+ finally :
143+ # Signal thread exit
144+ asyncio .run_coroutine_threadsafe (event_queue .put (("exit" , None )), loop )
127145
128- while True :
146+ def _advance_resource_version (self , rv : str | None ) -> bool :
147+ """Advance resource_version only when rv is strictly newer.
148+
149+ K8s resourceVersions are opaque strings but etcd encodes them as
150+ monotonically increasing integers. Prevents stale responses from
151+ rolling back the watch cursor.
152+
153+ Returns:
154+ True if version was updated, False if skipped (old version or error)
155+ """
156+ if not rv :
157+ return False
158+ with self ._resource_version_lock :
159+ if self ._resource_version is None :
160+ self ._resource_version = rv
161+ return True
129162 try :
163+ if int (rv ) > int (self ._resource_version ):
164+ self ._resource_version = rv
165+ return True
166+ return False
167+ except ValueError :
168+ # Non-integer resourceVersion — skip to avoid downgrade
169+ logger .error (f"Non-integer resourceVersion detected: rv={ rv } , current={ self ._resource_version } , skipping" )
170+ return False
171+
172+ async def _process_events (self , event_queue : asyncio .Queue , thread_future : Future ) -> str | None :
173+ """Process events from queue and update cache.
130174
131- def _watch_in_thread ():
132- w = watch .Watch ()
133- stream = w .stream (
134- self ._custom_api .list_namespaced_custom_object ,
135- group = self ._group ,
136- version = self ._version ,
137- namespace = self ._namespace ,
138- plural = self ._plural ,
139- resource_version = resource_version ,
140- timeout_seconds = self ._watch_timeout_seconds ,
141- )
142- events = []
143- for event in stream :
144- events .append (event )
145- return events
175+ Args:
176+ event_queue: Queue containing events from watch thread
177+ thread_future: Future representing the watch thread
178+
179+ Returns:
180+ Latest resource version or None if disconnected
181+ """
182+ while True :
183+ try :
184+ msg_type , event = await asyncio .wait_for (event_queue .get (), timeout = 1.0 )
146185
147- events = await asyncio .to_thread (_watch_in_thread )
186+ if msg_type == "exit" : # Watch thread exited
187+ break
148188
149- async with self ._cache_lock :
150- for event in events :
151- event_type = event ["type" ]
152- obj = event ["object" ]
153- name = obj .get ("metadata" , {}).get ("name" )
154- new_rv = obj .get ("metadata" , {}).get ("resourceVersion" )
189+ if msg_type == "event" :
190+ event_type = event ["type" ]
191+ obj = event ["object" ]
192+ name = obj .get ("metadata" , {}).get ("name" )
193+ new_rv = obj .get ("metadata" , {}).get ("resourceVersion" )
155194
156- if new_rv :
157- resource_version = new_rv
195+ # Skip event if no resourceVersion or version is stale
196+ if not new_rv :
197+ continue
198+ if not self ._advance_resource_version (new_rv ):
199+ continue
158200
159- if not name :
160- continue
201+ if not name :
202+ continue
161203
204+ async with self ._cache_lock :
162205 if event_type in ["ADDED" , "MODIFIED" ]:
163206 self ._cache [name ] = obj
164207 elif event_type == "DELETED" :
165208 self ._cache .pop (name , None )
166209
210+ logger .debug (f"Cache updated: { event_type } { name } , rv={ self ._resource_version } " )
211+
212+ except asyncio .TimeoutError :
213+ # Check if thread has exited
214+ if thread_future .done () and event_queue .empty ():
215+ break
216+ continue
217+ except asyncio .CancelledError :
218+ raise
219+
220+ return self ._resource_version
221+
222+ async def _watch_resources (self ):
223+ """Background task to watch resources and maintain cache.
224+
225+ Implements Kubernetes Informer pattern:
226+ 1. Initial list-and-sync to populate cache
227+ 2. Continuous watch for ADDED/MODIFIED/DELETED events (real-time)
228+ 3. Auto-reconnect on watch timeout or network failures
229+ 4. Re-sync on reconnect to avoid event loss
230+ """
231+ try :
232+ await self ._list_and_sync_cache ()
233+ logger .info (
234+ f"Initial cache populated with { len (self ._cache )} resources, resourceVersion={ self ._resource_version } "
235+ )
236+ except Exception as e :
237+ logger .error (f"Failed to populate initial cache: { e } " )
238+
239+ self ._event_queue = asyncio .Queue ()
240+ loop = asyncio .get_event_loop ()
241+
242+ while True :
243+ try :
244+ self ._stop_event .clear ()
245+
246+ # Start watch in background thread
247+ thread_future = self ._watch_executor .submit (
248+ self ._watch_in_thread ,
249+ self ._resource_version ,
250+ self ._event_queue ,
251+ loop
252+ )
253+
254+ # Process events in real-time
255+ await self ._process_events (self ._event_queue , thread_future )
256+
257+ # Wait for thread to complete
258+ if not thread_future .done ():
259+ self ._stop_event .set ()
260+ try :
261+ thread_future .result (timeout = 5.0 )
262+ except Exception :
263+ pass
264+
167265 except asyncio .CancelledError :
168266 logger .info ("Watch task cancelled" )
267+ self ._stop_event .set ()
169268 raise
170269 except Exception as e :
171270 logger .warning (f"Watch stream disconnected: { e } , reconnecting immediately..." )
172271 try :
173- resource_version = await self ._list_and_sync_cache ()
272+ await self ._list_and_sync_cache ()
174273 logger .info (
175- f"Re-synced cache with { len (self ._cache )} resources, resourceVersion={ resource_version } "
274+ f"Re-synced cache with { len (self ._cache )} resources, resourceVersion={ self . _resource_version } "
176275 )
177276 except Exception as list_err :
178277 logger .error (
0 commit comments