This is an automated email from the ASF dual-hosted git repository.
davsclaus pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/main by this push:
new f96b728a4d58 CAMEL-25273: camel-kubernetes - consumers must watch
again when the client closes their watch with an error (#27302)
f96b728a4d58 is described below
commit f96b728a4d589475b955970fcd7906269ceff47b
Author: allthingssecurity <[email protected]>
AuthorDate: Sun Oct 4 12:34:29 2026 +0530
CAMEL-25273: camel-kubernetes - consumers must watch again when the client
closes their watch with an error (#27302)
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../component/kubernetes/KubernetesHelper.java | 68 ++++++++++
.../config_maps/KubernetesConfigMapsConsumer.java | 8 +-
.../KubernetesCustomResourcesConsumer.java | 13 +-
.../deployments/KubernetesDeploymentsConsumer.java | 8 +-
.../events/KubernetesEventsConsumer.java | 8 +-
.../kubernetes/hpa/KubernetesHPAConsumer.java | 8 +-
.../namespaces/KubernetesNamespacesConsumer.java | 8 +-
.../kubernetes/nodes/KubernetesNodesConsumer.java | 8 +-
.../kubernetes/pods/KubernetesPodsConsumer.java | 8 +-
.../KubernetesReplicationControllersConsumer.java | 9 +-
.../services/KubernetesServicesConsumer.java | 8 +-
.../OpenshiftDeploymentConfigsConsumer.java | 9 +-
.../KubernetesPodsConsumerWatchClosedTest.java | 144 +++++++++++++++++++++
13 files changed, 295 insertions(+), 12 deletions(-)
diff --git
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/KubernetesHelper.java
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/KubernetesHelper.java
index 558e3706c64b..c647781cfb2b 100644
---
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/KubernetesHelper.java
+++
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/KubernetesHelper.java
@@ -16,6 +16,8 @@
*/
package org.apache.camel.component.kubernetes;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.RejectedExecutionException;
import java.util.function.Supplier;
import io.fabric8.kubernetes.client.Config;
@@ -25,6 +27,7 @@ import io.fabric8.kubernetes.client.KubernetesClientBuilder;
import io.fabric8.kubernetes.client.Watch;
import org.apache.camel.Exchange;
import org.apache.camel.support.MessageHelper;
+import org.apache.camel.support.service.ServiceSupport;
import org.apache.camel.util.ObjectHelper;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -90,6 +93,71 @@ public final class KubernetesHelper {
}
}
+ /**
+ * The delay before a consumer watches again after the Kubernetes client
closed its watch with an error.
+ */
+ static final long WATCH_AGAIN_DELAY_MILLIS = 1000;
+
+ /**
+ * The maximum delay between two attempts to watch again, when watching
again keeps failing (for example while the
+ * API server cannot be reached).
+ */
+ static final long WATCH_AGAIN_MAX_DELAY_MILLIS = 30000;
+
+ /**
+ * Watches again when the Kubernetes client closed the watch of a consumer
with an error. The client reconnects a
+ * watch by itself after transient errors, and only closes it with an
exception when it gives up: when the API
+ * server answers 410 Gone because the resource version of the watch is
too old (which happens to long-running
+ * watches), or when the reconnect limit is reached. The consumer would
then not receive any event anymore.
+ * <p>
+ * If watching again fails (the API server cannot be reached, answers 403,
...), it is tried again with a delay that
+ * doubles up to {@link #WATCH_AGAIN_MAX_DELAY_MILLIS}, until it succeeds
or the consumer is stopped.
+ *
+ * @param consumer the consumer of the watch
+ * @param executor the executor of the consumer
+ * @param task the task that creates the watch of the consumer
+ */
+ public static void watchAgain(ServiceSupport consumer, ExecutorService
executor, Runnable task) {
+ watchAgain(consumer, executor, task, WATCH_AGAIN_DELAY_MILLIS);
+ }
+
+ private static void watchAgain(ServiceSupport consumer, ExecutorService
executor, Runnable task, long delayMillis) {
+ if (!consumer.isRunAllowed() || executor == null ||
executor.isShutdown()) {
+ return;
+ }
+ LOG.info("Watching again for {} in {} ms", consumer, delayMillis);
+ try {
+ executor.submit(() -> {
+ // wait before watching again, so that an API server that
keeps closing the watch with an error or
+ // cannot be reached is not called in a tight loop
+ try {
+ Thread.sleep(delayMillis);
+ } catch (InterruptedException e) {
+ // the consumer is stopping (its executor is shut down)
+ Thread.currentThread().interrupt();
+ return;
+ }
+ if (!consumer.isRunAllowed()) {
+ return;
+ }
+ try {
+ task.run();
+ } catch (Exception e) {
+ if (Thread.currentThread().isInterrupted() ||
!consumer.isRunAllowed()) {
+ LOG.debug("Failed to watch again for {} as it is
stopping", consumer, e);
+ return;
+ }
+ long nextDelayMillis = Math.min(delayMillis * 2,
WATCH_AGAIN_MAX_DELAY_MILLIS);
+ LOG.warn("Failed to watch again for {}, trying again in {}
ms: {}", consumer, nextDelayMillis,
+ e.getMessage(), e);
+ watchAgain(consumer, executor, task, nextDelayMillis);
+ }
+ });
+ } catch (RejectedExecutionException e) {
+ LOG.debug("Cannot watch again for {} as it is stopping", consumer,
e);
+ }
+ }
+
public static String extractOperation(AbstractKubernetesEndpoint endpoint,
Exchange exchange) {
String operation;
diff --git
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/config_maps/KubernetesConfigMapsConsumer.java
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/config_maps/KubernetesConfigMapsConsumer.java
index 16172967e9ea..6fd6452f461c 100644
---
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/config_maps/KubernetesConfigMapsConsumer.java
+++
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/config_maps/KubernetesConfigMapsConsumer.java
@@ -82,7 +82,7 @@ public class KubernetesConfigMapsConsumer extends
DefaultConsumer {
class ConfigMapsConsumerTask implements Runnable {
- private Watch watch;
+ private volatile Watch watch;
@Override
public void run() {
@@ -139,9 +139,15 @@ public class KubernetesConfigMapsConsumer extends
DefaultConsumer {
public void onClose(WatcherException cause) {
if (cause != null) {
LOG.error(cause.getMessage(), cause);
+ // the client gave up the watch (410 Gone): watch again
+
KubernetesHelper.watchAgain(KubernetesConfigMapsConsumer.this, executor,
ConfigMapsConsumerTask.this);
}
}
});
+ if (!isRunAllowed()) {
+ // the consumer was stopped while the watch was being created,
so stopping could not close it
+ watch.close();
+ }
}
public Watch getWatch() {
diff --git
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/customresources/KubernetesCustomResourcesConsumer.java
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/customresources/KubernetesCustomResourcesConsumer.java
index b8cd2e3306b1..a429a2a12c4e 100644
---
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/customresources/KubernetesCustomResourcesConsumer.java
+++
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/customresources/KubernetesCustomResourcesConsumer.java
@@ -80,7 +80,7 @@ public class KubernetesCustomResourcesConsumer extends
DefaultConsumer {
class CustomResourcesConsumerTask implements Runnable {
- private Watch watch;
+ private volatile Watch watch;
@Override
public void run() {
@@ -89,7 +89,7 @@ public class KubernetesCustomResourcesConsumer extends
DefaultConsumer {
}
String namespace =
getEndpoint().getKubernetesConfiguration().getNamespace();
try {
- getEndpoint().getKubernetesClient()
+ watch = getEndpoint().getKubernetesClient()
.genericKubernetesResources(getCRDContext(getEndpoint().getKubernetesConfiguration()))
.inNamespace(namespace)
.watch(new Watcher<>() {
@@ -113,11 +113,20 @@ public class KubernetesCustomResourcesConsumer extends
DefaultConsumer {
public void onClose(WatcherException cause) {
if (cause != null) {
LOG.error(cause.getMessage(), cause);
+ // the client gave up the watch (410
Gone): watch again
+
KubernetesHelper.watchAgain(KubernetesCustomResourcesConsumer.this, executor,
+ CustomResourcesConsumerTask.this);
}
}
});
+ if (!isRunAllowed()) {
+ // the consumer was stopped while the watch was being
created, so stopping could not close it
+ watch.close();
+ }
} catch (Exception e) {
LOG.error("Exception in handling githubsource instance
change", e);
+ // so that watching again is tried again
+ throw e;
}
}
diff --git
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/deployments/KubernetesDeploymentsConsumer.java
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/deployments/KubernetesDeploymentsConsumer.java
index 10941a4beea7..8a5ab50196df 100644
---
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/deployments/KubernetesDeploymentsConsumer.java
+++
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/deployments/KubernetesDeploymentsConsumer.java
@@ -82,7 +82,7 @@ public class KubernetesDeploymentsConsumer extends
DefaultConsumer {
class DeploymentsConsumerTask implements Runnable {
- private Watch watch;
+ private volatile Watch watch;
@Override
public void run() {
@@ -140,10 +140,16 @@ public class KubernetesDeploymentsConsumer extends
DefaultConsumer {
public void onClose(WatcherException cause) {
if (cause != null) {
LOG.error(cause.getMessage(), cause);
+ // the client gave up the watch (410 Gone): watch again
+
KubernetesHelper.watchAgain(KubernetesDeploymentsConsumer.this, executor,
DeploymentsConsumerTask.this);
}
}
});
+ if (!isRunAllowed()) {
+ // the consumer was stopped while the watch was being created,
so stopping could not close it
+ watch.close();
+ }
}
public Watch getWatch() {
diff --git
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/events/KubernetesEventsConsumer.java
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/events/KubernetesEventsConsumer.java
index b9a3f329318f..44e6474eed92 100644
---
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/events/KubernetesEventsConsumer.java
+++
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/events/KubernetesEventsConsumer.java
@@ -82,7 +82,7 @@ public class KubernetesEventsConsumer extends DefaultConsumer
{
class EventsConsumerTask implements Runnable {
- private Watch watch;
+ private volatile Watch watch;
@Override
public void run() {
@@ -139,9 +139,15 @@ public class KubernetesEventsConsumer extends
DefaultConsumer {
public void onClose(WatcherException cause) {
if (cause != null) {
LOG.error(cause.getMessage(), cause);
+ // the client gave up the watch (410 Gone): watch again
+
KubernetesHelper.watchAgain(KubernetesEventsConsumer.this, executor,
EventsConsumerTask.this);
}
}
});
+ if (!isRunAllowed()) {
+ // the consumer was stopped while the watch was being created,
so stopping could not close it
+ watch.close();
+ }
}
public Watch getWatch() {
diff --git
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/hpa/KubernetesHPAConsumer.java
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/hpa/KubernetesHPAConsumer.java
index 927474cc68d7..e6b069f2e46b 100644
---
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/hpa/KubernetesHPAConsumer.java
+++
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/hpa/KubernetesHPAConsumer.java
@@ -82,7 +82,7 @@ public class KubernetesHPAConsumer extends DefaultConsumer {
class HpaConsumerTask implements Runnable {
- private Watch watch;
+ private volatile Watch watch;
@Override
public void run() {
@@ -141,10 +141,16 @@ public class KubernetesHPAConsumer extends
DefaultConsumer {
public void onClose(WatcherException cause) {
if (cause != null) {
LOG.error(cause.getMessage(), cause);
+ // the client gave up the watch (410 Gone): watch again
+
KubernetesHelper.watchAgain(KubernetesHPAConsumer.this, executor,
HpaConsumerTask.this);
}
}
});
+ if (!isRunAllowed()) {
+ // the consumer was stopped while the watch was being created,
so stopping could not close it
+ watch.close();
+ }
}
public Watch getWatch() {
diff --git
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/namespaces/KubernetesNamespacesConsumer.java
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/namespaces/KubernetesNamespacesConsumer.java
index bc920a21b8d3..300380459b72 100644
---
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/namespaces/KubernetesNamespacesConsumer.java
+++
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/namespaces/KubernetesNamespacesConsumer.java
@@ -81,7 +81,7 @@ public class KubernetesNamespacesConsumer extends
DefaultConsumer {
class NamespacesConsumerTask implements Runnable {
- private Watch watch;
+ private volatile Watch watch;
@Override
public void run() {
@@ -120,9 +120,15 @@ public class KubernetesNamespacesConsumer extends
DefaultConsumer {
public void onClose(WatcherException cause) {
if (cause != null) {
LOG.error(cause.getMessage(), cause);
+ // the client gave up the watch (410 Gone): watch again
+
KubernetesHelper.watchAgain(KubernetesNamespacesConsumer.this, executor,
NamespacesConsumerTask.this);
}
}
});
+ if (!isRunAllowed()) {
+ // the consumer was stopped while the watch was being created,
so stopping could not close it
+ watch.close();
+ }
}
public Watch getWatch() {
diff --git
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/nodes/KubernetesNodesConsumer.java
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/nodes/KubernetesNodesConsumer.java
index 35e07c5b70e8..e75b8b366c78 100644
---
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/nodes/KubernetesNodesConsumer.java
+++
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/nodes/KubernetesNodesConsumer.java
@@ -81,7 +81,7 @@ public class KubernetesNodesConsumer extends DefaultConsumer {
class NodesConsumerTask implements Runnable {
- private Watch watch;
+ private volatile Watch watch;
@Override
public void run() {
@@ -118,10 +118,16 @@ public class KubernetesNodesConsumer extends
DefaultConsumer {
public void onClose(WatcherException cause) {
if (cause != null) {
LOG.error(cause.getMessage(), cause);
+ // the client gave up the watch (410 Gone): watch again
+
KubernetesHelper.watchAgain(KubernetesNodesConsumer.this, executor,
NodesConsumerTask.this);
}
}
});
+ if (!isRunAllowed()) {
+ // the consumer was stopped while the watch was being created,
so stopping could not close it
+ watch.close();
+ }
}
public Watch getWatch() {
diff --git
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/pods/KubernetesPodsConsumer.java
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/pods/KubernetesPodsConsumer.java
index addcc543d536..2754f18dd0e9 100644
---
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/pods/KubernetesPodsConsumer.java
+++
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/pods/KubernetesPodsConsumer.java
@@ -82,7 +82,7 @@ public class KubernetesPodsConsumer extends DefaultConsumer {
class PodsConsumerTask implements Runnable {
- private Watch watch;
+ private volatile Watch watch;
@Override
public void run() {
@@ -139,10 +139,16 @@ public class KubernetesPodsConsumer extends
DefaultConsumer {
public void onClose(WatcherException cause) {
if (cause != null) {
LOG.error(cause.getMessage(), cause);
+ // the client gave up the watch (410 Gone): watch again
+
KubernetesHelper.watchAgain(KubernetesPodsConsumer.this, executor,
PodsConsumerTask.this);
}
}
});
+ if (!isRunAllowed()) {
+ // the consumer was stopped while the watch was being created,
so stopping could not close it
+ watch.close();
+ }
}
public Watch getWatch() {
diff --git
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/replication_controllers/KubernetesReplicationControllersConsumer.java
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/replication_controllers/KubernetesReplicationControllersConsumer.java
index 9785ecaf4941..1549971544f0 100644
---
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/replication_controllers/KubernetesReplicationControllersConsumer.java
+++
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/replication_controllers/KubernetesReplicationControllersConsumer.java
@@ -82,7 +82,7 @@ public class KubernetesReplicationControllersConsumer extends
DefaultConsumer {
class ReplicationControllersConsumerTask implements Runnable {
- private Watch watch;
+ private volatile Watch watch;
@Override
public void run() {
@@ -140,10 +140,17 @@ public class KubernetesReplicationControllersConsumer
extends DefaultConsumer {
public void onClose(WatcherException cause) {
if (cause != null) {
LOG.error(cause.getMessage(), cause);
+ // the client gave up the watch (410 Gone): watch again
+
KubernetesHelper.watchAgain(KubernetesReplicationControllersConsumer.this,
executor,
+ ReplicationControllersConsumerTask.this);
}
}
});
+ if (!isRunAllowed()) {
+ // the consumer was stopped while the watch was being created,
so stopping could not close it
+ watch.close();
+ }
}
public Watch getWatch() {
diff --git
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/services/KubernetesServicesConsumer.java
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/services/KubernetesServicesConsumer.java
index b92f165053cd..4a088ec31f53 100644
---
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/services/KubernetesServicesConsumer.java
+++
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/services/KubernetesServicesConsumer.java
@@ -82,7 +82,7 @@ public class KubernetesServicesConsumer extends
DefaultConsumer {
class ServicesConsumerTask implements Runnable {
- private Watch watch;
+ private volatile Watch watch;
@Override
public void run() {
@@ -140,10 +140,16 @@ public class KubernetesServicesConsumer extends
DefaultConsumer {
public void onClose(WatcherException cause) {
if (cause != null) {
LOG.error(cause.getMessage(), cause);
+ // the client gave up the watch (410 Gone): watch again
+
KubernetesHelper.watchAgain(KubernetesServicesConsumer.this, executor,
ServicesConsumerTask.this);
}
}
});
+ if (!isRunAllowed()) {
+ // the consumer was stopped while the watch was being created,
so stopping could not close it
+ watch.close();
+ }
}
public Watch getWatch() {
diff --git
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/openshift/deploymentconfigs/OpenshiftDeploymentConfigsConsumer.java
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/openshift/deploymentconfigs/OpenshiftDeploymentConfigsConsumer.java
index b66cb579d0cd..6ac0f8fc9a53 100644
---
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/openshift/deploymentconfigs/OpenshiftDeploymentConfigsConsumer.java
+++
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/openshift/deploymentconfigs/OpenshiftDeploymentConfigsConsumer.java
@@ -83,7 +83,7 @@ public class OpenshiftDeploymentConfigsConsumer extends
DefaultConsumer {
class DeploymentsConfigConsumerTask implements Runnable {
- private Watch watch;
+ private volatile Watch watch;
@Override
public void run() {
@@ -143,10 +143,17 @@ public class OpenshiftDeploymentConfigsConsumer extends
DefaultConsumer {
public void onClose(WatcherException cause) {
if (cause != null) {
LOG.error(cause.getMessage(), cause);
+ // the client gave up the watch (410 Gone): watch again
+
KubernetesHelper.watchAgain(OpenshiftDeploymentConfigsConsumer.this, executor,
+ DeploymentsConfigConsumerTask.this);
}
}
});
+ if (!isRunAllowed()) {
+ // the consumer was stopped while the watch was being created,
so stopping could not close it
+ watch.close();
+ }
}
public Watch getWatch() {
diff --git
a/components/camel-kubernetes/src/test/java/org/apache/camel/component/kubernetes/consumer/KubernetesPodsConsumerWatchClosedTest.java
b/components/camel-kubernetes/src/test/java/org/apache/camel/component/kubernetes/consumer/KubernetesPodsConsumerWatchClosedTest.java
new file mode 100644
index 000000000000..dcad0ef94c4c
--- /dev/null
+++
b/components/camel-kubernetes/src/test/java/org/apache/camel/component/kubernetes/consumer/KubernetesPodsConsumerWatchClosedTest.java
@@ -0,0 +1,144 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You 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 org.apache.camel.component.kubernetes.consumer;
+
+import io.fabric8.kubernetes.api.model.Pod;
+import io.fabric8.kubernetes.api.model.PodBuilder;
+import io.fabric8.kubernetes.api.model.StatusBuilder;
+import io.fabric8.kubernetes.api.model.WatchEvent;
+import io.fabric8.kubernetes.client.KubernetesClient;
+import io.fabric8.kubernetes.client.NamespacedKubernetesClient;
+import io.fabric8.kubernetes.client.server.mock.EnableKubernetesMockClient;
+import io.fabric8.kubernetes.client.server.mock.KubernetesMockServer;
+import org.apache.camel.BindToRegistry;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.kubernetes.KubernetesTestSupport;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+/**
+ * The Kubernetes client reconnects a watch after transient errors by itself,
but closes it for good, with an exception,
+ * when the API server answers 410 Gone (the resource version of the watch is
too old, which happens to long-running
+ * watches). The consumer must then watch again.
+ */
+@EnableKubernetesMockClient
+public class KubernetesPodsConsumerWatchClosedTest extends
KubernetesTestSupport {
+
+ private static final String WATCH_PATH =
"/api/v1/namespaces/test/pods?allowWatchBookmarks=true&watch=true";
+
+ KubernetesMockServer server;
+ NamespacedKubernetesClient client;
+
+ @BindToRegistry("kubernetesClient")
+ public KubernetesClient getClient() {
+ return client;
+ }
+
+ @Override
+ public boolean isUseRouteBuilder() {
+ return false;
+ }
+
+ @Test
+ public void testWatchAgainAfterTheWatchIsClosedWithAnError() throws
Exception {
+ server.expect().withPath(WATCH_PATH)
+ .andUpgradeToWebSocket()
+ .open()
+ .waitFor(10)
+ .andEmit(new WatchEvent(
+ new
StatusBuilder().withCode(410).withReason("Expired").withMessage("too old
resource version")
+ .build(),
+ "ERROR"))
+ .done()
+ .once();
+ server.expect().withPath(WATCH_PATH)
+ .andUpgradeToWebSocket()
+ .open()
+ .waitFor(10)
+ .andEmit(new WatchEvent(
+ new
PodBuilder().withNewMetadata().withName("pod1").withNamespace("test")
+
.withResourceVersion("2").endMetadata().build(),
+ "ADDED"))
+ .done()
+ .once();
+
+ MockEndpoint mock = getMockEndpoint("mock:result");
+ mock.expectedMessageCount(1);
+ mock.expectedMessagesMatches(e -> "pod1".equals(
+ e.getMessage().getBody(Pod.class).getMetadata().getName()));
+
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+
from("kubernetes-pods://kubernetes?kubernetesClient=#kubernetesClient&namespace=test")
+ .to("mock:result");
+ }
+ });
+ context.start();
+
+ mock.assertIsSatisfied();
+ }
+
+ @Test
+ public void testWatchAgainAfterAFailedAttemptToWatchAgain() throws
Exception {
+ server.expect().withPath(WATCH_PATH)
+ .andUpgradeToWebSocket()
+ .open()
+ .waitFor(10)
+ .andEmit(new WatchEvent(
+ new
StatusBuilder().withCode(410).withReason("Expired").withMessage("too old
resource version")
+ .build(),
+ "ERROR"))
+ .done()
+ .once();
+ // the first attempt to watch again fails
+ server.expect().withPath(WATCH_PATH)
+ .andReturn(403, new
StatusBuilder().withCode(403).withReason("Forbidden").withMessage("forbidden").build())
+ .once();
+ server.expect().withPath(WATCH_PATH)
+ .andUpgradeToWebSocket()
+ .open()
+ .waitFor(10)
+ .andEmit(new WatchEvent(
+ new
PodBuilder().withNewMetadata().withName("pod1").withNamespace("test")
+
.withResourceVersion("2").endMetadata().build(),
+ "ADDED"))
+ .done()
+ .once();
+
+ MockEndpoint mock = getMockEndpoint("mock:result");
+ mock.expectedMessageCount(1);
+ mock.expectedMessagesMatches(e -> "pod1".equals(
+ e.getMessage().getBody(Pod.class).getMetadata().getName()));
+
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+
from("kubernetes-pods://kubernetes?kubernetesClient=#kubernetesClient&namespace=test")
+ .to("mock:result");
+ }
+ });
+ context.start();
+
+ // watched again after 1 second, failed, and watched again after 2
more seconds
+ mock.setResultWaitTime(20000);
+ mock.assertIsSatisfied();
+ assertEquals(3, server.getRequestCount());
+ }
+}