nkokitkar commented on code in PR #25508:
URL: https://github.com/apache/camel/pull/25508#discussion_r3822386666
##########
components/camel-paho/src/main/java/org/apache/camel/component/paho/PahoConsumer.java:
##########
@@ -57,81 +57,121 @@ 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 = 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);
-
+ 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);
}
- 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;
+ closeOwnedClient(ownedClient, startException);
+ }
+ throw startException;
+ }
}
@Override
protected void doStop() throws Exception {
- super.doStop();
-
- if (stopClient && client != null && client.isConnected()) {
- String topic = getEndpoint().getTopic();
- // only unsubscribe if we are not durable
- if (getEndpoint().getConfiguration().isCleanSession()) {
- LOG.debug("Unsubscribing client: {} from topic: {}", clientId,
topic);
- client.unsubscribe(topic);
- } else {
- LOG.debug("Client: {} is durable so will not unsubscribe from
topic: {}", clientId, topic);
+ MqttClient ownedClient = stopClient ? client : null;
+ Exception stopException = null;
+ try {
+ super.doStop();
+
+ if (ownedClient != null && ownedClient.isConnected()) {
+ String topic = getEndpoint().getTopic();
+ // only unsubscribe if we are not durable
+ if (getEndpoint().getConfiguration().isCleanSession()) {
+ LOG.debug("Unsubscribing client: {} from topic: {}",
clientId, topic);
+ ownedClient.unsubscribe(topic);
+ } else {
+ LOG.debug("Client: {} is durable so will not unsubscribe
from topic: {}", clientId, topic);
+ }
+ LOG.debug("Disconnecting client: {} from broker: {}", clientId,
+ getEndpoint().getConfiguration().getBrokerUrl());
+ ownedClient.disconnect();
+ }
+ } catch (Exception e) {
+ stopException = e;
+ } finally {
+ client = null;
+ stopClient = false;
+ if (ownedClient != null) {
+ stopException = closeOwnedClient(ownedClient, stopException);
+ }
+ }
+ if (stopException != null) {
+ throw stopException;
+ }
+ }
+
+ MqttClient createClient() throws MqttException {
+ return new MqttClient(
+ getEndpoint().getConfiguration().getBrokerUrl(),
+ clientId,
+
PahoEndpoint.createMqttClientPersistence(getEndpoint().getConfiguration()));
+ }
+
+ private Exception closeOwnedClient(MqttClient ownedClient, Exception
primaryException) {
Review Comment:
Confirmed against Paho 1.2.5 that `close(true)` still rejects connected
clients. Commit `364767015edf` now disconnects connected owned clients before
close on failed startup and adds MQTT v3/v5 subscribe-failure coverage.
_GitHub Copilot on behalf of nkokitkar._
##########
components/camel-paho/src/main/java/org/apache/camel/component/paho/PahoConsumer.java:
##########
@@ -57,81 +57,121 @@ 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 = 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);
-
+ 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);
}
- 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;
+ closeOwnedClient(ownedClient, startException);
+ }
+ throw startException;
+ }
}
@Override
protected void doStop() throws Exception {
- super.doStop();
-
- if (stopClient && client != null && client.isConnected()) {
- String topic = getEndpoint().getTopic();
- // only unsubscribe if we are not durable
- if (getEndpoint().getConfiguration().isCleanSession()) {
- LOG.debug("Unsubscribing client: {} from topic: {}", clientId,
topic);
- client.unsubscribe(topic);
- } else {
- LOG.debug("Client: {} is durable so will not unsubscribe from
topic: {}", clientId, topic);
+ MqttClient ownedClient = stopClient ? client : null;
+ Exception stopException = null;
+ try {
+ super.doStop();
+
+ if (ownedClient != null && ownedClient.isConnected()) {
+ String topic = getEndpoint().getTopic();
+ // only unsubscribe if we are not durable
+ if (getEndpoint().getConfiguration().isCleanSession()) {
+ LOG.debug("Unsubscribing client: {} from topic: {}",
clientId, topic);
+ ownedClient.unsubscribe(topic);
+ } else {
+ LOG.debug("Client: {} is durable so will not unsubscribe
from topic: {}", clientId, topic);
+ }
+ LOG.debug("Disconnecting client: {} from broker: {}", clientId,
+ getEndpoint().getConfiguration().getBrokerUrl());
+ ownedClient.disconnect();
+ }
+ } catch (Exception e) {
+ stopException = e;
+ } finally {
+ client = null;
Review Comment:
Thanks. The shutdown cleanup remains covered, and the new change is limited
to ensuring the failed-start path also disconnects before close when startup
progressed past connect.
_GitHub Copilot on behalf of nkokitkar._
##########
components/camel-paho/src/test/java/org/apache/camel/component/paho/PahoConsumerLifecycleTest.java:
##########
@@ -0,0 +1,142 @@
+/*
+ * 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 {
Review Comment:
Added the missing subscribe-failure regression test for both MQTT v3 and v5.
Each suite now verifies that a connected owned client is disconnected and
closed while preserving the subscribe exception.
_GitHub Copilot on behalf of nkokitkar._
##########
components/camel-paho-mqtt5/src/main/java/org/apache/camel/component/paho/mqtt5/PahoMqtt5Consumer.java:
##########
@@ -58,91 +58,131 @@ 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());
+ 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;
+ closeOwnedClient(ownedClient, startException);
+ }
+ throw startException;
+ }
}
@Override
protected void doStop() throws Exception {
- super.doStop();
-
- if (stopClient && client != null && client.isConnected()) {
- String topic = getEndpoint().getTopic();
- // only unsubscribe if we are not durable
- if (getEndpoint().getConfiguration().isCleanStart()) {
- LOG.debug("Unsubscribing client: {} from topic: {}", clientId,
topic);
- client.unsubscribe(topic);
- } else {
- LOG.debug("Client: {} is durable so will not unsubscribe from
topic: {}", clientId, topic);
+ MqttClient ownedClient = stopClient ? client : null;
+ Exception stopException = null;
+ try {
+ super.doStop();
+
+ if (ownedClient != null && ownedClient.isConnected()) {
+ String topic = getEndpoint().getTopic();
+ // only unsubscribe if we are not durable
+ if (getEndpoint().getConfiguration().isCleanStart()) {
+ LOG.debug("Unsubscribing client: {} from topic: {}",
clientId, topic);
+ ownedClient.unsubscribe(topic);
+ } else {
+ LOG.debug("Client: {} is durable so will not unsubscribe
from topic: {}", clientId, topic);
+ }
+ LOG.debug("Disconnecting client: {} from broker: {}", clientId,
+ getEndpoint().getConfiguration().getBrokerUrl());
+ ownedClient.disconnect();
+ }
+ } catch (Exception e) {
+ stopException = e;
+ } finally {
+ client = null;
+ stopClient = false;
+ if (ownedClient != null) {
+ stopException = closeOwnedClient(ownedClient, stopException);
+ }
+ }
+ if (stopException != null) {
+ throw stopException;
+ }
+ }
+
+ MqttClient createClient() throws MqttException {
+ return new MqttClient(
+ getEndpoint().getConfiguration().getBrokerUrl(),
+ clientId,
+
PahoMqtt5Endpoint.createMqttClientPersistence(getEndpoint().getConfiguration()));
+ }
+
+ private Exception closeOwnedClient(MqttClient ownedClient, Exception
primaryException) {
Review Comment:
Applied the same disconnect-before-close cleanup to `PahoMqtt5Consumer` in
`364767015edf`, with symmetric subscribe-failure coverage.
_GitHub Copilot on behalf of nkokitkar._
--
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]