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
+        }
+    }
+}

Reply via email to