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 59de7b0e4c83 CAMEL-25275: camel-consul - key/value and event consumers 
must keep watching after a failed query (#27304)
59de7b0e4c83 is described below

commit 59de7b0e4c83dedf0d7156f37ff500463b79df74
Author: allthingssecurity <[email protected]>
AuthorDate: Sat Oct 3 13:46:25 2026 +0530

    CAMEL-25275: camel-consul - key/value and event consumers must keep 
watching after a failed query (#27304)
    
    The `consul:kv` and `consul:event` consumers watch with a chain of blocking 
queries: the answer to a query starts the next one. Their failure callbacks 
only reported the exception, so after one failed query (the Consul agent 
restarts, a network error, a 500 while the cluster has no leader) no further 
query was started and the consumer stopped watching for good, while the route 
stayed started.
    This change: after a failure the key/value consumer queries the key again 
after `blockSeconds` (at least one second, so that an agent that is down is not 
queried in a loop), on a single-thread scheduler of its own (started and shut 
down with the consumer, as `ConsulEventConsumer` already has), and the event 
consumer schedules its next query as after an answer (its `watch()` already 
waits `blockSeconds`, CAMEL-12418). Both only while the consumer runs. 
`mockito-core` is added as a test [...]
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 components/camel-consul/pom.xml                    |   6 +
 .../consul/endpoint/ConsulEventConsumer.java       |  11 +-
 .../consul/endpoint/ConsulKeyValueConsumer.java    |  27 +++++
 .../component/consul/ConsulEventWatchTest.java     | 132 +++++++++++++++++++++
 .../component/consul/ConsulKeyValueWatchTest.java  | 118 ++++++++++++++++++
 5 files changed, 293 insertions(+), 1 deletion(-)

diff --git a/components/camel-consul/pom.xml b/components/camel-consul/pom.xml
index 7da83e7c0082..ebc92f662f86 100644
--- a/components/camel-consul/pom.xml
+++ b/components/camel-consul/pom.xml
@@ -114,6 +114,12 @@
             <artifactId>assertj-core</artifactId>
             <scope>test</scope>
         </dependency>
+        <dependency>
+            <groupId>org.mockito</groupId>
+            <artifactId>mockito-core</artifactId>
+            <version>${mockito-version}</version>
+            <scope>test</scope>
+        </dependency>
         <dependency>
             <groupId>org.testcontainers</groupId>
             <artifactId>testcontainers-junit-jupiter</artifactId>
diff --git 
a/components/camel-consul/src/main/java/org/apache/camel/component/consul/endpoint/ConsulEventConsumer.java
 
b/components/camel-consul/src/main/java/org/apache/camel/component/consul/endpoint/ConsulEventConsumer.java
index 45e8c1b2cc0e..7f6a857d013d 100644
--- 
a/components/camel-consul/src/main/java/org/apache/camel/component/consul/endpoint/ConsulEventConsumer.java
+++ 
b/components/camel-consul/src/main/java/org/apache/camel/component/consul/endpoint/ConsulEventConsumer.java
@@ -78,9 +78,13 @@ public final class ConsulEventConsumer extends 
AbstractConsulConsumer<EventClien
 
         @Override
         public void watch(final EventClient client) {
+            query(client, configuration.getBlockSeconds());
+        }
+
+        private void query(final EventClient client, long delaySeconds) {
             Runnable runnable = () -> client.listEvents(key,
                     QueryOptions.blockSeconds(configuration.getBlockSeconds(), 
index.get()).build(), EventWatcher.this);
-            scheduledExecutorService.schedule(runnable, 
configuration.getBlockSeconds(), TimeUnit.SECONDS);
+            scheduledExecutorService.schedule(runnable, delaySeconds, 
TimeUnit.SECONDS);
         }
 
         @Override
@@ -98,6 +102,11 @@ public final class ConsulEventConsumer extends 
AbstractConsulConsumer<EventClien
         @Override
         public void onFailure(Throwable throwable) {
             onError(throwable);
+            // only an answer starts the next query: query again, or the 
events are not watched anymore. Wait at
+            // least one second, so that a Consul agent that is down is not 
queried in a loop
+            if (isRunAllowed()) {
+                query(client(), Math.max(1, configuration.getBlockSeconds()));
+            }
         }
 
         private void onEvent(Event event) {
diff --git 
a/components/camel-consul/src/main/java/org/apache/camel/component/consul/endpoint/ConsulKeyValueConsumer.java
 
b/components/camel-consul/src/main/java/org/apache/camel/component/consul/endpoint/ConsulKeyValueConsumer.java
index bcff5c9f52ef..8c9e141c28c5 100644
--- 
a/components/camel-consul/src/main/java/org/apache/camel/component/consul/endpoint/ConsulKeyValueConsumer.java
+++ 
b/components/camel-consul/src/main/java/org/apache/camel/component/consul/endpoint/ConsulKeyValueConsumer.java
@@ -18,6 +18,8 @@ package org.apache.camel.component.consul.endpoint;
 
 import java.util.List;
 import java.util.Optional;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
 
 import org.apache.camel.Exchange;
 import org.apache.camel.Message;
@@ -34,10 +36,29 @@ import org.kiwiproject.consul.option.QueryOptions;
 
 public final class ConsulKeyValueConsumer extends 
AbstractConsulConsumer<KeyValueClient> {
 
+    private ScheduledExecutorService scheduledExecutorService;
+
     public ConsulKeyValueConsumer(ConsulEndpoint endpoint, ConsulConfiguration 
configuration, Processor processor) {
         super(endpoint, configuration, processor, Consul::keyValueClient);
     }
 
+    @Override
+    protected void doStart() throws Exception {
+        // to watch the key again after a failed query
+        this.scheduledExecutorService = 
getEndpoint().getCamelContext().getExecutorServiceManager()
+                .newSingleThreadScheduledExecutor(this, 
"ConsulKeyValueConsumer");
+        super.doStart();
+    }
+
+    @Override
+    protected void doStop() throws Exception {
+        if (this.scheduledExecutorService != null) {
+            
getEndpoint().getCamelContext().getExecutorServiceManager().shutdownNow(scheduledExecutorService);
+            this.scheduledExecutorService = null;
+        }
+        super.doStop();
+    }
+
     @Override
     protected Runnable createWatcher(KeyValueClient client) throws Exception {
         return configuration.isRecursive() ? new RecursivePathWatcher(client) 
: new PathWatcher(client);
@@ -68,6 +89,12 @@ public final class ConsulKeyValueConsumer extends 
AbstractConsulConsumer<KeyValu
         @Override
         public void onFailure(Throwable throwable) {
             onError(throwable);
+            // only an answer starts the next query: query again, or the key 
is not watched anymore. Wait at
+            // least one second, so that a Consul agent that is down is not 
queried in a loop
+            ScheduledExecutorService executor = scheduledExecutorService;
+            if (isRunAllowed() && executor != null) {
+                executor.schedule(this, Math.max(1, 
configuration.getBlockSeconds()), TimeUnit.SECONDS);
+            }
         }
 
         protected void onValue(Value value) {
diff --git 
a/components/camel-consul/src/test/java/org/apache/camel/component/consul/ConsulEventWatchTest.java
 
b/components/camel-consul/src/test/java/org/apache/camel/component/consul/ConsulEventWatchTest.java
new file mode 100644
index 000000000000..7cb1ae3acd7e
--- /dev/null
+++ 
b/components/camel-consul/src/test/java/org/apache/camel/component/consul/ConsulEventWatchTest.java
@@ -0,0 +1,132 @@
+/*
+ * 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;
+
+import java.math.BigInteger;
+import java.util.List;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.TimeUnit;
+
+import org.apache.camel.BindToRegistry;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.junit.jupiter.api.Test;
+import org.kiwiproject.consul.Consul;
+import org.kiwiproject.consul.ConsulException;
+import org.kiwiproject.consul.EventClient;
+import org.kiwiproject.consul.async.EventResponseCallback;
+import org.kiwiproject.consul.model.EventResponse;
+import org.kiwiproject.consul.model.ImmutableEventResponse;
+import org.kiwiproject.consul.model.event.ImmutableEvent;
+import org.kiwiproject.consul.option.QueryOptions;
+
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.timeout;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * The event consumer watches the events with a chain of blocking queries: 
each answer schedules the next query.
+ */
+public class ConsulEventWatchTest extends CamelTestSupport {
+
+    private static final String EVENT = "camel-watch";
+    private static final String NO_BLOCK_EVENT = "camel-watch-no-block";
+
+    private final EventClient eventClient = mock(EventClient.class);
+    private final List<EventResponseCallback> queries = new 
CopyOnWriteArrayList<>();
+    private final List<Long> noBlockQueryTimes = new CopyOnWriteArrayList<>();
+
+    @BindToRegistry("consul")
+    public Consul consul() {
+        Consul consul = mock(Consul.class);
+        when(consul.eventClient()).thenReturn(eventClient);
+        // the consumer with blockSeconds=0 queries as soon as it starts: the 
first query fails, the next is pending
+        doAnswer(inv -> {
+            noBlockQueryTimes.add(System.nanoTime());
+            if (noBlockQueryTimes.size() == 1) {
+                inv.getArgument(2, EventResponseCallback.class).onFailure(new 
ConsulException("Consul is not available"));
+            }
+            return null;
+        }).when(eventClient).listEvents(eq(NO_BLOCK_EVENT), 
any(QueryOptions.class), any(EventResponseCallback.class));
+        return consul;
+    }
+
+    @Test
+    public void testWatchGoesOnAfterAFailedQuery() throws Exception {
+        // the first query fails (for example while the Consul agent 
restarts), the next ones are pending
+        doAnswer(inv -> {
+            EventResponseCallback callback = inv.getArgument(2);
+            queries.add(callback);
+            if (queries.size() == 1) {
+                callback.onFailure(new ConsulException("Consul is not 
available"));
+            }
+            return null;
+        }).when(eventClient).listEvents(eq(EVENT), any(QueryOptions.class), 
any(EventResponseCallback.class));
+
+        // the consumer must query the events again
+        verify(eventClient, timeout(5000).times(2)).listEvents(eq(EVENT), 
any(QueryOptions.class),
+                any(EventResponseCallback.class));
+
+        MockEndpoint mock = getMockEndpoint("mock:event");
+        mock.expectedBodiesReceived("bar");
+        queries.get(1).onComplete(response("bar"));
+        mock.assertIsSatisfied();
+    }
+
+    @Test
+    public void testFailedQueryIsRetriedAfterOneSecondWithoutBlockSeconds() {
+        verify(eventClient, 
timeout(5000).times(2)).listEvents(eq(NO_BLOCK_EVENT), any(QueryOptions.class),
+                any(EventResponseCallback.class));
+
+        // blockSeconds=0 must not query a Consul agent that is down in a loop
+        long delay = TimeUnit.NANOSECONDS.toMillis(noBlockQueryTimes.get(1) - 
noBlockQueryTimes.get(0));
+        assertTrue(delay >= 900, "The failed query was retried after " + delay 
+ " ms");
+    }
+
+    private static EventResponse response(String payload) {
+        return ImmutableEventResponse.builder()
+                .addEvents(ImmutableEvent.builder()
+                        .id("2a7d7cbc-d6ad-4d5e-8e3c-4b36a54f6bb7")
+                        .name(EVENT)
+                        .payload(payload)
+                        .version(1)
+                        .lTime(1L)
+                        .build())
+                .index(BigInteger.ONE)
+                .build();
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                
fromF("consul:event?key=%s&blockSeconds=1&consulClient=#consul", EVENT)
+                        .to("mock:event");
+
+                
fromF("consul:event?key=%s&blockSeconds=0&consulClient=#consul", NO_BLOCK_EVENT)
+                        .to("mock:noBlockEvent");
+            }
+        };
+    }
+}
diff --git 
a/components/camel-consul/src/test/java/org/apache/camel/component/consul/ConsulKeyValueWatchTest.java
 
b/components/camel-consul/src/test/java/org/apache/camel/component/consul/ConsulKeyValueWatchTest.java
new file mode 100644
index 000000000000..52b972344323
--- /dev/null
+++ 
b/components/camel-consul/src/test/java/org/apache/camel/component/consul/ConsulKeyValueWatchTest.java
@@ -0,0 +1,118 @@
+/*
+ * 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;
+
+import java.math.BigInteger;
+import java.util.Base64;
+import java.util.List;
+import java.util.Optional;
+import java.util.concurrent.CopyOnWriteArrayList;
+
+import org.apache.camel.BindToRegistry;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.apache.camel.test.junit6.CamelTestSupport;
+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.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.option.QueryOptions;
+
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.timeout;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * The key/value consumer watches a key with a chain of blocking queries: each 
answer starts the next query.
+ */
+public class ConsulKeyValueWatchTest extends CamelTestSupport {
+
+    private static final String KEY = "camel/watch";
+
+    private final KeyValueClient keyValueClient = mock(KeyValueClient.class);
+    private final List<ConsulResponseCallback<Optional<Value>>> queries = new 
CopyOnWriteArrayList<>();
+
+    @BindToRegistry("consul")
+    public Consul consul() {
+        Consul consul = mock(Consul.class);
+        when(consul.keyValueClient()).thenReturn(keyValueClient);
+        return consul;
+    }
+
+    @Override
+    public boolean isUseRouteBuilder() {
+        return false;
+    }
+
+    @Test
+    public void testWatchGoesOnAfterAFailedQuery() throws Exception {
+        // the first query fails (for example while the Consul agent 
restarts), the next ones are pending
+        doAnswer(inv -> {
+            ConsulResponseCallback<Optional<Value>> callback = 
inv.getArgument(2);
+            queries.add(callback);
+            if (queries.size() == 1) {
+                callback.onFailure(new ConsulException("Consul is not 
available"));
+            }
+            return null;
+        }).when(keyValueClient).getValue(eq(KEY), any(QueryOptions.class), 
any());
+
+        addRoute();
+        context.start();
+
+        // the consumer must query the key again
+        verify(keyValueClient, timeout(5000).times(2)).getValue(eq(KEY), 
any(QueryOptions.class), any());
+
+        MockEndpoint mock = getMockEndpoint("mock:kv");
+        mock.expectedBodiesReceived("bar");
+        lastQuery().onComplete(response("bar", 2));
+        mock.assertIsSatisfied();
+    }
+
+    private void addRoute() throws Exception {
+        context.addRoutes(new RouteBuilder() {
+            @Override
+            public void configure() {
+                
fromF("consul:kv?key=%s&valueAsString=true&blockSeconds=1&consulClient=#consul",
 KEY).routeId("kv")
+                        .to("mock:kv");
+            }
+        });
+    }
+
+    private ConsulResponseCallback<Optional<Value>> lastQuery() {
+        return queries.get(queries.size() - 1);
+    }
+
+    private static ConsulResponse<Optional<Value>> response(String value, long 
index) {
+        Value v = ImmutableValue.builder()
+                .key(KEY)
+                .value(Base64.getEncoder().encodeToString(value.getBytes()))
+                .createIndex(1)
+                .modifyIndex(index)
+                .lockIndex(0)
+                .flags(0)
+                .build();
+        return new ConsulResponse<>(Optional.of(v), 0, true, 
BigInteger.valueOf(index), (String) null, (String) null);
+    }
+}

Reply via email to