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