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 fabf0146f97d70598e8264d760ef8002aa5ffb1b Author: Thomas Strauss <[email protected]> AuthorDate: Tue Sep 29 23:40:56 2026 +0200 CAMEL-25030: camel-hivemq: add MQTT 3.1.1 protocol support Introduces a new mqttVersion option (MQTT_5_0, the default, or MQTT_3_1_1) so the HiveMQ component can connect using either MQTT protocol version, instead of MQTT 5 only. The hivemq-mqtt-client library exposes MQTT 3.1.1 and MQTT 5 as two structurally separate client APIs with no common publish/subscribe interface. To avoid duplicating connect/publish/subscribe/disconnect logic across HiveMQEndpoint, HiveMQConsumer and HiveMQProducer, this adds a small internal HiveMQClientAdapter interface with one implementation per protocol version (Mqtt3ClientAdapter, Mqtt5ClientAdapter). HiveMQEndpoint/Consumer/Producer depend only on this protocol-neutral adapter; HiveMQEndpoint.createClient()/connect() are package-private, since nothing outside the package needs them. HiveMQClientAdapter.getClient(Class<T>) gives callers who need functionality this component does not cover direct, type-safe access to the real underlying client object (Mqtt5AsyncClient or Mqtt3AsyncClient), exposed publicly via a same-named method on HiveMQConsumer and HiveMQProducer. Named to match the plain getClient() convention used by every other native-client-holding Camel component (camel-paho, camel-paho-mqtt5, camel-aws2-sqs, camel-elasticsearch and others), extended with the Class<T>/Optional<T> shape Camel's own Component.getExtension(Class<T>) already uses for "give me this as a specific type" lookups. Regenerates this component's own metadata/configurers and the two per-component DSL builder factories (Hivemq(Component|Endpoint) BuilderFactory) so the typed Java DSL exposes mqttVersion like every other option. Verified with a full, non-quickly `mvn install -DskipTests` build of this module (BUILD SUCCESS, formatter/checkstyle clean). AI-assisted contribution: authored together with Claude Sonnet 5 (Anthropic). Co-Authored-By: Claude Sonnet 5 <[email protected]> --- .../hivemq/HiveMQComponentConfigurer.java | 6 + .../component/hivemq/HiveMQEndpointConfigurer.java | 6 + .../component/hivemq/HiveMQEndpointUriFactory.java | 3 +- .../org/apache/camel/component/hivemq/hivemq.json | 44 +++--- .../component/hivemq/HiveMQClientAdapter.java | 63 +++++++++ .../component/hivemq/HiveMQConfiguration.java | 18 ++- .../camel/component/hivemq/HiveMQConsumer.java | 39 +++--- .../camel/component/hivemq/HiveMQEndpoint.java | 104 ++------------ .../camel/component/hivemq/HiveMQMessage.java | 25 ++++ .../camel/component/hivemq/HiveMQProducer.java | 27 ++-- .../camel/component/hivemq/Mqtt3ClientAdapter.java | 156 +++++++++++++++++++++ ...HiveMQEndpoint.java => Mqtt5ClientAdapter.java} | 146 +++++++------------ .../dsl/HivemqComponentBuilderFactory.java | 23 ++- .../endpoint/dsl/HiveMQEndpointBuilderFactory.java | 120 ++++++++++++++-- 14 files changed, 523 insertions(+), 257 deletions(-) diff --git a/components/camel-hivemq/src/generated/java/org/apache/camel/component/hivemq/HiveMQComponentConfigurer.java b/components/camel-hivemq/src/generated/java/org/apache/camel/component/hivemq/HiveMQComponentConfigurer.java index b8b64c1a7591..aea77890f75f 100644 --- a/components/camel-hivemq/src/generated/java/org/apache/camel/component/hivemq/HiveMQComponentConfigurer.java +++ b/components/camel-hivemq/src/generated/java/org/apache/camel/component/hivemq/HiveMQComponentConfigurer.java @@ -42,6 +42,8 @@ public class HiveMQComponentConfigurer extends PropertyConfigurerSupport impleme case "host": getOrCreateConfiguration(target).setHost(property(camelContext, java.lang.String.class, value)); return true; case "lazystartproducer": case "lazyStartProducer": target.setLazyStartProducer(property(camelContext, boolean.class, value)); return true; + case "mqttversion": + case "mqttVersion": getOrCreateConfiguration(target).setMqttVersion(property(camelContext, com.hivemq.client.mqtt.MqttVersion.class, value)); return true; case "password": getOrCreateConfiguration(target).setPassword(property(camelContext, java.lang.String.class, value)); return true; case "port": getOrCreateConfiguration(target).setPort(property(camelContext, int.class, value)); return true; case "qos": getOrCreateConfiguration(target).setQos(property(camelContext, com.hivemq.client.mqtt.datatypes.MqttQos.class, value)); return true; @@ -67,6 +69,8 @@ public class HiveMQComponentConfigurer extends PropertyConfigurerSupport impleme case "host": return java.lang.String.class; case "lazystartproducer": case "lazyStartProducer": return boolean.class; + case "mqttversion": + case "mqttVersion": return com.hivemq.client.mqtt.MqttVersion.class; case "password": return java.lang.String.class; case "port": return int.class; case "qos": return com.hivemq.client.mqtt.datatypes.MqttQos.class; @@ -93,6 +97,8 @@ public class HiveMQComponentConfigurer extends PropertyConfigurerSupport impleme case "host": return getOrCreateConfiguration(target).getHost(); case "lazystartproducer": case "lazyStartProducer": return target.isLazyStartProducer(); + case "mqttversion": + case "mqttVersion": return getOrCreateConfiguration(target).getMqttVersion(); case "password": return getOrCreateConfiguration(target).getPassword(); case "port": return getOrCreateConfiguration(target).getPort(); case "qos": return getOrCreateConfiguration(target).getQos(); diff --git a/components/camel-hivemq/src/generated/java/org/apache/camel/component/hivemq/HiveMQEndpointConfigurer.java b/components/camel-hivemq/src/generated/java/org/apache/camel/component/hivemq/HiveMQEndpointConfigurer.java index 672b100b8cb8..8c3b7146e6e1 100644 --- a/components/camel-hivemq/src/generated/java/org/apache/camel/component/hivemq/HiveMQEndpointConfigurer.java +++ b/components/camel-hivemq/src/generated/java/org/apache/camel/component/hivemq/HiveMQEndpointConfigurer.java @@ -36,6 +36,8 @@ public class HiveMQEndpointConfigurer extends PropertyConfigurerSupport implemen case "host": target.getConfiguration().setHost(property(camelContext, java.lang.String.class, value)); return true; case "lazystartproducer": case "lazyStartProducer": target.setLazyStartProducer(property(camelContext, boolean.class, value)); return true; + case "mqttversion": + case "mqttVersion": target.getConfiguration().setMqttVersion(property(camelContext, com.hivemq.client.mqtt.MqttVersion.class, value)); return true; case "password": target.getConfiguration().setPassword(property(camelContext, java.lang.String.class, value)); return true; case "port": target.getConfiguration().setPort(property(camelContext, int.class, value)); return true; case "qos": target.getConfiguration().setQos(property(camelContext, com.hivemq.client.mqtt.datatypes.MqttQos.class, value)); return true; @@ -62,6 +64,8 @@ public class HiveMQEndpointConfigurer extends PropertyConfigurerSupport implemen case "host": return java.lang.String.class; case "lazystartproducer": case "lazyStartProducer": return boolean.class; + case "mqttversion": + case "mqttVersion": return com.hivemq.client.mqtt.MqttVersion.class; case "password": return java.lang.String.class; case "port": return int.class; case "qos": return com.hivemq.client.mqtt.datatypes.MqttQos.class; @@ -89,6 +93,8 @@ public class HiveMQEndpointConfigurer extends PropertyConfigurerSupport implemen case "host": return target.getConfiguration().getHost(); case "lazystartproducer": case "lazyStartProducer": return target.isLazyStartProducer(); + case "mqttversion": + case "mqttVersion": return target.getConfiguration().getMqttVersion(); case "password": return target.getConfiguration().getPassword(); case "port": return target.getConfiguration().getPort(); case "qos": return target.getConfiguration().getQos(); diff --git a/components/camel-hivemq/src/generated/java/org/apache/camel/component/hivemq/HiveMQEndpointUriFactory.java b/components/camel-hivemq/src/generated/java/org/apache/camel/component/hivemq/HiveMQEndpointUriFactory.java index f385056ff92d..0a3ded932025 100644 --- a/components/camel-hivemq/src/generated/java/org/apache/camel/component/hivemq/HiveMQEndpointUriFactory.java +++ b/components/camel-hivemq/src/generated/java/org/apache/camel/component/hivemq/HiveMQEndpointUriFactory.java @@ -24,7 +24,7 @@ public class HiveMQEndpointUriFactory extends org.apache.camel.support.component private static final Set<String> ENDPOINT_IDENTITY_PROPERTY_NAMES; private static final Map<String, String> MULTI_VALUE_PREFIXES; static { - Set<String> props = new HashSet<>(14); + Set<String> props = new HashSet<>(15); props.add("bridgeErrorHandler"); props.add("cleanStart"); props.add("clientId"); @@ -32,6 +32,7 @@ public class HiveMQEndpointUriFactory extends org.apache.camel.support.component props.add("exchangePattern"); props.add("host"); props.add("lazyStartProducer"); + props.add("mqttVersion"); props.add("password"); props.add("port"); props.add("qos"); diff --git a/components/camel-hivemq/src/generated/resources/META-INF/org/apache/camel/component/hivemq/hivemq.json b/components/camel-hivemq/src/generated/resources/META-INF/org/apache/camel/component/hivemq/hivemq.json index 9b0950bfb59f..0dbbee22aad2 100644 --- a/components/camel-hivemq/src/generated/resources/META-INF/org/apache/camel/component/hivemq/hivemq.json +++ b/components/camel-hivemq/src/generated/resources/META-INF/org/apache/camel/component/hivemq/hivemq.json @@ -24,19 +24,20 @@ "remote": true }, "componentProperties": { - "cleanStart": { "index": 0, "kind": "property", "displayName": "Clean Start", "group": "common", "label": "", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "defaultValue": true, "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Whether to initiate a clean start (MQTT 5) upon connecting to the broker." }, + "cleanStart": { "index": 0, "kind": "property", "displayName": "Clean Start", "group": "common", "label": "", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "defaultValue": true, "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Whether to initiate a clean session upon connecting to the broker (called clean session in MQTT 3.1.1 a [...] "clientId": { "index": 1, "kind": "property", "displayName": "Client Id", "group": "common", "label": "", "required": false, "type": "string", "javaType": "java.lang.String", "deprecated": false, "autowired": false, "secret": false, "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Client identifier used when connecting to the HiveMQ broker." }, "configuration": { "index": 2, "kind": "property", "displayName": "Configuration", "group": "common", "label": "", "required": false, "type": "object", "javaType": "org.apache.camel.component.hivemq.HiveMQConfiguration", "deprecated": false, "autowired": false, "secret": false, "description": "Component configuration." }, "host": { "index": 3, "kind": "property", "displayName": "Host", "group": "common", "label": "", "required": false, "type": "string", "javaType": "java.lang.String", "deprecated": false, "autowired": false, "secret": false, "defaultValue": "localhost", "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Hostname or IP address of the HiveMQ MQTT broker." }, - "port": { "index": 4, "kind": "property", "displayName": "Port", "group": "common", "label": "", "required": false, "type": "integer", "javaType": "int", "deprecated": false, "autowired": false, "secret": false, "defaultValue": 1883, "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Port number of the HiveMQ MQTT broker." }, - "qos": { "index": 5, "kind": "property", "displayName": "Qos", "group": "common", "label": "", "required": false, "type": "enum", "javaType": "com.hivemq.client.mqtt.datatypes.MqttQos", "enum": [ "AT_MOST_ONCE", "AT_LEAST_ONCE", "EXACTLY_ONCE" ], "deprecated": false, "autowired": false, "secret": false, "defaultValue": "AT_LEAST_ONCE", "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Default Quality [...] - "retained": { "index": 6, "kind": "property", "displayName": "Retained", "group": "common", "label": "", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "defaultValue": false, "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Whether published messages should be retained by the MQTT broker." }, - "bridgeErrorHandler": { "index": 7, "kind": "property", "displayName": "Bridge Error Handler", "group": "consumer", "label": "consumer", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "defaultValue": false, "description": "Allows for bridging the consumer to the Camel routing Error Handler, which mean any exceptions (if possible) occurred while the Camel consumer is trying to pickup incoming messages, or the like [...] - "lazyStartProducer": { "index": 8, "kind": "property", "displayName": "Lazy Start Producer", "group": "producer", "label": "producer", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "defaultValue": false, "description": "Whether the producer should be started lazy (on the first message). By starting lazy you can use this to allow CamelContext and routes to startup in situations where a producer may otherwise fail [...] - "autowiredEnabled": { "index": 9, "kind": "property", "displayName": "Autowired Enabled", "group": "advanced", "label": "advanced", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "defaultValue": true, "description": "Whether autowiring is enabled. This is used for automatic autowiring options (the option must be marked as autowired) by looking up in the registry to find if there is a single instance of matching t [...] - "password": { "index": 10, "kind": "property", "displayName": "Password", "group": "security", "label": "security", "required": false, "type": "string", "javaType": "java.lang.String", "deprecated": false, "deprecationNote": "", "autowired": false, "secret": true, "security": "secret", "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Password for authentication with the HiveMQ broker." }, - "ssl": { "index": 11, "kind": "property", "displayName": "Ssl", "group": "security", "label": "security", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "security": "insecure:ssl", "insecureValue": "false", "defaultValue": false, "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Whether to enable SSL\/TLS encryption for the broker [...] - "username": { "index": 12, "kind": "property", "displayName": "Username", "group": "security", "label": "security", "required": false, "type": "string", "javaType": "java.lang.String", "deprecated": false, "deprecationNote": "", "autowired": false, "secret": false, "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Username for authentication with the HiveMQ broker." } + "mqttVersion": { "index": 4, "kind": "property", "displayName": "Mqtt Version", "group": "common", "label": "", "required": false, "type": "enum", "javaType": "com.hivemq.client.mqtt.MqttVersion", "enum": [ "MQTT_3_1_1", "MQTT_5_0" ], "deprecated": false, "autowired": false, "secret": false, "defaultValue": "MQTT_5_0", "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "The MQTT protocol version to use [...] + "port": { "index": 5, "kind": "property", "displayName": "Port", "group": "common", "label": "", "required": false, "type": "integer", "javaType": "int", "deprecated": false, "autowired": false, "secret": false, "defaultValue": 1883, "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Port number of the HiveMQ MQTT broker." }, + "qos": { "index": 6, "kind": "property", "displayName": "Qos", "group": "common", "label": "", "required": false, "type": "enum", "javaType": "com.hivemq.client.mqtt.datatypes.MqttQos", "enum": [ "AT_MOST_ONCE", "AT_LEAST_ONCE", "EXACTLY_ONCE" ], "deprecated": false, "autowired": false, "secret": false, "defaultValue": "AT_LEAST_ONCE", "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Default Quality [...] + "retained": { "index": 7, "kind": "property", "displayName": "Retained", "group": "common", "label": "", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "defaultValue": false, "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Whether published messages should be retained by the MQTT broker." }, + "bridgeErrorHandler": { "index": 8, "kind": "property", "displayName": "Bridge Error Handler", "group": "consumer", "label": "consumer", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "defaultValue": false, "description": "Allows for bridging the consumer to the Camel routing Error Handler, which mean any exceptions (if possible) occurred while the Camel consumer is trying to pickup incoming messages, or the like [...] + "lazyStartProducer": { "index": 9, "kind": "property", "displayName": "Lazy Start Producer", "group": "producer", "label": "producer", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "defaultValue": false, "description": "Whether the producer should be started lazy (on the first message). By starting lazy you can use this to allow CamelContext and routes to startup in situations where a producer may otherwise fail [...] + "autowiredEnabled": { "index": 10, "kind": "property", "displayName": "Autowired Enabled", "group": "advanced", "label": "advanced", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "defaultValue": true, "description": "Whether autowiring is enabled. This is used for automatic autowiring options (the option must be marked as autowired) by looking up in the registry to find if there is a single instance of matching [...] + "password": { "index": 11, "kind": "property", "displayName": "Password", "group": "security", "label": "security", "required": false, "type": "string", "javaType": "java.lang.String", "deprecated": false, "deprecationNote": "", "autowired": false, "secret": true, "security": "secret", "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Password for authentication with the HiveMQ broker." }, + "ssl": { "index": 12, "kind": "property", "displayName": "Ssl", "group": "security", "label": "security", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "security": "insecure:ssl", "insecureValue": "false", "defaultValue": false, "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Whether to enable SSL\/TLS encryption for the broker [...] + "username": { "index": 13, "kind": "property", "displayName": "Username", "group": "security", "label": "security", "required": false, "type": "string", "javaType": "java.lang.String", "deprecated": false, "deprecationNote": "", "autowired": false, "secret": false, "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Username for authentication with the HiveMQ broker." } }, "headers": { "CamelHiveMQTopic": { "index": 0, "kind": "header", "displayName": "", "group": "common", "label": "", "required": false, "javaType": "String", "deprecated": false, "deprecationNote": "", "autowired": false, "secret": false, "description": "The topic to publish\/subscribe to.", "constantName": "org.apache.camel.component.hivemq.HiveMQConstants#MQTT_TOPIC" }, @@ -46,18 +47,19 @@ }, "properties": { "topic": { "index": 0, "kind": "path", "displayName": "Topic", "group": "common", "label": "", "required": true, "type": "string", "javaType": "java.lang.String", "deprecated": false, "deprecationNote": "", "autowired": false, "secret": false, "description": "The MQTT topic name or pattern to subscribe to or publish on." }, - "cleanStart": { "index": 1, "kind": "parameter", "displayName": "Clean Start", "group": "common", "label": "", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "defaultValue": true, "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Whether to initiate a clean start (MQTT 5) upon connecting to the broker." }, + "cleanStart": { "index": 1, "kind": "parameter", "displayName": "Clean Start", "group": "common", "label": "", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "defaultValue": true, "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Whether to initiate a clean session upon connecting to the broker (called clean session in MQTT 3.1.1 [...] "clientId": { "index": 2, "kind": "parameter", "displayName": "Client Id", "group": "common", "label": "", "required": false, "type": "string", "javaType": "java.lang.String", "deprecated": false, "autowired": false, "secret": false, "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Client identifier used when connecting to the HiveMQ broker." }, "host": { "index": 3, "kind": "parameter", "displayName": "Host", "group": "common", "label": "", "required": false, "type": "string", "javaType": "java.lang.String", "deprecated": false, "autowired": false, "secret": false, "defaultValue": "localhost", "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Hostname or IP address of the HiveMQ MQTT broker." }, - "port": { "index": 4, "kind": "parameter", "displayName": "Port", "group": "common", "label": "", "required": false, "type": "integer", "javaType": "int", "deprecated": false, "autowired": false, "secret": false, "defaultValue": 1883, "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Port number of the HiveMQ MQTT broker." }, - "qos": { "index": 5, "kind": "parameter", "displayName": "Qos", "group": "common", "label": "", "required": false, "type": "enum", "javaType": "com.hivemq.client.mqtt.datatypes.MqttQos", "enum": [ "AT_MOST_ONCE", "AT_LEAST_ONCE", "EXACTLY_ONCE" ], "deprecated": false, "autowired": false, "secret": false, "defaultValue": "AT_LEAST_ONCE", "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Default Quality [...] - "retained": { "index": 6, "kind": "parameter", "displayName": "Retained", "group": "common", "label": "", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "defaultValue": false, "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Whether published messages should be retained by the MQTT broker." }, - "bridgeErrorHandler": { "index": 7, "kind": "parameter", "displayName": "Bridge Error Handler", "group": "consumer (advanced)", "label": "consumer,advanced", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "defaultValue": false, "description": "Allows for bridging the consumer to the Camel routing Error Handler, which mean any exceptions (if possible) occurred while the Camel consumer is trying to pickup incoming [...] - "exceptionHandler": { "index": 8, "kind": "parameter", "displayName": "Exception Handler", "group": "consumer (advanced)", "label": "consumer,advanced", "required": false, "type": "object", "javaType": "org.apache.camel.spi.ExceptionHandler", "optionalPrefix": "consumer.", "deprecated": false, "autowired": false, "secret": false, "description": "To let the consumer use a custom ExceptionHandler. Notice if the option bridgeErrorHandler is enabled then this option is not in use. By def [...] - "exchangePattern": { "index": 9, "kind": "parameter", "displayName": "Exchange Pattern", "group": "consumer (advanced)", "label": "consumer,advanced", "required": false, "type": "enum", "javaType": "org.apache.camel.ExchangePattern", "enum": [ "InOnly", "InOut" ], "deprecated": false, "autowired": false, "secret": false, "description": "Sets the exchange pattern when the consumer creates an exchange." }, - "lazyStartProducer": { "index": 10, "kind": "parameter", "displayName": "Lazy Start Producer", "group": "producer (advanced)", "label": "producer,advanced", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "defaultValue": false, "description": "Whether the producer should be started lazy (on the first message). By starting lazy you can use this to allow CamelContext and routes to startup in situations where a produ [...] - "password": { "index": 11, "kind": "parameter", "displayName": "Password", "group": "security", "label": "security", "required": false, "type": "string", "javaType": "java.lang.String", "deprecated": false, "deprecationNote": "", "autowired": false, "secret": true, "security": "secret", "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Password for authentication with the HiveMQ broker." }, - "ssl": { "index": 12, "kind": "parameter", "displayName": "Ssl", "group": "security", "label": "security", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "security": "insecure:ssl", "insecureValue": "false", "defaultValue": false, "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Whether to enable SSL\/TLS encryption for the broke [...] - "username": { "index": 13, "kind": "parameter", "displayName": "Username", "group": "security", "label": "security", "required": false, "type": "string", "javaType": "java.lang.String", "deprecated": false, "deprecationNote": "", "autowired": false, "secret": false, "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Username for authentication with the HiveMQ broker." } + "mqttVersion": { "index": 4, "kind": "parameter", "displayName": "Mqtt Version", "group": "common", "label": "", "required": false, "type": "enum", "javaType": "com.hivemq.client.mqtt.MqttVersion", "enum": [ "MQTT_3_1_1", "MQTT_5_0" ], "deprecated": false, "autowired": false, "secret": false, "defaultValue": "MQTT_5_0", "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "The MQTT protocol version to use [...] + "port": { "index": 5, "kind": "parameter", "displayName": "Port", "group": "common", "label": "", "required": false, "type": "integer", "javaType": "int", "deprecated": false, "autowired": false, "secret": false, "defaultValue": 1883, "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Port number of the HiveMQ MQTT broker." }, + "qos": { "index": 6, "kind": "parameter", "displayName": "Qos", "group": "common", "label": "", "required": false, "type": "enum", "javaType": "com.hivemq.client.mqtt.datatypes.MqttQos", "enum": [ "AT_MOST_ONCE", "AT_LEAST_ONCE", "EXACTLY_ONCE" ], "deprecated": false, "autowired": false, "secret": false, "defaultValue": "AT_LEAST_ONCE", "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Default Quality [...] + "retained": { "index": 7, "kind": "parameter", "displayName": "Retained", "group": "common", "label": "", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "defaultValue": false, "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Whether published messages should be retained by the MQTT broker." }, + "bridgeErrorHandler": { "index": 8, "kind": "parameter", "displayName": "Bridge Error Handler", "group": "consumer (advanced)", "label": "consumer,advanced", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "defaultValue": false, "description": "Allows for bridging the consumer to the Camel routing Error Handler, which mean any exceptions (if possible) occurred while the Camel consumer is trying to pickup incoming [...] + "exceptionHandler": { "index": 9, "kind": "parameter", "displayName": "Exception Handler", "group": "consumer (advanced)", "label": "consumer,advanced", "required": false, "type": "object", "javaType": "org.apache.camel.spi.ExceptionHandler", "optionalPrefix": "consumer.", "deprecated": false, "autowired": false, "secret": false, "description": "To let the consumer use a custom ExceptionHandler. Notice if the option bridgeErrorHandler is enabled then this option is not in use. By def [...] + "exchangePattern": { "index": 10, "kind": "parameter", "displayName": "Exchange Pattern", "group": "consumer (advanced)", "label": "consumer,advanced", "required": false, "type": "enum", "javaType": "org.apache.camel.ExchangePattern", "enum": [ "InOnly", "InOut" ], "deprecated": false, "autowired": false, "secret": false, "description": "Sets the exchange pattern when the consumer creates an exchange." }, + "lazyStartProducer": { "index": 11, "kind": "parameter", "displayName": "Lazy Start Producer", "group": "producer (advanced)", "label": "producer,advanced", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "defaultValue": false, "description": "Whether the producer should be started lazy (on the first message). By starting lazy you can use this to allow CamelContext and routes to startup in situations where a produ [...] + "password": { "index": 12, "kind": "parameter", "displayName": "Password", "group": "security", "label": "security", "required": false, "type": "string", "javaType": "java.lang.String", "deprecated": false, "deprecationNote": "", "autowired": false, "secret": true, "security": "secret", "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Password for authentication with the HiveMQ broker." }, + "ssl": { "index": 13, "kind": "parameter", "displayName": "Ssl", "group": "security", "label": "security", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "security": "insecure:ssl", "insecureValue": "false", "defaultValue": false, "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Whether to enable SSL\/TLS encryption for the broke [...] + "username": { "index": 14, "kind": "parameter", "displayName": "Username", "group": "security", "label": "security", "required": false, "type": "string", "javaType": "java.lang.String", "deprecated": false, "deprecationNote": "", "autowired": false, "secret": false, "configurationClass": "org.apache.camel.component.hivemq.HiveMQConfiguration", "configurationField": "configuration", "description": "Username for authentication with the HiveMQ broker." } } } diff --git a/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQClientAdapter.java b/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQClientAdapter.java new file mode 100644 index 000000000000..17a504fa603b --- /dev/null +++ b/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQClientAdapter.java @@ -0,0 +1,63 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.camel.component.hivemq; + +import java.util.Optional; +import java.util.concurrent.CompletableFuture; +import java.util.function.Consumer; + +import com.hivemq.client.mqtt.datatypes.MqttQos; + +/** + * Adapts an underlying HiveMQ MQTT client, hiding whether it speaks MQTT 3.1.1 or MQTT 5 behind a single API so that + * {@link HiveMQEndpoint}, {@link HiveMQConsumer} and {@link HiveMQProducer} do not need to know which protocol version + * is in use. + */ +interface HiveMQClientAdapter { + + /** + * Sends the CONNECT packet. The returned future completes when the CONNACK is received. + */ + CompletableFuture<?> connect(boolean cleanStart); + + /** + * Cancels automatic reconnect and disconnects if currently connected. Safe to call from any client state. + */ + void stop(); + + boolean isConnected(); + + /** + * Whether the client is connected or automatically retrying a connection, i.e. not fully settled into a + * disconnected state. Used to verify that {@link #stop()} really stopped the automatic-reconnect loop. + */ + boolean isConnectedOrReconnecting(); + + CompletableFuture<?> subscribe(String topicFilter, MqttQos qos, Consumer<HiveMQMessage> callback); + + CompletableFuture<?> unsubscribe(String topicFilter); + + CompletableFuture<?> publish(String topic, byte[] payload, MqttQos qos, boolean retained); + + /** + * Gives direct access to the underlying HiveMQ MQTT Client library client (e.g. {@code Mqtt5AsyncClient} or + * {@code Mqtt3AsyncClient}), for callers who need functionality this adapter does not expose. Empty if the + * underlying client is not an instance of {@code clazz}, e.g. because this connection uses the other protocol + * version. + */ + <T> Optional<T> getClient(Class<T> clazz); +} diff --git a/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQConfiguration.java b/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQConfiguration.java index 47bd662534e4..158c7fff9541 100644 --- a/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQConfiguration.java +++ b/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQConfiguration.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.apache.camel.RuntimeCamelException; import org.apache.camel.spi.Metadata; @@ -31,6 +32,12 @@ public class HiveMQConfiguration implements Cloneable { @UriParam(defaultValue = HiveMQConstants.DEFAULT_HOST) private String host = HiveMQConstants.DEFAULT_HOST; + /** + * The MQTT protocol version to use when connecting to the broker. + */ + @UriParam(defaultValue = "MQTT_5_0") + private MqttVersion mqttVersion = MqttVersion.MQTT_5_0; + /** * Port number of the HiveMQ MQTT broker. */ @@ -56,7 +63,8 @@ public class HiveMQConfiguration implements Cloneable { private boolean retained; /** - * Whether to initiate a clean start (MQTT 5) upon connecting to the broker. + * Whether to initiate a clean session upon connecting to the broker (called "clean session" in MQTT 3.1.1 and + * "clean start" in MQTT 5). */ @UriParam(defaultValue = "true") private boolean cleanStart = true; @@ -89,6 +97,14 @@ public class HiveMQConfiguration implements Cloneable { this.host = host; } + public MqttVersion getMqttVersion() { + return mqttVersion; + } + + public void setMqttVersion(MqttVersion mqttVersion) { + this.mqttVersion = mqttVersion; + } + public int getPort() { return port; } diff --git a/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQConsumer.java b/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQConsumer.java index 7d2d378e3b04..62ce0d50b99e 100644 --- a/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQConsumer.java +++ b/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQConsumer.java @@ -16,12 +16,11 @@ */ package org.apache.camel.component.hivemq; +import java.util.Optional; import java.util.concurrent.ExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; -import com.hivemq.client.mqtt.mqtt5.Mqtt5AsyncClient; -import com.hivemq.client.mqtt.mqtt5.message.publish.Mqtt5Publish; import org.apache.camel.Exchange; import org.apache.camel.Processor; import org.apache.camel.support.DefaultConsumer; @@ -32,7 +31,7 @@ public class HiveMQConsumer extends DefaultConsumer { private static final Logger LOG = LoggerFactory.getLogger(HiveMQConsumer.class); private final HiveMQEndpoint endpoint; - private Mqtt5AsyncClient client; + private HiveMQClientAdapter client; private ExecutorService executor; public HiveMQConsumer(HiveMQEndpoint endpoint, Processor processor) { @@ -47,27 +46,25 @@ public class HiveMQConsumer extends DefaultConsumer { client = endpoint.createClient(); endpoint.connect(client); - client.subscribeWith() - .topicFilter(endpoint.getTopic()) - .qos(endpoint.getConfiguration().getQos()) - .callback(this::onMessage) - .send() + client.subscribe(endpoint.getTopic(), endpoint.getConfiguration().getQos(), this::onMessage) .orTimeout(HiveMQConstants.DEFAULT_CONNECT_TIMEOUT_SECONDS, TimeUnit.SECONDS) .join(); } @Override protected void doStop() throws Exception { - if (client != null && client.getState().isConnected()) { + if (client != null && client.isConnected()) { try { - client.unsubscribeWith().topicFilter(endpoint.getTopic()).send() + client.unsubscribe(endpoint.getTopic()) .orTimeout(5, TimeUnit.SECONDS).join(); } catch (Exception e) { // Best-effort unsubscribe before cancelling reconnect / disconnect LOG.debug("Failed to unsubscribe from topic {} during shutdown", endpoint.getTopic(), e); } } - endpoint.stopClient(client); + if (client != null) { + client.stop(); + } client = null; if (executor != null) { endpoint.getCamelContext().getExecutorServiceManager().shutdownNow(executor); @@ -76,12 +73,22 @@ public class HiveMQConsumer extends DefaultConsumer { super.doStop(); } - private void onMessage(Mqtt5Publish publish) { + /** + * Gives direct access to the underlying HiveMQ MQTT Client library client (e.g. {@code Mqtt5AsyncClient} or + * {@code Mqtt3AsyncClient}, depending on the {@code mqttVersion} this consumer is connected with) for use cases + * this component does not cover. Empty before the consumer has started, or if {@code clazz} does not match the + * protocol version in use. + */ + public <T> Optional<T> getClient(Class<T> clazz) { + return client == null ? Optional.empty() : client.getClient(clazz); + } + + private void onMessage(HiveMQMessage message) { Exchange exchange = createExchange(false); - exchange.getIn().setBody(publish.getPayloadAsBytes()); - exchange.getIn().setHeader(HiveMQConstants.MQTT_TOPIC, publish.getTopic().toString()); - exchange.getIn().setHeader(HiveMQConstants.MQTT_QOS, publish.getQos()); - exchange.getIn().setHeader(HiveMQConstants.MQTT_RETAINED, publish.isRetain()); + exchange.getIn().setBody(message.payload()); + exchange.getIn().setHeader(HiveMQConstants.MQTT_TOPIC, message.topic()); + exchange.getIn().setHeader(HiveMQConstants.MQTT_QOS, message.qos()); + exchange.getIn().setHeader(HiveMQConstants.MQTT_RETAINED, message.retained()); ExecutorService worker = executor; if (worker == null) { diff --git a/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQEndpoint.java b/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQEndpoint.java index d20c4b02f3b6..9ec1f0bf5f1f 100644 --- a/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQEndpoint.java +++ b/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQEndpoint.java @@ -16,21 +16,10 @@ */ package org.apache.camel.component.hivemq; -import java.nio.charset.StandardCharsets; -import java.util.Map; import java.util.concurrent.CompletionException; -import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; -import java.util.concurrent.atomic.AtomicBoolean; -import java.util.concurrent.atomic.AtomicReference; - -import com.hivemq.client.mqtt.MqttClient; -import com.hivemq.client.mqtt.MqttClientState; -import com.hivemq.client.mqtt.mqtt5.Mqtt5AsyncClient; -import com.hivemq.client.mqtt.mqtt5.Mqtt5ClientBuilder; -import com.hivemq.client.mqtt.mqtt5.message.auth.Mqtt5SimpleAuth; -import com.hivemq.client.mqtt.mqtt5.message.auth.Mqtt5SimpleAuthBuilder; + import org.apache.camel.Category; import org.apache.camel.Consumer; import org.apache.camel.Processor; @@ -42,15 +31,11 @@ import org.apache.camel.spi.UriEndpoint; import org.apache.camel.spi.UriParam; import org.apache.camel.spi.UriPath; import org.apache.camel.support.DefaultEndpoint; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; @UriEndpoint(firstVersion = "4.23.0", scheme = "hivemq", title = "HiveMQ", syntax = "hivemq:topic", category = { Category.MESSAGING, Category.IOT }, headersClass = HiveMQConstants.class) public class HiveMQEndpoint extends DefaultEndpoint implements EndpointServiceLocation { - private static final Logger LOG = LoggerFactory.getLogger(HiveMQEndpoint.class); - /** * The MQTT topic name or pattern to subscribe to or publish on. */ @@ -65,8 +50,6 @@ public class HiveMQEndpoint extends DefaultEndpoint implements EndpointServiceLo @Metadata(description = "To use a custom HiveMQConfiguration") private HiveMQConfiguration configuration; - private final Map<Mqtt5AsyncClient, AtomicBoolean> reconnectCancellations = new ConcurrentHashMap<>(); - public HiveMQEndpoint(String uri, HiveMQComponent component, HiveMQConfiguration configuration, String topic) { super(uri, component); this.configuration = configuration; @@ -85,93 +68,24 @@ public class HiveMQEndpoint extends DefaultEndpoint implements EndpointServiceLo return consumer; } - public Mqtt5AsyncClient createClient() { - AtomicBoolean cancelReconnect = new AtomicBoolean(); - AtomicReference<Mqtt5AsyncClient> clientRef = new AtomicReference<>(); - Mqtt5ClientBuilder builder = MqttClient.builder() - .serverHost(configuration.getHost()) - .serverPort(configuration.getPort()) - .automaticReconnectWithDefaultConfig() - .addDisconnectedListener(context -> { - // Initial connect() does not complete while auto-reconnect keeps retrying (HiveMQ #302). - // Also honour an explicit stop so DISCONNECTED_RECONNECT / CONNECTING_RECONNECT are cancelled. - if (cancelReconnect.get() || context.getClientConfig().getState() == MqttClientState.CONNECTING) { - context.getReconnector().reconnect(false); - } - }) - .addConnectedListener(context -> { - // HiveMQ schedules reconnect after listeners return; cancelReconnect cannot abort that delay. - // If a reconnect succeeds after Camel stop, disconnect immediately (USER source skips auto-reconnect). - if (cancelReconnect.get()) { - Mqtt5AsyncClient started = clientRef.get(); - if (started != null && started.getState().isConnected()) { - try { - started.disconnect(); - } catch (Exception e) { - // Already disconnecting or not connected - } - } - } - }) - .useMqttVersion5(); - - if (configuration.getClientId() != null) { - builder.identifier(configuration.getClientId()); - } - - if (configuration.isSsl()) { - builder.sslWithDefaultConfig(); - } - - if (configuration.getUsername() != null) { - Mqtt5SimpleAuthBuilder.Complete authBuilder - = Mqtt5SimpleAuth.builder().username(configuration.getUsername()); - if (configuration.getPassword() != null) { - authBuilder.password(configuration.getPassword().getBytes(StandardCharsets.UTF_8)); - } - builder.simpleAuth(authBuilder.build()); - } - - Mqtt5AsyncClient client = builder.buildAsync(); - clientRef.set(client); - reconnectCancellations.put(client, cancelReconnect); - return client; + HiveMQClientAdapter createClient() { + return switch (configuration.getMqttVersion()) { + case MQTT_3_1_1 -> new Mqtt3ClientAdapter(configuration); + case MQTT_5_0 -> new Mqtt5ClientAdapter(configuration); + }; } - public void connect(Mqtt5AsyncClient client) { + void connect(HiveMQClientAdapter client) { try { - client.connectWith() - .cleanStart(configuration.isCleanStart()) - .send() + client.connect(configuration.isCleanStart()) .orTimeout(HiveMQConstants.DEFAULT_CONNECT_TIMEOUT_SECONDS, TimeUnit.SECONDS) .join(); } catch (CompletionException e) { - stopClient(client); + client.stop(); throw unwrapConnectFailure(e); } } - /** - * Stops automatic reconnect and disconnects if currently connected. Safe to call from any client state. - */ - public void stopClient(Mqtt5AsyncClient client) { - if (client == null) { - return; - } - AtomicBoolean cancelReconnect = reconnectCancellations.remove(client); - if (cancelReconnect != null) { - cancelReconnect.set(true); - } - try { - if (client.getState().isConnected()) { - client.disconnect().orTimeout(5, TimeUnit.SECONDS).join(); - } - } catch (Exception e) { - // Not connected, already disconnecting, or reconnecting: the disconnected listener cancels reconnect. - LOG.debug("Failed to disconnect HiveMQ client during shutdown", e); - } - } - private static RuntimeCamelException unwrapConnectFailure(CompletionException e) { Throwable cause = e.getCause() != null ? e.getCause() : e; if (cause instanceof TimeoutException) { diff --git a/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQMessage.java b/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQMessage.java new file mode 100644 index 000000000000..0500259e8d76 --- /dev/null +++ b/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQMessage.java @@ -0,0 +1,25 @@ +/* + * 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; + +/** + * A received MQTT publish, independent of the MQTT protocol version (3.1.1 or 5) it was received over. + */ +record HiveMQMessage(String topic, byte[] payload, MqttQos qos, boolean retained) { +} diff --git a/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQProducer.java b/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQProducer.java index 6d8bba768156..f233bb6e7398 100644 --- a/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQProducer.java +++ b/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQProducer.java @@ -16,9 +16,9 @@ */ package org.apache.camel.component.hivemq; +import java.util.Optional; + import com.hivemq.client.mqtt.datatypes.MqttQos; -import com.hivemq.client.mqtt.mqtt5.Mqtt5AsyncClient; -import com.hivemq.client.mqtt.mqtt5.message.publish.Mqtt5Publish; import org.apache.camel.AsyncCallback; import org.apache.camel.Exchange; import org.apache.camel.support.DefaultAsyncProducer; @@ -26,7 +26,7 @@ import org.apache.camel.support.DefaultAsyncProducer; public class HiveMQProducer extends DefaultAsyncProducer { private final HiveMQEndpoint endpoint; - private Mqtt5AsyncClient client; + private HiveMQClientAdapter client; public HiveMQProducer(HiveMQEndpoint endpoint) { super(endpoint); @@ -42,11 +42,23 @@ public class HiveMQProducer extends DefaultAsyncProducer { @Override protected void doStop() throws Exception { - endpoint.stopClient(client); + if (client != null) { + client.stop(); + } client = null; super.doStop(); } + /** + * Gives direct access to the underlying HiveMQ MQTT Client library client (e.g. {@code Mqtt5AsyncClient} or + * {@code Mqtt3AsyncClient}, depending on the {@code mqttVersion} this producer is connected with) for use cases + * this component does not cover. Empty before the producer has started, or if {@code clazz} does not match the + * protocol version in use. + */ + public <T> Optional<T> getClient(Class<T> clazz) { + return client == null ? Optional.empty() : client.getClient(clazz); + } + @Override public boolean process(Exchange exchange, AsyncCallback callback) { String targetTopic = exchange.getIn().getHeader(HiveMQConstants.OVERRIDE_TOPIC, endpoint.getTopic(), String.class); @@ -59,12 +71,7 @@ public class HiveMQProducer extends DefaultAsyncProducer { payload = new byte[0]; } - client.publish(Mqtt5Publish.builder() - .topic(targetTopic) - .qos(qos) - .retain(retained) - .payload(payload) - .build()) + client.publish(targetTopic, payload, qos, retained) .whenComplete((publishResult, throwable) -> { if (throwable != null) { exchange.setException(throwable); diff --git a/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/Mqtt3ClientAdapter.java b/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/Mqtt3ClientAdapter.java new file mode 100644 index 000000000000..648ca916d92e --- /dev/null +++ b/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/Mqtt3ClientAdapter.java @@ -0,0 +1,156 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.camel.component.hivemq; + +import java.nio.charset.StandardCharsets; +import java.util.Optional; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Consumer; + +import com.hivemq.client.mqtt.MqttClient; +import com.hivemq.client.mqtt.MqttClientState; +import com.hivemq.client.mqtt.datatypes.MqttQos; +import com.hivemq.client.mqtt.mqtt3.Mqtt3AsyncClient; +import com.hivemq.client.mqtt.mqtt3.Mqtt3ClientBuilder; +import com.hivemq.client.mqtt.mqtt3.message.auth.Mqtt3SimpleAuth; +import com.hivemq.client.mqtt.mqtt3.message.auth.Mqtt3SimpleAuthBuilder; +import com.hivemq.client.mqtt.mqtt3.message.publish.Mqtt3Publish; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +final class Mqtt3ClientAdapter implements HiveMQClientAdapter { + + private static final Logger LOG = LoggerFactory.getLogger(Mqtt3ClientAdapter.class); + + private final Mqtt3AsyncClient client; + private final AtomicBoolean cancelReconnect = new AtomicBoolean(); + + Mqtt3ClientAdapter(HiveMQConfiguration configuration) { + AtomicReference<Mqtt3AsyncClient> clientRef = new AtomicReference<>(); + Mqtt3ClientBuilder builder = MqttClient.builder() + .serverHost(configuration.getHost()) + .serverPort(configuration.getPort()) + .automaticReconnectWithDefaultConfig() + .addDisconnectedListener(context -> { + // Initial connect() does not complete while auto-reconnect keeps retrying (HiveMQ #302). + // Also honour an explicit stop so DISCONNECTED_RECONNECT / CONNECTING_RECONNECT are cancelled. + if (cancelReconnect.get() || context.getClientConfig().getState() == MqttClientState.CONNECTING) { + context.getReconnector().reconnect(false); + } + }) + .addConnectedListener(context -> { + // HiveMQ schedules reconnect after listeners return; cancelReconnect cannot abort that delay. + // If a reconnect succeeds after Camel stop, disconnect immediately (USER source skips auto-reconnect). + if (cancelReconnect.get()) { + Mqtt3AsyncClient started = clientRef.get(); + if (started != null && started.getState().isConnected()) { + try { + started.disconnect(); + } catch (Exception e) { + // Already disconnecting or not connected + } + } + } + }) + .useMqttVersion3(); + + if (configuration.getClientId() != null) { + builder.identifier(configuration.getClientId()); + } + + if (configuration.isSsl()) { + builder.sslWithDefaultConfig(); + } + + if (configuration.getUsername() != null) { + Mqtt3SimpleAuthBuilder.Complete authBuilder + = Mqtt3SimpleAuth.builder().username(configuration.getUsername()); + if (configuration.getPassword() != null) { + authBuilder.password(configuration.getPassword().getBytes(StandardCharsets.UTF_8)); + } + builder.simpleAuth(authBuilder.build()); + } + + client = builder.buildAsync(); + clientRef.set(client); + } + + @Override + public CompletableFuture<?> connect(boolean cleanStart) { + return client.connectWith().cleanSession(cleanStart).send(); + } + + @Override + public void stop() { + cancelReconnect.set(true); + try { + if (client.getState().isConnected()) { + client.disconnect().orTimeout(5, TimeUnit.SECONDS).join(); + } + } catch (Exception e) { + // Not connected, already disconnecting, or reconnecting: the disconnected listener cancels reconnect. + LOG.debug("Failed to disconnect HiveMQ MQTT 3.1.1 client during shutdown", e); + } + } + + @Override + public boolean isConnected() { + return client.getState().isConnected(); + } + + @Override + public boolean isConnectedOrReconnecting() { + return client.getState().isConnectedOrReconnect(); + } + + @Override + public CompletableFuture<?> subscribe(String topicFilter, MqttQos qos, Consumer<HiveMQMessage> callback) { + return client.subscribeWith() + .topicFilter(topicFilter) + .qos(qos) + .callback(publish -> callback.accept(toMessage(publish))) + .send(); + } + + @Override + public CompletableFuture<?> unsubscribe(String topicFilter) { + return client.unsubscribeWith().topicFilter(topicFilter).send(); + } + + @Override + public CompletableFuture<?> publish(String topic, byte[] payload, MqttQos qos, boolean retained) { + return client.publish(Mqtt3Publish.builder() + .topic(topic) + .qos(qos) + .retain(retained) + .payload(payload) + .build()); + } + + private static HiveMQMessage toMessage(Mqtt3Publish publish) { + return new HiveMQMessage( + publish.getTopic().toString(), publish.getPayloadAsBytes(), publish.getQos(), publish.isRetain()); + } + + @Override + public <T> Optional<T> getClient(Class<T> clazz) { + return clazz.isInstance(client) ? Optional.of(clazz.cast(client)) : Optional.empty(); + } +} diff --git a/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQEndpoint.java b/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/Mqtt5ClientAdapter.java similarity index 50% copy from components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQEndpoint.java copy to components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/Mqtt5ClientAdapter.java index d20c4b02f3b6..1be32bb431e0 100644 --- a/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/HiveMQEndpoint.java +++ b/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/Mqtt5ClientAdapter.java @@ -17,76 +17,32 @@ package org.apache.camel.component.hivemq; import java.nio.charset.StandardCharsets; -import java.util.Map; -import java.util.concurrent.CompletionException; -import java.util.concurrent.ConcurrentHashMap; +import java.util.Optional; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; -import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Consumer; import com.hivemq.client.mqtt.MqttClient; import com.hivemq.client.mqtt.MqttClientState; +import com.hivemq.client.mqtt.datatypes.MqttQos; import com.hivemq.client.mqtt.mqtt5.Mqtt5AsyncClient; import com.hivemq.client.mqtt.mqtt5.Mqtt5ClientBuilder; import com.hivemq.client.mqtt.mqtt5.message.auth.Mqtt5SimpleAuth; import com.hivemq.client.mqtt.mqtt5.message.auth.Mqtt5SimpleAuthBuilder; -import org.apache.camel.Category; -import org.apache.camel.Consumer; -import org.apache.camel.Processor; -import org.apache.camel.Producer; -import org.apache.camel.RuntimeCamelException; -import org.apache.camel.spi.EndpointServiceLocation; -import org.apache.camel.spi.Metadata; -import org.apache.camel.spi.UriEndpoint; -import org.apache.camel.spi.UriParam; -import org.apache.camel.spi.UriPath; -import org.apache.camel.support.DefaultEndpoint; +import com.hivemq.client.mqtt.mqtt5.message.publish.Mqtt5Publish; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -@UriEndpoint(firstVersion = "4.23.0", scheme = "hivemq", title = "HiveMQ", syntax = "hivemq:topic", - category = { Category.MESSAGING, Category.IOT }, headersClass = HiveMQConstants.class) -public class HiveMQEndpoint extends DefaultEndpoint implements EndpointServiceLocation { +final class Mqtt5ClientAdapter implements HiveMQClientAdapter { - private static final Logger LOG = LoggerFactory.getLogger(HiveMQEndpoint.class); + private static final Logger LOG = LoggerFactory.getLogger(Mqtt5ClientAdapter.class); - /** - * The MQTT topic name or pattern to subscribe to or publish on. - */ - @UriPath - @Metadata(required = true) - private String topic; + private final Mqtt5AsyncClient client; + private final AtomicBoolean cancelReconnect = new AtomicBoolean(); - /** - * The HiveMQ component configuration options. - */ - @UriParam - @Metadata(description = "To use a custom HiveMQConfiguration") - private HiveMQConfiguration configuration; - - private final Map<Mqtt5AsyncClient, AtomicBoolean> reconnectCancellations = new ConcurrentHashMap<>(); - - public HiveMQEndpoint(String uri, HiveMQComponent component, HiveMQConfiguration configuration, String topic) { - super(uri, component); - this.configuration = configuration; - this.topic = topic; - } - - @Override - public Producer createProducer() throws Exception { - return new HiveMQProducer(this); - } - - @Override - public Consumer createConsumer(Processor processor) throws Exception { - HiveMQConsumer consumer = new HiveMQConsumer(this, processor); - configureConsumer(consumer); - return consumer; - } - - public Mqtt5AsyncClient createClient() { - AtomicBoolean cancelReconnect = new AtomicBoolean(); + Mqtt5ClientAdapter(HiveMQConfiguration configuration) { AtomicReference<Mqtt5AsyncClient> clientRef = new AtomicReference<>(); Mqtt5ClientBuilder builder = MqttClient.builder() .serverHost(configuration.getHost()) @@ -132,77 +88,69 @@ public class HiveMQEndpoint extends DefaultEndpoint implements EndpointServiceLo builder.simpleAuth(authBuilder.build()); } - Mqtt5AsyncClient client = builder.buildAsync(); + client = builder.buildAsync(); clientRef.set(client); - reconnectCancellations.put(client, cancelReconnect); - return client; } - public void connect(Mqtt5AsyncClient client) { - try { - client.connectWith() - .cleanStart(configuration.isCleanStart()) - .send() - .orTimeout(HiveMQConstants.DEFAULT_CONNECT_TIMEOUT_SECONDS, TimeUnit.SECONDS) - .join(); - } catch (CompletionException e) { - stopClient(client); - throw unwrapConnectFailure(e); - } + @Override + public CompletableFuture<?> connect(boolean cleanStart) { + return client.connectWith().cleanStart(cleanStart).send(); } - /** - * Stops automatic reconnect and disconnects if currently connected. Safe to call from any client state. - */ - public void stopClient(Mqtt5AsyncClient client) { - if (client == null) { - return; - } - AtomicBoolean cancelReconnect = reconnectCancellations.remove(client); - if (cancelReconnect != null) { - cancelReconnect.set(true); - } + @Override + public void stop() { + cancelReconnect.set(true); try { if (client.getState().isConnected()) { client.disconnect().orTimeout(5, TimeUnit.SECONDS).join(); } } catch (Exception e) { // Not connected, already disconnecting, or reconnecting: the disconnected listener cancels reconnect. - LOG.debug("Failed to disconnect HiveMQ client during shutdown", e); + LOG.debug("Failed to disconnect HiveMQ MQTT 5 client during shutdown", e); } } - private static RuntimeCamelException unwrapConnectFailure(CompletionException e) { - Throwable cause = e.getCause() != null ? e.getCause() : e; - if (cause instanceof TimeoutException) { - return new RuntimeCamelException("Timed out connecting to the HiveMQ broker", cause); - } - return new RuntimeCamelException("Failed to connect to the HiveMQ broker", cause); + @Override + public boolean isConnected() { + return client.getState().isConnected(); } - public String getTopic() { - return topic; + @Override + public boolean isConnectedOrReconnecting() { + return client.getState().isConnectedOrReconnect(); } - public void setTopic(String topic) { - this.topic = topic; + @Override + public CompletableFuture<?> subscribe(String topicFilter, MqttQos qos, Consumer<HiveMQMessage> callback) { + return client.subscribeWith() + .topicFilter(topicFilter) + .qos(qos) + .callback(publish -> callback.accept(toMessage(publish))) + .send(); } - public HiveMQConfiguration getConfiguration() { - return configuration; + @Override + public CompletableFuture<?> unsubscribe(String topicFilter) { + return client.unsubscribeWith().topicFilter(topicFilter).send(); } - public void setConfiguration(HiveMQConfiguration configuration) { - this.configuration = configuration; + @Override + public CompletableFuture<?> publish(String topic, byte[] payload, MqttQos qos, boolean retained) { + return client.publish(Mqtt5Publish.builder() + .topic(topic) + .qos(qos) + .retain(retained) + .payload(payload) + .build()); } - @Override - public String getServiceUrl() { - return configuration.getHost() + ":" + configuration.getPort(); + private static HiveMQMessage toMessage(Mqtt5Publish publish) { + return new HiveMQMessage( + publish.getTopic().toString(), publish.getPayloadAsBytes(), publish.getQos(), publish.isRetain()); } @Override - public String getServiceProtocol() { - return "mqtt"; + public <T> Optional<T> getClient(Class<T> clazz) { + return clazz.isInstance(client) ? Optional.of(clazz.cast(client)) : Optional.empty(); } } diff --git a/dsl/camel-componentdsl/src/generated/java/org/apache/camel/builder/component/dsl/HivemqComponentBuilderFactory.java b/dsl/camel-componentdsl/src/generated/java/org/apache/camel/builder/component/dsl/HivemqComponentBuilderFactory.java index 01a5ac2a0b7e..b13da7ad1ec2 100644 --- a/dsl/camel-componentdsl/src/generated/java/org/apache/camel/builder/component/dsl/HivemqComponentBuilderFactory.java +++ b/dsl/camel-componentdsl/src/generated/java/org/apache/camel/builder/component/dsl/HivemqComponentBuilderFactory.java @@ -52,8 +52,8 @@ public interface HivemqComponentBuilderFactory { /** - * Whether to initiate a clean start (MQTT 5) upon connecting to the - * broker. + * Whether to initiate a clean session upon connecting to the broker + * (called clean session in MQTT 3.1.1 and clean start in MQTT 5). * * The option is a: <code>boolean</code> type. * @@ -117,6 +117,24 @@ public interface HivemqComponentBuilderFactory { } + /** + * The MQTT protocol version to use when connecting to the broker. + * + * The option is a: + * <code>com.hivemq.client.mqtt.MqttVersion</code> type. + * + * Default: MQTT_5_0 + * Group: common + * + * @param mqttVersion the value to set + * @return the dsl builder + */ + default HivemqComponentBuilder mqttVersion(com.hivemq.client.mqtt.MqttVersion mqttVersion) { + doSetProperty("mqttVersion", mqttVersion); + return this; + } + + /** * Port number of the HiveMQ MQTT broker. * @@ -315,6 +333,7 @@ public interface HivemqComponentBuilderFactory { case "clientId": getOrCreateConfiguration((HiveMQComponent) component).setClientId((java.lang.String) value); return true; case "configuration": ((HiveMQComponent) component).setConfiguration((org.apache.camel.component.hivemq.HiveMQConfiguration) value); return true; case "host": getOrCreateConfiguration((HiveMQComponent) component).setHost((java.lang.String) value); return true; + case "mqttVersion": getOrCreateConfiguration((HiveMQComponent) component).setMqttVersion((com.hivemq.client.mqtt.MqttVersion) value); return true; case "port": getOrCreateConfiguration((HiveMQComponent) component).setPort((int) value); return true; case "qos": getOrCreateConfiguration((HiveMQComponent) component).setQos((com.hivemq.client.mqtt.datatypes.MqttQos) value); return true; case "retained": getOrCreateConfiguration((HiveMQComponent) component).setRetained((boolean) value); return true; diff --git a/dsl/camel-endpointdsl/src/generated/java/org/apache/camel/builder/endpoint/dsl/HiveMQEndpointBuilderFactory.java b/dsl/camel-endpointdsl/src/generated/java/org/apache/camel/builder/endpoint/dsl/HiveMQEndpointBuilderFactory.java index 80a520c3918c..1ebea7ecc53d 100644 --- a/dsl/camel-endpointdsl/src/generated/java/org/apache/camel/builder/endpoint/dsl/HiveMQEndpointBuilderFactory.java +++ b/dsl/camel-endpointdsl/src/generated/java/org/apache/camel/builder/endpoint/dsl/HiveMQEndpointBuilderFactory.java @@ -44,8 +44,8 @@ public interface HiveMQEndpointBuilderFactory { return (AdvancedHiveMQEndpointConsumerBuilder) this; } /** - * Whether to initiate a clean start (MQTT 5) upon connecting to the - * broker. + * Whether to initiate a clean session upon connecting to the broker + * (called clean session in MQTT 3.1.1 and clean start in MQTT 5). * * The option is a: <code>boolean</code> type. * @@ -60,8 +60,8 @@ public interface HiveMQEndpointBuilderFactory { return this; } /** - * Whether to initiate a clean start (MQTT 5) upon connecting to the - * broker. + * Whether to initiate a clean session upon connecting to the broker + * (called clean session in MQTT 3.1.1 and clean start in MQTT 5). * * The option will be converted to a <code>boolean</code> type. * @@ -104,6 +104,38 @@ public interface HiveMQEndpointBuilderFactory { doSetProperty("host", host); return this; } + /** + * The MQTT protocol version to use when connecting to the broker. + * + * The option is a: <code>com.hivemq.client.mqtt.MqttVersion</code> + * type. + * + * Default: MQTT_5_0 + * Group: common + * + * @param mqttVersion the value to set + * @return the dsl builder + */ + default HiveMQEndpointConsumerBuilder mqttVersion(com.hivemq.client.mqtt.MqttVersion mqttVersion) { + doSetProperty("mqttVersion", mqttVersion); + return this; + } + /** + * The MQTT protocol version to use when connecting to the broker. + * + * The option will be converted to a + * <code>com.hivemq.client.mqtt.MqttVersion</code> type. + * + * Default: MQTT_5_0 + * Group: common + * + * @param mqttVersion the value to set + * @return the dsl builder + */ + default HiveMQEndpointConsumerBuilder mqttVersion(String mqttVersion) { + doSetProperty("mqttVersion", mqttVersion); + return this; + } /** * Port number of the HiveMQ MQTT broker. * @@ -395,8 +427,8 @@ public interface HiveMQEndpointBuilderFactory { } /** - * Whether to initiate a clean start (MQTT 5) upon connecting to the - * broker. + * Whether to initiate a clean session upon connecting to the broker + * (called clean session in MQTT 3.1.1 and clean start in MQTT 5). * * The option is a: <code>boolean</code> type. * @@ -411,8 +443,8 @@ public interface HiveMQEndpointBuilderFactory { return this; } /** - * Whether to initiate a clean start (MQTT 5) upon connecting to the - * broker. + * Whether to initiate a clean session upon connecting to the broker + * (called clean session in MQTT 3.1.1 and clean start in MQTT 5). * * The option will be converted to a <code>boolean</code> type. * @@ -455,6 +487,38 @@ public interface HiveMQEndpointBuilderFactory { doSetProperty("host", host); return this; } + /** + * The MQTT protocol version to use when connecting to the broker. + * + * The option is a: <code>com.hivemq.client.mqtt.MqttVersion</code> + * type. + * + * Default: MQTT_5_0 + * Group: common + * + * @param mqttVersion the value to set + * @return the dsl builder + */ + default HiveMQEndpointProducerBuilder mqttVersion(com.hivemq.client.mqtt.MqttVersion mqttVersion) { + doSetProperty("mqttVersion", mqttVersion); + return this; + } + /** + * The MQTT protocol version to use when connecting to the broker. + * + * The option will be converted to a + * <code>com.hivemq.client.mqtt.MqttVersion</code> type. + * + * Default: MQTT_5_0 + * Group: common + * + * @param mqttVersion the value to set + * @return the dsl builder + */ + default HiveMQEndpointProducerBuilder mqttVersion(String mqttVersion) { + doSetProperty("mqttVersion", mqttVersion); + return this; + } /** * Port number of the HiveMQ MQTT broker. * @@ -675,8 +739,8 @@ public interface HiveMQEndpointBuilderFactory { } /** - * Whether to initiate a clean start (MQTT 5) upon connecting to the - * broker. + * Whether to initiate a clean session upon connecting to the broker + * (called clean session in MQTT 3.1.1 and clean start in MQTT 5). * * The option is a: <code>boolean</code> type. * @@ -691,8 +755,8 @@ public interface HiveMQEndpointBuilderFactory { return this; } /** - * Whether to initiate a clean start (MQTT 5) upon connecting to the - * broker. + * Whether to initiate a clean session upon connecting to the broker + * (called clean session in MQTT 3.1.1 and clean start in MQTT 5). * * The option will be converted to a <code>boolean</code> type. * @@ -735,6 +799,38 @@ public interface HiveMQEndpointBuilderFactory { doSetProperty("host", host); return this; } + /** + * The MQTT protocol version to use when connecting to the broker. + * + * The option is a: <code>com.hivemq.client.mqtt.MqttVersion</code> + * type. + * + * Default: MQTT_5_0 + * Group: common + * + * @param mqttVersion the value to set + * @return the dsl builder + */ + default HiveMQEndpointBuilder mqttVersion(com.hivemq.client.mqtt.MqttVersion mqttVersion) { + doSetProperty("mqttVersion", mqttVersion); + return this; + } + /** + * The MQTT protocol version to use when connecting to the broker. + * + * The option will be converted to a + * <code>com.hivemq.client.mqtt.MqttVersion</code> type. + * + * Default: MQTT_5_0 + * Group: common + * + * @param mqttVersion the value to set + * @return the dsl builder + */ + default HiveMQEndpointBuilder mqttVersion(String mqttVersion) { + doSetProperty("mqttVersion", mqttVersion); + return this; + } /** * Port number of the HiveMQ MQTT broker. *
