From c6f6bb9424e3827db9626b5328dd92cd007ab9f4 Mon Sep 17 00:00:00 2001 From: Brendan Burns <5751682+brendandburns@users.noreply.github.com> Date: Tue, 1 Sep 2026 19:19:49 +0000 Subject: [PATCH] Add predicate-aware informer event handlers Add predicate-aware event handler overloads, transition handling, malformed-object hardening, and coverage. --- .../e2e/informer/NamespaceInformerTest.java | 61 ++++++++++ .../FilteringResourceEventHandler.java | 80 +++++++++++++ .../client/informer/SharedInformer.java | 20 ++++ .../impl/DefaultSharedIndexInformer.java | 21 +++- .../FilteringResourceEventHandlerTest.java | 107 ++++++++++++++++++ 5 files changed, 287 insertions(+), 2 deletions(-) create mode 100644 util/src/main/java/io/kubernetes/client/informer/FilteringResourceEventHandler.java create mode 100644 util/src/test/java/io/kubernetes/client/informer/FilteringResourceEventHandlerTest.java diff --git a/e2e/src/test/java/io/kubernetes/client/e2e/informer/NamespaceInformerTest.java b/e2e/src/test/java/io/kubernetes/client/e2e/informer/NamespaceInformerTest.java index bbf92916b2..1f74c7a4b4 100644 --- a/e2e/src/test/java/io/kubernetes/client/e2e/informer/NamespaceInformerTest.java +++ b/e2e/src/test/java/io/kubernetes/client/e2e/informer/NamespaceInformerTest.java @@ -17,12 +17,18 @@ import io.kubernetes.client.informer.SharedIndexInformer; import io.kubernetes.client.informer.SharedInformerFactory; +import io.kubernetes.client.informer.ResourceEventHandler; import io.kubernetes.client.informer.cache.Lister; import io.kubernetes.client.openapi.ApiClient; +import io.kubernetes.client.openapi.apis.CoreV1Api; import io.kubernetes.client.openapi.models.V1Namespace; import io.kubernetes.client.openapi.models.V1NamespaceList; +import io.kubernetes.client.openapi.models.V1ObjectMeta; import io.kubernetes.client.util.ClientBuilder; import io.kubernetes.client.util.generic.GenericKubernetesApi; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; import org.junit.jupiter.api.Test; class NamespaceInformerTest { @@ -57,4 +63,59 @@ void listWatchingNamespaces() throws Exception { informerFactory.stopAllRegisteredInformers(true); } } + + @Test + void listWatchingNamespacesWithPredicateHandler() throws Exception { + ApiClient client = ClientBuilder.defaultClient(); + CoreV1Api coreV1Api = new CoreV1Api(client); + SharedInformerFactory informerFactory = new SharedInformerFactory(client); + String selectedNamespace = "e2e-filtered-selected"; + String ignoredNamespace = "e2e-filtered-ignored"; + + coreV1Api.createNamespace(new V1Namespace().metadata(new V1ObjectMeta().name(selectedNamespace))).execute(); + coreV1Api.createNamespace(new V1Namespace().metadata(new V1ObjectMeta().name(ignoredNamespace))).execute(); + + GenericKubernetesApi api = + new GenericKubernetesApi<>(V1Namespace.class, V1NamespaceList.class, "", "v1", "namespaces", client); + + SharedIndexInformer nsInformer = + informerFactory.sharedIndexInformerFor(api, V1Namespace.class, 0); + CountDownLatch selectedSeen = new CountDownLatch(1); + AtomicBoolean ignoredSeen = new AtomicBoolean(false); + AtomicBoolean selectedSeenByHandler = new AtomicBoolean(false); + try { + nsInformer.addEventHandler( + new ResourceEventHandler() { + @Override + public void onAdd(V1Namespace obj) { + String name = obj.getMetadata().getName(); + if (selectedNamespace.equals(name)) { + selectedSeenByHandler.set(true); + selectedSeen.countDown(); + } + if (ignoredNamespace.equals(name)) { + ignoredSeen.set(true); + } + } + + @Override + public void onUpdate(V1Namespace oldObj, V1Namespace newObj) {} + + @Override + public void onDelete(V1Namespace obj, boolean deletedFinalStateUnknown) {} + }, + ns -> selectedNamespace.equals(ns.getMetadata().getName())); + + informerFactory.startAllRegisteredInformers(); + + await().untilAsserted(() -> assertThat(nsInformer.hasSynced()).isTrue()); + assertThat(selectedSeen.await(30, TimeUnit.SECONDS)).isTrue(); + assertThat(selectedSeenByHandler.get()).isTrue(); + assertThat(ignoredSeen.get()).isFalse(); + } finally { + informerFactory.stopAllRegisteredInformers(true); + coreV1Api.deleteNamespace(selectedNamespace).execute(); + coreV1Api.deleteNamespace(ignoredNamespace).execute(); + } + } } diff --git a/util/src/main/java/io/kubernetes/client/informer/FilteringResourceEventHandler.java b/util/src/main/java/io/kubernetes/client/informer/FilteringResourceEventHandler.java new file mode 100644 index 0000000000..4b64a2a84f --- /dev/null +++ b/util/src/main/java/io/kubernetes/client/informer/FilteringResourceEventHandler.java @@ -0,0 +1,80 @@ +/* +Copyright 2026 The Kubernetes Authors. +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. +*/ +package io.kubernetes.client.informer; + +import io.kubernetes.client.common.KubernetesObject; +import java.util.Objects; +import java.util.function.Predicate; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * FilteringResourceEventHandler dispatches events only when objects match a predicate. + */ +public class FilteringResourceEventHandler + implements ResourceEventHandler { + + private static final Logger log = LoggerFactory.getLogger(FilteringResourceEventHandler.class); + + private final ResourceEventHandler delegate; + private final Predicate filter; + + public FilteringResourceEventHandler( + ResourceEventHandler delegate, Predicate filter) { + this.delegate = Objects.requireNonNull(delegate); + this.filter = Objects.requireNonNull(filter); + } + + @Override + public void onAdd(ApiType obj) { + if (matches(obj)) { + delegate.onAdd(obj); + } + } + + @Override + public void onUpdate(ApiType oldObj, ApiType newObj) { + boolean oldMatched = matches(oldObj); + boolean newMatched = matches(newObj); + if (oldMatched && newMatched) { + delegate.onUpdate(oldObj, newObj); + return; + } + if (!oldMatched && newMatched) { + delegate.onAdd(newObj); + return; + } + if (oldMatched) { + delegate.onDelete(oldObj, false); + } + } + + @Override + public void onDelete(ApiType obj, boolean deletedFinalStateUnknown) { + if (matches(obj)) { + delegate.onDelete(obj, deletedFinalStateUnknown); + } + } + + private boolean matches(ApiType obj) { + if (obj == null) { + return false; + } + try { + return filter.test(obj); + } catch (RuntimeException e) { + log.warn("Predicate threw an exception; dropping informer event", e); + return false; + } + } +} diff --git a/util/src/main/java/io/kubernetes/client/informer/SharedInformer.java b/util/src/main/java/io/kubernetes/client/informer/SharedInformer.java index 8c656fa2b7..2a7a4c3a51 100644 --- a/util/src/main/java/io/kubernetes/client/informer/SharedInformer.java +++ b/util/src/main/java/io/kubernetes/client/informer/SharedInformer.java @@ -13,6 +13,7 @@ package io.kubernetes.client.informer; import io.kubernetes.client.common.KubernetesObject; +import java.util.function.Predicate; /* * SharedInformer defines basic methods of a informer. @@ -26,6 +27,14 @@ public interface SharedInformer { */ void addEventHandler(ResourceEventHandler handler); + /** + * Add event handler with predicate filter. + * + * @param handler the handler + * @param filter the object filter + */ + void addEventHandler(ResourceEventHandler handler, Predicate filter); + /** * addEventHandlerWithResyncPeriod adds an event handler to the shared informer using the * specified resync period. Events to a single handler are delivered sequentially, but there is no @@ -36,6 +45,17 @@ public interface SharedInformer { */ void addEventHandlerWithResyncPeriod(ResourceEventHandler handler, long resyncPeriod); + /** + * addEventHandlerWithResyncPeriod adds an event handler with the specified resync period and + * object filter. + * + * @param handler the event handler + * @param resyncPeriod the specific resync period + * @param filter the object filter + */ + void addEventHandlerWithResyncPeriod( + ResourceEventHandler handler, long resyncPeriod, Predicate filter); + /** run starts the shared informer, which will be stopped until stop() is called. */ void run(); diff --git a/util/src/main/java/io/kubernetes/client/informer/impl/DefaultSharedIndexInformer.java b/util/src/main/java/io/kubernetes/client/informer/impl/DefaultSharedIndexInformer.java index 4fd1db93a3..7c1690b66f 100644 --- a/util/src/main/java/io/kubernetes/client/informer/impl/DefaultSharedIndexInformer.java +++ b/util/src/main/java/io/kubernetes/client/informer/impl/DefaultSharedIndexInformer.java @@ -18,6 +18,7 @@ import io.kubernetes.client.informer.ResourceEventHandler; import io.kubernetes.client.informer.SharedIndexInformer; import io.kubernetes.client.informer.TransformFunc; +import io.kubernetes.client.informer.FilteringResourceEventHandler; import io.kubernetes.client.informer.cache.Cache; import io.kubernetes.client.informer.cache.Controller; import io.kubernetes.client.informer.cache.DeltaFIFO; @@ -31,6 +32,7 @@ import java.util.concurrent.ThreadFactory; import java.util.function.BiConsumer; import java.util.function.Function; +import java.util.function.Predicate; import org.apache.commons.collections4.CollectionUtils; import org.apache.commons.lang3.tuple.MutablePair; import org.slf4j.Logger; @@ -152,13 +154,27 @@ public DefaultSharedIndexInformer( /** add event callback */ @Override public void addEventHandler(ResourceEventHandler handler) { - addEventHandlerWithResyncPeriod(handler, defaultEventHandlerResyncPeriod); + addEventHandler(handler, obj -> true); + } + + @Override + public void addEventHandler( + ResourceEventHandler handler, Predicate filterPredicate) { + addEventHandlerWithResyncPeriod(handler, defaultEventHandlerResyncPeriod, filterPredicate); } /** add event callback with a resync period */ @Override public void addEventHandlerWithResyncPeriod( ResourceEventHandler handler, long resyncPeriodMillis) { + addEventHandlerWithResyncPeriod(handler, resyncPeriodMillis, obj -> true); + } + + @Override + public void addEventHandlerWithResyncPeriod( + ResourceEventHandler handler, + long resyncPeriodMillis, + Predicate filterPredicate) { if (stopped) { log.info( "DefaultSharedIndexInformer#Handler was not added to shared informer because it has stopped already"); @@ -195,7 +211,8 @@ public void addEventHandlerWithResyncPeriod( ProcessorListener listener = new ProcessorListener( - handler, determineResyncPeriod(resyncCheckPeriodMillis, this.resyncCheckPeriodMillis)); + new FilteringResourceEventHandler<>(handler, filterPredicate), + determineResyncPeriod(resyncCheckPeriodMillis, this.resyncCheckPeriodMillis)); if (!started) { this.processor.addListener(listener); return; diff --git a/util/src/test/java/io/kubernetes/client/informer/FilteringResourceEventHandlerTest.java b/util/src/test/java/io/kubernetes/client/informer/FilteringResourceEventHandlerTest.java new file mode 100644 index 0000000000..3a6a78ba33 --- /dev/null +++ b/util/src/test/java/io/kubernetes/client/informer/FilteringResourceEventHandlerTest.java @@ -0,0 +1,107 @@ +/* +Copyright 2026 The Kubernetes Authors. +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. +*/ +package io.kubernetes.client.informer; + +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoInteractions; + +import io.kubernetes.client.openapi.models.V1ObjectMeta; +import io.kubernetes.client.openapi.models.V1Pod; +import org.junit.jupiter.api.Test; +import org.mockito.Mockito; + +class FilteringResourceEventHandlerTest { + + private static V1Pod pod(String name) { + return new V1Pod().metadata(new V1ObjectMeta().namespace("default").name(name)); + } + + @Test + void dropsAddForFilteredOutObject() { + ResourceEventHandler delegate = Mockito.mock(ResourceEventHandler.class); + FilteringResourceEventHandler handler = + new FilteringResourceEventHandler<>(delegate, obj -> "selected".equals(obj.getMetadata().getName())); + + handler.onAdd(pod("ignored")); + + verifyNoInteractions(delegate); + } + + @Test + void convertsUpdateTransitionIntoAdd() { + ResourceEventHandler delegate = Mockito.mock(ResourceEventHandler.class); + FilteringResourceEventHandler handler = + new FilteringResourceEventHandler<>(delegate, obj -> "selected".equals(obj.getMetadata().getName())); + V1Pod oldObj = pod("ignored"); + V1Pod newObj = pod("selected"); + + handler.onUpdate(oldObj, newObj); + + verify(delegate).onAdd(newObj); + verify(delegate, never()).onUpdate(oldObj, newObj); + verify(delegate, never()).onDelete(oldObj, false); + } + + @Test + void convertsUpdateTransitionIntoDelete() { + ResourceEventHandler delegate = Mockito.mock(ResourceEventHandler.class); + FilteringResourceEventHandler handler = + new FilteringResourceEventHandler<>(delegate, obj -> "selected".equals(obj.getMetadata().getName())); + V1Pod oldObj = pod("selected"); + V1Pod newObj = pod("ignored"); + + handler.onUpdate(oldObj, newObj); + + verify(delegate).onDelete(oldObj, false); + verify(delegate, never()).onUpdate(oldObj, newObj); + verify(delegate, never()).onAdd(newObj); + } + + @Test + void preservesUpdateWhenOldAndNewMatch() { + ResourceEventHandler delegate = Mockito.mock(ResourceEventHandler.class); + FilteringResourceEventHandler handler = + new FilteringResourceEventHandler<>(delegate, obj -> "selected".equals(obj.getMetadata().getName())); + V1Pod oldObj = pod("selected"); + V1Pod newObj = pod("selected"); + + handler.onUpdate(oldObj, newObj); + + verify(delegate).onUpdate(oldObj, newObj); + verify(delegate, never()).onAdd(newObj); + verify(delegate, never()).onDelete(oldObj, false); + } + + @Test + void dropsDeleteForFilteredOutObject() { + ResourceEventHandler delegate = Mockito.mock(ResourceEventHandler.class); + FilteringResourceEventHandler handler = + new FilteringResourceEventHandler<>(delegate, obj -> "selected".equals(obj.getMetadata().getName())); + + handler.onDelete(pod("ignored"), false); + + verifyNoInteractions(delegate); + } + + @Test + void shouldIgnoreDeleteWhenPredicateThrows() { + ResourceEventHandler delegate = Mockito.mock(ResourceEventHandler.class); + FilteringResourceEventHandler handler = + new FilteringResourceEventHandler<>(delegate, obj -> "selected".equals(obj.getMetadata().getName())); + + handler.onDelete(new V1Pod(), true); + + verifyNoInteractions(delegate); + } +}