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 3485b7c94d..7b86d1bddd 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 @@ -51,9 +51,15 @@ import org.springframework.util.Assert; public class MqttPahoMessageDrivenChannelAdapter extends AbstractMqttMessageDrivenChannelAdapter implements MqttCallback, ApplicationEventPublisherAware { + /** + * The default completion timeout in milliseconds. + */ public static final long DEFAULT_COMPLETION_TIMEOUT = 30_000L; - public static final long STOP_COMPLETION_TIMEOUT = 5_000L; + /** + * The default disconnect completion timeout in milliseconds. + */ + public static final long DISCONNECT_COMPLETION_TIMEOUT = 5_000L; private static final int DEFAULT_RECOVERY_INTERVAL = 10_000; @@ -63,7 +69,7 @@ public class MqttPahoMessageDrivenChannelAdapter extends AbstractMqttMessageDriv private long completionTimeout = DEFAULT_COMPLETION_TIMEOUT; - private long stopCompletionTimeout = STOP_COMPLETION_TIMEOUT; + private long disconnectCompletionTimeout = DISCONNECT_COMPLETION_TIMEOUT; private volatile IMqttClient client; @@ -127,13 +133,13 @@ public class MqttPahoMessageDrivenChannelAdapter extends AbstractMqttMessageDriv } /** - * Set the completion timeout wnen stopping. Not settable using the namespace. - * Default {@value #STOP_COMPLETION_TIMEOUT} milliseconds. + * Set the completion timeout when disconnecting. Not settable using the namespace. + * Default {@value #DISCONNECT_COMPLETION_TIMEOUT} milliseconds. * @param completionTimeout The timeout. - * @since 5.2.5 + * @since 5.1.10 */ - public void setStopCompletionTimeout(long completionTimeout) { - this.stopCompletionTimeout = completionTimeout; + public void setDisconnectCompletionTimeout(long completionTimeout) { + this.disconnectCompletionTimeout = completionTimeout; } /** @@ -184,7 +190,7 @@ public class MqttPahoMessageDrivenChannelAdapter extends AbstractMqttMessageDriv logger.error("Exception while unsubscribing", e); } try { - this.client.disconnectForcibly(this.stopCompletionTimeout); + this.client.disconnectForcibly(this.disconnectCompletionTimeout); } catch (MqttException e) { logger.error("Exception while disconnecting", e); @@ -276,7 +282,7 @@ public class MqttPahoMessageDrivenChannelAdapter extends AbstractMqttMessageDriv this.applicationEventPublisher.publishEvent(new MqttConnectionFailedEvent(this, e)); } logger.error("Error connecting or subscribing to " + Arrays.toString(topics), e); - this.client.disconnectForcibly(this.completionTimeout); + this.client.disconnectForcibly(this.disconnectCompletionTimeout); try { this.client.setCallback(null); this.client.close(); 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 0d2690e99a..fb93564233 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-2019 the original author or authors. + * Copyright 2002-2020 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. @@ -50,10 +50,17 @@ public class MqttPahoMessageHandler extends AbstractMqttMessageHandler /** * The default completion timeout in milliseconds. */ - public static final long DEFAULT_COMPLETION_TIMEOUT = 30000L; + public static final long DEFAULT_COMPLETION_TIMEOUT = 30_000L; + + /** + * The default disconnect completion timeout in milliseconds. + */ + public static final long DISCONNECT_COMPLETION_TIMEOUT = 5_000L; private long completionTimeout = DEFAULT_COMPLETION_TIMEOUT; + private long disconnectCompletionTimeout = DISCONNECT_COMPLETION_TIMEOUT; + private final MqttPahoClientFactory clientFactory; private boolean async; @@ -131,6 +138,16 @@ public class MqttPahoMessageHandler extends AbstractMqttMessageHandler this.completionTimeout = completionTimeout; } + /** + * Set the completion timeout when disconnecting. Not settable using the namespace. + * Default {@value #DISCONNECT_COMPLETION_TIMEOUT} milliseconds. + * @param completionTimeout The timeout. + * @since 5.1.10 + */ + public void setDisconnectCompletionTimeout(long completionTimeout) { + this.disconnectCompletionTimeout = completionTimeout; + } + @Override public void setApplicationEventPublisher(ApplicationEventPublisher applicationEventPublisher) { this.applicationEventPublisher = applicationEventPublisher; @@ -152,7 +169,7 @@ public class MqttPahoMessageHandler extends AbstractMqttMessageHandler try { IMqttAsyncClient theClient = this.client; if (theClient != null) { - theClient.disconnect().waitForCompletion(this.completionTimeout); + theClient.disconnect().waitForCompletion(this.disconnectCompletionTimeout); theClient.close(); this.client = null; }