e0e362d700
SDK Tests / changes (push) Successful in 2m29s
Real E2E Tests / changes (push) Successful in 2m29s
Deploy Docs Pages / build (push) Has been cancelled
Deploy Docs Pages / deploy (push) Has been cancelled
Real E2E Tests / JavaScript E2E (docker bridge) (push) Has been cancelled
Real E2E Tests / Python E2E (docker bridge) (push) Has been cancelled
Real E2E Tests / Java E2E (docker bridge) (push) Has been cancelled
Real E2E Tests / C# E2E (docker bridge) (push) Has been cancelled
Real E2E Tests / Go E2E (docker bridge) (push) Has been cancelled
Real E2E Tests / Real E2E CI (push) Has been cancelled
SDK Tests / SDK CI (push) Has been cancelled
SDK Tests / CLI Tests (push) Has been cancelled
SDK Tests / Python SDK Quality (code-interpreter) (push) Has been cancelled
SDK Tests / Python SDK Quality (sandbox) (push) Has been cancelled
SDK Tests / Python SDK Tests (code-interpreter) (push) Has been cancelled
SDK Tests / JavaScript SDK Quality And Tests (code-interpreter) (push) Has been cancelled
SDK Tests / JavaScript SDK Quality And Tests (sandbox) (push) Has been cancelled
SDK Tests / Python SDK Tests (sandbox) (push) Has been cancelled
SDK Tests / CLI Quality (push) Has been cancelled
SDK Tests / Kotlin SDK Quality And Tests (sandbox) (push) Has been cancelled
SDK Tests / Kotlin SDK Quality And Tests (code-interpreter) (push) Has been cancelled
SDK Tests / C# SDK Quality And Tests (code-interpreter) (push) Has been cancelled
SDK Tests / C# SDK Quality And Tests (sandbox) (push) Has been cancelled
SDK Tests / Go SDK Quality And Tests (push) Has been cancelled
224 lines
8.0 KiB
Python
224 lines
8.0 KiB
Python
# Copyright 2025 Alibaba Group Holding Ltd.
|
|
#
|
|
# Licensed under the Apache License, Version 2.0 (the "License");
|
|
# you may not use this file except in compliance with the License.
|
|
# You may obtain a copy of the License at
|
|
#
|
|
# http://www.apache.org/licenses/LICENSE-2.0
|
|
#
|
|
# Unless required by applicable law or agreed to in writing, software
|
|
# distributed under the License is distributed on an "AS IS" BASIS,
|
|
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
# See the License for the specific language governing permissions and
|
|
# limitations under the License.
|
|
|
|
"""Lightweight informer-style cache for namespaced custom resources."""
|
|
|
|
import logging
|
|
import threading
|
|
from typing import Any, Callable, Dict, List, Optional
|
|
|
|
from kubernetes import watch
|
|
from kubernetes.client import ApiException
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
class WorkloadInformer:
|
|
"""Maintain an in-memory cache of a namespaced custom resource via watch."""
|
|
|
|
def __init__(
|
|
self,
|
|
list_fn: Callable[..., Any],
|
|
resync_period_seconds: int = 300,
|
|
watch_timeout_seconds: int = 60,
|
|
enable_watch: bool = True,
|
|
thread_name: str = "workload-informer",
|
|
):
|
|
"""
|
|
Args:
|
|
list_fn: Callable that lists the custom resource, with signature
|
|
``list_fn(**kwargs) -> dict``. Typically a bound method
|
|
like ``custom_api.list_namespaced_custom_object``.
|
|
resync_period_seconds: Full-resync interval when watch is disabled.
|
|
watch_timeout_seconds: Per-stream watch timeout before restart.
|
|
enable_watch: When False only the initial list is performed.
|
|
thread_name: Name for the background thread, used in stack traces
|
|
and debuggers. Should be unique per informer instance.
|
|
"""
|
|
self.list_fn = list_fn
|
|
self.resync_period_seconds = resync_period_seconds
|
|
self.watch_timeout_seconds = watch_timeout_seconds
|
|
self.enable_watch = enable_watch
|
|
self._thread_name = thread_name
|
|
|
|
self._cache: Dict[str, Dict[str, Any]] = {}
|
|
self._lock = threading.RLock()
|
|
self._resource_version: Optional[str] = None
|
|
self._has_synced = False
|
|
self._stop_event = threading.Event()
|
|
self._thread: Optional[threading.Thread] = None
|
|
|
|
@property
|
|
def has_synced(self) -> bool:
|
|
"""Return True once an initial list has completed."""
|
|
return self._has_synced
|
|
|
|
def start(self) -> None:
|
|
"""Start the background watch thread if not already running."""
|
|
if self._thread and self._thread.is_alive():
|
|
return
|
|
|
|
self._thread = threading.Thread(
|
|
target=self._run,
|
|
name=self._thread_name,
|
|
daemon=True,
|
|
)
|
|
self._thread.start()
|
|
|
|
def stop(self) -> None:
|
|
"""Stop the background watch thread."""
|
|
self._stop_event.set()
|
|
|
|
def get(self, name: str) -> Optional[Dict[str, Any]]:
|
|
"""Return cached object by name, if present."""
|
|
with self._lock:
|
|
return self._cache.get(name)
|
|
|
|
def list(self) -> List[Dict[str, Any]]:
|
|
"""Return a snapshot of every cached object."""
|
|
with self._lock:
|
|
return list(self._cache.values())
|
|
|
|
def update_cache(self, obj: Dict[str, Any]) -> None:
|
|
"""Upsert a single object into the cache.
|
|
|
|
Only advances ``_resource_version`` if the incoming version is strictly
|
|
newer, preventing a stale API response from rolling back the watch cursor.
|
|
"""
|
|
metadata = obj.get("metadata", {})
|
|
name = metadata.get("name")
|
|
if not name:
|
|
return
|
|
|
|
with self._lock:
|
|
self._cache[name] = obj
|
|
self._advance_resource_version(metadata.get("resourceVersion"))
|
|
|
|
def delete_from_cache(self, name: str) -> None:
|
|
"""Evict a single object from the cache by name."""
|
|
with self._lock:
|
|
self._cache.pop(name, None)
|
|
|
|
def _advance_resource_version(self, rv: Optional[str]) -> None:
|
|
"""Advance ``_resource_version`` only when *rv* is strictly newer.
|
|
|
|
K8s resourceVersions are opaque strings but etcd encodes them as
|
|
monotonically increasing integers. If the conversion fails we skip the
|
|
update (conservative: keep the current, newer cursor).
|
|
|
|
Must be called with ``self._lock`` already held.
|
|
"""
|
|
if not rv:
|
|
return
|
|
if self._resource_version is None:
|
|
self._resource_version = rv
|
|
return
|
|
try:
|
|
if int(rv) > int(self._resource_version):
|
|
self._resource_version = rv
|
|
except ValueError:
|
|
# Non-integer resourceVersion — skip to avoid downgrade.
|
|
pass
|
|
|
|
def _run(self) -> None:
|
|
backoff = 1.0
|
|
while not self._stop_event.is_set():
|
|
try:
|
|
if not self._has_synced:
|
|
self._full_resync()
|
|
backoff = 1.0
|
|
|
|
if not self.enable_watch:
|
|
self._stop_event.wait(self.resync_period_seconds)
|
|
self._has_synced = False # trigger a fresh list on next loop
|
|
continue
|
|
|
|
self._run_watch_loop()
|
|
backoff = 1.0
|
|
except ApiException as exc:
|
|
if exc.status == 410:
|
|
# Resource version too old; force a fresh list on next loop.
|
|
self._resource_version = None
|
|
self._has_synced = False
|
|
else:
|
|
logger.warning(f"Informer watch error: {exc}", exc_info=True)
|
|
self._has_synced = False
|
|
self._stop_event.wait(min(backoff, 30.0))
|
|
backoff = min(backoff * 2, 30.0)
|
|
except Exception as exc: # pragma: no cover - defensive
|
|
logger.warning(f"Unexpected informer error: {exc}", exc_info=True)
|
|
self._has_synced = False
|
|
self._stop_event.wait(min(backoff, 30.0))
|
|
backoff = min(backoff * 2, 30.0)
|
|
|
|
def _full_resync(self) -> None:
|
|
"""Perform a full list to refresh the cache."""
|
|
resp = self.list_fn()
|
|
|
|
# list response is a dict for CustomObjectsApi
|
|
items = resp.get("items", []) if isinstance(resp, dict) else []
|
|
metadata = resp.get("metadata", {}) if isinstance(resp, dict) else {}
|
|
resource_version = metadata.get("resourceVersion")
|
|
|
|
# Build new cache outside the lock to avoid blocking readers
|
|
new_cache: Dict[str, Dict[str, Any]] = {}
|
|
for item in items:
|
|
name = item.get("metadata", {}).get("name")
|
|
if name:
|
|
new_cache[name] = item
|
|
|
|
with self._lock:
|
|
self._cache = new_cache
|
|
self._advance_resource_version(resource_version)
|
|
self._has_synced = True
|
|
|
|
def _run_watch_loop(self) -> None:
|
|
"""Stream watch events to keep the cache fresh."""
|
|
w = watch.Watch()
|
|
try:
|
|
for event in w.stream(
|
|
self.list_fn,
|
|
resource_version=self._resource_version,
|
|
timeout_seconds=self.watch_timeout_seconds,
|
|
):
|
|
if self._stop_event.is_set():
|
|
break
|
|
self._handle_event(event)
|
|
finally:
|
|
w.stop()
|
|
|
|
def _handle_event(self, event: Dict[str, Any]) -> None:
|
|
obj = event.get("object")
|
|
if obj is None:
|
|
return
|
|
|
|
if not isinstance(obj, dict):
|
|
try:
|
|
obj = obj.to_dict()
|
|
except Exception:
|
|
return
|
|
|
|
metadata = obj.get("metadata", {})
|
|
name = metadata.get("name")
|
|
if not name:
|
|
return
|
|
|
|
event_type = event.get("type")
|
|
with self._lock:
|
|
if event_type == "DELETED":
|
|
self._cache.pop(name, None)
|
|
else:
|
|
self._cache[name] = obj
|
|
self._advance_resource_version(metadata.get("resourceVersion"))
|