atiaomar1978-hub commented on code in PR #25508:
URL: https://github.com/apache/camel/pull/25508#discussion_r3824979346
##########
components/camel-paho/src/main/java/org/apache/camel/component/paho/PahoConsumer.java:
##########
@@ -57,81 +57,128 @@ public void setClient(MqttClient client) {
protected void doStart() throws Exception {
super.doStart();
- connectOptions =
PahoEndpoint.createMqttConnectOptions(getEndpoint().getConfiguration());
-
- if (client == null) {
- clientId = getEndpoint().getConfiguration().getClientId();
- if (clientId == null) {
- clientId = "camel-" + MqttClient.generateClientId();
+ stopClient = client == null;
+ try {
+ connectOptions =
PahoEndpoint.createMqttConnectOptions(getEndpoint().getConfiguration());
+
+ if (stopClient) {
+ clientId = getEndpoint().getConfiguration().getClientId();
+ if (clientId == null) {
+ clientId = "camel-" + MqttClient.generateClientId();
+ }
+ client = createClient();
+ LOG.debug("Connecting client: {} to broker: {}", clientId,
getEndpoint().getConfiguration().getBrokerUrl());
+ if (getEndpoint().getConfiguration().isManualAcksEnabled()) {
+ client.setManualAcks(true);
+ }
+ client.connect(connectOptions);
}
- stopClient = true;
- client = new MqttClient(
- getEndpoint().getConfiguration().getBrokerUrl(),
- clientId,
-
PahoEndpoint.createMqttClientPersistence(getEndpoint().getConfiguration()));
- LOG.debug("Connecting client: {} to broker: {}", clientId,
getEndpoint().getConfiguration().getBrokerUrl());
- if (getEndpoint().getConfiguration().isManualAcksEnabled()) {
- client.setManualAcks(true);
- }
- client.connect(connectOptions);
- }
+ client.setCallback(new MqttCallbackExtended() {
- client.setCallback(new MqttCallbackExtended() {
-
- @Override
- public void connectComplete(boolean reconnect, String serverURI) {
- if (reconnect) {
- try {
- client.subscribe(getEndpoint().getTopic(),
getEndpoint().getConfiguration().getQos());
- } catch (MqttException e) {
- LOG.error("MQTT resubscribe failed {}",
e.getMessage(), e);
+ @Override
+ public void connectComplete(boolean reconnect, String
serverURI) {
+ if (reconnect) {
+ try {
+ client.subscribe(getEndpoint().getTopic(),
getEndpoint().getConfiguration().getQos());
+ } catch (MqttException e) {
+ LOG.error("MQTT resubscribe failed {}",
e.getMessage(), e);
+ }
}
}
- }
- @Override
- public void connectionLost(Throwable cause) {
- LOG.debug("MQTT broker connection lost due {}",
cause.getMessage(), cause);
- }
+ @Override
+ public void connectionLost(Throwable cause) {
+ LOG.debug("MQTT broker connection lost due {}",
cause.getMessage(), cause);
+ }
- @Override
- public void messageArrived(String topic, MqttMessage message)
throws Exception {
- LOG.debug("Message arrived on topic: {} -> {}", topic,
message);
- Exchange exchange = createExchange(message, topic);
+ @Override
+ public void messageArrived(String topic, MqttMessage message)
throws Exception {
+ LOG.debug("Message arrived on topic: {} -> {}", topic,
message);
+ Exchange exchange = createExchange(message, topic);
- // use default consumer callback
- AsyncCallback cb = defaultConsumerCallback(exchange, true);
- getAsyncProcessor().process(exchange, cb);
- }
+ // use default consumer callback
+ AsyncCallback cb = defaultConsumerCallback(exchange, true);
+ getAsyncProcessor().process(exchange, cb);
+ }
- @Override
- public void deliveryComplete(IMqttDeliveryToken token) {
- LOG.debug("Delivery complete. Token: {}", token);
- }
- });
+ @Override
+ public void deliveryComplete(IMqttDeliveryToken token) {
+ LOG.debug("Delivery complete. Token: {}", token);
+ }
+ });
- LOG.debug("Subscribing client: {} to topic: {}", clientId,
getEndpoint().getTopic());
- client.subscribe(getEndpoint().getTopic(),
getEndpoint().getConfiguration().getQos());
+ LOG.debug("Subscribing client: {} to topic: {}", clientId,
getEndpoint().getTopic());
+ client.subscribe(getEndpoint().getTopic(),
getEndpoint().getConfiguration().getQos());
+ } catch (Exception startException) {
+ MqttClient ownedClient = stopClient ? client : null;
+ if (ownedClient != null) {
+ client = null;
+ stopClient = false;
+ if (ownedClient.isConnected()) {
Review Comment:
**Resolved — disconnect before close on failed startup**
`doStart` catch now disconnects the owned client when `isConnected()` before
`closeOwnedClient()`, with disconnect failures suppressed on the primary
startup exception. Confirmed against Paho 1.2.5 behavior noted in your reply —
this is the correct fix.
##########
components/camel-paho/src/test/java/org/apache/camel/component/paho/PahoConsumerLifecycleTest.java:
##########
@@ -0,0 +1,158 @@
+/*
+ * 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.paho;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.ExtendedCamelContext;
+import org.apache.camel.Processor;
+import org.apache.camel.spi.ExchangeFactory;
+import org.eclipse.paho.client.mqttv3.MqttClient;
+import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
+import org.eclipse.paho.client.mqttv3.MqttException;
+import org.junit.jupiter.api.Test;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.catchThrowableOfType;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+class PahoConsumerLifecycleTest {
+
+ @Test
+ void failedStartForceClosesOwnedClient() throws Exception {
+ MqttClient client = mock(MqttClient.class);
+ MqttException connectException = new
MqttException(MqttException.REASON_CODE_CLIENT_EXCEPTION);
+
doThrow(connectException).when(client).connect(any(MqttConnectOptions.class));
+ PahoConsumer consumer = createConsumer(new PahoConfiguration(),
client);
+
+ MqttException thrown = catchThrowableOfType(MqttException.class,
consumer::doStart);
+
+ assertThat(thrown).isSameAs(connectException);
+ verify(client).close(true);
+ }
+
+ @Test
+ void failedStartSuppressesCloseFailure() throws Exception {
+ MqttClient client = mock(MqttClient.class);
+ MqttException connectException = new
MqttException(MqttException.REASON_CODE_CLIENT_EXCEPTION);
+ MqttException closeException = new
MqttException(MqttException.REASON_CODE_CLIENT_DISCONNECTING);
+
doThrow(connectException).when(client).connect(any(MqttConnectOptions.class));
+ doThrow(closeException).when(client).close(true);
+ PahoConsumer consumer = createConsumer(new PahoConfiguration(),
client);
+
+ MqttException thrown = catchThrowableOfType(MqttException.class,
consumer::doStart);
+
+ assertThat(thrown).isSameAs(connectException);
+ assertThat(thrown.getSuppressed()).containsExactly(closeException);
+ }
+
+ @Test
+ void failedStartAfterConnectDisconnectsAndClosesOwnedClient() throws
Exception {
Review Comment:
**Resolved — subscribe-failure regression test**
`failedStartAfterConnectDisconnectsAndClosesOwnedClient` verifies connect
succeeds, subscribe throws, and cleanup calls `disconnect()` + `close(true)`.
Same test mirrored in MQTT5 suite — good coverage.
##########
components/camel-paho-mqtt5/src/main/java/org/apache/camel/component/paho/mqtt5/PahoMqtt5Consumer.java:
##########
@@ -58,91 +58,138 @@ public void setClient(MqttClient client) {
protected void doStart() throws Exception {
super.doStart();
- connectionOptions = getEndpoint().createMqttConnectionOptions();
-
- if (client == null) {
- clientId = getEndpoint().getConfiguration().getClientId();
- if (clientId == null) {
- clientId = PahoMqtt5Endpoint.generateClientId();
- }
- stopClient = true;
- client = new MqttClient(
- getEndpoint().getConfiguration().getBrokerUrl(),
- clientId,
-
PahoMqtt5Endpoint.createMqttClientPersistence(getEndpoint().getConfiguration()));
- LOG.debug("Connecting client: {} to broker: {}", clientId,
getEndpoint().getConfiguration().getBrokerUrl());
- if (getEndpoint().getConfiguration().isManualAcksEnabled()) {
- client.setManualAcks(true);
-
+ stopClient = client == null;
+ try {
+ connectionOptions = getEndpoint().createMqttConnectionOptions();
+
+ if (stopClient) {
+ clientId = getEndpoint().getConfiguration().getClientId();
+ if (clientId == null) {
+ clientId = PahoMqtt5Endpoint.generateClientId();
+ }
+ client = createClient();
+ LOG.debug("Connecting client: {} to broker: {}", clientId,
getEndpoint().getConfiguration().getBrokerUrl());
+ if (getEndpoint().getConfiguration().isManualAcksEnabled()) {
+ client.setManualAcks(true);
+ }
+ client.connect(connectionOptions);
}
- client.connect(connectionOptions);
- }
- client.setCallback(new MqttCallback() {
+ client.setCallback(new MqttCallback() {
- @Override
- public void connectComplete(boolean reconnect, String serverURI) {
- if (reconnect) {
- try {
- client.subscribe(getEndpoint().getTopic(),
getEndpoint().getConfiguration().getQos());
- } catch (MqttException e) {
- LOG.error("MQTT resubscribe failed {}",
e.getMessage(), e);
+ @Override
+ public void connectComplete(boolean reconnect, String
serverURI) {
+ if (reconnect) {
+ try {
+ client.subscribe(getEndpoint().getTopic(),
getEndpoint().getConfiguration().getQos());
+ } catch (MqttException e) {
+ LOG.error("MQTT resubscribe failed {}",
e.getMessage(), e);
+ }
}
}
- }
- @Override
- public void authPacketArrived(int reasonCode, MqttProperties
properties) {
- LOG.debug("Auth packet arrived {} {}", reasonCode, properties);
- }
+ @Override
+ public void authPacketArrived(int reasonCode, MqttProperties
properties) {
+ LOG.debug("Auth packet arrived {} {}", reasonCode,
properties);
+ }
- @Override
- public void disconnected(MqttDisconnectResponse response) {
- LOG.debug("MQTT broker disconnected due {}",
response.getReasonString(), response.getException());
- }
+ @Override
+ public void disconnected(MqttDisconnectResponse response) {
+ LOG.debug("MQTT broker disconnected due {}",
response.getReasonString(), response.getException());
+ }
- @Override
- public void mqttErrorOccurred(MqttException exception) {
- LOG.debug("Error occurred {}", exception.getMessage(),
exception);
- }
+ @Override
+ public void mqttErrorOccurred(MqttException exception) {
+ LOG.debug("Error occurred {}", exception.getMessage(),
exception);
+ }
+
+ @Override
+ public void messageArrived(String topic, MqttMessage message)
throws Exception {
+ LOG.debug("Message arrived on topic: {} -> {}", topic,
message);
+ Exchange exchange = createExchange(message, topic);
- @Override
- public void messageArrived(String topic, MqttMessage message)
throws Exception {
- LOG.debug("Message arrived on topic: {} -> {}", topic,
message);
- Exchange exchange = createExchange(message, topic);
+ // use default consumer callback
+ AsyncCallback cb = defaultConsumerCallback(exchange, true);
+ getAsyncProcessor().process(exchange, cb);
+ }
- // use default consumer callback
- AsyncCallback cb = defaultConsumerCallback(exchange, true);
- getAsyncProcessor().process(exchange, cb);
- }
+ @Override
+ public void deliveryComplete(IMqttToken token) {
+ LOG.debug("Delivery complete. Token: {}", token);
+ }
+ });
- @Override
- public void deliveryComplete(IMqttToken token) {
- LOG.debug("Delivery complete. Token: {}", token);
+ LOG.debug("Subscribing client: {} to topic: {}", clientId,
getEndpoint().getTopic());
+ client.subscribe(getEndpoint().getTopic(),
getEndpoint().getConfiguration().getQos());
+ } catch (Exception startException) {
+ MqttClient ownedClient = stopClient ? client : null;
+ if (ownedClient != null) {
+ client = null;
+ stopClient = false;
+ if (ownedClient.isConnected()) {
Review Comment:
**Resolved — symmetric MQTT5 fix**
Same disconnect-before-close cleanup in the `doStart` catch block. Both
consumers now share identical lifecycle semantics.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]