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 363d2ef35fb5ed4c3faa764b973280536dc7b825
Author: Thomas Strauss <[email protected]>
AuthorDate: Tue Sep 29 23:41:29 2026 +0200

    CAMEL-25030: camel-hivemq: add test coverage for MQTT 3.1.1 support
    
    Adds HiveMQMqtt311PubSubIT, a Testcontainers-based pub/sub IT mirroring
    the existing MQTT 5 IT but with mqttVersion=MQTT_3_1_1 (verified
    against a live HiveMQ broker), and HiveMQClientAccessTest, covering
    HiveMQConsumer/HiveMQProducer.getClient(Class<T>) for both protocol
    versions (empty before start, correct type present, other type empty).
    
    Updates the white-box unit tests that previously reached into the raw
    Mqtt5AsyncClient (HiveMQEndpointConnectTest, HiveMQEndpointAuthTest) to
    go through the HiveMQClientAdapter/getClient(Class<T>) API instead, and
    extends HiveMQEndpointAuthTest to cover simple-auth handling for both
    MQTT 3.1.1 and MQTT 5. HiveMQConfigurationTest now also verifies the
    mqttVersion default and that it survives configuration copying.
    
    Verified: all unit tests and all 6 ITs (pub/sub, QoS/retain,
    empty-payload, topic-override, send-dynamic, MQTT 3.1.1 pub/sub) 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]>
---
 .../component/hivemq/HiveMQClientAccessTest.java   | 99 ++++++++++++++++++++++
 .../component/hivemq/HiveMQConfigurationTest.java  |  4 +
 .../component/hivemq/HiveMQEndpointAuthTest.java   | 68 ++++++++++++---
 .../hivemq/HiveMQEndpointConnectTest.java          | 20 ++---
 .../component/hivemq/HiveMQMqtt311PubSubIT.java    | 72 ++++++++++++++++
 5 files changed, 237 insertions(+), 26 deletions(-)

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
new file mode 100644
index 000000000000..97ff93e2eab0
--- /dev/null
+++ 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQClientAccessTest.java
@@ -0,0 +1,99 @@
+/*
+ * 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.mqtt3.Mqtt3AsyncClient;
+import com.hivemq.client.mqtt.mqtt5.Mqtt5AsyncClient;
+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 static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+class HiveMQClientAccessTest {
+
+    private DefaultCamelContext camelContext;
+    private HiveMQComponent component;
+
+    @BeforeEach
+    void setUp() {
+        camelContext = new DefaultCamelContext();
+        component = new HiveMQComponent();
+        component.setCamelContext(camelContext);
+    }
+
+    @AfterEach
+    void tearDown() {
+        camelContext.stop();
+    }
+
+    @Test
+    @DisplayName("Producer exposes no client implementation before it has 
started")
+    void producerHasNoClientBeforeStart() {
+        HiveMQProducer producer = new HiveMQProducer(newEndpoint(new 
HiveMQConfiguration()));
+
+        assertThat(producer.getClient(Mqtt5AsyncClient.class)).isEmpty();
+    }
+
+    @Test
+    @DisplayName("Consumer exposes no client implementation before it has 
started")
+    void consumerHasNoClientBeforeStart() {
+        HiveMQConsumer consumer = new HiveMQConsumer(newEndpoint(new 
HiveMQConfiguration()), exchange -> {
+        });
+
+        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()));
+
+        assertThatThrownBy(producer::doStart);
+
+        assertThat(producer.getClient(Mqtt5AsyncClient.class)).isPresent();
+        assertThat(producer.getClient(Mqtt3AsyncClient.class)).isEmpty();
+    }
+
+    @Test
+    @DisplayName("Consumer exposes the underlying Mqtt5AsyncClient once 
created, not the Mqtt3 type")
+    void consumerExposesMatchingClientTypeAfterConnectAttempt() {
+        HiveMQConsumer consumer = new 
HiveMQConsumer(newEndpoint(unreachableConfiguration()), exchange -> {
+        });
+
+        assertThatThrownBy(consumer::doStart);
+
+        assertThat(consumer.getClient(Mqtt5AsyncClient.class)).isPresent();
+        assertThat(consumer.getClient(Mqtt3AsyncClient.class)).isEmpty();
+    }
+
+    private static HiveMQConfiguration unreachableConfiguration() {
+        HiveMQConfiguration configuration = new HiveMQConfiguration();
+        configuration.setHost("127.0.0.1");
+        configuration.setPort(1);
+        return configuration;
+    }
+
+    private HiveMQEndpoint newEndpoint(HiveMQConfiguration configuration) {
+        HiveMQEndpoint endpoint = new HiveMQEndpoint("hivemq:test", component, 
configuration, "test");
+        endpoint.setCamelContext(camelContext);
+        return endpoint;
+    }
+}
diff --git 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQConfigurationTest.java
 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQConfigurationTest.java
index e125f64b912a..f0d20bd30d07 100644
--- 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQConfigurationTest.java
+++ 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQConfigurationTest.java
@@ -16,6 +16,7 @@
  */
 package org.apache.camel.component.hivemq;
 
+import com.hivemq.client.mqtt.MqttVersion;
 import com.hivemq.client.mqtt.datatypes.MqttQos;
 import org.junit.jupiter.api.DisplayName;
 import org.junit.jupiter.api.Test;
@@ -31,6 +32,7 @@ class HiveMQConfigurationTest {
 
         assertThat(config.getHost()).isEqualTo(HiveMQConstants.DEFAULT_HOST);
         assertThat(config.getPort()).isEqualTo(HiveMQConstants.DEFAULT_PORT);
+        assertThat(config.getMqttVersion()).isEqualTo(MqttVersion.MQTT_5_0);
         assertThat(config.getQos()).isEqualTo(MqttQos.AT_LEAST_ONCE);
         assertThat(config.isRetained()).isFalse();
         assertThat(config.isCleanStart()).isTrue();
@@ -43,6 +45,7 @@ class HiveMQConfigurationTest {
         HiveMQConfiguration original = new HiveMQConfiguration();
         original.setHost("broker.hivemq.com");
         original.setPort(8883);
+        original.setMqttVersion(MqttVersion.MQTT_3_1_1);
         original.setQos(MqttQos.EXACTLY_ONCE);
         original.setRetained(true);
         original.setUsername("admin");
@@ -53,6 +56,7 @@ class HiveMQConfigurationTest {
         assertThat(copy).isNotSameAs(original);
         assertThat(copy.getHost()).isEqualTo("broker.hivemq.com");
         assertThat(copy.getPort()).isEqualTo(8883);
+        assertThat(copy.getMqttVersion()).isEqualTo(MqttVersion.MQTT_3_1_1);
         assertThat(copy.getQos()).isEqualTo(MqttQos.EXACTLY_ONCE);
         assertThat(copy.isRetained()).isTrue();
         assertThat(copy.getUsername()).isEqualTo("admin");
diff --git 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQEndpointAuthTest.java
 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQEndpointAuthTest.java
index 69910044ff9f..8a2e1ad29fd2 100644
--- 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQEndpointAuthTest.java
+++ 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQEndpointAuthTest.java
@@ -16,8 +16,11 @@
  */
 package org.apache.camel.component.hivemq;
 
+import java.nio.ByteBuffer;
 import java.nio.charset.StandardCharsets;
 
+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;
 import org.junit.jupiter.api.AfterEach;
@@ -45,35 +48,74 @@ class HiveMQEndpointAuthTest {
     }
 
     @Test
-    @DisplayName("Username without password builds MQTT simple auth and does 
not NPE")
-    void usernameWithoutPasswordDoesNotThrow() {
+    @DisplayName("MQTT 5: username without password builds MQTT simple auth 
and does not NPE")
+    void mqtt5UsernameWithoutPasswordDoesNotThrow() {
         HiveMQConfiguration configuration = new HiveMQConfiguration();
+        configuration.setMqttVersion(MqttVersion.MQTT_5_0);
         configuration.setUsername("mqtt-user");
         configuration.setPassword(null);
 
-        HiveMQEndpoint endpoint = new HiveMQEndpoint("hivemq:test", component, 
configuration, "test");
-        Mqtt5AsyncClient client = endpoint.createClient();
+        Mqtt5AsyncClient client = createEndpoint(configuration).createClient()
+                .getClient(Mqtt5AsyncClient.class).orElseThrow();
 
         assertThat(client.getConfig().getSimpleAuth()).isPresent();
         
assertThat(client.getConfig().getSimpleAuth().orElseThrow().getPassword()).isEmpty();
     }
 
     @Test
-    @DisplayName("Username and password are both applied to MQTT simple auth")
-    void usernameWithPasswordSetsPassword() {
+    @DisplayName("MQTT 5: username and password are both applied to MQTT 
simple auth")
+    void mqtt5UsernameWithPasswordSetsPassword() {
         HiveMQConfiguration configuration = new HiveMQConfiguration();
+        configuration.setMqttVersion(MqttVersion.MQTT_5_0);
         configuration.setUsername("mqtt-user");
         configuration.setPassword("secret");
 
-        HiveMQEndpoint endpoint = new HiveMQEndpoint("hivemq:test", component, 
configuration, "test");
-        Mqtt5AsyncClient client = endpoint.createClient();
+        Mqtt5AsyncClient client = createEndpoint(configuration).createClient()
+                .getClient(Mqtt5AsyncClient.class).orElseThrow();
 
         assertThat(client.getConfig().getSimpleAuth()).isPresent();
         
assertThat(client.getConfig().getSimpleAuth().orElseThrow().getPassword())
-                .hasValueSatisfying(buffer -> {
-                    byte[] bytes = new byte[buffer.remaining()];
-                    buffer.get(bytes);
-                    
assertThat(bytes).isEqualTo("secret".getBytes(StandardCharsets.UTF_8));
-                });
+                .hasValueSatisfying(buffer -> 
assertThat(toBytes(buffer)).isEqualTo("secret".getBytes(StandardCharsets.UTF_8)));
+    }
+
+    @Test
+    @DisplayName("MQTT 3.1.1: username without password builds MQTT simple 
auth and does not NPE")
+    void mqtt311UsernameWithoutPasswordDoesNotThrow() {
+        HiveMQConfiguration configuration = new HiveMQConfiguration();
+        configuration.setMqttVersion(MqttVersion.MQTT_3_1_1);
+        configuration.setUsername("mqtt-user");
+        configuration.setPassword(null);
+
+        Mqtt3AsyncClient client = createEndpoint(configuration).createClient()
+                .getClient(Mqtt3AsyncClient.class).orElseThrow();
+
+        assertThat(client.getConfig().getSimpleAuth()).isPresent();
+        
assertThat(client.getConfig().getSimpleAuth().orElseThrow().getPassword()).isEmpty();
+    }
+
+    @Test
+    @DisplayName("MQTT 3.1.1: username and password are both applied to MQTT 
simple auth")
+    void mqtt311UsernameWithPasswordSetsPassword() {
+        HiveMQConfiguration configuration = new HiveMQConfiguration();
+        configuration.setMqttVersion(MqttVersion.MQTT_3_1_1);
+        configuration.setUsername("mqtt-user");
+        configuration.setPassword("secret");
+
+        Mqtt3AsyncClient client = createEndpoint(configuration).createClient()
+                .getClient(Mqtt3AsyncClient.class).orElseThrow();
+
+        assertThat(client.getConfig().getSimpleAuth()).isPresent();
+        
assertThat(client.getConfig().getSimpleAuth().orElseThrow().getPassword())
+                .hasValueSatisfying(buffer -> 
assertThat(toBytes(buffer)).isEqualTo("secret".getBytes(StandardCharsets.UTF_8)));
+    }
+
+    private HiveMQEndpoint createEndpoint(HiveMQConfiguration configuration) {
+        return new HiveMQEndpoint("hivemq:test", component, configuration, 
"test");
+    }
+
+    private static byte[] toBytes(ByteBuffer buffer) {
+        byte[] bytes = new byte[buffer.remaining()];
+        buffer.get(bytes);
+        return bytes;
     }
 }
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 c96568983caa..ba7407b74c1b 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,8 +18,6 @@ package org.apache.camel.component.hivemq;
 
 import java.util.concurrent.TimeUnit;
 
-import com.hivemq.client.mqtt.MqttClientState;
-import com.hivemq.client.mqtt.mqtt5.Mqtt5AsyncClient;
 import org.apache.camel.RuntimeCamelException;
 import org.apache.camel.impl.DefaultCamelContext;
 import org.junit.jupiter.api.AfterEach;
@@ -56,7 +54,7 @@ class HiveMQEndpointConnectTest {
         configuration.setPort(1);
 
         HiveMQEndpoint endpoint = new HiveMQEndpoint("hivemq:test", component, 
configuration, "test");
-        Mqtt5AsyncClient client = endpoint.createClient();
+        HiveMQClientAdapter client = endpoint.createClient();
 
         long started = System.nanoTime();
         assertThatThrownBy(() -> 
endpoint.connect(client)).isInstanceOf(RuntimeCamelException.class);
@@ -64,27 +62,23 @@ class HiveMQEndpointConnectTest {
 
         
assertThat(elapsedMs).isLessThan(TimeUnit.SECONDS.toMillis(HiveMQConstants.DEFAULT_CONNECT_TIMEOUT_SECONDS));
         await().atMost(5, TimeUnit.SECONDS)
-                .untilAsserted(() -> 
assertThat(client.getState().isConnectedOrReconnect()).isFalse());
-        assertThat(client.getState()).isEqualTo(MqttClientState.DISCONNECTED);
+                .untilAsserted(() -> 
assertThat(client.isConnectedOrReconnecting()).isFalse());
     }
 
     @Test
-    @DisplayName("stopClient cancels automatic reconnect when the client is 
not CONNECTED")
+    @DisplayName("stop() cancels automatic reconnect when the client is not 
CONNECTED")
     void stopClientCancelsReconnect() {
         HiveMQConfiguration configuration = new HiveMQConfiguration();
         configuration.setHost("127.0.0.1");
         configuration.setPort(1);
 
         HiveMQEndpoint endpoint = new HiveMQEndpoint("hivemq:test", component, 
configuration, "test");
-        Mqtt5AsyncClient client = endpoint.createClient();
+        HiveMQClientAdapter client = endpoint.createClient();
 
-        client.connectWith().cleanStart(true).send();
-        endpoint.stopClient(client);
+        client.connect(true);
+        client.stop();
 
         await().atMost(5, TimeUnit.SECONDS)
-                .untilAsserted(() -> {
-                    
assertThat(client.getState().isConnectedOrReconnect()).isFalse();
-                    
assertThat(client.getState()).isEqualTo(MqttClientState.DISCONNECTED);
-                });
+                .untilAsserted(() -> 
assertThat(client.isConnectedOrReconnecting()).isFalse());
     }
 }
diff --git 
a/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQMqtt311PubSubIT.java
 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQMqtt311PubSubIT.java
new file mode 100644
index 000000000000..6ef6afbc72f4
--- /dev/null
+++ 
b/components/camel-hivemq/src/test/java/org/apache/camel/component/hivemq/HiveMQMqtt311PubSubIT.java
@@ -0,0 +1,72 @@
+/*
+ * 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.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 HiveMQMqtt311PubSubIT extends CamelTestSupport {
+
+    @RegisterExtension
+    public static HiveMQService HIVEMQ_SERVICE = 
HiveMQServiceFactory.createService();
+
+    @EndpointInject("mock:result")
+    private MockEndpoint mockResult;
+
+    @Produce
+    private ProducerTemplate template;
+
+    @Test
+    @DisplayName("Publish string payload and receive it over MQTT 3.1.1")
+    public void testBasicPubSub() throws Exception {
+        mockResult.expectedBodiesReceived("Hello HiveMQ MQTT 3.1.1!");
+        mockResult.expectedHeaderReceived(HiveMQConstants.MQTT_TOPIC, 
"test/mqtt311");
+        mockResult.expectedHeaderReceived(HiveMQConstants.MQTT_QOS, 
MqttQos.AT_LEAST_ONCE);
+
+        template.sendBody("direct:startMqtt311", "Hello HiveMQ MQTT 3.1.1!");
+
+        mockResult.assertIsSatisfied();
+    }
+
+    @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:startMqtt311")
+                        
.toF("hivemq:test/mqtt311?host=%s&port=%d&mqttVersion=MQTT_3_1_1", host, port);
+
+                
fromF("hivemq:test/mqtt311?host=%s&port=%d&mqttVersion=MQTT_3_1_1", host, port)
+                        .to("mock:result");
+            }
+        };
+    }
+}

Reply via email to