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