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