|
|
@ -45,6 +45,7 @@ import io.netty.handler.timeout.IdleStateHandler; |
|
|
import io.netty.util.concurrent.DefaultPromise; |
|
|
import io.netty.util.concurrent.DefaultPromise; |
|
|
import io.netty.util.concurrent.Future; |
|
|
import io.netty.util.concurrent.Future; |
|
|
import io.netty.util.concurrent.Promise; |
|
|
import io.netty.util.concurrent.Promise; |
|
|
|
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
|
|
|
|
|
|
import java.util.Collections; |
|
|
import java.util.Collections; |
|
|
import java.util.HashSet; |
|
|
import java.util.HashSet; |
|
|
@ -60,6 +61,7 @@ import java.util.concurrent.atomic.AtomicInteger; |
|
|
* Represents an MqttClientImpl connected to a single MQTT server. Will try to keep the connection going at all times |
|
|
* Represents an MqttClientImpl connected to a single MQTT server. Will try to keep the connection going at all times |
|
|
*/ |
|
|
*/ |
|
|
@SuppressWarnings({"WeakerAccess", "unused"}) |
|
|
@SuppressWarnings({"WeakerAccess", "unused"}) |
|
|
|
|
|
@Slf4j |
|
|
final class MqttClientImpl implements MqttClient { |
|
|
final class MqttClientImpl implements MqttClient { |
|
|
|
|
|
|
|
|
private final Set<String> serverSubscriptions = new HashSet<>(); |
|
|
private final Set<String> serverSubscriptions = new HashSet<>(); |
|
|
@ -131,6 +133,7 @@ final class MqttClientImpl implements MqttClient { |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private Future<MqttConnectResult> connect(String host, int port, boolean reconnect) { |
|
|
private Future<MqttConnectResult> connect(String host, int port, boolean reconnect) { |
|
|
|
|
|
log.trace("[{}] Connecting to server, isReconnect - {}", channel != null ? channel.id() : "UNKNOWN", reconnect); |
|
|
if (this.eventLoop == null) { |
|
|
if (this.eventLoop == null) { |
|
|
this.eventLoop = new NioEventLoopGroup(); |
|
|
this.eventLoop = new NioEventLoopGroup(); |
|
|
} |
|
|
} |
|
|
@ -147,10 +150,12 @@ final class MqttClientImpl implements MqttClient { |
|
|
future.addListener((ChannelFutureListener) f -> { |
|
|
future.addListener((ChannelFutureListener) f -> { |
|
|
if (f.isSuccess()) { |
|
|
if (f.isSuccess()) { |
|
|
MqttClientImpl.this.channel = f.channel(); |
|
|
MqttClientImpl.this.channel = f.channel(); |
|
|
|
|
|
log.debug("[{}][{}] Connected successfully {}!", host, port, this.channel.id()); |
|
|
MqttClientImpl.this.channel.closeFuture().addListener((ChannelFutureListener) channelFuture -> { |
|
|
MqttClientImpl.this.channel.closeFuture().addListener((ChannelFutureListener) channelFuture -> { |
|
|
if (isConnected()) { |
|
|
if (isConnected()) { |
|
|
return; |
|
|
return; |
|
|
} |
|
|
} |
|
|
|
|
|
log.debug("[{}][{}] Channel is closed {}!", host, port, this.channel.id()); |
|
|
ChannelClosedException e = new ChannelClosedException("Channel is closed!"); |
|
|
ChannelClosedException e = new ChannelClosedException("Channel is closed!"); |
|
|
if (callback != null) { |
|
|
if (callback != null) { |
|
|
callback.connectionLost(e); |
|
|
callback.connectionLost(e); |
|
|
@ -169,6 +174,7 @@ final class MqttClientImpl implements MqttClient { |
|
|
scheduleConnectIfRequired(host, port, true); |
|
|
scheduleConnectIfRequired(host, port, true); |
|
|
}); |
|
|
}); |
|
|
} else { |
|
|
} else { |
|
|
|
|
|
log.debug("[{}][{}] Connect failed, trying reconnect!", host, port); |
|
|
scheduleConnectIfRequired(host, port, reconnect); |
|
|
scheduleConnectIfRequired(host, port, reconnect); |
|
|
} |
|
|
} |
|
|
}); |
|
|
}); |
|
|
@ -176,6 +182,7 @@ final class MqttClientImpl implements MqttClient { |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private void scheduleConnectIfRequired(String host, int port, boolean reconnect) { |
|
|
private void scheduleConnectIfRequired(String host, int port, boolean reconnect) { |
|
|
|
|
|
log.trace("[{}] Scheduling connect to server, isReconnect - {}", channel != null ? channel.id() : "UNKNOWN", reconnect); |
|
|
if (clientConfig.isReconnect() && !disconnected) { |
|
|
if (clientConfig.isReconnect() && !disconnected) { |
|
|
if (reconnect) { |
|
|
if (reconnect) { |
|
|
this.reconnect = true; |
|
|
this.reconnect = true; |
|
|
@ -191,6 +198,7 @@ final class MqttClientImpl implements MqttClient { |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
public Future<MqttConnectResult> reconnect() { |
|
|
public Future<MqttConnectResult> reconnect() { |
|
|
|
|
|
log.trace("[{}] Reconnecting to server, isReconnect - {}", channel != null ? channel.id() : "UNKNOWN", reconnect); |
|
|
if (host == null) { |
|
|
if (host == null) { |
|
|
throw new IllegalStateException("Cannot reconnect. Call connect() first"); |
|
|
throw new IllegalStateException("Cannot reconnect. Call connect() first"); |
|
|
} |
|
|
} |
|
|
@ -281,6 +289,7 @@ final class MqttClientImpl implements MqttClient { |
|
|
*/ |
|
|
*/ |
|
|
@Override |
|
|
@Override |
|
|
public Future<Void> off(String topic, MqttHandler handler) { |
|
|
public Future<Void> off(String topic, MqttHandler handler) { |
|
|
|
|
|
log.trace("[{}] Unsubscribing from {}", channel != null ? channel.id() : "UNKNOWN", topic); |
|
|
Promise<Void> future = new DefaultPromise<>(this.eventLoop.next()); |
|
|
Promise<Void> future = new DefaultPromise<>(this.eventLoop.next()); |
|
|
for (MqttSubscription subscription : this.handlerToSubscription.get(handler)) { |
|
|
for (MqttSubscription subscription : this.handlerToSubscription.get(handler)) { |
|
|
this.subscriptions.remove(topic, subscription); |
|
|
this.subscriptions.remove(topic, subscription); |
|
|
@ -299,6 +308,7 @@ final class MqttClientImpl implements MqttClient { |
|
|
*/ |
|
|
*/ |
|
|
@Override |
|
|
@Override |
|
|
public Future<Void> off(String topic) { |
|
|
public Future<Void> off(String topic) { |
|
|
|
|
|
log.trace("[{}] Unsubscribing from {}", channel != null ? channel.id() : "UNKNOWN", topic); |
|
|
Promise<Void> future = new DefaultPromise<>(this.eventLoop.next()); |
|
|
Promise<Void> future = new DefaultPromise<>(this.eventLoop.next()); |
|
|
ImmutableSet<MqttSubscription> subscriptions = ImmutableSet.copyOf(this.subscriptions.get(topic)); |
|
|
ImmutableSet<MqttSubscription> subscriptions = ImmutableSet.copyOf(this.subscriptions.get(topic)); |
|
|
for (MqttSubscription subscription : subscriptions) { |
|
|
for (MqttSubscription subscription : subscriptions) { |
|
|
@ -360,6 +370,7 @@ final class MqttClientImpl implements MqttClient { |
|
|
*/ |
|
|
*/ |
|
|
@Override |
|
|
@Override |
|
|
public Future<Void> publish(String topic, ByteBuf payload, MqttQoS qos, boolean retain) { |
|
|
public Future<Void> publish(String topic, ByteBuf payload, MqttQoS qos, boolean retain) { |
|
|
|
|
|
log.trace("[{}] Publishing message to {}", channel != null ? channel.id() : "UNKNOWN", topic); |
|
|
Promise<Void> future = new DefaultPromise<>(this.eventLoop.next()); |
|
|
Promise<Void> future = new DefaultPromise<>(this.eventLoop.next()); |
|
|
MqttFixedHeader fixedHeader = new MqttFixedHeader(MqttMessageType.PUBLISH, false, qos, retain, 0); |
|
|
MqttFixedHeader fixedHeader = new MqttFixedHeader(MqttMessageType.PUBLISH, false, qos, retain, 0); |
|
|
MqttPublishVariableHeader variableHeader = new MqttPublishVariableHeader(topic, getNewMessageId().messageId()); |
|
|
MqttPublishVariableHeader variableHeader = new MqttPublishVariableHeader(topic, getNewMessageId().messageId()); |
|
|
@ -404,6 +415,7 @@ final class MqttClientImpl implements MqttClient { |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
public void disconnect() { |
|
|
public void disconnect() { |
|
|
|
|
|
log.trace("[{}] Disconnecting from server", channel != null ? channel.id() : "UNKNOWN"); |
|
|
disconnected = true; |
|
|
disconnected = true; |
|
|
if (this.channel != null) { |
|
|
if (this.channel != null) { |
|
|
MqttMessage message = new MqttMessage(new MqttFixedHeader(MqttMessageType.DISCONNECT, false, MqttQoS.AT_MOST_ONCE, false, 0)); |
|
|
MqttMessage message = new MqttMessage(new MqttFixedHeader(MqttMessageType.DISCONNECT, false, MqttQoS.AT_MOST_ONCE, false, 0)); |
|
|
@ -435,6 +447,7 @@ final class MqttClientImpl implements MqttClient { |
|
|
return null; |
|
|
return null; |
|
|
} |
|
|
} |
|
|
if (this.channel.isActive()) { |
|
|
if (this.channel.isActive()) { |
|
|
|
|
|
log.trace("[{}] Sending message {}", channel != null ? channel.id() : "UNKNOWN", message); |
|
|
return this.channel.writeAndFlush(message); |
|
|
return this.channel.writeAndFlush(message); |
|
|
} |
|
|
} |
|
|
return this.channel.newFailedFuture(new ChannelClosedException("Channel is closed!")); |
|
|
return this.channel.newFailedFuture(new ChannelClosedException("Channel is closed!")); |
|
|
@ -450,6 +463,7 @@ final class MqttClientImpl implements MqttClient { |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private Future<Void> createSubscription(String topic, MqttHandler handler, boolean once, MqttQoS qos) { |
|
|
private Future<Void> createSubscription(String topic, MqttHandler handler, boolean once, MqttQoS qos) { |
|
|
|
|
|
log.trace("[{}] Creating subscription to {}", channel != null ? channel.id() : "UNKNOWN", topic); |
|
|
if (this.pendingSubscribeTopics.contains(topic)) { |
|
|
if (this.pendingSubscribeTopics.contains(topic)) { |
|
|
Optional<Map.Entry<Integer, MqttPendingSubscription>> subscriptionEntry = this.pendingSubscriptions.entrySet().stream().filter((e) -> e.getValue().getTopic().equals(topic)).findAny(); |
|
|
Optional<Map.Entry<Integer, MqttPendingSubscription>> subscriptionEntry = this.pendingSubscriptions.entrySet().stream().filter((e) -> e.getValue().getTopic().equals(topic)).findAny(); |
|
|
if (subscriptionEntry.isPresent()) { |
|
|
if (subscriptionEntry.isPresent()) { |
|
|
|