This is an automated email from the ASF dual-hosted git repository.

apupier pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git

commit 68979efc222c6a73d5cf5f2cd17bd3ac571c201d
Author: Thomas Strauss <[email protected]>
AuthorDate: Wed Sep 30 14:29:46 2026 +0200

    CAMEL-25030: camel-hivemq: expand MQTT 3.1.1 test coverage
    
    Addresses review feedback (@Croway, @davsclaus) that this PR's test
    coverage was thinner than the split-component alternative, which
    duplicated the entire test suite per protocol version.
    
    Parameterizes HiveMQQosAndRetainIT, HiveMQEmptyPayloadIT,
    HiveMQOverrideTopicIT and HiveMQSendDynamicIT over MqttVersion, via an
    AbstractHiveMQ*IT base per test (Camel's existing convention for
    IT variants - e.g. AbstractTransactionIT, AbstractRAGIT - Failsafe
    skips abstract classes so only the concrete MQTT_5_0/MQTT_3_1_1
    subclasses run). CamelTestSupport's createRouteBuilder() runs before a
    JUnit5 @ParameterizedTest method would receive its argument, so
    per-class subclassing is used instead of @ParameterizedTest for these;
    each concrete subclass is a one-line override of the abstract
    mqttVersion() method.
    
    HiveMQEndpointConnectTest and HiveMQClientAccessTest are plain unit
    tests (no route-builder lifecycle constraint), so these are
    parameterized directly with @ParameterizedTest + 
@EnumSource(MqttVersion.class).
    This specifically closes the gap davsclaus flagged: the reconnect-cancel
    listener logic was previously only exercised via Mqtt5ClientAdapter,
    leaving Mqtt3ClientAdapter's identical logic (a mistake in which could
    hang a 3.1.1 route on start or keep reconnecting after stop) completely
    untested.
    
    Verified: all unit tests pass (HiveMQEndpointConnectTest and
    HiveMQClientAccessTest now run twice, once per protocol version), and
    all 10 ITs (5 original + 5 new MQTT 3.1.1 variants) pass against a live
    HiveMQ CE broker.
    
    AI-assisted contribution: authored together with Claude Sonnet 5 
(Anthropic).
    
    Co-Authored-By: Claude Sonnet 5 <[email protected]>
---
 ...adIT.java => AbstractHiveMQEmptyPayloadIT.java} |  9 ++-
 ...cIT.java => AbstractHiveMQOverrideTopicIT.java} |  9 ++-
 ...inIT.java => AbstractHiveMQQosAndRetainIT.java} | 10 +--
 ...micIT.java => AbstractHiveMQSendDynamicIT.java} | 19 ++++--
 .../component/hivemq/HiveMQClientAccessTest.java   | 40 ++++++++----
 .../component/hivemq/HiveMQEmptyPayloadIT.java     | 54 ++--------------
 .../hivemq/HiveMQEndpointConnectTest.java          | 16 +++--
 .../hivemq/HiveMQMqtt311EmptyPayloadIT.java        | 27 ++++++++
 .../hivemq/HiveMQMqtt311OverrideTopicIT.java       | 27 ++++++++
 .../hivemq/HiveMQMqtt311QosAndRetainIT.java        | 27 ++++++++
 .../hivemq/HiveMQMqtt311SendDynamicIT.java         | 27 ++++++++
 .../component/hivemq/HiveMQOverrideTopicIT.java    | 53 ++--------------
 .../component/hivemq/HiveMQQosAndRetainIT.java     | 74 ++--------------------
 .../component/hivemq/HiveMQSendDynamicIT.java      | 51 ++-------------
 14 files changed, 193 insertions(+), 250 deletions(-)

diff --git 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQEmptyPayloadIT.java
 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/AbstractHiveMQEmptyPayloadIT.java
similarity index 87%
copy from 
components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQEmptyPayloadIT.java
copy to 
components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/AbstractHiveMQEmptyPayloadIT.java
index fd90d24ee709..6c3bdfa681c7 100644
--- 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQEmptyPayloadIT.java
+++ 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/AbstractHiveMQEmptyPayloadIT.java
@@ -16,6 +16,7 @@
  */
 package org.apache.camel.component.hivemq;
 
+import com.hivemq.client.mqtt.MqttVersion;
 import org.apache.camel.EndpointInject;
 import org.apache.camel.Produce;
 import org.apache.camel.ProducerTemplate;
@@ -30,7 +31,7 @@ import org.junit.jupiter.api.extension.RegisterExtension;
 
 import static org.assertj.core.api.Assertions.assertThat;
 
-public class HiveMQEmptyPayloadIT extends CamelTestSupport {
+public abstract class AbstractHiveMQEmptyPayloadIT extends CamelTestSupport {
 
     @RegisterExtension
     public static HiveMQService HIVEMQ_SERVICE = 
HiveMQServiceFactory.createService();
@@ -41,6 +42,8 @@ public class HiveMQEmptyPayloadIT extends CamelTestSupport {
     @Produce
     private ProducerTemplate template;
 
+    protected abstract MqttVersion mqttVersion();
+
     @Test
     @DisplayName("Null Camel body is published as an MQTT empty payload")
     public void testNullBodyPublishesEmptyPayload() throws Exception {
@@ -63,9 +66,9 @@ public class HiveMQEmptyPayloadIT extends CamelTestSupport {
                 int port = HIVEMQ_SERVICE.getMqttPort();
 
                 from("direct:startEmpty")
-                        .toF("hivemq:test/empty?host=%s&port=%d", host, port);
+                        
.toF("hivemq:test/empty?host=%s&port=%d&mqttVersion=%s", host, port, 
mqttVersion());
 
-                fromF("hivemq:test/empty?host=%s&port=%d", host, port)
+                fromF("hivemq:test/empty?host=%s&port=%d&mqttVersion=%s", 
host, port, mqttVersion())
                         .to("mock:emptyResult");
             }
         };
diff --git 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQOverrideTopicIT.java
 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/AbstractHiveMQOverrideTopicIT.java
similarity index 89%
copy from 
components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQOverrideTopicIT.java
copy to 
components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/AbstractHiveMQOverrideTopicIT.java
index e66505bf1e9d..534285963b35 100644
--- 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQOverrideTopicIT.java
+++ 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/AbstractHiveMQOverrideTopicIT.java
@@ -16,6 +16,7 @@
  */
 package org.apache.camel.component.hivemq;
 
+import com.hivemq.client.mqtt.MqttVersion;
 import org.apache.camel.EndpointInject;
 import org.apache.camel.Produce;
 import org.apache.camel.ProducerTemplate;
@@ -28,7 +29,7 @@ import org.junit.jupiter.api.DisplayName;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.extension.RegisterExtension;
 
-public class HiveMQOverrideTopicIT extends CamelTestSupport {
+public abstract class AbstractHiveMQOverrideTopicIT extends CamelTestSupport {
 
     @RegisterExtension
     public static HiveMQService HIVEMQ_SERVICE = 
HiveMQServiceFactory.createService();
@@ -39,6 +40,8 @@ public class HiveMQOverrideTopicIT extends CamelTestSupport {
     @Produce
     private ProducerTemplate template;
 
+    protected abstract MqttVersion mqttVersion();
+
     @Test
     @DisplayName("Publish to default endpoint topic but route to override 
topic via header")
     public void testTopicOverride() throws Exception {
@@ -62,9 +65,9 @@ public class HiveMQOverrideTopicIT extends CamelTestSupport {
 
                 // Endpoint points to "orders/default", but producer will 
override to "orders/processed"
                 from("direct:startOverride")
-                        .toF("hivemq:orders/default?host=%s&port=%d", host, 
port);
+                        
.toF("hivemq:orders/default?host=%s&port=%d&mqttVersion=%s", host, port, 
mqttVersion());
 
-                fromF("hivemq:orders/processed?host=%s&port=%d", host, port)
+                
fromF("hivemq:orders/processed?host=%s&port=%d&mqttVersion=%s", host, port, 
mqttVersion())
                         .to("mock:overrideResult");
             }
         };
diff --git 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQQosAndRetainIT.java
 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/AbstractHiveMQQosAndRetainIT.java
similarity index 92%
copy from 
components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQQosAndRetainIT.java
copy to 
components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/AbstractHiveMQQosAndRetainIT.java
index 717332bda325..9d4f1188598e 100644
--- 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQQosAndRetainIT.java
+++ 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/AbstractHiveMQQosAndRetainIT.java
@@ -14,12 +14,12 @@
  * See the License for the specific language governing permissions and
  * limitations under the License.
  */
-
 package org.apache.camel.component.hivemq;
 
 import java.util.HashMap;
 import java.util.Map;
 
+import com.hivemq.client.mqtt.MqttVersion;
 import com.hivemq.client.mqtt.datatypes.MqttQos;
 import org.apache.camel.EndpointInject;
 import org.apache.camel.Produce;
@@ -33,7 +33,7 @@ import org.junit.jupiter.api.DisplayName;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.extension.RegisterExtension;
 
-public class HiveMQQosAndRetainIT extends CamelTestSupport {
+public abstract class AbstractHiveMQQosAndRetainIT extends CamelTestSupport {
 
     @RegisterExtension
     public static HiveMQService HIVEMQ_SERVICE = 
HiveMQServiceFactory.createService();
@@ -44,6 +44,8 @@ public class HiveMQQosAndRetainIT extends CamelTestSupport {
     @Produce
     private ProducerTemplate template;
 
+    protected abstract MqttVersion mqttVersion();
+
     @Test
     @DisplayName("Publish with QoS 2 (EXACTLY_ONCE) and retain flag set to 
true")
     public void testQosAndRetainHeaders() throws Exception {
@@ -67,7 +69,7 @@ public class HiveMQQosAndRetainIT extends CamelTestSupport {
         context.addRoutes(new RouteBuilder() {
             @Override
             public void configure() {
-                fromF("hivemq:system/alerts?host=%s&port=%d&qos=EXACTLY_ONCE", 
host, port)
+                
fromF("hivemq:system/alerts?host=%s&port=%d&qos=EXACTLY_ONCE&mqttVersion=%s", 
host, port, mqttVersion())
                         .to("mock:qosResult");
             }
         });
@@ -86,7 +88,7 @@ public class HiveMQQosAndRetainIT extends CamelTestSupport {
                 int port = HIVEMQ_SERVICE.getMqttPort();
 
                 from("direct:startQos")
-                        .toF("hivemq:system/alerts?host=%s&port=%d", host, 
port);
+                        
.toF("hivemq:system/alerts?host=%s&port=%d&mqttVersion=%s", host, port, 
mqttVersion());
             }
         };
     }
diff --git 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQSendDynamicIT.java
 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/AbstractHiveMQSendDynamicIT.java
similarity index 78%
copy from 
components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQSendDynamicIT.java
copy to 
components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/AbstractHiveMQSendDynamicIT.java
index 4cb5e364802c..d2a315058628 100644
--- 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQSendDynamicIT.java
+++ 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/AbstractHiveMQSendDynamicIT.java
@@ -16,6 +16,7 @@
  */
 package org.apache.camel.component.hivemq;
 
+import com.hivemq.client.mqtt.MqttVersion;
 import org.apache.camel.builder.RouteBuilder;
 import org.apache.camel.test.infra.hivemq.services.HiveMQService;
 import org.apache.camel.test.infra.hivemq.services.HiveMQServiceFactory;
@@ -26,17 +27,20 @@ import org.junit.jupiter.api.extension.RegisterExtension;
 
 import static org.assertj.core.api.Assertions.assertThat;
 
-public class HiveMQSendDynamicIT extends CamelTestSupport {
+public abstract class AbstractHiveMQSendDynamicIT extends CamelTestSupport {
 
     @RegisterExtension
     public static HiveMQService HIVEMQ_SERVICE = 
HiveMQServiceFactory.createService();
 
+    protected abstract MqttVersion mqttVersion();
+
     @Test
     @DisplayName("toD with multiple dynamic topics reuses a single hivemq 
endpoint")
     @SuppressWarnings("deprecation")
     public void testToDReusesSingleEndpoint() {
-        template.sendBodyAndHeader("direct:start", "Hello bar", "where", 
"HiveMQSendDynamicIT-bar");
-        template.sendBodyAndHeader("direct:start", "Hello beer", "where", 
"HiveMQSendDynamicIT-beer");
+        String topicPrefix = getClass().getSimpleName();
+        template.sendBodyAndHeader("direct:start", "Hello bar", "where", 
topicPrefix + "-bar");
+        template.sendBodyAndHeader("direct:start", "Hello beer", "where", 
topicPrefix + "-beer");
 
         long count = context.getEndpoints().stream()
                 .filter(e -> e.getEndpointUri().startsWith("hivemq:"))
@@ -46,10 +50,12 @@ public class HiveMQSendDynamicIT extends CamelTestSupport {
         String host = HIVEMQ_SERVICE.getMqttHost();
         int port = HIVEMQ_SERVICE.getMqttPort();
         String out = consumer.receiveBody(
-                
String.format("hivemq:HiveMQSendDynamicIT-bar?host=%s&port=%d", host, port), 
5000, String.class);
+                String.format("hivemq:%s-bar?host=%s&port=%d&mqttVersion=%s", 
topicPrefix, host, port, mqttVersion()),
+                5000, String.class);
         assertThat(out).isEqualTo("Hello bar");
         out = consumer.receiveBody(
-                
String.format("hivemq:HiveMQSendDynamicIT-beer?host=%s&port=%d", host, port), 
5000, String.class);
+                String.format("hivemq:%s-beer?host=%s&port=%d&mqttVersion=%s", 
topicPrefix, host, port, mqttVersion()),
+                5000, String.class);
         assertThat(out).isEqualTo("Hello beer");
     }
 
@@ -63,7 +69,8 @@ public class HiveMQSendDynamicIT extends CamelTestSupport {
                 int port = HIVEMQ_SERVICE.getMqttPort();
 
                 from("direct:start")
-                        
.toD(String.format("hivemq:${header.where}?host=%s&port=%d&retained=true", 
host, port));
+                        
.toD(String.format("hivemq:${header.where}?host=%s&port=%d&retained=true&mqttVersion=%s",
 host, port,
+                                mqttVersion()));
             }
         };
     }
diff --git 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQClientAccessTest.java
 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQClientAccessTest.java
index 97ff93e2eab0..f26e6dc3aaa3 100644
--- 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQClientAccessTest.java
+++ 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQClientAccessTest.java
@@ -16,6 +16,7 @@
  */
 package org.apache.camel.component.hivemq;
 
+import com.hivemq.client.mqtt.MqttVersion;
 import com.hivemq.client.mqtt.mqtt3.Mqtt3AsyncClient;
 import com.hivemq.client.mqtt.mqtt5.Mqtt5AsyncClient;
 import org.apache.camel.impl.DefaultCamelContext;
@@ -23,6 +24,8 @@ import org.junit.jupiter.api.AfterEach;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.DisplayName;
 import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.EnumSource;
 
 import static org.assertj.core.api.Assertions.assertThat;
 import static org.assertj.core.api.Assertions.assertThatThrownBy;
@@ -61,31 +64,42 @@ class HiveMQClientAccessTest {
         assertThat(consumer.getClient(Mqtt5AsyncClient.class)).isEmpty();
     }
 
-    @Test
-    @DisplayName("Producer exposes the underlying Mqtt5AsyncClient once 
created, not the Mqtt3 type")
-    void producerExposesMatchingClientTypeAfterConnectAttempt() {
-        HiveMQProducer producer = new 
HiveMQProducer(newEndpoint(unreachableConfiguration()));
+    @ParameterizedTest(name = "{0}")
+    @EnumSource(MqttVersion.class)
+    @DisplayName("Producer exposes the client matching its configured protocol 
version, not the other one")
+    void producerExposesMatchingClientTypeAfterConnectAttempt(MqttVersion 
mqttVersion) {
+        HiveMQProducer producer = new 
HiveMQProducer(newEndpoint(unreachableConfiguration(mqttVersion)));
 
         assertThatThrownBy(producer::doStart);
 
-        assertThat(producer.getClient(Mqtt5AsyncClient.class)).isPresent();
-        assertThat(producer.getClient(Mqtt3AsyncClient.class)).isEmpty();
+        
assertThat(producer.getClient(matchingClientType(mqttVersion))).isPresent();
+        assertThat(producer.getClient(otherClientType(mqttVersion))).isEmpty();
     }
 
-    @Test
-    @DisplayName("Consumer exposes the underlying Mqtt5AsyncClient once 
created, not the Mqtt3 type")
-    void consumerExposesMatchingClientTypeAfterConnectAttempt() {
-        HiveMQConsumer consumer = new 
HiveMQConsumer(newEndpoint(unreachableConfiguration()), exchange -> {
+    @ParameterizedTest(name = "{0}")
+    @EnumSource(MqttVersion.class)
+    @DisplayName("Consumer exposes the client matching its configured protocol 
version, not the other one")
+    void consumerExposesMatchingClientTypeAfterConnectAttempt(MqttVersion 
mqttVersion) {
+        HiveMQConsumer consumer = new 
HiveMQConsumer(newEndpoint(unreachableConfiguration(mqttVersion)), exchange -> {
         });
 
         assertThatThrownBy(consumer::doStart);
 
-        assertThat(consumer.getClient(Mqtt5AsyncClient.class)).isPresent();
-        assertThat(consumer.getClient(Mqtt3AsyncClient.class)).isEmpty();
+        
assertThat(consumer.getClient(matchingClientType(mqttVersion))).isPresent();
+        assertThat(consumer.getClient(otherClientType(mqttVersion))).isEmpty();
+    }
+
+    private static Class<?> matchingClientType(MqttVersion mqttVersion) {
+        return mqttVersion == MqttVersion.MQTT_5_0 ? Mqtt5AsyncClient.class : 
Mqtt3AsyncClient.class;
+    }
+
+    private static Class<?> otherClientType(MqttVersion mqttVersion) {
+        return mqttVersion == MqttVersion.MQTT_5_0 ? Mqtt3AsyncClient.class : 
Mqtt5AsyncClient.class;
     }
 
-    private static HiveMQConfiguration unreachableConfiguration() {
+    private static HiveMQConfiguration unreachableConfiguration(MqttVersion 
mqttVersion) {
         HiveMQConfiguration configuration = new HiveMQConfiguration();
+        configuration.setMqttVersion(mqttVersion);
         configuration.setHost("127.0.0.1");
         configuration.setPort(1);
         return configuration;
diff --git 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQEmptyPayloadIT.java
 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQEmptyPayloadIT.java
index fd90d24ee709..02f9164f3522 100644
--- 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQEmptyPayloadIT.java
+++ 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQEmptyPayloadIT.java
@@ -16,58 +16,12 @@
  */
 package org.apache.camel.component.hivemq;
 
-import org.apache.camel.EndpointInject;
-import org.apache.camel.Produce;
-import org.apache.camel.ProducerTemplate;
-import org.apache.camel.builder.RouteBuilder;
-import org.apache.camel.component.mock.MockEndpoint;
-import org.apache.camel.test.infra.hivemq.services.HiveMQService;
-import org.apache.camel.test.infra.hivemq.services.HiveMQServiceFactory;
-import org.apache.camel.test.junit6.CamelTestSupport;
-import org.junit.jupiter.api.DisplayName;
-import org.junit.jupiter.api.Test;
-import org.junit.jupiter.api.extension.RegisterExtension;
+import com.hivemq.client.mqtt.MqttVersion;
 
-import static org.assertj.core.api.Assertions.assertThat;
-
-public class HiveMQEmptyPayloadIT extends CamelTestSupport {
-
-    @RegisterExtension
-    public static HiveMQService HIVEMQ_SERVICE = 
HiveMQServiceFactory.createService();
-
-    @EndpointInject("mock:emptyResult")
-    private MockEndpoint mockEmptyResult;
-
-    @Produce
-    private ProducerTemplate template;
-
-    @Test
-    @DisplayName("Null Camel body is published as an MQTT empty payload")
-    public void testNullBodyPublishesEmptyPayload() throws Exception {
-        mockEmptyResult.expectedMessageCount(1);
-
-        template.sendBody("direct:startEmpty", null);
-
-        mockEmptyResult.assertIsSatisfied();
-        byte[] payload = 
mockEmptyResult.getExchanges().get(0).getIn().getBody(byte[].class);
-        assertThat(payload).isEmpty();
-    }
+public class HiveMQEmptyPayloadIT extends AbstractHiveMQEmptyPayloadIT {
 
     @Override
-    @SuppressWarnings("deprecation")
-    protected RouteBuilder createRouteBuilder() {
-        return new RouteBuilder() {
-            @Override
-            public void configure() {
-                String host = HIVEMQ_SERVICE.getMqttHost();
-                int port = HIVEMQ_SERVICE.getMqttPort();
-
-                from("direct:startEmpty")
-                        .toF("hivemq:test/empty?host=%s&port=%d", host, port);
-
-                fromF("hivemq:test/empty?host=%s&port=%d", host, port)
-                        .to("mock:emptyResult");
-            }
-        };
+    protected MqttVersion mqttVersion() {
+        return MqttVersion.MQTT_5_0;
     }
 }
diff --git 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQEndpointConnectTest.java
 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQEndpointConnectTest.java
index ba7407b74c1b..2bd3736f7145 100644
--- 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQEndpointConnectTest.java
+++ 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQEndpointConnectTest.java
@@ -18,12 +18,14 @@ package org.apache.camel.component.hivemq;
 
 import java.util.concurrent.TimeUnit;
 
+import com.hivemq.client.mqtt.MqttVersion;
 import org.apache.camel.RuntimeCamelException;
 import org.apache.camel.impl.DefaultCamelContext;
 import org.junit.jupiter.api.AfterEach;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.DisplayName;
-import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.EnumSource;
 
 import static org.assertj.core.api.Assertions.assertThat;
 import static org.assertj.core.api.Assertions.assertThatThrownBy;
@@ -46,10 +48,12 @@ class HiveMQEndpointConnectTest {
         camelContext.stop();
     }
 
-    @Test
+    @ParameterizedTest(name = "{0}")
+    @EnumSource(MqttVersion.class)
     @DisplayName("Unreachable broker fails connect() within a bounded time 
instead of hanging")
-    void connectToUnreachableBrokerFailsBounded() {
+    void connectToUnreachableBrokerFailsBounded(MqttVersion mqttVersion) {
         HiveMQConfiguration configuration = new HiveMQConfiguration();
+        configuration.setMqttVersion(mqttVersion);
         configuration.setHost("127.0.0.1");
         configuration.setPort(1);
 
@@ -65,10 +69,12 @@ class HiveMQEndpointConnectTest {
                 .untilAsserted(() -> 
assertThat(client.isConnectedOrReconnecting()).isFalse());
     }
 
-    @Test
+    @ParameterizedTest(name = "{0}")
+    @EnumSource(MqttVersion.class)
     @DisplayName("stop() cancels automatic reconnect when the client is not 
CONNECTED")
-    void stopClientCancelsReconnect() {
+    void stopClientCancelsReconnect(MqttVersion mqttVersion) {
         HiveMQConfiguration configuration = new HiveMQConfiguration();
+        configuration.setMqttVersion(mqttVersion);
         configuration.setHost("127.0.0.1");
         configuration.setPort(1);
 
diff --git 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQMqtt311EmptyPayloadIT.java
 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQMqtt311EmptyPayloadIT.java
new file mode 100644
index 000000000000..dd6590387443
--- /dev/null
+++ 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQMqtt311EmptyPayloadIT.java
@@ -0,0 +1,27 @@
+/*
+ * 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 com.hivemq.client.mqtt.MqttVersion;
+
+public class HiveMQMqtt311EmptyPayloadIT extends AbstractHiveMQEmptyPayloadIT {
+
+    @Override
+    protected MqttVersion mqttVersion() {
+        return MqttVersion.MQTT_3_1_1;
+    }
+}
diff --git 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQMqtt311OverrideTopicIT.java
 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQMqtt311OverrideTopicIT.java
new file mode 100644
index 000000000000..cb1842268005
--- /dev/null
+++ 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQMqtt311OverrideTopicIT.java
@@ -0,0 +1,27 @@
+/*
+ * 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 com.hivemq.client.mqtt.MqttVersion;
+
+public class HiveMQMqtt311OverrideTopicIT extends 
AbstractHiveMQOverrideTopicIT {
+
+    @Override
+    protected MqttVersion mqttVersion() {
+        return MqttVersion.MQTT_3_1_1;
+    }
+}
diff --git 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQMqtt311QosAndRetainIT.java
 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQMqtt311QosAndRetainIT.java
new file mode 100644
index 000000000000..cc23f8e086b7
--- /dev/null
+++ 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQMqtt311QosAndRetainIT.java
@@ -0,0 +1,27 @@
+/*
+ * 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 com.hivemq.client.mqtt.MqttVersion;
+
+public class HiveMQMqtt311QosAndRetainIT extends AbstractHiveMQQosAndRetainIT {
+
+    @Override
+    protected MqttVersion mqttVersion() {
+        return MqttVersion.MQTT_3_1_1;
+    }
+}
diff --git 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQMqtt311SendDynamicIT.java
 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQMqtt311SendDynamicIT.java
new file mode 100644
index 000000000000..65f6c53e0b49
--- /dev/null
+++ 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQMqtt311SendDynamicIT.java
@@ -0,0 +1,27 @@
+/*
+ * 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 com.hivemq.client.mqtt.MqttVersion;
+
+public class HiveMQMqtt311SendDynamicIT extends AbstractHiveMQSendDynamicIT {
+
+    @Override
+    protected MqttVersion mqttVersion() {
+        return MqttVersion.MQTT_3_1_1;
+    }
+}
diff --git 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQOverrideTopicIT.java
 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQOverrideTopicIT.java
index e66505bf1e9d..f29babb53e6a 100644
--- 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQOverrideTopicIT.java
+++ 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQOverrideTopicIT.java
@@ -16,57 +16,12 @@
  */
 package org.apache.camel.component.hivemq;
 
-import org.apache.camel.EndpointInject;
-import org.apache.camel.Produce;
-import org.apache.camel.ProducerTemplate;
-import org.apache.camel.builder.RouteBuilder;
-import org.apache.camel.component.mock.MockEndpoint;
-import org.apache.camel.test.infra.hivemq.services.HiveMQService;
-import org.apache.camel.test.infra.hivemq.services.HiveMQServiceFactory;
-import org.apache.camel.test.junit6.CamelTestSupport;
-import org.junit.jupiter.api.DisplayName;
-import org.junit.jupiter.api.Test;
-import org.junit.jupiter.api.extension.RegisterExtension;
+import com.hivemq.client.mqtt.MqttVersion;
 
-public class HiveMQOverrideTopicIT extends CamelTestSupport {
-
-    @RegisterExtension
-    public static HiveMQService HIVEMQ_SERVICE = 
HiveMQServiceFactory.createService();
-
-    @EndpointInject("mock:overrideResult")
-    private MockEndpoint mockOverrideResult;
-
-    @Produce
-    private ProducerTemplate template;
-
-    @Test
-    @DisplayName("Publish to default endpoint topic but route to override 
topic via header")
-    public void testTopicOverride() throws Exception {
-        mockOverrideResult.expectedBodiesReceived("Routed to custom topic");
-        mockOverrideResult.expectedHeaderReceived(HiveMQConstants.MQTT_TOPIC, 
"orders/processed");
-
-        template.sendBodyAndHeader("direct:startOverride", "Routed to custom 
topic",
-                HiveMQConstants.OVERRIDE_TOPIC, "orders/processed");
-
-        mockOverrideResult.assertIsSatisfied();
-    }
+public class HiveMQOverrideTopicIT extends AbstractHiveMQOverrideTopicIT {
 
     @Override
-    @SuppressWarnings("deprecation")
-    protected RouteBuilder createRouteBuilder() {
-        return new RouteBuilder() {
-            @Override
-            public void configure() {
-                String host = HIVEMQ_SERVICE.getMqttHost();
-                int port = HIVEMQ_SERVICE.getMqttPort();
-
-                // Endpoint points to "orders/default", but producer will 
override to "orders/processed"
-                from("direct:startOverride")
-                        .toF("hivemq:orders/default?host=%s&port=%d", host, 
port);
-
-                fromF("hivemq:orders/processed?host=%s&port=%d", host, port)
-                        .to("mock:overrideResult");
-            }
-        };
+    protected MqttVersion mqttVersion() {
+        return MqttVersion.MQTT_5_0;
     }
 }
diff --git 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQQosAndRetainIT.java
 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQQosAndRetainIT.java
index 717332bda325..ee197fcd8b39 100644
--- 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQQosAndRetainIT.java
+++ 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQQosAndRetainIT.java
@@ -14,80 +14,14 @@
  * See the License for the specific language governing permissions and
  * limitations under the License.
  */
-
 package org.apache.camel.component.hivemq;
 
-import java.util.HashMap;
-import java.util.Map;
-
-import com.hivemq.client.mqtt.datatypes.MqttQos;
-import org.apache.camel.EndpointInject;
-import org.apache.camel.Produce;
-import org.apache.camel.ProducerTemplate;
-import org.apache.camel.builder.RouteBuilder;
-import org.apache.camel.component.mock.MockEndpoint;
-import org.apache.camel.test.infra.hivemq.services.HiveMQService;
-import org.apache.camel.test.infra.hivemq.services.HiveMQServiceFactory;
-import org.apache.camel.test.junit6.CamelTestSupport;
-import org.junit.jupiter.api.DisplayName;
-import org.junit.jupiter.api.Test;
-import org.junit.jupiter.api.extension.RegisterExtension;
-
-public class HiveMQQosAndRetainIT extends CamelTestSupport {
-
-    @RegisterExtension
-    public static HiveMQService HIVEMQ_SERVICE = 
HiveMQServiceFactory.createService();
-
-    @EndpointInject("mock:qosResult")
-    private MockEndpoint mockQosResult;
-
-    @Produce
-    private ProducerTemplate template;
-
-    @Test
-    @DisplayName("Publish with QoS 2 (EXACTLY_ONCE) and retain flag set to 
true")
-    public void testQosAndRetainHeaders() throws Exception {
-        @SuppressWarnings("deprecation")
-        String host = HIVEMQ_SERVICE.getMqttHost();
-        @SuppressWarnings("deprecation")
-        int port = HIVEMQ_SERVICE.getMqttPort();
+import com.hivemq.client.mqtt.MqttVersion;
 
-        mockQosResult.expectedBodiesReceived("Retained Critical Message");
-        mockQosResult.expectedHeaderReceived(HiveMQConstants.MQTT_QOS, 
MqttQos.EXACTLY_ONCE);
-        mockQosResult.expectedHeaderReceived(HiveMQConstants.MQTT_RETAINED, 
true);
-
-        Map<String, Object> headers = new HashMap<>();
-        headers.put(HiveMQConstants.MQTT_QOS, MqttQos.EXACTLY_ONCE);
-        headers.put(HiveMQConstants.MQTT_RETAINED, true);
-
-        // 1. Publish retained message to broker
-        template.sendBodyAndHeaders("direct:startQos", "Retained Critical 
Message", headers);
-
-        // 2. Start dynamic consumer route AFTER publishing so it fetches the 
retained message
-        context.addRoutes(new RouteBuilder() {
-            @Override
-            public void configure() {
-                fromF("hivemq:system/alerts?host=%s&port=%d&qos=EXACTLY_ONCE", 
host, port)
-                        .to("mock:qosResult");
-            }
-        });
-
-        mockQosResult.assertIsSatisfied();
-    }
+public class HiveMQQosAndRetainIT extends AbstractHiveMQQosAndRetainIT {
 
     @Override
-    protected RouteBuilder createRouteBuilder() {
-        return new RouteBuilder() {
-            @Override
-            public void configure() {
-                @SuppressWarnings("deprecation")
-                String host = HIVEMQ_SERVICE.getMqttHost();
-                @SuppressWarnings("deprecation")
-                int port = HIVEMQ_SERVICE.getMqttPort();
-
-                from("direct:startQos")
-                        .toF("hivemq:system/alerts?host=%s&port=%d", host, 
port);
-            }
-        };
+    protected MqttVersion mqttVersion() {
+        return MqttVersion.MQTT_5_0;
     }
 }
diff --git 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQSendDynamicIT.java
 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQSendDynamicIT.java
index 4cb5e364802c..ed054f072b7e 100644
--- 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQSendDynamicIT.java
+++ 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQSendDynamicIT.java
@@ -16,55 +16,12 @@
  */
 package org.apache.camel.component.hivemq;
 
-import org.apache.camel.builder.RouteBuilder;
-import org.apache.camel.test.infra.hivemq.services.HiveMQService;
-import org.apache.camel.test.infra.hivemq.services.HiveMQServiceFactory;
-import org.apache.camel.test.junit6.CamelTestSupport;
-import org.junit.jupiter.api.DisplayName;
-import org.junit.jupiter.api.Test;
-import org.junit.jupiter.api.extension.RegisterExtension;
+import com.hivemq.client.mqtt.MqttVersion;
 
-import static org.assertj.core.api.Assertions.assertThat;
-
-public class HiveMQSendDynamicIT extends CamelTestSupport {
-
-    @RegisterExtension
-    public static HiveMQService HIVEMQ_SERVICE = 
HiveMQServiceFactory.createService();
-
-    @Test
-    @DisplayName("toD with multiple dynamic topics reuses a single hivemq 
endpoint")
-    @SuppressWarnings("deprecation")
-    public void testToDReusesSingleEndpoint() {
-        template.sendBodyAndHeader("direct:start", "Hello bar", "where", 
"HiveMQSendDynamicIT-bar");
-        template.sendBodyAndHeader("direct:start", "Hello beer", "where", 
"HiveMQSendDynamicIT-beer");
-
-        long count = context.getEndpoints().stream()
-                .filter(e -> e.getEndpointUri().startsWith("hivemq:"))
-                .count();
-        assertThat(count).as("There should only be 1 hivemq 
endpoint").isEqualTo(1);
-
-        String host = HIVEMQ_SERVICE.getMqttHost();
-        int port = HIVEMQ_SERVICE.getMqttPort();
-        String out = consumer.receiveBody(
-                
String.format("hivemq:HiveMQSendDynamicIT-bar?host=%s&port=%d", host, port), 
5000, String.class);
-        assertThat(out).isEqualTo("Hello bar");
-        out = consumer.receiveBody(
-                
String.format("hivemq:HiveMQSendDynamicIT-beer?host=%s&port=%d", host, port), 
5000, String.class);
-        assertThat(out).isEqualTo("Hello beer");
-    }
+public class HiveMQSendDynamicIT extends AbstractHiveMQSendDynamicIT {
 
     @Override
-    @SuppressWarnings("deprecation")
-    protected RouteBuilder createRouteBuilder() {
-        return new RouteBuilder() {
-            @Override
-            public void configure() {
-                String host = HIVEMQ_SERVICE.getMqttHost();
-                int port = HIVEMQ_SERVICE.getMqttPort();
-
-                from("direct:start")
-                        
.toD(String.format("hivemq:${header.where}?host=%s&port=%d&retained=true", 
host, port));
-            }
-        };
+    protected MqttVersion mqttVersion() {
+        return MqttVersion.MQTT_5_0;
     }
 }

Reply via email to