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