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 1f77f911259b002108f6c1e30739f70f81f16cf1 Author: Thomas Strauss <[email protected]> AuthorDate: Wed Sep 30 14:22:24 2026 +0200 CAMEL-25030: camel-hivemq: extract shared reconnect-cancellation listener Mqtt3ClientAdapter and Mqtt5ClientAdapter each had an identical ~20-line block registering the HiveMQ #302 reconnect-cancellation workaround (addDisconnectedListener/addConnectedListener + the cancelReconnect AtomicBoolean dance), flagged in review as duplicated code that would need any future fix made twice. Everything in that block is actually version-neutral in the HiveMQ MQTT Client library: MqttClientBuilderBase (the common ancestor of Mqtt3ClientBuilder/Mqtt5ClientBuilder), the MqttClientConnectedListener/ MqttClientDisconnectedListener types and their contexts (all in the shared com.hivemq.client.mqtt.lifecycle package), and MqttClient#getState() are all shared - verified against the library's actual type hierarchy, not assumed. Only the final disconnect() call is genuinely protocol-specific (MQTT 5's DISCONNECT carries reason codes/properties MQTT 3.1.1 doesn't have), so it isn't part of this class - it's supplied by each adapter as a one-line Runnable. Extracts this into a new package-private MqttReconnectCancellation helper, parameterized over MqttClientBuilderBase<B>. Both adapters now build the version-neutral MqttClient.builder() with host/port/reconnect config, register the shared listener via MqttReconnectCancellation.apply(...), and only then call useMqttVersion3()/useMqttVersion5() to get their version-specific builder - no behavioural change, same listener logic, now defined once. Verified with a full, non-quickly `mvn install -DskipTests` build (BUILD SUCCESS) and all unit tests plus all 6 ITs (including the MQTT 3.1.1 pub/sub IT) passing against a live HiveMQ CE broker. AI-assisted contribution: authored together with Claude Sonnet 5 (Anthropic). Co-Authored-By: Claude Sonnet 5 <[email protected]> --- .../camel/component/hivemq/Mqtt3ClientAdapter.java | 32 +++-------- .../camel/component/hivemq/Mqtt5ClientAdapter.java | 32 +++-------- .../hivemq/MqttReconnectCancellation.java | 63 ++++++++++++++++++++++ 3 files changed, 75 insertions(+), 52 deletions(-) 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 index 648ca916d92e..8c31fc191bba 100644 --- 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 @@ -25,7 +25,6 @@ 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; @@ -44,31 +43,12 @@ final class Mqtt3ClientAdapter implements HiveMQClientAdapter { 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 - } - } - } - }) + Mqtt3ClientBuilder builder = MqttReconnectCancellation.apply( + MqttClient.builder() + .serverHost(configuration.getHost()) + .serverPort(configuration.getPort()) + .automaticReconnectWithDefaultConfig(), + cancelReconnect, clientRef::get, () -> clientRef.get().disconnect()) .useMqttVersion3(); if (configuration.getClientId() != null) { diff --git a/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/Mqtt5ClientAdapter.java b/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/Mqtt5ClientAdapter.java index 1be32bb431e0..e1744497e951 100644 --- a/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/Mqtt5ClientAdapter.java +++ b/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/Mqtt5ClientAdapter.java @@ -25,7 +25,6 @@ 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; @@ -44,31 +43,12 @@ final class Mqtt5ClientAdapter implements HiveMQClientAdapter { Mqtt5ClientAdapter(HiveMQConfiguration configuration) { 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 - } - } - } - }) + Mqtt5ClientBuilder builder = MqttReconnectCancellation.apply( + MqttClient.builder() + .serverHost(configuration.getHost()) + .serverPort(configuration.getPort()) + .automaticReconnectWithDefaultConfig(), + cancelReconnect, clientRef::get, () -> clientRef.get().disconnect()) .useMqttVersion5(); if (configuration.getClientId() != null) { diff --git a/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/MqttReconnectCancellation.java b/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/MqttReconnectCancellation.java new file mode 100644 index 000000000000..c589dc2745fa --- /dev/null +++ b/components/camel-hivemq/src/main/java/org/apache/camel/component/hivemq/MqttReconnectCancellation.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.concurrent.atomic.AtomicBoolean; +import java.util.function.Supplier; + +import com.hivemq.client.mqtt.MqttClient; +import com.hivemq.client.mqtt.MqttClientBuilderBase; +import com.hivemq.client.mqtt.MqttClientState; + +/** + * Registers the HiveMQ #302 reconnect-cancellation workaround on an {@link MqttClientBuilderBase}, shared between + * {@link Mqtt3ClientAdapter} and {@link Mqtt5ClientAdapter}. Everything involved - the builder base, the listener + * types, and {@link MqttClient#getState()} - is version-neutral in the HiveMQ MQTT Client library; only the actual + * {@code disconnect()} call is protocol-specific (MQTT 5's DISCONNECT carries reason codes/properties that MQTT 3.1.1 + * does not have), which is why it is supplied by the caller instead of being part of this class. + */ +final class MqttReconnectCancellation { + + private MqttReconnectCancellation() { + } + + static <B extends MqttClientBuilderBase<B>> B apply( + B builder, AtomicBoolean cancelReconnect, Supplier<? extends MqttClient> client, Runnable disconnect) { + return builder + .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()) { + MqttClient started = client.get(); + if (started != null && started.getState().isConnected()) { + try { + disconnect.run(); + } catch (Exception e) { + // Already disconnecting or not connected + } + } + } + }); + } +}
