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 174a4cd750e4 CAMEL-25212: camel-infinispan - cluster view must keep
refreshing the leadership after an error, and give it up on stop (#27161)
174a4cd750e4 is described below
commit 174a4cd750e41f7ff131cd785447ff8395f406bc
Author: allthingssecurity <[email protected]>
AuthorDate: Thu Oct 1 21:10:36 2026 +0530
CAMEL-25212: camel-infinispan - cluster view must keep refreshing the
leadership after an error, and give it up on stop (#27161)
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../infinispan/cluster/InfinispanClusterView.java | 10 +-
.../cluster/InfinispanEmbeddedClusterView.java | 46 ++++-
...nfinispanEmbeddedClusterViewLeadershipTest.java | 210 +++++++++++++++++++++
.../cluster/InfinispanRemoteClusterView.java | 56 ++++--
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 13 ++
5 files changed, 316 insertions(+), 19 deletions(-)
diff --git
a/components/camel-infinispan/camel-infinispan-common/src/main/java/org/apache/camel/component/infinispan/cluster/InfinispanClusterView.java
b/components/camel-infinispan/camel-infinispan-common/src/main/java/org/apache/camel/component/infinispan/cluster/InfinispanClusterView.java
index e7e9e7eb3e5e..c3cfe1119d81 100644
---
a/components/camel-infinispan/camel-infinispan-common/src/main/java/org/apache/camel/component/infinispan/cluster/InfinispanClusterView.java
+++
b/components/camel-infinispan/camel-infinispan-common/src/main/java/org/apache/camel/component/infinispan/cluster/InfinispanClusterView.java
@@ -56,7 +56,15 @@ public abstract class InfinispanClusterView extends
AbstractCamelClusterView {
if (!master && this.leader.compareAndSet(true, false)) {
LOGGER.debug("Leadership lost for id: {}", id);
- fireLeadershipChangedEvent(getLeader().orElse(null));
+ CamelClusterMember newLeader = null;
+ try {
+ // the leader key may still hold this member (it is
removed after the leadership is given up on stop)
+ newLeader = getLeader().filter(m -> !Objects.equals(id,
m.getId())).orElse(null);
+ } catch (Exception e) {
+ // the event must be fired even if the current leader
cannot be looked up
+ LOGGER.debug("Cannot look up the current leader", e);
+ }
+ fireLeadershipChangedEvent(newLeader);
return;
}
}
diff --git
a/components/camel-infinispan/camel-infinispan-embedded/src/main/java/org/apache/camel/component/infinispan/embedded/cluster/InfinispanEmbeddedClusterView.java
b/components/camel-infinispan/camel-infinispan-embedded/src/main/java/org/apache/camel/component/infinispan/embedded/cluster/InfinispanEmbeddedClusterView.java
index 4d5d34c6a28d..fe6cd0b47352 100644
---
a/components/camel-infinispan/camel-infinispan-embedded/src/main/java/org/apache/camel/component/infinispan/embedded/cluster/InfinispanEmbeddedClusterView.java
+++
b/components/camel-infinispan/camel-infinispan-embedded/src/main/java/org/apache/camel/component/infinispan/embedded/cluster/InfinispanEmbeddedClusterView.java
@@ -182,11 +182,17 @@ public class InfinispanEmbeddedClusterView extends
InfinispanClusterView {
getCamelContext().getExecutorServiceManager().shutdownGraceful(executorService);
- if (cache != null) {
- cache.remove(InfinispanClusterService.LEADER_KEY,
getClusterService().getId());
+ try {
+ // tell the listeners (such as clustered routes) that the
local member is no longer the leader, before
+ // the leader key is removed and another member can take over
the leadership
+ setLeader(false);
+ } finally {
+ if (cache != null) {
+ cache.remove(InfinispanClusterService.LEADER_KEY,
getClusterService().getId());
- LOGGER.info("Removing local member, key={}",
getLocalMember().getId());
- cache.remove(getLocalMember().getId());
+ LOGGER.info("Removing local member, key={}",
getLocalMember().getId());
+ cache.remove(getLocalMember().getId());
+ }
}
}
@@ -205,6 +211,15 @@ public class InfinispanEmbeddedClusterView extends
InfinispanClusterView {
return;
}
+ refreshLeadership();
+ refreshMembership();
+ } finally {
+ lock.unlock();
+ }
+ }
+
+ private void refreshLeadership() {
+ try {
final String leaderKey = InfinispanClusterService.LEADER_KEY;
final String localId = getLocalMember().getId();
@@ -240,12 +255,29 @@ public class InfinispanEmbeddedClusterView extends
InfinispanClusterView {
setLeader(false);
}
}
+ } catch (Exception e) {
+ // an exception must not end the periodic refresh of the
leadership (as it would with an exception
+ // thrown out of this task), and as the leadership could not
be refreshed, give it up until the next run
+ LOGGER.warn("Error while refreshing the leadership of id={}
(will try again): {}",
+ getLocalMember().getId(), e.getMessage());
+ LOGGER.debug("Error while refreshing the leadership", e);
+ try {
+ setLeader(false);
+ } catch (Exception ex) {
+ LOGGER.debug("Error while giving up the leadership", ex);
+ }
+ }
+ }
- // refresh local membership
+ private void refreshMembership() {
+ try {
cache.put(getLocalMember().getId(), isLeader() ? "true" :
"false", configuration.getLifespan(),
configuration.getLifespanTimeUnit());
- } finally {
- lock.unlock();
+ } catch (Exception e) {
+ // the membership entry does not hold the leadership, so keep
it and try again on the next run
+ LOGGER.warn("Error while refreshing the membership of id={}
(will try again): {}",
+ getLocalMember().getId(), e.getMessage());
+ LOGGER.debug("Error while refreshing the membership", e);
}
}
diff --git
a/components/camel-infinispan/camel-infinispan-embedded/src/test/java/org/apache/camel/component/infinispan/embedded/cluster/InfinispanEmbeddedClusterViewLeadershipTest.java
b/components/camel-infinispan/camel-infinispan-embedded/src/test/java/org/apache/camel/component/infinispan/embedded/cluster/InfinispanEmbeddedClusterViewLeadershipTest.java
new file mode 100644
index 000000000000..f02a45b7ac8a
--- /dev/null
+++
b/components/camel-infinispan/camel-infinispan-embedded/src/test/java/org/apache/camel/component/infinispan/embedded/cluster/InfinispanEmbeddedClusterViewLeadershipTest.java
@@ -0,0 +1,210 @@
+/*
+ * 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.infinispan.embedded.cluster;
+
+import java.util.List;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.camel.cluster.CamelClusterEventListener;
+import org.apache.camel.cluster.CamelClusterMember;
+import org.apache.camel.cluster.CamelClusterView;
+import org.apache.camel.component.infinispan.cluster.InfinispanClusterService;
+import org.apache.camel.impl.DefaultCamelContext;
+import org.infinispan.Cache;
+import org.infinispan.cache.impl.AbstractDelegatingCache;
+import org.infinispan.commons.CacheException;
+import org.infinispan.manager.DefaultCacheManager;
+import org.infinispan.manager.EmbeddedCacheManager;
+import org.infinispan.manager.impl.AbstractDelegatingEmbeddedCacheManager;
+import org.junit.jupiter.api.Test;
+
+import static
org.apache.camel.component.infinispan.embedded.cluster.InfinispanEmbeddedClusteredTestSupport.createCache;
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.awaitility.Awaitility.await;
+
+public class InfinispanEmbeddedClusterViewLeadershipTest {
+ private static final String VIEW_NAME = "myView";
+ private static final String NODE_ID = "node";
+
+ @Test
+ public void leadershipIsGivenUpAndRefreshedAfterALeaderKeyError() throws
Exception {
+ try (DefaultCacheManager cacheContainer = new DefaultCacheManager()) {
+ createCache(cacheContainer, VIEW_NAME);
+ FailingCacheManager cacheManager = new
FailingCacheManager(cacheContainer);
+
+ try (DefaultCamelContext context = new DefaultCamelContext()) {
+ context.disableJMX();
+ InfinispanEmbeddedClusterService clusterService =
createClusterService(cacheManager, 1000);
+ context.addService(clusterService);
+ context.start();
+
+ CamelClusterView view = clusterService.getView(VIEW_NAME);
+ List<Boolean> events = new CopyOnWriteArrayList<>();
+ view.addEventListener((CamelClusterEventListener.Leadership)
(v, leader) -> events
+ .add(v.getLocalMember().isLeader()));
+ await().atMost(10, TimeUnit.SECONDS).until(() ->
view.getLocalMember().isLeader());
+ await().atMost(10, TimeUnit.SECONDS).until(() ->
cacheManager.replaced.get() > 0);
+ int received = events.size();
+
+ // the next refresh of the leader key fails
+ cacheManager.failNextReplace.set(true);
+ await().atMost(10, TimeUnit.SECONDS).until(() ->
!cacheManager.failNextReplace.get());
+ int replaced = cacheManager.replaced.get();
+
+ // the leadership is given up, and as it is still refreshed
afterwards the node is the leader again
+ await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> {
+
assertThat(cacheManager.replaced.get()).isGreaterThan(replaced + 1);
+ assertThat(view.getLocalMember().isLeader()).isTrue();
+ });
+ assertThat(events.subList(received,
events.size())).containsSubsequence(false, true);
+ }
+ }
+ }
+
+ @Test
+ public void leadershipIsKeptAfterAMembershipError() throws Exception {
+ try (DefaultCacheManager cacheContainer = new DefaultCacheManager()) {
+ createCache(cacheContainer, VIEW_NAME);
+ FailingCacheManager cacheManager = new
FailingCacheManager(cacheContainer);
+
+ try (DefaultCamelContext context = new DefaultCamelContext()) {
+ context.disableJMX();
+ // a longer lifespan, so that a slow refresh does not let the
leader key expire and change the leadership
+ InfinispanEmbeddedClusterService clusterService =
createClusterService(cacheManager, 2000);
+ context.addService(clusterService);
+ context.start();
+
+ CamelClusterView view = clusterService.getView(VIEW_NAME);
+ List<Boolean> events = new CopyOnWriteArrayList<>();
+ view.addEventListener((CamelClusterEventListener.Leadership)
(v, leader) -> events
+ .add(v.getLocalMember().isLeader()));
+ await().atMost(10, TimeUnit.SECONDS).until(() ->
view.getLocalMember().isLeader());
+ int received = events.size();
+
+ // the next refresh of the membership fails
+ cacheManager.failNextMembershipPut.set(true);
+ await().atMost(10, TimeUnit.SECONDS).until(() ->
!cacheManager.failNextMembershipPut.get());
+ int replaced = cacheManager.replaced.get();
+ int membershipPut = cacheManager.membershipPut.get();
+
+ // the leadership and the membership are still refreshed
afterwards
+ await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> {
+
assertThat(cacheManager.replaced.get()).isGreaterThan(replaced + 1);
+
assertThat(cacheManager.membershipPut.get()).isGreaterThan(membershipPut + 1);
+ });
+
+ // and the leadership has not been given up meanwhile
+ assertThat(view.getLocalMember().isLeader()).isTrue();
+ assertThat(events).hasSize(received);
+ }
+ }
+ }
+
+ @Test
+ public void
stoppingTheViewGivesUpTheLeadershipBeforeReleasingTheLeaderKey() throws
Exception {
+ try (DefaultCacheManager cacheContainer = new DefaultCacheManager()) {
+ createCache(cacheContainer, VIEW_NAME);
+ Cache<String, String> cache = cacheContainer.getCache(VIEW_NAME);
+
+ try (DefaultCamelContext context = new DefaultCamelContext()) {
+ context.disableJMX();
+ InfinispanEmbeddedClusterService clusterService =
createClusterService(cacheContainer, 1000);
+ context.addService(clusterService);
+ context.start();
+
+ CamelClusterView view = clusterService.getView(VIEW_NAME);
+ List<Boolean> events = new CopyOnWriteArrayList<>();
+ List<CamelClusterMember> leaders = new
CopyOnWriteArrayList<>();
+ List<String> leaderKeys = new CopyOnWriteArrayList<>();
+ view.addEventListener((CamelClusterEventListener.Leadership)
(v, leader) -> {
+ events.add(v.getLocalMember().isLeader());
+ leaders.add(leader);
+
leaderKeys.add(String.valueOf(cache.get(InfinispanClusterService.LEADER_KEY)));
+ });
+ await().atMost(10, TimeUnit.SECONDS).until(() ->
view.getLocalMember().isLeader());
+ int received = events.size();
+
+ view.stop();
+
+ // the listeners are told that the local member is no longer
the leader
+ assertThat(view.getLocalMember().isLeader()).isFalse();
+ assertThat(events).hasSize(received + 1);
+ assertThat(events.get(received)).isFalse();
+ assertThat(leaders.get(received)).isNull();
+
+ // while the leader key was still held, so that no other
member could take over the leadership yet
+ assertThat(leaderKeys.get(received)).isEqualTo(NODE_ID);
+
+ // which is released afterwards
+
assertThat(cache.get(InfinispanClusterService.LEADER_KEY)).isNull();
+ }
+ }
+ }
+
+ private static InfinispanEmbeddedClusterService
createClusterService(EmbeddedCacheManager cacheManager, long lifespan) {
+ InfinispanEmbeddedClusterService clusterService = new
InfinispanEmbeddedClusterService();
+ clusterService.setCacheContainer(cacheManager);
+ clusterService.setId(NODE_ID);
+ clusterService.setLifespan(lifespan);
+ clusterService.setLifespanTimeUnit(TimeUnit.MILLISECONDS);
+ return clusterService;
+ }
+
+ /**
+ * A cache manager whose caches count the refreshes of the leader key and
of the membership, and fail the next one
+ * when requested.
+ */
+ private static final class FailingCacheManager extends
AbstractDelegatingEmbeddedCacheManager {
+ private final AtomicBoolean failNextReplace = new AtomicBoolean();
+ private final AtomicBoolean failNextMembershipPut = new
AtomicBoolean();
+ private final AtomicInteger replaced = new AtomicInteger();
+ private final AtomicInteger membershipPut = new AtomicInteger();
+
+ private FailingCacheManager(EmbeddedCacheManager cm) {
+ super(cm);
+ }
+
+ @Override
+ public <K, V> Cache<K, V> getCache(String cacheName) {
+ Cache<K, V> cache = super.getCache(cacheName);
+ return new AbstractDelegatingCache<>(cache) {
+ @Override
+ public boolean replace(K key, V oldValue, V value, long
lifespan, TimeUnit unit) {
+ if (failNextReplace.compareAndSet(true, false)) {
+ throw new CacheException("Simulated failure while
refreshing the leader key");
+ }
+ replaced.incrementAndGet();
+ return super.replace(key, oldValue, value, lifespan, unit);
+ }
+
+ @Override
+ public V put(K key, V value, long lifespan, TimeUnit unit) {
+ if (NODE_ID.equals(key)) {
+ if (failNextMembershipPut.compareAndSet(true, false)) {
+ throw new CacheException("Simulated failure while
refreshing the membership");
+ }
+ membershipPut.incrementAndGet();
+ }
+ return super.put(key, value, lifespan, unit);
+ }
+ };
+ }
+ }
+}
diff --git
a/components/camel-infinispan/camel-infinispan/src/main/java/org/apache/camel/component/infinispan/remote/cluster/InfinispanRemoteClusterView.java
b/components/camel-infinispan/camel-infinispan/src/main/java/org/apache/camel/component/infinispan/remote/cluster/InfinispanRemoteClusterView.java
index f7a634a199e7..46417b936a5b 100644
---
a/components/camel-infinispan/camel-infinispan/src/main/java/org/apache/camel/component/infinispan/remote/cluster/InfinispanRemoteClusterView.java
+++
b/components/camel-infinispan/camel-infinispan/src/main/java/org/apache/camel/component/infinispan/remote/cluster/InfinispanRemoteClusterView.java
@@ -190,16 +190,24 @@ public class InfinispanRemoteClusterView extends
InfinispanClusterView {
getCamelContext().getExecutorServiceManager().shutdownGraceful(executorService);
- if (cache != null) {
- if (this.version != null) {
-
cache.removeWithVersion(InfinispanClusterService.LEADER_KEY, this.version);
+ try {
+ // tell the listeners (such as clustered routes) that the
local member is no longer the leader, before
+ // the leader key is removed and another member can take over
the leadership
+ setLeader(false);
+ } finally {
+ try {
+ if (cache != null) {
+ if (this.version != null) {
+
cache.removeWithVersion(InfinispanClusterService.LEADER_KEY, this.version);
+ }
+
+ LOGGER.info("Removing local member, key={}",
getLocalMember().getId());
+ cache.remove(getLocalMember().getId());
+ }
+ } finally {
+ this.version = null;
}
-
- LOGGER.info("Removing local member, key={}",
getLocalMember().getId());
- cache.remove(getLocalMember().getId());
}
-
- this.version = null;
}
private boolean isLeader() {
@@ -217,6 +225,15 @@ public class InfinispanRemoteClusterView extends
InfinispanClusterView {
return;
}
+ refreshLeadership();
+ refreshMembership();
+ } finally {
+ lock.unlock();
+ }
+ }
+
+ private void refreshLeadership() {
+ try {
final String leaderKey = InfinispanClusterService.LEADER_KEY;
final String localId = getLocalMember().getId();
@@ -265,12 +282,29 @@ public class InfinispanRemoteClusterView extends
InfinispanClusterView {
setLeader(false);
}
}
+ } catch (Exception e) {
+ // an exception must not end the periodic refresh of the
leadership (as it would with an exception
+ // thrown out of this task), and as the leadership could not
be refreshed, give it up until the next run
+ LOGGER.warn("Error while refreshing the leadership of id={}
(will try again): {}",
+ getLocalMember().getId(), e.getMessage());
+ LOGGER.debug("Error while refreshing the leadership", e);
+ try {
+ setLeader(false);
+ } catch (Exception ex) {
+ LOGGER.debug("Error while giving up the leadership", ex);
+ }
+ }
+ }
- // refresh local membership
+ private void refreshMembership() {
+ try {
cache.put(getLocalMember().getId(), isLeader() ? "true" :
"false", configuration.getLifespan(),
configuration.getLifespanTimeUnit());
- } finally {
- lock.unlock();
+ } catch (Exception e) {
+ // the membership entry does not hold the leadership, so keep
it and try again on the next run
+ LOGGER.warn("Error while refreshing the membership of id={}
(will try again): {}",
+ getLocalMember().getId(), e.getMessage());
+ LOGGER.debug("Error while refreshing the membership", e);
}
}
diff --git
a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
index 2df785d4225b..7fbdd1071cb2 100644
--- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
+++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
@@ -3676,6 +3676,19 @@ delay). Previously the routes were then started right
away.
The `acquireLockInterval` of the file lock cluster service must now be at
least 1 millisecond.
+=== camel-infinispan - clustered routes stop when the Infinispan cluster view
is stopped
+
+When an Infinispan cluster view (`camel-infinispan` or
`camel-infinispan-embedded`) is stopped while this member is
+the leader, for example with the `stopView` JMX operation, the stop of the
cluster service or the stop of the
+clustered route controller, the listeners are now told that the leadership is
lost, before the leader key is removed
+from the cache. The clustered routes (`ClusteredRoutePolicy`) and the `master`
consumers are then stopped. Previously
+they kept running while the leader key was removed, so another member could
take the leadership and both members ran
+the routes.
+
+An error while refreshing the leader key no longer stops the periodic refresh
of the leadership. The leadership is
+given up (so the clustered routes are stopped) until the next refresh
succeeds, as this member can no longer be sure
+that it holds the leader key. An error while refreshing only the membership
entry of this member keeps the leadership.
+
=== camel-pqc - FileBasedKeyLifecycleManager restricts keyId to a flat file
name
`FileBasedKeyLifecycleManager` now confines every key file to its configured
key directory. A `keyId`