Skip to content
Open
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
98 changes: 94 additions & 4 deletions robusta_krr/core/integrations/kubernetes/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,8 @@ def __init__(self, cluster: Optional[str] = None):
self.autoscaling_v2 = client.AutoscalingV2Api(api_client=self.api_client)

self.__kind_available: defaultdict[KindLiteral, bool] = defaultdict(lambda: True)
self._strimzipodset_api_version: Optional[str] = None
self._strimzipodset_api_version_checked = False

self.__jobs_for_cronjobs: dict[str, list[V1Job]] = {}
self.__jobs_loading_locks: defaultdict[str, asyncio.Lock] = defaultdict(asyncio.Lock)
Expand Down Expand Up @@ -471,23 +473,111 @@ async def _extract_containers(item: Any) -> list[V1Container]:
extract_containers=_extract_containers,
)

def _list_strimzipodsets(self) -> list[K8sObjectData]:
def _probe_strimzipodset_api_version(self, version: str) -> None:
probe_kwargs = {"limit": 1}
group = "core.strimzi.io"
plural = "strimzipodsets"

if self.namespaces == "*":
self.custom_objects.list_cluster_custom_object(
group=group,
version=version,
plural=plural,
**probe_kwargs,
)
return

if not self.namespaces:
logger.debug("No namespaces configured, skipping StrimziPodSet version probe")
raise ApiException(status=404, reason="No namespaces configured")

last_error: Optional[ApiException] = None
for namespace in self.namespaces:
try:
self.custom_objects.list_namespaced_custom_object(
group=group,
version=version,
plural=plural,
namespace=namespace,
**probe_kwargs,
)
return
except ApiException as e:
last_error = e
if e.status == 404:
logger.debug(
f"StrimziPodSet API version {version} not available in namespace {namespace}"
)
else:
logger.debug(
f"Skipping namespace {namespace} while probing StrimziPodSet API version {version}: {e.reason}"
)
continue

if last_error is not None:
raise last_error

def _resolve_strimzipodset_api_version(self) -> Optional[str]:
if self._strimzipodset_api_version_checked:
return self._strimzipodset_api_version

if self.namespaces != "*" and not self.namespaces:
logger.debug("No namespaces configured, skipping StrimziPodSet API detection")
self._strimzipodset_api_version_checked = True
self._strimzipodset_api_version = None
return None

# Strimzi >= 1.0 serves v1 only; older operators still expose v1beta2.
for version in ("v1", "v1beta2"):
try:
self._probe_strimzipodset_api_version(version)
except ApiException as e:
if e.status == 404:
continue
logger.exception(
f"Error {e.status} probing StrimziPodSet API version {version} in cluster {self.cluster}: {e.reason}"
)
continue
except Exception:
logger.exception(
f"Unexpected error probing StrimziPodSet API version {version} in cluster {self.cluster}"
)
return None
else:
self._strimzipodset_api_version = version
self._strimzipodset_api_version_checked = True
logger.debug(f"StrimziPodSet API version {version} available in {self.cluster}")
return version

self._strimzipodset_api_version = None
self._strimzipodset_api_version_checked = True
return None
Comment thread
coderabbitai[bot] marked this conversation as resolved.

async def _list_strimzipodsets(self) -> list[K8sObjectData]:
# NOTE: Using custom objects API returns dicts, but all other APIs return objects
# We need to handle this difference using a small wrapper
return self._list_scannable_objects(
loop = asyncio.get_running_loop()
version = await loop.run_in_executor(self.executor, self._resolve_strimzipodset_api_version)
if version is None:
if self.__kind_available["StrimziPodSet"]:
logger.debug(f"StrimziPodSet API not available in {self.cluster}")
self.__kind_available["StrimziPodSet"] = False
return []

return await self._list_scannable_objects(
kind="StrimziPodSet",
all_namespaces_request=lambda **kwargs: ObjectLikeDict(
self.custom_objects.list_cluster_custom_object(
group="core.strimzi.io",
version="v1beta2",
version=version,
plural="strimzipodsets",
**kwargs,
)
),
namespaced_request=lambda **kwargs: ObjectLikeDict(
self.custom_objects.list_namespaced_custom_object(
group="core.strimzi.io",
version="v1beta2",
version=version,
plural="strimzipodsets",
**kwargs,
)
Expand Down
149 changes: 149 additions & 0 deletions tests/test_strimzipodset_api_version.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,149 @@
from unittest.mock import MagicMock, patch

import pytest
from kubernetes.client import ApiException

from robusta_krr.core.integrations.kubernetes import ClusterLoader
from robusta_krr.core.models.config import Config


@pytest.fixture
def mock_config():
config = MagicMock(spec=Config)
config.max_workers = 4
config.get_kube_client = MagicMock()
config.resources = "*"
config.selector = None
config.namespaces = "*"
return config


@pytest.fixture
def loader(mock_config):
with patch("robusta_krr.core.integrations.kubernetes.settings", mock_config):
yield ClusterLoader(cluster="test-cluster")


@pytest.fixture
def namespaced_loader(mock_config):
mock_config.namespaces = ["sentry"]
with patch("robusta_krr.core.integrations.kubernetes.settings", mock_config):
yield ClusterLoader(cluster="test-cluster")


def test_resolve_strimzipodset_api_version_prefers_v1(loader):
loader.custom_objects = MagicMock()

version = loader._resolve_strimzipodset_api_version()

assert version == "v1"
loader.custom_objects.list_cluster_custom_object.assert_called_once_with(
group="core.strimzi.io",
version="v1",
plural="strimzipodsets",
limit=1,
)
loader.custom_objects.list_namespaced_custom_object.assert_not_called()


def test_resolve_strimzipodset_api_version_falls_back_to_v1beta2(loader):
loader.custom_objects = MagicMock()
loader.custom_objects.list_cluster_custom_object.side_effect = [
ApiException(status=404, reason="Not Found"),
{"items": []},
]

version = loader._resolve_strimzipodset_api_version()

assert version == "v1beta2"
assert loader.custom_objects.list_cluster_custom_object.call_count == 2
loader.custom_objects.list_cluster_custom_object.assert_any_call(
group="core.strimzi.io",
version="v1",
plural="strimzipodsets",
limit=1,
)
loader.custom_objects.list_cluster_custom_object.assert_any_call(
group="core.strimzi.io",
version="v1beta2",
plural="strimzipodsets",
limit=1,
)


def test_resolve_strimzipodset_api_version_returns_none_when_unavailable(loader):
loader.custom_objects = MagicMock()
loader.custom_objects.list_cluster_custom_object.side_effect = ApiException(status=404, reason="Not Found")

version = loader._resolve_strimzipodset_api_version()

assert version is None
assert loader._strimzipodset_api_version_checked
assert loader.custom_objects.list_cluster_custom_object.call_count == 2


def test_resolve_strimzipodset_api_version_is_cached(loader):
loader.custom_objects = MagicMock()

assert loader._resolve_strimzipodset_api_version() == "v1"
assert loader._resolve_strimzipodset_api_version() == "v1"
loader.custom_objects.list_cluster_custom_object.assert_called_once()


def test_resolve_strimzipodset_api_version_uses_namespaced_probe(namespaced_loader):
namespaced_loader.custom_objects = MagicMock()

version = namespaced_loader._resolve_strimzipodset_api_version()

assert version == "v1"
namespaced_loader.custom_objects.list_cluster_custom_object.assert_not_called()
namespaced_loader.custom_objects.list_namespaced_custom_object.assert_called_once_with(
group="core.strimzi.io",
version="v1",
plural="strimzipodsets",
namespace="sentry",
limit=1,
)


def test_resolve_strimzipodset_api_version_tries_v1beta2_after_v1_non_404(loader):
loader.custom_objects = MagicMock()
loader.custom_objects.list_cluster_custom_object.side_effect = [
ApiException(status=403, reason="Forbidden"),
{"items": []},
]

version = loader._resolve_strimzipodset_api_version()

assert version == "v1beta2"
assert loader.custom_objects.list_cluster_custom_object.call_count == 2


def test_resolve_strimzipodset_api_version_does_not_cache_transient_errors(loader):
loader.custom_objects = MagicMock()
loader.custom_objects.list_cluster_custom_object.side_effect = RuntimeError("connection reset")

assert loader._resolve_strimzipodset_api_version() is None
assert not loader._strimzipodset_api_version_checked

loader.custom_objects.list_cluster_custom_object.side_effect = None
loader.custom_objects.list_cluster_custom_object.return_value = {"items": []}

assert loader._resolve_strimzipodset_api_version() == "v1"
assert loader._strimzipodset_api_version_checked


def test_resolve_strimzipodset_api_version_probes_until_namespace_succeeds(mock_config):
mock_config.namespaces = ["blocked", "sentry"]
with patch("robusta_krr.core.integrations.kubernetes.settings", mock_config):
loader = ClusterLoader(cluster="test-cluster")
loader.custom_objects = MagicMock()
loader.custom_objects.list_namespaced_custom_object.side_effect = [
ApiException(status=403, reason="Forbidden"),
{"items": []},
]

version = loader._resolve_strimzipodset_api_version()

assert version == "v1"
assert loader.custom_objects.list_namespaced_custom_object.call_count == 2