diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/aot/MqttRuntimeHints.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/aot/MqttRuntimeHints.java new file mode 100644 index 0000000000..e591b80378 --- /dev/null +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/aot/MqttRuntimeHints.java @@ -0,0 +1,57 @@ +/* + * Copyright 2024 the original author or authors. + * + * Licensed 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 + * + * https://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.springframework.integration.mqtt.aot; + +import java.util.stream.Stream; + +import org.springframework.aot.hint.ExecutableMode; +import org.springframework.aot.hint.ReflectionHints; +import org.springframework.aot.hint.RuntimeHints; +import org.springframework.aot.hint.RuntimeHintsRegistrar; +import org.springframework.util.ClassUtils; +import org.springframework.util.ReflectionUtils; + +/** + * {@link RuntimeHintsRegistrar} for Spring Integration MQTT module. + * + * @author Artem Bilan + * + * @since 6.1.9 + */ +class MqttRuntimeHints implements RuntimeHintsRegistrar { + + @Override + public void registerHints(RuntimeHints hints, ClassLoader classLoader) { + ReflectionHints reflectionHints = hints.reflection(); + // TODO until the real fix in Paho library. + Stream.of("org.eclipse.paho.client.mqttv3.MqttAsyncClient", "org.eclipse.paho.mqttv5.client.MqttAsyncClient") + .filter((typeName) -> ClassUtils.isPresent(typeName, classLoader)) + .map((typeName) -> loadClassByName(typeName, classLoader)) + .flatMap((type) -> Stream.ofNullable(ReflectionUtils.findMethod(type, "stopReconnectCycle"))) + .forEach(method -> reflectionHints.registerMethod(method, ExecutableMode.INVOKE)); + } + + private static Class loadClassByName(String typeName, ClassLoader classLoader) { + try { + return ClassUtils.forName(typeName, classLoader); + } + catch (ClassNotFoundException ex) { + throw new IllegalArgumentException(ex); + } + } + +} diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/aot/package-info.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/aot/package-info.java new file mode 100644 index 0000000000..69f0b3b5ae --- /dev/null +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/aot/package-info.java @@ -0,0 +1,6 @@ +/** + * Provides classes to support Spring AOT. + */ +@org.springframework.lang.NonNullApi +@org.springframework.lang.NonNullFields +package org.springframework.integration.mqtt.aot; diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/core/Mqttv3ClientManager.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/core/Mqttv3ClientManager.java index 6a9a0249e1..8fb3f99182 100644 --- a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/core/Mqttv3ClientManager.java +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/core/Mqttv3ClientManager.java @@ -1,5 +1,5 @@ /* - * Copyright 2022-2023 the original author or authors. + * Copyright 2022-2024 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -26,6 +26,7 @@ import org.eclipse.paho.client.mqttv3.MqttException; import org.eclipse.paho.client.mqttv3.MqttMessage; import org.springframework.integration.mqtt.event.MqttConnectionFailedEvent; +import org.springframework.integration.mqtt.support.MqttUtils; import org.springframework.util.Assert; /** @@ -149,6 +150,9 @@ public class Mqttv3ClientManager } try { client.disconnectForcibly(getDisconnectCompletionTimeout()); + if (getConnectionInfo().isAutomaticReconnect()) { + MqttUtils.stopClientReconnectCycle(client); + } } catch (MqttException e) { logger.error("Could not disconnect from the client", e); diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/core/Mqttv5ClientManager.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/core/Mqttv5ClientManager.java index a89b34aa60..36bc028b70 100644 --- a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/core/Mqttv5ClientManager.java +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/core/Mqttv5ClientManager.java @@ -1,5 +1,5 @@ /* - * Copyright 2022-2023 the original author or authors. + * Copyright 2022-2024 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -28,6 +28,7 @@ import org.eclipse.paho.mqttv5.common.MqttMessage; import org.eclipse.paho.mqttv5.common.packet.MqttProperties; import org.springframework.integration.mqtt.event.MqttConnectionFailedEvent; +import org.springframework.integration.mqtt.support.MqttUtils; import org.springframework.util.Assert; /** @@ -151,6 +152,9 @@ public class Mqttv5ClientManager try { client.disconnectForcibly(getDisconnectCompletionTimeout()); + if (getConnectionInfo().isAutomaticReconnect()) { + MqttUtils.stopClientReconnectCycle(client); + } } catch (MqttException e) { logger.error("Could not disconnect from the client", e); diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/MqttPahoMessageDrivenChannelAdapter.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/MqttPahoMessageDrivenChannelAdapter.java index d7a77ac9d7..6df2802ace 100644 --- a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/MqttPahoMessageDrivenChannelAdapter.java +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/MqttPahoMessageDrivenChannelAdapter.java @@ -228,6 +228,9 @@ public class MqttPahoMessageDrivenChannelAdapter try { this.client.disconnectForcibly(getDisconnectCompletionTimeout()); + if (getConnectionInfo().isAutomaticReconnect()) { + MqttUtils.stopClientReconnectCycle(this.client); + } } catch (MqttException ex) { logger.error(ex, "Exception while disconnecting"); diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/Mqttv5PahoMessageDrivenChannelAdapter.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/Mqttv5PahoMessageDrivenChannelAdapter.java index e71b53a9c0..3d015b8713 100644 --- a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/Mqttv5PahoMessageDrivenChannelAdapter.java +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/Mqttv5PahoMessageDrivenChannelAdapter.java @@ -53,6 +53,7 @@ import org.springframework.integration.mqtt.event.MqttSubscribedEvent; import org.springframework.integration.mqtt.support.MqttHeaderMapper; import org.springframework.integration.mqtt.support.MqttHeaders; import org.springframework.integration.mqtt.support.MqttMessageConverter; +import org.springframework.integration.mqtt.support.MqttUtils; import org.springframework.lang.Nullable; import org.springframework.messaging.Message; import org.springframework.messaging.MessageHeaders; @@ -299,6 +300,9 @@ public class Mqttv5PahoMessageDrivenChannelAdapter } if (getClientManager() == null) { this.mqttClient.disconnectForcibly(getDisconnectCompletionTimeout()); + if (getConnectionInfo().isAutomaticReconnect()) { + MqttUtils.stopClientReconnectCycle(this.mqttClient); + } } } } diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/MqttPahoMessageHandler.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/MqttPahoMessageHandler.java index 317a9b9604..3a565d76cc 100644 --- a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/MqttPahoMessageHandler.java +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/MqttPahoMessageHandler.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2023 the original author or authors. + * Copyright 2002-2024 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -176,6 +176,9 @@ public class MqttPahoMessageHandler extends AbstractMqttMessageHandler method = new AtomicReference<>(); ReflectionUtils.doWithMethods(MqttPahoMessageDrivenChannelAdapter.class, m -> { m.setAccessible(true);