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 ff898996c7b1 CAMEL-25286: camel-consul - cluster view must keep taking
part in the leader election after a failed query or an invalidated session
(#27311)
ff898996c7b1 is described below
commit ff898996c7b19fc4dd496551a161766958c0a888
Author: allthingssecurity <[email protected]>
AuthorDate: Sun Oct 4 12:34:20 2026 +0530
CAMEL-25286: camel-consul - cluster view must keep taking part in the
leader election after a failed query or an invalidated session (#27311)
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../consul/cluster/ConsulClusterView.java | 124 ++++++--
.../cluster/ConsulClusterViewRecoveryTest.java | 318 +++++++++++++++++++++
2 files changed, 416 insertions(+), 26 deletions(-)
diff --git
a/components/camel-consul/src/main/java/org/apache/camel/component/consul/cluster/ConsulClusterView.java
b/components/camel-consul/src/main/java/org/apache/camel/component/consul/cluster/ConsulClusterView.java
index cfa1fae71f43..ca8277173723 100644
---
a/components/camel-consul/src/main/java/org/apache/camel/component/consul/cluster/ConsulClusterView.java
+++
b/components/camel-consul/src/main/java/org/apache/camel/component/consul/cluster/ConsulClusterView.java
@@ -20,6 +20,8 @@ import java.math.BigInteger;
import java.util.Collections;
import java.util.List;
import java.util.Optional;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicReference;
import java.util.concurrent.locks.Lock;
@@ -30,6 +32,7 @@ import org.apache.camel.cluster.CamelClusterMember;
import org.apache.camel.support.cluster.AbstractCamelClusterView;
import org.apache.camel.util.ObjectHelper;
import org.kiwiproject.consul.Consul;
+import org.kiwiproject.consul.ConsulException;
import org.kiwiproject.consul.KeyValueClient;
import org.kiwiproject.consul.SessionClient;
import org.kiwiproject.consul.async.ConsulResponseCallback;
@@ -53,6 +56,7 @@ final class ConsulClusterView extends
AbstractCamelClusterView {
private Consul client;
private SessionClient sessionClient;
private KeyValueClient keyValueClient;
+ private ScheduledExecutorService executorService;
private String path;
ConsulClusterView(ConsulClusterService service, ConsulClusterConfiguration
configuration, String namespace) {
@@ -96,22 +100,33 @@ final class ConsulClusterView extends
AbstractCamelClusterView {
sessionClient = client.sessionClient();
keyValueClient = client.keyValueClient();
- sessionId.set(sessionClient
-
.createSession(ImmutableSession.builder().name(getNamespace()).ttl(configuration.getSessionTtl()
+ "s")
- .lockDelay(configuration.getSessionLockDelay() +
"s").build())
- .getId());
-
+ sessionId.set(createSession());
LOGGER.debug("Acquired session with id '{}'", sessionId.get());
- boolean lock = acquireLock();
- LOGGER.debug("Acquire lock on path '{}' with id '{}' result '{}'",
path, sessionId.get(), lock);
- localMember.setMaster(lock);
- watcher.watch();
+ // to watch again after a failed query. Created once the session
exists, as a view that fails to start is
+ // not stopped
+ executorService =
getCamelContext().getExecutorServiceManager().newSingleThreadScheduledExecutor(this,
+ "ConsulClusterView");
+ try {
+ boolean lock = acquireLock();
+ LOGGER.debug("Acquire lock on path '{}' with id '{}' result
'{}'", path, sessionId.get(), lock);
+
+ localMember.setMaster(lock);
+ watcher.watch();
+ } catch (Exception e) {
+
getCamelContext().getExecutorServiceManager().shutdownNow(executorService);
+ executorService = null;
+ throw e;
+ }
}
}
@Override
protected void doStop() throws Exception {
+ if (executorService != null) {
+
getCamelContext().getExecutorServiceManager().shutdownNow(executorService);
+ executorService = null;
+ }
if (sessionId.get() != null) {
if (keyValueClient.releaseLock(this.path, sessionId.get())) {
LOGGER.debug("Successfully released lock on path '{}' with id
'{}'", path, sessionId.get());
@@ -126,6 +141,55 @@ final class ConsulClusterView extends
AbstractCamelClusterView {
}
}
+ private String createSession() {
+ return sessionClient
+
.createSession(ImmutableSession.builder().name(getNamespace()).ttl(configuration.getSessionTtl()
+ "s")
+ .lockDelay(configuration.getSessionLockDelay() +
"s").build())
+ .getId();
+ }
+
+ private void renewSession(String sid) {
+ try {
+ if (sessionClient.renewSession(sid).isPresent()) {
+ return;
+ }
+ } catch (ConsulException e) {
+ if (!e.hasCode() || e.getCode() != 404) {
+ // for example Consul cannot be reached: the session is
renewed again with the next query
+ LOGGER.debug("Failed to renew session with id '{}': {}", sid,
e.getMessage(), e);
+ return;
+ }
+ }
+
+ // the session does not exist anymore: Consul invalidated it (its TTL
expired, the health check of the agent
+ // failed, Consul lost its data) and released the lock. Create a new
session, or this node can never take the
+ // leadership again
+ sessionIdLock.lock();
+ try {
+ if ((isStarting() || isStarted()) && sid.equals(sessionId.get())) {
+ localMember.setMaster(false);
+ sessionId.set(createSession());
+ LOGGER.info("Session with id '{}' was invalidated by Consul,
created session with id '{}'", sid,
+ sessionId.get());
+ }
+ } catch (Exception e) {
+ // tried again with the next query
+ LOGGER.debug("Failed to create a session to replace session with
id '{}': {}", sid, e.getMessage(), e);
+ } finally {
+ sessionIdLock.unlock();
+ }
+ }
+
+ private CamelClusterMember currentLeader() {
+ try {
+ return getLeader().orElse(null);
+ } catch (Exception e) {
+ // for example Consul cannot be reached
+ LOGGER.debug("Failed to get the leader on path '{}': {}", path,
e.getMessage(), e);
+ return null;
+ }
+ }
+
private boolean acquireLock() {
sessionIdLock.lock();
try {
@@ -153,7 +217,7 @@ final class ConsulClusterView extends
AbstractCamelClusterView {
}
if (!master && this.master.compareAndSet(true, false)) {
LOGGER.debug("Leadership lost for session id {}",
sessionId.get());
- fireLeadershipChangedEvent(getLeader().orElse(null));
+ fireLeadershipChangedEvent(currentLeader());
}
}
@@ -239,12 +303,12 @@ final class ConsulClusterView extends
AbstractCamelClusterView {
@Override
public void onComplete(ConsulResponse<Optional<Value>> consulResponse)
{
if (isStarting() || isStarted()) {
- Optional<Value> value = consulResponse.getResponse();
- if (value.isPresent()) {
- Optional<String> sid = value.get().getSession();
+ index.set(consulResponse.getIndex());
+ try {
+ Optional<String> sid =
consulResponse.getResponse().flatMap(Value::getSession);
if (!sid.isPresent()) {
- // If the key is not held by any session, try acquire a
- // lock (become leader)
+ // If the key is not held by any session (or does not
exist, for
+ // example after Consul lost its data), try acquire a
lock (become leader)
boolean lock = acquireLock();
LOGGER.debug("Try to acquire lock on path '{}' with id
'{}', result '{}'", path, sessionId.get(), lock);
@@ -258,9 +322,12 @@ final class ConsulClusterView extends
AbstractCamelClusterView {
localMember.setMaster(sid.get().equals(sessionId.get()));
}
+ } catch (Exception e) {
+ // for example Consul cannot be reached anymore: the
leadership cannot be confirmed
+ LOGGER.debug("Failed to update the leadership on path
'{}': {}", path, e.getMessage(), e);
+ localMember.setMaster(false);
}
- index.set(consulResponse.getIndex());
watch();
}
}
@@ -269,16 +336,23 @@ final class ConsulClusterView extends
AbstractCamelClusterView {
public void onFailure(Throwable throwable) {
LOGGER.debug("{}", throwable.getMessage(), throwable);
- if (sessionId.get() != null) {
- keyValueClient.releaseLock(configuration.getRootPath(),
sessionId.get());
- }
-
+ // the leadership cannot be confirmed: give it up locally, which
can only lead to no leader, never to two.
+ // The lock is kept: releasing it explicitly skips the lock-delay
of Consul, so another node could take the
+ // leadership while the clustered routes of this node are still
stopping. If this node really is cut off
+ // from Consul, its session expires and Consul releases the lock
and applies the lock-delay
localMember.setMaster(false);
- watch();
+
+ // keep watching, the leadership is taken again when Consul
answers. Wait, so that a Consul agent that
+ // cannot be reached is not queried in a loop
+ ScheduledExecutorService executor = executorService;
+ if ((isStarting() || isStarted()) && executor != null) {
+ executor.schedule(this::watch, Math.max(1,
configuration.getSessionRefreshInterval()), TimeUnit.SECONDS);
+ }
}
public void watch() {
- if (sessionId.get() == null) {
+ String sid = sessionId.get();
+ if (sid == null) {
return;
}
@@ -287,10 +361,8 @@ final class ConsulClusterView extends
AbstractCamelClusterView {
keyValueClient.getValue(path,
QueryOptions.blockSeconds(configuration.getSessionRefreshInterval(),
index.get()).build(), this);
- if (sessionId.get() != null) {
- // Refresh session
- sessionClient.renewSession(sessionId.get());
- }
+ // Refresh session
+ renewSession(sid);
}
}
}
diff --git
a/components/camel-consul/src/test/java/org/apache/camel/component/consul/cluster/ConsulClusterViewRecoveryTest.java
b/components/camel-consul/src/test/java/org/apache/camel/component/consul/cluster/ConsulClusterViewRecoveryTest.java
new file mode 100644
index 000000000000..8130f2eea3af
--- /dev/null
+++
b/components/camel-consul/src/test/java/org/apache/camel/component/consul/cluster/ConsulClusterViewRecoveryTest.java
@@ -0,0 +1,318 @@
+/*
+ * 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.consul.cluster;
+
+import java.math.BigInteger;
+import java.net.ConnectException;
+import java.util.List;
+import java.util.Optional;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import okhttp3.Request;
+import okhttp3.ResponseBody;
+import org.apache.camel.CamelContext;
+import org.apache.camel.cluster.CamelClusterView;
+import org.apache.camel.impl.DefaultCamelContext;
+import org.apache.camel.spi.ExecutorServiceManager;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.kiwiproject.consul.Consul;
+import org.kiwiproject.consul.ConsulException;
+import org.kiwiproject.consul.KeyValueClient;
+import org.kiwiproject.consul.SessionClient;
+import org.kiwiproject.consul.async.ConsulResponseCallback;
+import org.kiwiproject.consul.model.ConsulResponse;
+import org.kiwiproject.consul.model.kv.ImmutableValue;
+import org.kiwiproject.consul.model.kv.Value;
+import org.kiwiproject.consul.model.session.ImmutableSessionCreatedResponse;
+import org.kiwiproject.consul.model.session.Session;
+import org.kiwiproject.consul.model.session.SessionInfo;
+import org.kiwiproject.consul.option.QueryOptions;
+import retrofit2.Call;
+import retrofit2.Response;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.spy;
+import static org.mockito.Mockito.timeout;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * The cluster view must keep taking part in the leader election after Consul
could not be reached and after Consul
+ * invalidated its session, as it happens when the Consul agent restarts.
Consul is simulated by mocked clients.
+ */
+class ConsulClusterViewRecoveryTest {
+
+ private static final String NAMESPACE = "my-ns";
+ private static final String PATH = "/camel/" + NAMESPACE;
+
+ private final SessionClient sessionClient = mock(SessionClient.class);
+ private final KeyValueClient keyValueClient = mock(KeyValueClient.class);
+ private final SessionInfo sessionInfo = mock(SessionInfo.class);
+ private final List<ConsulResponseCallback<Optional<Value>>> queries = new
CopyOnWriteArrayList<>();
+
+ // the state of the simulated Consul
+ private final Set<String> sessions = ConcurrentHashMap.newKeySet();
+ private final AtomicInteger createdSessions = new AtomicInteger();
+ private volatile boolean reachable = true;
+ private volatile boolean keyExists;
+ private volatile String lockHolder;
+
+ private ConsulClusterConfiguration configuration;
+ private CamelContext context;
+ private CamelClusterView view;
+
+ @BeforeEach
+ void setUp() throws Exception {
+ Consul consul = mock(Consul.class);
+ when(consul.sessionClient()).thenReturn(sessionClient);
+ when(consul.keyValueClient()).thenReturn(keyValueClient);
+ simulateConsul();
+
+ configuration = new ConsulClusterConfiguration() {
+ @Override
+ public Consul createConsulClient(CamelContext camelContext) {
+ return consul;
+ }
+ };
+ configuration.setSessionRefreshInterval(1);
+
+ ConsulClusterService service = new ConsulClusterService(configuration);
+ service.setId("node-1");
+
+ context = new DefaultCamelContext();
+ context.addService(service);
+ context.start();
+
+ view = service.getView(NAMESPACE);
+ assertTrue(view.getLocalMember().isLeader());
+ }
+
+ @AfterEach
+ void tearDown() {
+ if (context != null) {
+ context.stop();
+ }
+ }
+
+ @Test
+ void keepsWatchingAfterConsulCouldNotBeReached() {
+ // the agent cannot be reached: the pending query fails
+ reachable = false;
+ deliver(() -> lastQuery().onFailure(notReachable()));
+
+ // the leadership cannot be confirmed anymore (the session is not
renewed): give it up
+ assertFalse(view.getLocalMember().isLeader());
+
+ // and query again, the agent is back
+ reachable = true;
+ verify(keyValueClient, timeout(5000).times(2)).getValue(eq(PATH),
any(QueryOptions.class), any());
+
+ // the lock was held by the session all along: the node is leader again
+ deliver(() -> lastQuery().onComplete(keyValue()));
+ assertTrue(view.getLocalMember().isLeader());
+ }
+
+ @Test
+ void keepsTheLockAfterAFailedQuery() {
+ String session = lockHolder;
+
+ deliver(() -> lastQuery().onFailure(notReachable()));
+
+ // the node steps down locally, but keeps the lock: an explicit
release would skip the lock-delay of Consul and
+ // let another node take the leadership while this node is still
stopping its clustered routes
+ assertFalse(view.getLocalMember().isLeader());
+ verify(keyValueClient, never()).releaseLock(anyString(), anyString());
+ assertEquals(session, lockHolder);
+ assertFalse(keyValueClient.acquireLock(PATH,
"session-of-another-node"));
+
+ // the next answer confirms that the session still holds the lock: the
node is leader again
+ verify(keyValueClient, timeout(5000).times(2)).getValue(eq(PATH),
any(QueryOptions.class), any());
+ deliver(() -> lastQuery().onComplete(keyValue()));
+ assertTrue(view.getLocalMember().isLeader());
+ }
+
+ @Test
+ void doesNotLeaveAnExecutorBehindWhenTheSessionCannotBeCreated() throws
Exception {
+ CamelContext otherContext = new DefaultCamelContext();
+ ExecutorServiceManager executorServiceManager =
spy(otherContext.getExecutorServiceManager());
+ otherContext.setExecutorServiceManager(executorServiceManager);
+ try {
+ ConsulClusterService service = new
ConsulClusterService(configuration);
+ service.setId("node-2");
+ otherContext.addService(service);
+ otherContext.start();
+
+ // the session cannot be created: the view fails to start, and a
failed view is not stopped
+ reachable = false;
+ assertThrows(Exception.class, () -> service.getView(NAMESPACE));
+
+ verify(executorServiceManager,
never()).newSingleThreadScheduledExecutor(any(), anyString());
+ } finally {
+ reachable = true;
+ otherContext.stop();
+ }
+ }
+
+ @Test
+ void createsANewSessionWhenConsulInvalidatedIt() {
+ // the agent restarts: Consul invalidates the session and releases its
lock
+ sessions.clear();
+ lockHolder = null;
+
+ deliver(() -> lastQuery().onComplete(keyValue()));
+ assertFalse(view.getLocalMember().isLeader());
+
+ // the next answer: the key is still free and the node takes the
leadership with a new session
+ deliver(() -> lastQuery().onComplete(keyValue()));
+ assertEquals(2, createdSessions.get());
+ assertEquals("session-2", lockHolder);
+ assertTrue(view.getLocalMember().isLeader());
+ }
+
+ @Test
+ void acquiresTheLockWhenTheKeyDoesNotExist() {
+ // Consul restarts without its data: no session and no key anymore
+ sessions.clear();
+ lockHolder = null;
+ keyExists = false;
+
+ // the node does not hold a lock anymore
+ deliver(() -> lastQuery().onComplete(keyValue()));
+ assertFalse(view.getLocalMember().isLeader());
+
+ // the next answer: the node creates the key and takes the leadership
with a new session
+ deliver(() -> lastQuery().onComplete(keyValue()));
+ assertEquals("session-2", lockHolder);
+ assertTrue(view.getLocalMember().isLeader());
+ }
+
+ private void simulateConsul() {
+ doAnswer(inv -> {
+ ensureReachable();
+ String id = "session-" + createdSessions.incrementAndGet();
+ sessions.add(id);
+ return ImmutableSessionCreatedResponse.builder().id(id).build();
+ }).when(sessionClient).createSession(any(Session.class));
+
+ doAnswer(inv -> {
+ ensureReachable();
+ return sessions.contains(inv.<String> getArgument(0)) ?
Optional.of(sessionInfo) : Optional.empty();
+ }).when(sessionClient).getSessionInfo(anyString());
+
+ doAnswer(inv -> {
+ ensureReachable();
+ String id = inv.getArgument(0);
+ if (!sessions.contains(id)) {
+ // what Consul answers for a session that does not exist
+ throw new ConsulException(
+ renewCall(id), Response.error(404,
ResponseBody.create("Session id '" + id + "' not found", null)));
+ }
+ return Optional.of(sessionInfo);
+ }).when(sessionClient).renewSession(anyString());
+
+ doAnswer(inv -> {
+ ensureReachable();
+ String id = inv.getArgument(1);
+ if (sessions.contains(id) && (lockHolder == null ||
lockHolder.equals(id))) {
+ lockHolder = id;
+ keyExists = true;
+ return true;
+ }
+ return false;
+ }).when(keyValueClient).acquireLock(eq(PATH), anyString());
+
+ doAnswer(inv -> {
+ ensureReachable();
+ if (PATH.equals(inv.getArgument(0)) &&
inv.getArgument(1).equals(lockHolder)) {
+ lockHolder = null;
+ return true;
+ }
+ return false;
+ }).when(keyValueClient).releaseLock(anyString(), anyString());
+
+ doAnswer(inv -> {
+ ensureReachable();
+ return Optional.ofNullable(lockHolder);
+ }).when(keyValueClient).getSession(PATH);
+
+ // blocking queries: answered by the test
+ doAnswer(inv -> {
+ queries.add(inv.getArgument(2));
+ return null;
+ }).when(keyValueClient).getValue(eq(PATH), any(QueryOptions.class),
any());
+ }
+
+ private void ensureReachable() {
+ if (!reachable) {
+ throw notReachable();
+ }
+ }
+
+ private static Call<?> renewCall(String id) {
+ Call<?> call = mock(Call.class);
+ when(call.request()).thenReturn(new
Request.Builder().url("http://localhost:8500/v1/session/renew/" + id).build());
+ return call;
+ }
+
+ private static ConsulException notReachable() {
+ return new ConsulException("Error connecting to Consul", new
ConnectException("Connection refused"));
+ }
+
+ private ConsulResponseCallback<Optional<Value>> lastQuery() {
+ return queries.get(queries.size() - 1);
+ }
+
+ private ConsulResponse<Optional<Value>> keyValue() {
+ Optional<Value> value = Optional.empty();
+ if (keyExists) {
+ value = Optional.of(ImmutableValue.builder()
+ .key(PATH)
+ .session(Optional.ofNullable(lockHolder))
+ .createIndex(1)
+ .modifyIndex(queries.size())
+ .lockIndex(1)
+ .flags(0)
+ .build());
+ }
+ return new ConsulResponse<>(value, 0, true,
BigInteger.valueOf(queries.size()), (String) null, (String) null);
+ }
+
+ /**
+ * Calls the callback of a query like the HTTP client does: an exception
thrown by the callback is only logged.
+ */
+ private static void deliver(Runnable callback) {
+ try {
+ callback.run();
+ } catch (RuntimeException e) {
+ // ignored, as by the HTTP client
+ }
+ }
+}