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());
+    }
+}

Reply via email to