|
|
|
@ -99,6 +99,8 @@ public class EdgeGrpcClient implements EdgeRpcClient { |
|
|
|
|
|
|
|
private volatile boolean connected; |
|
|
|
|
|
|
|
private volatile boolean streamActive; |
|
|
|
|
|
|
|
private static final ReentrantLock uplinkMsgLock = new ReentrantLock(); |
|
|
|
|
|
|
|
@Override |
|
|
|
@ -109,6 +111,7 @@ public class EdgeGrpcClient implements EdgeRpcClient { |
|
|
|
Consumer<DownlinkMsg> onDownlink, |
|
|
|
Consumer<Exception> onError) { |
|
|
|
connected = false; |
|
|
|
streamActive = false; |
|
|
|
NettyChannelBuilder builder = NettyChannelBuilder.forAddress(rpcHost, rpcPort) |
|
|
|
.eventLoopGroup(workerGroup) |
|
|
|
.channelType(channelType()) |
|
|
|
@ -147,6 +150,7 @@ public class EdgeGrpcClient implements EdgeRpcClient { |
|
|
|
EdgeRpcServiceGrpc.EdgeRpcServiceStub stub = EdgeRpcServiceGrpc.newStub(channel); |
|
|
|
log.info("[{}] Sending a connect request to the TB!", edgeKey); |
|
|
|
this.inputStream = stub.withCompression("gzip").handleMsgs(initOutputStream(edgeKey, onUplinkResponse, onEdgeUpdate, onDownlink, onError)); |
|
|
|
streamActive = true; |
|
|
|
this.inputStream.onNext(RequestMsg.newBuilder() |
|
|
|
.setMsgType(RequestMsgType.CONNECT_RPC_MESSAGE) |
|
|
|
.setConnectRequestMsg(ConnectRequestMsg.newBuilder() |
|
|
|
@ -219,6 +223,7 @@ public class EdgeGrpcClient implements EdgeRpcClient { |
|
|
|
@Override |
|
|
|
public void onError(Throwable t) { |
|
|
|
connected = false; |
|
|
|
streamActive = false; |
|
|
|
log.warn("[{}] Stream was terminated due to error:", edgeKey, t); |
|
|
|
try { |
|
|
|
EdgeGrpcClient.this.disconnect(true); |
|
|
|
@ -231,6 +236,7 @@ public class EdgeGrpcClient implements EdgeRpcClient { |
|
|
|
@Override |
|
|
|
public void onCompleted() { |
|
|
|
connected = false; |
|
|
|
streamActive = false; |
|
|
|
log.info("[{}] Stream was closed and completed successfully!", edgeKey); |
|
|
|
} |
|
|
|
}; |
|
|
|
@ -239,6 +245,7 @@ public class EdgeGrpcClient implements EdgeRpcClient { |
|
|
|
@Override |
|
|
|
public void disconnect(boolean onError) throws InterruptedException { |
|
|
|
connected = false; |
|
|
|
streamActive = false; |
|
|
|
if (!onError) { |
|
|
|
try { |
|
|
|
if (inputStream != null) { |
|
|
|
@ -280,7 +287,7 @@ public class EdgeGrpcClient implements EdgeRpcClient { |
|
|
|
public void sendUplinkMsg(UplinkMsg msg) { |
|
|
|
uplinkMsgLock.lock(); |
|
|
|
try { |
|
|
|
if (!connected) { |
|
|
|
if (!streamActive) { |
|
|
|
log.debug("Uplink msg is skipped, the cloud session is not established: {}", msg); |
|
|
|
return; |
|
|
|
} |
|
|
|
@ -297,7 +304,7 @@ public class EdgeGrpcClient implements EdgeRpcClient { |
|
|
|
public void sendSyncRequestMsg(boolean fullSyncRequired) { |
|
|
|
uplinkMsgLock.lock(); |
|
|
|
try { |
|
|
|
if (!connected) { |
|
|
|
if (!streamActive) { |
|
|
|
log.debug("Sync request msg is skipped, the cloud session is not established"); |
|
|
|
return; |
|
|
|
} |
|
|
|
@ -317,7 +324,7 @@ public class EdgeGrpcClient implements EdgeRpcClient { |
|
|
|
public void sendDownlinkResponseMsg(DownlinkResponseMsg downlinkResponseMsg) { |
|
|
|
uplinkMsgLock.lock(); |
|
|
|
try { |
|
|
|
if (!connected) { |
|
|
|
if (!streamActive) { |
|
|
|
log.debug("Downlink response msg is skipped, the cloud session is not established: {}", downlinkResponseMsg); |
|
|
|
return; |
|
|
|
} |
|
|
|
|