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 97e434aa597d CAMEL-25209: camel-kubernetes - leadership controller 
must keep running after a failed lease renewal (#27158)
97e434aa597d is described below

commit 97e434aa597d58a37c9307554fbcc6732fc4f60a
Author: allthingssecurity <[email protected]>
AuthorDate: Thu Oct 1 12:24:34 2026 +0530

    CAMEL-25209: camel-kubernetes - leadership controller must keep running 
after a failed lease renewal (#27158)
    
    KubernetesLeadershipController runs its state machine as a chain of tasks
    on a single thread: each run schedules the next one. All Kubernetes calls
    catch their exceptions, except the renewal of the Lease in the LEADER
    state (refreshLeaseRenewTime, a PUT with the resource version, used with
    the default Lease resource type). When that PUT fails (an API server outage
    that outlasts the client's own retries, or an error that is not retried,
    such as a 409 conflict), the exception ends the task, the next run is
    never scheduled and nothing is logged.
    
    The leader then stops refreshing its local leadership, so after the
    renew deadline its TimedLeaderNotifier reports that there is no leader
    and the pod stops its master routes. The other pods keep seeing the
    Lease held by a running and ready pod, so none of them takes over: no pod
    of the group leads until the former leader pod is restarted.
    
    refreshStatus now catches an exception, logs it and schedules the next
    run, as the other failure paths do.
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../lock/KubernetesLeadershipController.java       | 23 ++++++++++++++++++
 .../cluster/KubernetesClusterServiceTest.java      | 27 ++++++++++++++++++++++
 .../kubernetes/cluster/utils/LockTestServer.java   | 20 ++++++++++++++++
 3 files changed, 70 insertions(+)

diff --git 
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/cluster/lock/KubernetesLeadershipController.java
 
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/cluster/lock/KubernetesLeadershipController.java
index 00bfdce5e23a..0b81d0863822 100644
--- 
a/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/cluster/lock/KubernetesLeadershipController.java
+++ 
b/components/camel-kubernetes/src/main/java/org/apache/camel/component/kubernetes/cluster/lock/KubernetesLeadershipController.java
@@ -23,6 +23,7 @@ import java.util.List;
 import java.util.Objects;
 import java.util.Optional;
 import java.util.Set;
+import java.util.concurrent.RejectedExecutionException;
 import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.TimeUnit;
 import java.util.stream.Collectors;
@@ -131,6 +132,28 @@ public class KubernetesLeadershipController implements 
Service {
     }
 
     private void refreshStatus() {
+        try {
+            doRefreshStatus();
+        } catch (Exception e) {
+            ScheduledExecutorService executor = this.serializedExecutor;
+            if (executor == null || executor.isShutdown()) {
+                LOG.debug("{} Exception thrown while refreshing the leadership 
status during stop", logPrefix, e);
+                return;
+            }
+            // a failure must not end the refresh loop: the next run pulls the 
current state from the cluster again
+            LOG.warn("{} Error while refreshing the leadership status (will 
try again): {}", logPrefix, e.getMessage());
+            LOG.debug("{} Exception thrown while refreshing the leadership 
status", logPrefix, e);
+            try {
+                executor.schedule(this::refreshStatus,
+                        jitter(this.lockConfiguration.getRetryPeriodMillis(), 
this.lockConfiguration.getJitterFactor()),
+                        TimeUnit.MILLISECONDS);
+            } catch (RejectedExecutionException ex) {
+                LOG.debug("{} Cannot reschedule the leadership status refresh 
as the controller is stopping", logPrefix);
+            }
+        }
+    }
+
+    private void doRefreshStatus() {
         switch (currentState) {
             case NOT_LEADER:
                 refreshStatusNotLeader();
diff --git 
a/components/camel-kubernetes/src/test/java/org/apache/camel/component/kubernetes/cluster/KubernetesClusterServiceTest.java
 
b/components/camel-kubernetes/src/test/java/org/apache/camel/component/kubernetes/cluster/KubernetesClusterServiceTest.java
index 1f6e986c31d9..2e5da6289142 100644
--- 
a/components/camel-kubernetes/src/test/java/org/apache/camel/component/kubernetes/cluster/KubernetesClusterServiceTest.java
+++ 
b/components/camel-kubernetes/src/test/java/org/apache/camel/component/kubernetes/cluster/KubernetesClusterServiceTest.java
@@ -248,6 +248,33 @@ public class KubernetesClusterServiceTest extends 
CamelTestSupport {
         assertEquals(firstLeader, mypod2.getCurrentLeader());
     }
 
+    @Test
+    public void testLeaderKeepsLeadershipAfterFailedLeaseRenewal() {
+        LeaderRecorder mypod1 = addMember("mypod1", LeaseResourceType.Lease);
+        LeaderRecorder mypod2 = addMember("mypod2", LeaseResourceType.Lease);
+        context.start();
+
+        mypod1.waitForAnyLeader(5, TimeUnit.SECONDS);
+        mypod2.waitForAnyLeader(5, TimeUnit.SECONDS);
+
+        String leader = mypod1.getCurrentLeader();
+        assertNotNull(leader);
+        assertEquals(leader, mypod2.getCurrentLeader());
+
+        // the leader can read the lease but the renewal of the lease fails 
once
+        withLockServer(leader, server -> server.setRefuseUpdateRequests(true));
+        await().atMost(5, TimeUnit.SECONDS)
+                .untilAsserted(() -> withLockServer(leader, server -> 
assertTrue(server.getRefusedUpdateRequests() > 0)));
+        withLockServer(leader, server -> 
server.setRefuseUpdateRequests(false));
+
+        // the leader keeps renewing the lease, so both pods keep seeing it as 
the leader
+        await().during(LEASE_TIME_MILLIS, TimeUnit.MILLISECONDS).atMost(3 * 
LEASE_TIME_MILLIS, TimeUnit.MILLISECONDS)
+                .untilAsserted(() -> {
+                    assertEquals(leader, mypod1.getCurrentLeader());
+                    assertEquals(leader, mypod2.getCurrentLeader());
+                });
+    }
+
     @Test
     public void testSharedConfigMap() {
         LeaderRecorder a1 = addMember("a1", LeaseResourceType.ConfigMap);
diff --git 
a/components/camel-kubernetes/src/test/java/org/apache/camel/component/kubernetes/cluster/utils/LockTestServer.java
 
b/components/camel-kubernetes/src/test/java/org/apache/camel/component/kubernetes/cluster/utils/LockTestServer.java
index 662ba3bc0aaf..6a8e6549bd09 100644
--- 
a/components/camel-kubernetes/src/test/java/org/apache/camel/component/kubernetes/cluster/utils/LockTestServer.java
+++ 
b/components/camel-kubernetes/src/test/java/org/apache/camel/component/kubernetes/cluster/utils/LockTestServer.java
@@ -23,6 +23,7 @@ import java.util.HashMap;
 import java.util.Map;
 import java.util.Set;
 import java.util.TreeSet;
+import java.util.concurrent.atomic.AtomicInteger;
 import java.util.stream.Collectors;
 
 import com.fasterxml.jackson.databind.ObjectMapper;
@@ -57,6 +58,10 @@ public class LockTestServer<T extends HasMetadata> extends 
KubernetesMockServer
 
     private boolean refuseRequests;
 
+    private volatile boolean refuseUpdateRequests;
+
+    private final AtomicInteger refusedUpdateRequests = new AtomicInteger();
+
     private Long delayRequests;
 
     private Set<String> pods;
@@ -241,6 +246,10 @@ public class LockTestServer<T extends HasMetadata> extends 
KubernetesMockServer
                         if (refuseRequests) {
                             return 500;
                         }
+                        if (refuseUpdateRequests) {
+                            refusedUpdateRequests.incrementAndGet();
+                            return 500;
+                        }
 
                         T resource;
                         try {
@@ -288,6 +297,17 @@ public class LockTestServer<T extends HasMetadata> extends 
KubernetesMockServer
         this.refuseRequests = refuseRequests;
     }
 
+    /**
+     * Refuses (with an internal server error) only the requests that update 
the lock resource.
+     */
+    public void setRefuseUpdateRequests(boolean refuseUpdateRequests) {
+        this.refuseUpdateRequests = refuseUpdateRequests;
+    }
+
+    public int getRefusedUpdateRequests() {
+        return refusedUpdateRequests.get();
+    }
+
     public synchronized Collection<String> getCurrentPods() {
         return new TreeSet<>(this.pods);
     }

Reply via email to