@ -182,13 +182,15 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
void processRpcRequest ( TbActorCtx context , ToDeviceRpcRequestActorMsg msg ) {
ToDeviceRpcRequest request = msg . getMsg ( ) ;
UUID rpcId = request . getId ( ) ;
log . debug ( "[{}][{}] Received rpc request to process ..." , deviceId , rpcId ) ;
ToDeviceRpcRequestMsg rpcRequest = creteToDeviceRpcRequestMsg ( request ) ;
long timeout = request . getExpirationTime ( ) - System . currentTimeMillis ( ) ;
boolean persisted = request . isPersisted ( ) ;
if ( timeout < = 0 ) {
log . debug ( "[{}][{}] Ignoring message due to exp time reached, {}" , deviceId , request . get Id ( ) , request . getExpirationTime ( ) ) ;
log . debug ( "[{}][{}] Ignoring message due to exp time reached, {}" , deviceId , rpc Id , request . getExpirationTime ( ) ) ;
if ( persisted ) {
createRpc ( request , RpcStatus . EXPIRED ) ;
}
@ -198,8 +200,9 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
}
boolean sent = false ;
int requestId = rpcRequest . getRequestId ( ) ;
if ( systemContext . isEdgesEnabled ( ) & & edgeId ! = null ) {
log . debug ( "[{}][{}] device is related to edge [{}]. Saving RPC request to edge queue" , tenantId , deviceId , edgeId . getId ( ) ) ;
log . debug ( "[{}][{}] device is related to edge: [{}]. Saving rpc request: [{}][{}] to edge queue" , tenantId , deviceId , edgeId . getId ( ) , rpcId , requestId ) ;
try {
saveRpcRequestToEdgeQueue ( request , rpcRequest . getRequestId ( ) ) . get ( ) ;
sent = true ;
@ -209,10 +212,11 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
} else if ( isSendNewRpcAvailable ( ) ) {
sent = rpcSubscriptions . size ( ) > 0 ;
Set < UUID > syncSessionSet = new HashSet < > ( ) ;
rpcSubscriptions . forEach ( ( key , value ) - > {
sendToTransport ( rpcRequest , key , value . getNodeId ( ) ) ;
if ( SessionType . SYNC = = value . getType ( ) ) {
syncSessionSet . add ( key ) ;
rpcSubscriptions . forEach ( ( sessionId , sessionInfo ) - > {
log . debug ( "[{}][{}][{}][{}] send rpc request to transport ..." , deviceId , sessionId , rpcId , requestId ) ;
sendToTransport ( rpcRequest , sessionId , sessionInfo . getNodeId ( ) ) ;
if ( SessionType . SYNC = = sessionInfo . getType ( ) ) {
syncSessionSet . add ( sessionId ) ;
}
} ) ;
log . trace ( "Rpc syncSessionSet [{}] subscription after sent [{}]" , syncSessionSet , rpcSubscriptions ) ;
@ -221,28 +225,32 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
if ( persisted ) {
ObjectNode response = JacksonUtil . newObjectNode ( ) ;
response . put ( "rpcId" , request . get Id ( ) . toString ( ) ) ;
response . put ( "rpcId" , rpc Id . toString ( ) ) ;
systemContext . getTbCoreDeviceRpcService ( ) . processRpcResponseFromDeviceActor ( new FromDeviceRpcResponse ( msg . getMsg ( ) . getId ( ) , JacksonUtil . toString ( response ) , null ) ) ;
}
if ( ! persisted & & request . isOneway ( ) & & sent ) {
log . debug ( "[{}] Rpc command response sent [{}]!" , deviceId , request . getId ( ) ) ;
log . debug ( "[{}] Rpc command response sent [{}][{}] !" , deviceId , rpcId , requestId ) ;
systemContext . getTbCoreDeviceRpcService ( ) . processRpcResponseFromDeviceActor ( new FromDeviceRpcResponse ( msg . getMsg ( ) . getId ( ) , null , null ) ) ;
} else {
registerPendingRpcRequest ( context , msg , sent , rpcRequest , timeout ) ;
}
if ( sent ) {
log . debug ( "[{}] RPC request {} is sent!" , deviceId , request . getId ( ) ) ;
log . debug ( "[{}][{}][{}] Rpc request is sent!" , deviceId , rpcId , requestId ) ;
} else {
log . debug ( "[{}] RPC request {} is NOT sent!" , deviceId , request . getId ( ) ) ;
log . debug ( "[{}][{}][{}] Rpc request is NOT sent!" , deviceId , rpcId , requestId ) ;
}
}
private UUID getRpcIdFromRequest ( ToDeviceRpcRequestMsg request ) {
return new UUID ( request . getRequestIdMSB ( ) , request . getRequestIdLSB ( ) ) ;
}
private boolean isSendNewRpcAvailable ( ) {
return ! rpcSequential | | toDeviceRpcPendingMap . values ( ) . stream ( ) . filter ( md - > ! md . isDelivered ( ) ) . findAny ( ) . isEmpty ( ) ;
}
private Rpc createRpc ( ToDeviceRpcRequest request , RpcStatus status ) {
private void createRpc ( ToDeviceRpcRequest request , RpcStatus status ) {
Rpc rpc = new Rpc ( new RpcId ( request . getId ( ) ) ) ;
rpc . setCreatedTime ( System . currentTimeMillis ( ) ) ;
rpc . setTenantId ( tenantId ) ;
@ -251,7 +259,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
rpc . setRequest ( JacksonUtil . valueToTree ( request ) ) ;
rpc . setStatus ( status ) ;
rpc . setAdditionalInfo ( JacksonUtil . toJsonNode ( request . getAdditionalInfo ( ) ) ) ;
return systemContext . getTbRpcService ( ) . save ( tenantId , rpc ) ;
systemContext . getTbRpcService ( ) . save ( tenantId , rpc ) ;
}
private ToDeviceRpcRequestMsg creteToDeviceRpcRequestMsg ( ToDeviceRpcRequest request ) {
@ -280,7 +288,8 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
}
void processRemoveRpc ( TbActorCtx context , RemoveRpcActorMsg msg ) {
log . debug ( "[{}] Processing remove rpc command" , msg . getRequestId ( ) ) ;
UUID requestId = msg . getRequestId ( ) ;
log . debug ( "[{}][{}] Received remove rpc request ..." , deviceId , requestId ) ;
Map . Entry < Integer , ToDeviceRpcRequestMetadata > entry = null ;
for ( Map . Entry < Integer , ToDeviceRpcRequestMetadata > e : toDeviceRpcPendingMap . entrySet ( ) ) {
if ( e . getValue ( ) . getMsg ( ) . getMsg ( ) . getId ( ) . equals ( msg . getRequestId ( ) ) ) {
@ -290,36 +299,42 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
}
if ( entry ! = null ) {
Integer key = entry . getKey ( ) ;
if ( entry . getValue ( ) . isDelivered ( ) ) {
toDeviceRpcPendingMap . remove ( entry . getKey ( ) ) ;
toDeviceRpcPendingMap . remove ( key ) ;
} else {
Optional < Map . Entry < Integer , ToDeviceRpcRequestMetadata > > firstRpc = getFirstRpc ( ) ;
if ( firstRpc . isPresent ( ) & & entry . getKey ( ) . equals ( firstRpc . get ( ) . getKey ( ) ) ) {
toDeviceRpcPendingMap . remove ( entry . getKey ( ) ) ;
if ( firstRpc . isPresent ( ) & & key . equals ( firstRpc . get ( ) . getKey ( ) ) ) {
toDeviceRpcPendingMap . remove ( key ) ;
log . debug ( "[{}][{}][{}] Removed pending rpc! Going to send next pending request ..." , deviceId , requestId , key ) ;
sendNextPendingRequest ( context ) ;
} else {
toDeviceRpcPendingMap . remove ( entry . getKey ( ) ) ;
toDeviceRpcPendingMap . remove ( key ) ;
}
}
}
}
private void registerPendingRpcRequest ( TbActorCtx context , ToDeviceRpcRequestActorMsg msg , boolean sent , ToDeviceRpcRequestMsg rpcRequest , long timeout ) {
log . debug ( "[{}][{}][{}] Registering pending rpc request..." , deviceId , getRpcIdFromRequest ( rpcRequest ) , rpcRequest . getRequestId ( ) ) ;
toDeviceRpcPendingMap . put ( rpcRequest . getRequestId ( ) , new ToDeviceRpcRequestMetadata ( msg , sent ) ) ;
DeviceActorServerSideRpcTimeoutMsg timeoutMsg = new DeviceActorServerSideRpcTimeoutMsg ( rpcRequest . getRequestId ( ) , timeout ) ;
scheduleMsgWithDelay ( context , timeoutMsg , timeoutMsg . getTimeout ( ) ) ;
}
void processServerSideRpcTimeout ( TbActorCtx context , DeviceActorServerSideRpcTimeoutMsg msg ) {
ToDeviceRpcRequestMetadata requestMd = toDeviceRpcPendingMap . remove ( msg . getId ( ) ) ;
Integer requestId = msg . getId ( ) ;
ToDeviceRpcRequestMetadata requestMd = toDeviceRpcPendingMap . remove ( requestId ) ;
if ( requestMd ! = null ) {
log . debug ( "[{}] RPC request [{}] timeout detected!" , deviceId , msg . getId ( ) ) ;
UUID rpcId = requestMd . getMsg ( ) . getMsg ( ) . getId ( ) ;
log . debug ( "[{}][{}][{}] Rpc request timeout detected!" , deviceId , rpcId , requestId ) ;
if ( requestMd . getMsg ( ) . getMsg ( ) . isPersisted ( ) ) {
systemContext . getTbRpcService ( ) . save ( tenantId , new RpcId ( requestMd . getMsg ( ) . getMsg ( ) . get Id ( ) ) , RpcStatus . EXPIRED , null ) ;
systemContext . getTbRpcService ( ) . save ( tenantId , new RpcId ( rpc Id ) , RpcStatus . EXPIRED , null ) ;
}
systemContext . getTbCoreDeviceRpcService ( ) . processRpcResponseFromDeviceActor ( new FromDeviceRpcResponse ( requestMd . getMsg ( ) . getMsg ( ) . get Id ( ) ,
systemContext . getTbCoreDeviceRpcService ( ) . processRpcResponseFromDeviceActor ( new FromDeviceRpcResponse ( rpc Id ,
null , requestMd . isSent ( ) ? RpcError . TIMEOUT : RpcError . NO_ACTIVE_CONNECTION ) ) ;
if ( ! requestMd . isDelivered ( ) ) {
log . debug ( "[{}][{}][{}] Pending rpc timeout detected! Going to send next pending request ..." , deviceId , rpcId , requestId ) ;
sendNextPendingRequest ( context ) ;
}
}
@ -328,13 +343,13 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
private void sendPendingRequests ( TbActorCtx context , UUID sessionId , String nodeId ) {
SessionType sessionType = getSessionType ( sessionId ) ;
if ( ! toDeviceRpcPendingMap . isEmpty ( ) ) {
log . debug ( "[{}] Pushing {} pending RPC messages to new async session [{}] " , deviceId , toDeviceRpcPendingMap . size ( ) , sessionId ) ;
log . debug ( "[{}][{}] Pushing {} pending rpc messages to new async session! " , deviceId , sessionId , toDeviceRpcPendingMap . size ( ) ) ;
if ( sessionType = = SessionType . SYNC ) {
log . debug ( "[{}] Cleanup sync rpc session [{}]" , deviceId , sessionId ) ;
rpcSubscriptions . remove ( sessionId ) ;
}
} else {
log . debug ( "[{}] No pending RPC messages for new async session [{}]" , deviceId , sessionId ) ;
log . debug ( "[{}] No pending rpc messages for new async session [{}]" , deviceId , sessionId ) ;
}
Set < Integer > sentOneWayIds = new HashSet < > ( ) ;
@ -363,12 +378,13 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
return entry - > {
ToDeviceRpcRequest request = entry . getValue ( ) . getMsg ( ) . getMsg ( ) ;
ToDeviceRpcRequestBody body = request . getBody ( ) ;
Integer requestId = entry . getKey ( ) ;
if ( request . isOneway ( ) & & ! rpcSequential ) {
sentOneWayIds . add ( entry . getKey ( ) ) ;
sentOneWayIds . add ( requestId ) ;
systemContext . getTbCoreDeviceRpcService ( ) . processRpcResponseFromDeviceActor ( new FromDeviceRpcResponse ( request . getId ( ) , null , null ) ) ;
}
ToDeviceRpcRequestMsg rpcRequest = ToDeviceRpcRequestMsg . newBuilder ( )
. setRequestId ( entry . getKey ( ) )
. setRequestId ( requestId )
. setMethodName ( body . getMethod ( ) )
. setParams ( body . getParams ( ) )
. setExpirationTime ( request . getExpirationTime ( ) )
@ -377,6 +393,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
. setOneway ( request . isOneway ( ) )
. setPersisted ( request . isPersisted ( ) )
. build ( ) ;
log . debug ( "[{}][{}][{}][{}] Send pending rpc request to transport ..." , deviceId , sessionId , getRpcIdFromRequest ( rpcRequest ) , requestId ) ;
sendToTransport ( rpcRequest , sessionId , nodeId ) ;
} ;
}
@ -569,18 +586,26 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
private void processRpcResponses ( TbActorCtx context , SessionInfoProto sessionInfo , ToDeviceRpcResponseMsg responseMsg ) {
UUID sessionId = getSessionId ( sessionInfo ) ;
log . debug ( "[{}] Processing rpc command response [ {}] " , deviceId , sessionId ) ;
log . debug ( "[{}][{}] Processing rpc command response: {}" , deviceId , sessionId , responseMsg ) ;
ToDeviceRpcRequestMetadata requestMd = toDeviceRpcPendingMap . remove ( responseMsg . getRequestId ( ) ) ;
boolean success = requestMd ! = null ;
if ( success ) {
boolean delivered = requestMd . isDelivered ( ) ;
boolean hasError = StringUtils . isNotEmpty ( responseMsg . getError ( ) ) ;
try {
String payload = hasError ? responseMsg . getError ( ) : responseMsg . getPayload ( ) ;
String payload ;
if ( hasError ) {
payload = responseMsg . getError ( ) ;
} else if ( delivered ) {
payload = responseMsg . getPayload ( ) ;
} else {
payload = "Received response for undelivered rpc: " + responseMsg . getPayload ( ) ;
}
systemContext . getTbCoreDeviceRpcService ( ) . processRpcResponseFromDeviceActor (
new FromDeviceRpcResponse ( requestMd . getMsg ( ) . getMsg ( ) . getId ( ) ,
payload , null ) ) ;
if ( requestMd . getMsg ( ) . getMsg ( ) . isPersisted ( ) ) {
RpcStatus status = hasError ? RpcStatus . FAILED : RpcStatus . SUCCESSFUL ;
RpcStatus status = hasError | | ! delivered ? RpcStatus . FAILED : RpcStatus . SUCCESSFUL ;
JsonNode response ;
try {
response = JacksonUtil . toJsonNode ( payload ) ;
@ -590,25 +615,30 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
systemContext . getTbRpcService ( ) . save ( tenantId , new RpcId ( requestMd . getMsg ( ) . getMsg ( ) . getId ( ) ) , status , response ) ;
}
} finally {
if ( hasError & & ! requestMd . isDelivered ( ) ) {
if ( ! delivered ) {
String errorResponse = hasError ? "error" : "" ;
log . debug ( "[{}][{}][{}] Received {} response for undelivered rpc! Going to send next pending request ..." , deviceId , sessionId , responseMsg . getRequestId ( ) , errorResponse ) ;
sendNextPendingRequest ( context ) ;
}
}
} else {
log . debug ( "[{}] Rpc command response [{}] is stale!" , deviceId , responseMsg . getRequestId ( ) ) ;
log . debug ( "[{}][{}][{}] Rpc command response is stale!" , deviceId , session Id , responseMsg . getRequestId ( ) ) ;
}
}
private void processRpcResponseStatus ( TbActorCtx context , SessionInfoProto sessionInfo , ToDeviceRpcResponseStatusMsg responseMsg ) {
UUID rpcId = new UUID ( responseMsg . getRequestIdMSB ( ) , responseMsg . getRequestIdLSB ( ) ) ;
RpcStatus status = RpcStatus . valueOf ( responseMsg . getStatus ( ) ) ;
ToDeviceRpcRequestMetadata md = toDeviceRpcPendingMap . get ( responseMsg . getRequestId ( ) ) ;
UUID sessionId = getSessionId ( sessionInfo ) ;
int requestId = responseMsg . getRequestId ( ) ;
log . debug ( "[{}][{}][{}][{}] Processing rpc command response status: [{}]" , deviceId , sessionId , rpcId , requestId , status ) ;
ToDeviceRpcRequestMetadata md = toDeviceRpcPendingMap . get ( requestId ) ;
if ( md ! = null ) {
JsonNode response = null ;
if ( status . equals ( RpcStatus . DELIVERED ) ) {
if ( md . getMsg ( ) . getMsg ( ) . isOneway ( ) ) {
toDeviceRpcPendingMap . remove ( responseMsg . getRe questId ( ) ) ;
toDeviceRpcPendingMap . remove ( requestId ) ;
if ( rpcSequential ) {
systemContext . getTbCoreDeviceRpcService ( ) . processRpcResponseFromDeviceActor ( new FromDeviceRpcResponse ( rpcId , null , null ) ) ;
}
@ -619,7 +649,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
Integer maxRpcRetries = md . getMsg ( ) . getMsg ( ) . getRetries ( ) ;
maxRpcRetries = maxRpcRetries = = null ? systemContext . getMaxRpcRetries ( ) : Math . min ( maxRpcRetries , systemContext . getMaxRpcRetries ( ) ) ;
if ( maxRpcRetries < = md . getRetries ( ) ) {
toDeviceRpcPendingMap . remove ( responseMsg . getRe questId ( ) ) ;
toDeviceRpcPendingMap . remove ( requestId ) ;
status = RpcStatus . FAILED ;
response = JacksonUtil . newObjectNode ( ) . put ( "error" , "There was a Timeout and all retry attempts have been exhausted. Retry attempts set: " + maxRpcRetries ) ;
} else {
@ -631,10 +661,11 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
systemContext . getTbRpcService ( ) . save ( tenantId , new RpcId ( rpcId ) , status , response ) ;
}
if ( status ! = RpcStatus . SENT ) {
log . debug ( "[{}][{}][{}][{}] Rpc was {}! Going to send next pending request ..." , deviceId , sessionId , rpcId , requestId , status . name ( ) . toLowerCase ( ) ) ;
sendNextPendingRequest ( context ) ;
}
} else {
log . info ( "[{}][{}] Rpc has already removed from pending map." , deviceId , rpcId ) ;
log . warn ( "[{}][{}][{}][{}] Rpc has already been removed from pending map." , deviceId , sessionId , rpcId , request Id ) ;
}
}
@ -662,7 +693,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
private void processSubscriptionCommands ( TbActorCtx context , SessionInfoProto sessionInfo , SubscribeToRPCMsg subscribeCmd ) {
UUID sessionId = getSessionId ( sessionInfo ) ;
if ( subscribeCmd . getUnsubscribe ( ) ) {
log . debug ( "[{}] Canceling rpc subscription for session [{}]" , deviceId , sessionId ) ;
log . debug ( "[{}] Canceling rpc subscription for session: [{}]" , deviceId , sessionId ) ;
rpcSubscriptions . remove ( sessionId ) ;
} else {
SessionInfoMetaData sessionMD = sessions . get ( sessionId ) ;
@ -670,7 +701,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
sessionMD = new SessionInfoMetaData ( new SessionInfo ( subscribeCmd . getSessionType ( ) , sessionInfo . getNodeId ( ) ) ) ;
}
sessionMD . setSubscribedToRPC ( true ) ;
log . debug ( "[{}] Registering rpc subscription for session [{}] " , deviceId , sessionId ) ;
log . debug ( "[{}] Registered rpc subscription for session: [{}] Going to check for pending requests ... " , deviceId , sessionId ) ;
rpcSubscriptions . put ( sessionId , sessionMD . getSessionInfo ( ) ) ;
sendPendingRequests ( context , sessionId , sessionInfo . getNodeId ( ) ) ;
dumpSessions ( ) ;