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 f57419ff40b8 CAMEL-25293: camel-hivemq - the consumer must process the
messages it received before it stops (#27324)
f57419ff40b8 is described below
commit f57419ff40b85ca901f7c4c0504474d4aaa96194
Author: allthingssecurity <[email protected]>
AuthorDate: Sun Oct 4 12:33:21 2026 +0530
CAMEL-25293: camel-hivemq - the consumer must process the messages it
received before it stops (#27324)
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../camel/component/hivemq/HiveMQConsumer.java | 3 +-
.../component/hivemq/HiveMQConsumerStopTest.java | 183 +++++++++++++++++++++
2 files changed, 185 insertions(+), 1 deletion(-)
diff --git
a/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQConsumer.java
b/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQConsumer.java
index 62ce0d50b99e..051e288ece51 100644
---
a/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQConsumer.java
+++
b/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQConsumer.java
@@ -67,7 +67,8 @@ public class HiveMQConsumer extends DefaultConsumer {
}
client = null;
if (executor != null) {
-
endpoint.getCamelContext().getExecutorServiceManager().shutdownNow(executor);
+ // the client acknowledged the messages that are queued or being
processed, so let them complete
+
endpoint.getCamelContext().getExecutorServiceManager().shutdownGraceful(executor);
executor = null;
}
super.doStop();
diff --git
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQConsumerStopTest.java
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQConsumerStopTest.java
new file mode 100644
index 000000000000..7aa0eec537ed
--- /dev/null
+++
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQConsumerStopTest.java
@@ -0,0 +1,183 @@
+/*
+ * 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.hivemq;
+
+import java.nio.charset.StandardCharsets;
+import java.util.List;
+import java.util.Optional;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.ThreadFactory;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+import java.util.function.Consumer;
+
+import com.hivemq.client.mqtt.datatypes.MqttQos;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.impl.DefaultCamelContext;
+import org.apache.camel.spi.ThreadPoolProfile;
+import org.apache.camel.support.DefaultThreadPoolFactory;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/**
+ * The client acknowledges a message when the consumer's callback returns, and
the consumer then processes it on its own
+ * thread pool: when the route stops, the messages the consumer already
received must still be processed.
+ */
+class HiveMQConsumerStopTest {
+
+ // released when the consumer shuts its thread pool down, so the first
message is being processed at that time
+ private final CountDownLatch poolShutdown = new CountDownLatch(1);
+ private final CountDownLatch firstStarted = new CountDownLatch(1);
+ private final List<String> processed = new CopyOnWriteArrayList<>();
+ private final FakeClient client = new FakeClient();
+ private DefaultCamelContext camelContext;
+
+ @BeforeEach
+ void setUp() throws Exception {
+ camelContext = new DefaultCamelContext();
+ // a single consumer thread, so the second and third message wait in
the queue of the pool
+ camelContext.getExecutorServiceManager().setThreadPoolFactory(new
SingleThreadPoolFactory());
+
+ HiveMQComponent component = new HiveMQComponent();
+ component.setCamelContext(camelContext);
+ HiveMQEndpoint endpoint = new HiveMQEndpoint("hivemq:test", component,
new HiveMQConfiguration(), "test") {
+ @Override
+ HiveMQClientAdapter createClient() {
+ return client;
+ }
+ };
+ endpoint.setCamelContext(camelContext);
+
+ camelContext.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ from(endpoint).routeId("hivemq")
+ .process(exchange -> {
+ String body =
exchange.getIn().getBody(String.class);
+ if ("1".equals(body)) {
+ firstStarted.countDown();
+ assertThat(poolShutdown.await(20,
TimeUnit.SECONDS)).isTrue();
+ }
+ processed.add(body);
+ });
+ }
+ });
+ camelContext.start();
+ }
+
+ @AfterEach
+ void tearDown() {
+ camelContext.stop();
+ }
+
+ @Test
+ void receivedMessagesAreProcessedWhenTheRouteStops() throws Exception {
+ client.deliver("1");
+ client.deliver("2");
+ client.deliver("3");
+ assertThat(firstStarted.await(10, TimeUnit.SECONDS)).isTrue();
+
+ camelContext.getRouteController().stopRoute("hivemq");
+
+ assertThat(client.unsubscribed).isTrue();
+ assertThat(processed).containsExactly("1", "2", "3");
+ }
+
+ private final class SingleThreadPoolFactory extends
DefaultThreadPoolFactory {
+
+ @Override
+ public ExecutorService newThreadPool(ThreadPoolProfile profile,
ThreadFactory factory) {
+ if (!Boolean.TRUE.equals(profile.isDefaultProfile())) {
+ return super.newThreadPool(profile, factory);
+ }
+ return new ThreadPoolExecutor(1, 1, 0, TimeUnit.SECONDS, new
LinkedBlockingQueue<>(), factory) {
+ @Override
+ public void shutdown() {
+ super.shutdown();
+ poolShutdown.countDown();
+ }
+
+ @Override
+ public List<Runnable> shutdownNow() {
+ List<Runnable> notStarted = super.shutdownNow();
+ poolShutdown.countDown();
+ return notStarted;
+ }
+ };
+ }
+ }
+
+ private static final class FakeClient implements HiveMQClientAdapter {
+
+ private volatile Consumer<HiveMQMessage> callback;
+ private volatile boolean unsubscribed;
+
+ void deliver(String payload) {
+ callback.accept(new HiveMQMessage(
+ "test", payload.getBytes(StandardCharsets.UTF_8),
MqttQos.AT_LEAST_ONCE, false));
+ }
+
+ @Override
+ public CompletableFuture<?> connect(boolean cleanStart) {
+ return CompletableFuture.completedFuture(null);
+ }
+
+ @Override
+ public void stop() {
+ // noop
+ }
+
+ @Override
+ public boolean isConnected() {
+ return true;
+ }
+
+ @Override
+ public boolean isConnectedOrReconnecting() {
+ return true;
+ }
+
+ @Override
+ public CompletableFuture<?> subscribe(String topicFilter, MqttQos qos,
Consumer<HiveMQMessage> callback) {
+ this.callback = callback;
+ return CompletableFuture.completedFuture(null);
+ }
+
+ @Override
+ public CompletableFuture<?> unsubscribe(String topicFilter) {
+ unsubscribed = true;
+ return CompletableFuture.completedFuture(null);
+ }
+
+ @Override
+ public CompletableFuture<?> publish(String topic, byte[] payload,
MqttQos qos, boolean retained) {
+ return CompletableFuture.completedFuture(null);
+ }
+
+ @Override
+ public <T> Optional<T> getClient(Class<T> clazz) {
+ return Optional.empty();
+ }
+ }
+}