Browse Source

Added correct onComplete in case onError

pull/2436/head
Volodymyr Babak 6 years ago
parent
commit
5613baf066
  1. 10
      common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java
  2. 2
      common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java

10
common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java

@ -40,6 +40,7 @@ import org.thingsboard.server.gen.edge.UplinkResponseMsg;
import javax.net.ssl.SSLException;
import java.io.File;
import java.net.URISyntaxException;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.ReentrantLock;
import java.util.function.Consumer;
@ -107,7 +108,7 @@ public class EdgeGrpcClient implements EdgeRpcClient {
} else {
log.error("[{}] Failed to establish the connection! Code: {}. Error message: {}.", edgeKey, connectResponseMsg.getResponseCode(), connectResponseMsg.getErrorMsg());
try {
EdgeGrpcClient.this.disconnect();
EdgeGrpcClient.this.disconnect(true);
} catch (InterruptedException e) {
log.error("[{}] Got interruption during disconnect!", edgeKey, e);
}
@ -136,7 +137,12 @@ public class EdgeGrpcClient implements EdgeRpcClient {
}
@Override
public void disconnect() throws InterruptedException {
public void disconnect(boolean onError) throws InterruptedException {
if (!onError) {
try {
inputStream.onCompleted();
} catch (Exception ignored) {}
}
if (channel != null) {
channel.shutdown().awaitTermination(timeoutSecs, TimeUnit.SECONDS);
}

2
common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java

@ -32,7 +32,7 @@ public interface EdgeRpcClient {
Consumer<DownlinkMsg> onDownlink,
Consumer<Exception> onError);
void disconnect() throws InterruptedException;
void disconnect(boolean onError) throws InterruptedException;
void sendSyncRequestMsg();

Loading…
Cancel
Save