|
|
@ -103,6 +103,7 @@ import java.util.LinkedHashMap; |
|
|
import java.util.List; |
|
|
import java.util.List; |
|
|
import java.util.Map; |
|
|
import java.util.Map; |
|
|
import java.util.Objects; |
|
|
import java.util.Objects; |
|
|
|
|
|
import java.util.Optional; |
|
|
import java.util.Set; |
|
|
import java.util.Set; |
|
|
import java.util.UUID; |
|
|
import java.util.UUID; |
|
|
import java.util.function.Consumer; |
|
|
import java.util.function.Consumer; |
|
|
@ -232,15 +233,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private boolean isSendNewRpcAvailable() { |
|
|
private boolean isSendNewRpcAvailable() { |
|
|
if (rpcSequenceEnabled) { |
|
|
return !rpcSequenceEnabled || toDeviceRpcPendingMap.values().stream().filter(md -> !md.isDelivered()).findAny().isEmpty(); |
|
|
for (ToDeviceRpcRequestMetadata rpc : toDeviceRpcPendingMap.values()) { |
|
|
|
|
|
if (!rpc.isDelivered()) { |
|
|
|
|
|
return false; |
|
|
|
|
|
} |
|
|
|
|
|
} |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
return true; |
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private Rpc createRpc(ToDeviceRpcRequest request, RpcStatus status) { |
|
|
private Rpc createRpc(ToDeviceRpcRequest request, RpcStatus status) { |
|
|
@ -282,16 +275,26 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { |
|
|
|
|
|
|
|
|
void processRemoveRpc(TbActorCtx context, RemoveRpcActorMsg msg) { |
|
|
void processRemoveRpc(TbActorCtx context, RemoveRpcActorMsg msg) { |
|
|
log.debug("[{}] Processing remove rpc command", msg.getRequestId()); |
|
|
log.debug("[{}] Processing remove rpc command", msg.getRequestId()); |
|
|
Integer requestId = null; |
|
|
Map.Entry<Integer, ToDeviceRpcRequestMetadata> entry = null; |
|
|
for (Map.Entry<Integer, ToDeviceRpcRequestMetadata> entry : toDeviceRpcPendingMap.entrySet()) { |
|
|
for (Map.Entry<Integer, ToDeviceRpcRequestMetadata> e : toDeviceRpcPendingMap.entrySet()) { |
|
|
if (entry.getValue().getMsg().getMsg().getId().equals(msg.getRequestId())) { |
|
|
if (e.getValue().getMsg().getMsg().getId().equals(msg.getRequestId())) { |
|
|
requestId = entry.getKey(); |
|
|
entry = e; |
|
|
break; |
|
|
break; |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
if (requestId != null) { |
|
|
if (entry != null) { |
|
|
toDeviceRpcPendingMap.remove(requestId); |
|
|
if (entry.getValue().isDelivered()) { |
|
|
|
|
|
toDeviceRpcPendingMap.remove(entry.getKey()); |
|
|
|
|
|
} else { |
|
|
|
|
|
Optional<Map.Entry<Integer, ToDeviceRpcRequestMetadata>> firstRpc = getFirstRpc(); |
|
|
|
|
|
if (firstRpc.isPresent() && entry.getKey().equals(firstRpc.get().getKey())) { |
|
|
|
|
|
toDeviceRpcPendingMap.remove(entry.getKey()); |
|
|
|
|
|
sendNextPendingRequest(context); |
|
|
|
|
|
} else { |
|
|
|
|
|
toDeviceRpcPendingMap.remove(entry.getKey()); |
|
|
|
|
|
} |
|
|
|
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@ -330,7 +333,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { |
|
|
Set<Integer> sentOneWayIds = new HashSet<>(); |
|
|
Set<Integer> sentOneWayIds = new HashSet<>(); |
|
|
|
|
|
|
|
|
if (rpcSequenceEnabled) { |
|
|
if (rpcSequenceEnabled) { |
|
|
toDeviceRpcPendingMap.entrySet().stream().filter(e -> !e.getValue().isDelivered()).findFirst().ifPresent(processPendingRpc(context, sessionId, nodeId, sentOneWayIds)); |
|
|
getFirstRpc().ifPresent(processPendingRpc(context, sessionId, nodeId, sentOneWayIds)); |
|
|
} else if (sessionType == SessionType.ASYNC) { |
|
|
} else if (sessionType == SessionType.ASYNC) { |
|
|
toDeviceRpcPendingMap.entrySet().forEach(processPendingRpc(context, sessionId, nodeId, sentOneWayIds)); |
|
|
toDeviceRpcPendingMap.entrySet().forEach(processPendingRpc(context, sessionId, nodeId, sentOneWayIds)); |
|
|
} else { |
|
|
} else { |
|
|
@ -340,6 +343,10 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { |
|
|
sentOneWayIds.stream().filter(id -> !toDeviceRpcPendingMap.get(id).getMsg().getMsg().isPersisted()).forEach(toDeviceRpcPendingMap::remove); |
|
|
sentOneWayIds.stream().filter(id -> !toDeviceRpcPendingMap.get(id).getMsg().getMsg().isPersisted()).forEach(toDeviceRpcPendingMap::remove); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private Optional<Map.Entry<Integer, ToDeviceRpcRequestMetadata>> getFirstRpc() { |
|
|
|
|
|
return toDeviceRpcPendingMap.entrySet().stream().filter(e -> !e.getValue().isDelivered()).findFirst(); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
private void sendNextPendingRequest(TbActorCtx context) { |
|
|
private void sendNextPendingRequest(TbActorCtx context) { |
|
|
if (rpcSequenceEnabled) { |
|
|
if (rpcSequenceEnabled) { |
|
|
rpcSubscriptions.forEach((id, s) -> sendPendingRequests(context, id, s.getNodeId())); |
|
|
rpcSubscriptions.forEach((id, s) -> sendPendingRequests(context, id, s.getNodeId())); |
|
|
@ -599,7 +606,9 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { |
|
|
md.setDelivered(true); |
|
|
md.setDelivered(true); |
|
|
} |
|
|
} |
|
|
} else if (status.equals(RpcStatus.TIMEOUT)) { |
|
|
} else if (status.equals(RpcStatus.TIMEOUT)) { |
|
|
if (systemContext.getMaxPersistentRpcRetries() <= md.getRetries()) { |
|
|
Integer maxRpcRetries = md.getMsg().getMsg().getRetries(); |
|
|
|
|
|
maxRpcRetries = maxRpcRetries == null ? systemContext.getMaxRpcRetries() : Math.min(maxRpcRetries, systemContext.getMaxRpcRetries()); |
|
|
|
|
|
if (maxRpcRetries <= md.getRetries()) { |
|
|
toDeviceRpcPendingMap.remove(responseMsg.getRequestId()); |
|
|
toDeviceRpcPendingMap.remove(responseMsg.getRequestId()); |
|
|
status = RpcStatus.FAILED; |
|
|
status = RpcStatus.FAILED; |
|
|
} else { |
|
|
} else { |
|
|
|