@ -23,6 +23,7 @@ import com.google.common.util.concurrent.MoreExecutors;
import com.google.protobuf.InvalidProtocolBufferException ;
import com.google.protobuf.InvalidProtocolBufferException ;
import lombok.extern.slf4j.Slf4j ;
import lombok.extern.slf4j.Slf4j ;
import org.apache.commons.collections.CollectionUtils ;
import org.apache.commons.collections.CollectionUtils ;
import org.thingsboard.common.util.JacksonUtil ;
import org.thingsboard.rule.engine.api.RpcError ;
import org.thingsboard.rule.engine.api.RpcError ;
import org.thingsboard.rule.engine.api.msg.DeviceAttributesEventNotificationMsg ;
import org.thingsboard.rule.engine.api.msg.DeviceAttributesEventNotificationMsg ;
import org.thingsboard.rule.engine.api.msg.DeviceCredentialsUpdateNotificationMsg ;
import org.thingsboard.rule.engine.api.msg.DeviceCredentialsUpdateNotificationMsg ;
@ -38,12 +39,17 @@ import org.thingsboard.server.common.data.edge.EdgeEventActionType;
import org.thingsboard.server.common.data.edge.EdgeEventType ;
import org.thingsboard.server.common.data.edge.EdgeEventType ;
import org.thingsboard.server.common.data.id.DeviceId ;
import org.thingsboard.server.common.data.id.DeviceId ;
import org.thingsboard.server.common.data.id.EdgeId ;
import org.thingsboard.server.common.data.id.EdgeId ;
import org.thingsboard.server.common.data.id.RpcId ;
import org.thingsboard.server.common.data.id.TenantId ;
import org.thingsboard.server.common.data.id.TenantId ;
import org.thingsboard.server.common.data.kv.AttributeKey ;
import org.thingsboard.server.common.data.kv.AttributeKey ;
import org.thingsboard.server.common.data.kv.AttributeKvEntry ;
import org.thingsboard.server.common.data.kv.AttributeKvEntry ;
import org.thingsboard.server.common.data.kv.KvEntry ;
import org.thingsboard.server.common.data.kv.KvEntry ;
import org.thingsboard.server.common.data.page.PageData ;
import org.thingsboard.server.common.data.page.PageLink ;
import org.thingsboard.server.common.data.relation.EntityRelation ;
import org.thingsboard.server.common.data.relation.EntityRelation ;
import org.thingsboard.server.common.data.relation.RelationTypeGroup ;
import org.thingsboard.server.common.data.relation.RelationTypeGroup ;
import org.thingsboard.server.common.data.rpc.Rpc ;
import org.thingsboard.server.common.data.rpc.RpcStatus ;
import org.thingsboard.server.common.data.rpc.ToDeviceRpcRequestBody ;
import org.thingsboard.server.common.data.rpc.ToDeviceRpcRequestBody ;
import org.thingsboard.server.common.data.security.DeviceCredentials ;
import org.thingsboard.server.common.data.security.DeviceCredentials ;
import org.thingsboard.server.common.data.security.DeviceCredentialsType ;
import org.thingsboard.server.common.data.security.DeviceCredentialsType ;
@ -52,8 +58,8 @@ import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.common.msg.queue.TbCallback ;
import org.thingsboard.server.common.msg.queue.TbCallback ;
import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequest ;
import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequest ;
import org.thingsboard.server.common.msg.timeout.DeviceActorServerSideRpcTimeoutMsg ;
import org.thingsboard.server.common.msg.timeout.DeviceActorServerSideRpcTimeoutMsg ;
import org.thingsboard.server.gen.transport.TransportProtos ;
import org.thingsboard.server.gen.transport.TransportProtos.AttributeUpdateNotificationMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.AttributeUpdateNotificationMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.ClaimDeviceMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.DeviceSessionsCacheEntry ;
import org.thingsboard.server.gen.transport.TransportProtos.DeviceSessionsCacheEntry ;
import org.thingsboard.server.gen.transport.TransportProtos.GetAttributeRequestMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.GetAttributeRequestMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.GetAttributeResponseMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.GetAttributeResponseMsg ;
@ -68,10 +74,12 @@ import org.thingsboard.server.gen.transport.TransportProtos.SessionType;
import org.thingsboard.server.gen.transport.TransportProtos.SubscribeToAttributeUpdatesMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.SubscribeToAttributeUpdatesMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.SubscribeToRPCMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.SubscribeToRPCMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.SubscriptionInfoProto ;
import org.thingsboard.server.gen.transport.TransportProtos.SubscriptionInfoProto ;
import org.thingsboard.server.gen.transport.TransportProtos.ToDevicePersistedRpcResponseMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.ToDeviceRpcRequestMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.ToDeviceRpcRequestMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.ToDeviceRpcResponseMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.ToDeviceRpcResponseMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.ToServerRpcResponseMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.ToServerRpcResponseMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.ToTransportUpdateCredentialsProto ;
import org.thingsboard.server.gen.transport.TransportProtos.TransportToDeviceActorMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.TransportToDeviceActorMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.TsKvProto ;
import org.thingsboard.server.gen.transport.TransportProtos.TsKvProto ;
import org.thingsboard.server.service.rpc.FromDeviceRpcResponse ;
import org.thingsboard.server.service.rpc.FromDeviceRpcResponse ;
@ -162,20 +170,19 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
void processRpcRequest ( TbActorCtx context , ToDeviceRpcRequestActorMsg msg ) {
void processRpcRequest ( TbActorCtx context , ToDeviceRpcRequestActorMsg msg ) {
ToDeviceRpcRequest request = msg . getMsg ( ) ;
ToDeviceRpcRequest request = msg . getMsg ( ) ;
ToDeviceRpcRequestBody body = request . getBody ( ) ;
ToDeviceRpcRequestMsg rpcRequest = creteToDeviceRpcRequestMsg ( request ) ;
ToDeviceRpcRequestMsg rpcRequest = ToDeviceRpcRequestMsg . newBuilder ( )
. setRequestId ( rpcSeq + + )
. setMethodName ( body . getMethod ( ) )
. setParams ( body . getParams ( ) )
. setExpirationTime ( request . getExpirationTime ( ) )
. setRequestIdMSB ( request . getId ( ) . getMostSignificantBits ( ) )
. setRequestIdLSB ( request . getId ( ) . getLeastSignificantBits ( ) )
. build ( ) ;
long timeout = request . getExpirationTime ( ) - System . currentTimeMillis ( ) ;
long timeout = request . getExpirationTime ( ) - System . currentTimeMillis ( ) ;
boolean persisted = request . isPersisted ( ) ;
if ( timeout < = 0 ) {
if ( timeout < = 0 ) {
log . debug ( "[{}][{}] Ignoring message due to exp time reached, {}" , deviceId , request . getId ( ) , request . getExpirationTime ( ) ) ;
log . debug ( "[{}][{}] Ignoring message due to exp time reached, {}" , deviceId , request . getId ( ) , request . getExpirationTime ( ) ) ;
if ( persisted ) {
createRpc ( request , RpcStatus . TIMEOUT ) ;
}
return ;
return ;
} else if ( persisted ) {
createRpc ( request , RpcStatus . QUEUED ) ;
}
}
boolean sent ;
boolean sent ;
@ -192,10 +199,16 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
syncSessionSet . add ( key ) ;
syncSessionSet . add ( key ) ;
}
}
} ) ;
} ) ;
log . trace ( "46) Rpc syncSessionSet [{}] subscription after sent [{}]" , syncSessionSet , rpcSubscriptions ) ;
log . trace ( "46) Rpc syncSessionSet [{}] subscription after sent [{}]" , syncSessionSet , rpcSubscriptions ) ;
syncSessionSet . forEach ( rpcSubscriptions : : remove ) ;
syncSessionSet . forEach ( rpcSubscriptions : : remove ) ;
}
}
if ( persisted & & ! ( sent | | request . isOneway ( ) ) ) {
ObjectNode response = JacksonUtil . newObjectNode ( ) ;
response . put ( "rpcId" , request . getId ( ) . toString ( ) ) ;
systemContext . getTbCoreDeviceRpcService ( ) . processRpcResponseFromDeviceActor ( new FromDeviceRpcResponse ( msg . getMsg ( ) . getId ( ) , JacksonUtil . toString ( response ) , null ) ) ;
}
if ( request . isOneway ( ) & & sent ) {
if ( request . isOneway ( ) & & sent ) {
log . debug ( "[{}] Rpc command response sent [{}]!" , deviceId , request . getId ( ) ) ;
log . debug ( "[{}] Rpc command response sent [{}]!" , deviceId , request . getId ( ) ) ;
systemContext . getTbCoreDeviceRpcService ( ) . processRpcResponseFromDeviceActor ( new FromDeviceRpcResponse ( msg . getMsg ( ) . getId ( ) , null , null ) ) ;
systemContext . getTbCoreDeviceRpcService ( ) . processRpcResponseFromDeviceActor ( new FromDeviceRpcResponse ( msg . getMsg ( ) . getId ( ) , null , null ) ) ;
@ -209,6 +222,32 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
}
}
}
}
private Rpc createRpc ( ToDeviceRpcRequest request , RpcStatus status ) {
Rpc rpc = new Rpc ( new RpcId ( request . getId ( ) ) ) ;
rpc . setCreatedTime ( System . currentTimeMillis ( ) ) ;
rpc . setTenantId ( tenantId ) ;
rpc . setDeviceId ( deviceId ) ;
rpc . setExpirationTime ( request . getExpirationTime ( ) ) ;
rpc . setRequest ( JacksonUtil . valueToTree ( request ) ) ;
rpc . setStatus ( status ) ;
systemContext . getTbRpcService ( ) . save ( tenantId , rpc ) ;
return systemContext . getTbRpcService ( ) . save ( tenantId , rpc ) ;
}
private ToDeviceRpcRequestMsg creteToDeviceRpcRequestMsg ( ToDeviceRpcRequest request ) {
ToDeviceRpcRequestBody body = request . getBody ( ) ;
return ToDeviceRpcRequestMsg . newBuilder ( )
. setRequestId ( rpcSeq + + )
. setMethodName ( body . getMethod ( ) )
. setParams ( body . getParams ( ) )
. setExpirationTime ( request . getExpirationTime ( ) )
. setRequestIdMSB ( request . getId ( ) . getMostSignificantBits ( ) )
. setRequestIdLSB ( request . getId ( ) . getLeastSignificantBits ( ) )
. setOneway ( request . isOneway ( ) )
. setPersisted ( request . isPersisted ( ) )
. build ( ) ;
}
void processRpcResponsesFromEdge ( TbActorCtx context , FromDeviceRpcResponseActorMsg responseMsg ) {
void processRpcResponsesFromEdge ( TbActorCtx context , FromDeviceRpcResponseActorMsg responseMsg ) {
log . debug ( "[{}] Processing rpc command response from edge session" , deviceId ) ;
log . debug ( "[{}] Processing rpc command response from edge session" , deviceId ) ;
ToDeviceRpcRequestMetadata requestMd = toDeviceRpcPendingMap . remove ( responseMsg . getRequestId ( ) ) ;
ToDeviceRpcRequestMetadata requestMd = toDeviceRpcPendingMap . remove ( responseMsg . getRequestId ( ) ) ;
@ -230,6 +269,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
ToDeviceRpcRequestMetadata requestMd = toDeviceRpcPendingMap . remove ( msg . getId ( ) ) ;
ToDeviceRpcRequestMetadata requestMd = toDeviceRpcPendingMap . remove ( msg . getId ( ) ) ;
if ( requestMd ! = null ) {
if ( requestMd ! = null ) {
log . debug ( "[{}] RPC request [{}] timeout detected!" , deviceId , msg . getId ( ) ) ;
log . debug ( "[{}] RPC request [{}] timeout detected!" , deviceId , msg . getId ( ) ) ;
systemContext . getTbRpcService ( ) . save ( tenantId , new RpcId ( requestMd . getMsg ( ) . getMsg ( ) . getId ( ) ) , RpcStatus . TIMEOUT , null ) ;
systemContext . getTbCoreDeviceRpcService ( ) . processRpcResponseFromDeviceActor ( new FromDeviceRpcResponse ( requestMd . getMsg ( ) . getMsg ( ) . getId ( ) ,
systemContext . getTbCoreDeviceRpcService ( ) . processRpcResponseFromDeviceActor ( new FromDeviceRpcResponse ( requestMd . getMsg ( ) . getMsg ( ) . getId ( ) ,
null , requestMd . isSent ( ) ? RpcError . TIMEOUT : RpcError . NO_ACTIVE_CONNECTION ) ) ;
null , requestMd . isSent ( ) ? RpcError . TIMEOUT : RpcError . NO_ACTIVE_CONNECTION ) ) ;
}
}
@ -271,7 +311,10 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
. setExpirationTime ( request . getExpirationTime ( ) )
. setExpirationTime ( request . getExpirationTime ( ) )
. setRequestIdMSB ( request . getId ( ) . getMostSignificantBits ( ) )
. setRequestIdMSB ( request . getId ( ) . getMostSignificantBits ( ) )
. setRequestIdLSB ( request . getId ( ) . getLeastSignificantBits ( ) )
. setRequestIdLSB ( request . getId ( ) . getLeastSignificantBits ( ) )
. setOneway ( request . isOneway ( ) )
. setPersisted ( request . isPersisted ( ) )
. build ( ) ;
. build ( ) ;
sendToTransport ( rpcRequest , sessionId , nodeId ) ;
sendToTransport ( rpcRequest , sessionId , nodeId ) ;
} ;
} ;
}
}
@ -279,31 +322,39 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
void process ( TbActorCtx context , TransportToDeviceActorMsgWrapper wrapper ) {
void process ( TbActorCtx context , TransportToDeviceActorMsgWrapper wrapper ) {
TransportToDeviceActorMsg msg = wrapper . getMsg ( ) ;
TransportToDeviceActorMsg msg = wrapper . getMsg ( ) ;
TbCallback callback = wrapper . getCallback ( ) ;
TbCallback callback = wrapper . getCallback ( ) ;
var sessionInfo = msg . getSessionInfo ( ) ;
if ( msg . hasSessionEvent ( ) ) {
if ( msg . hasSessionEvent ( ) ) {
processSessionStateMsgs ( msg . getSessionInfo ( ) , msg . getSessionEvent ( ) ) ;
processSessionStateMsgs ( sessionInfo , msg . getSessionEvent ( ) ) ;
}
}
if ( msg . hasSubscribeToAttributes ( ) ) {
if ( msg . hasSubscribeToAttributes ( ) ) {
processSubscriptionCommands ( context , m sg . getS essionInfo( ) , msg . getSubscribeToAttributes ( ) ) ;
processSubscriptionCommands ( context , sessionInfo , msg . getSubscribeToAttributes ( ) ) ;
}
}
if ( msg . hasSubscribeToRPC ( ) ) {
if ( msg . hasSubscribeToRPC ( ) ) {
processSubscriptionCommands ( context , msg . getSessionInfo ( ) , msg . getSubscribeToRPC ( ) ) ;
processSubscriptionCommands ( context , sessionInfo , msg . getSubscribeToRPC ( ) ) ;
}
if ( msg . hasSendPendingRPC ( ) ) {
sendPendingRequests ( context , getSessionId ( sessionInfo ) , sessionInfo ) ;
}
}
if ( msg . hasGetAttributes ( ) ) {
if ( msg . hasGetAttributes ( ) ) {
handleGetAttributesRequest ( context , msg . getSessionInfo ( ) , msg . getGetAttributes ( ) ) ;
handleGetAttributesRequest ( context , sessionInfo , msg . getGetAttributes ( ) ) ;
}
}
if ( msg . hasToDeviceRPCCallResponse ( ) ) {
if ( msg . hasToDeviceRPCCallResponse ( ) ) {
processRpcResponses ( context , m sg . getS essionInfo( ) , msg . getToDeviceRPCCallResponse ( ) ) ;
processRpcResponses ( context , sessionInfo , msg . getToDeviceRPCCallResponse ( ) ) ;
}
}
if ( msg . hasSubscriptionInfo ( ) ) {
if ( msg . hasSubscriptionInfo ( ) ) {
handleSessionActivity ( context , m sg . getS essionInfo( ) , msg . getSubscriptionInfo ( ) ) ;
handleSessionActivity ( context , sessionInfo , msg . getSubscriptionInfo ( ) ) ;
}
}
if ( msg . hasClaimDevice ( ) ) {
if ( msg . hasClaimDevice ( ) ) {
handleClaimDeviceMsg ( context , msg . getSessionInfo ( ) , msg . getClaimDevice ( ) ) ;
handleClaimDeviceMsg ( context , sessionInfo , msg . getClaimDevice ( ) ) ;
}
if ( msg . hasPersistedRpcResponseMsg ( ) ) {
processPersistedRpcResponses ( context , sessionInfo , msg . getPersistedRpcResponseMsg ( ) ) ;
}
}
callback . onSuccess ( ) ;
callback . onSuccess ( ) ;
}
}
private void handleClaimDeviceMsg ( TbActorCtx context , SessionInfoProto sessionInfo , TransportProtos . ClaimDeviceMsg msg ) {
private void handleClaimDeviceMsg ( TbActorCtx context , SessionInfoProto sessionInfo , ClaimDeviceMsg msg ) {
DeviceId deviceId = new DeviceId ( new UUID ( msg . getDeviceIdMSB ( ) , msg . getDeviceIdLSB ( ) ) ) ;
DeviceId deviceId = new DeviceId ( new UUID ( msg . getDeviceIdMSB ( ) , msg . getDeviceIdLSB ( ) ) ) ;
systemContext . getClaimDevicesService ( ) . registerClaimingInfo ( tenantId , deviceId , msg . getSecretKey ( ) , msg . getDurationMs ( ) ) ;
systemContext . getClaimDevicesService ( ) . registerClaimingInfo ( tenantId , deviceId , msg . getSecretKey ( ) , msg . getDurationMs ( ) ) ;
}
}
@ -442,11 +493,22 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
if ( success ) {
if ( success ) {
systemContext . getTbCoreDeviceRpcService ( ) . processRpcResponseFromDeviceActor ( new FromDeviceRpcResponse ( requestMd . getMsg ( ) . getMsg ( ) . getId ( ) ,
systemContext . getTbCoreDeviceRpcService ( ) . processRpcResponseFromDeviceActor ( new FromDeviceRpcResponse ( requestMd . getMsg ( ) . getMsg ( ) . getId ( ) ,
responseMsg . getPayload ( ) , null ) ) ;
responseMsg . getPayload ( ) , null ) ) ;
if ( requestMd . getMsg ( ) . getMsg ( ) . isPersisted ( ) ) {
systemContext . getTbRpcService ( ) . save ( tenantId , new RpcId ( requestMd . getMsg ( ) . getMsg ( ) . getId ( ) ) , RpcStatus . SUCCESSFUL , JacksonUtil . toJsonNode ( responseMsg . getPayload ( ) ) ) ;
}
} else {
} else {
log . debug ( "[{}] Rpc command response [{}] is stale!" , deviceId , responseMsg . getRequestId ( ) ) ;
log . debug ( "[{}] Rpc command response [{}] is stale!" , deviceId , responseMsg . getRequestId ( ) ) ;
if ( requestMd . getMsg ( ) . getMsg ( ) . isPersisted ( ) ) {
systemContext . getTbRpcService ( ) . save ( tenantId , new RpcId ( requestMd . getMsg ( ) . getMsg ( ) . getId ( ) ) , RpcStatus . FAILED , JacksonUtil . toJsonNode ( responseMsg . getPayload ( ) ) ) ;
}
}
}
}
}
private void processPersistedRpcResponses ( TbActorCtx context , SessionInfoProto sessionInfo , ToDevicePersistedRpcResponseMsg responseMsg ) {
UUID rpcId = new UUID ( responseMsg . getRequestIdMSB ( ) , responseMsg . getRequestIdLSB ( ) ) ;
systemContext . getTbRpcService ( ) . save ( tenantId , new RpcId ( rpcId ) , RpcStatus . valueOf ( responseMsg . getStatus ( ) ) , null ) ;
}
private void processSubscriptionCommands ( TbActorCtx context , SessionInfoProto sessionInfo , SubscribeToAttributeUpdatesMsg subscribeCmd ) {
private void processSubscriptionCommands ( TbActorCtx context , SessionInfoProto sessionInfo , SubscribeToAttributeUpdatesMsg subscribeCmd ) {
UUID sessionId = getSessionId ( sessionInfo ) ;
UUID sessionId = getSessionId ( sessionInfo ) ;
if ( subscribeCmd . getUnsubscribe ( ) ) {
if ( subscribeCmd . getUnsubscribe ( ) ) {
@ -565,7 +627,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
void notifyTransportAboutProfileUpdate ( UUID sessionId , SessionInfoMetaData sessionMd , DeviceCredentials deviceCredentials ) {
void notifyTransportAboutProfileUpdate ( UUID sessionId , SessionInfoMetaData sessionMd , DeviceCredentials deviceCredentials ) {
log . info ( "2) LwM2Mtype: " ) ;
log . info ( "2) LwM2Mtype: " ) ;
TransportProtos . T oTransportUpdateCredentialsProto . Builder notification = TransportProtos . ToTransportUpdateCredentialsProto . newBuilder ( ) ;
ToTransportUpdateCredentialsProto . Builder notification = ToTransportUpdateCredentialsProto . newBuilder ( ) ;
notification . addCredentialsId ( deviceCredentials . getCredentialsId ( ) ) ;
notification . addCredentialsId ( deviceCredentials . getCredentialsId ( ) ) ;
notification . addCredentialsValue ( deviceCredentials . getCredentialsValue ( ) ) ;
notification . addCredentialsValue ( deviceCredentials . getCredentialsValue ( ) ) ;
ToTransportMsg msg = ToTransportMsg . newBuilder ( )
ToTransportMsg msg = ToTransportMsg . newBuilder ( )
@ -640,7 +702,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
ListenableFuture < EdgeEvent > future = systemContext . getEdgeEventService ( ) . saveAsync ( edgeEvent ) ;
ListenableFuture < EdgeEvent > future = systemContext . getEdgeEventService ( ) . saveAsync ( edgeEvent ) ;
Futures . addCallback ( future , new FutureCallback < EdgeEvent > ( ) {
Futures . addCallback ( future , new FutureCallback < EdgeEvent > ( ) {
@Override
@Override
public void onSuccess ( EdgeEvent result ) {
public void onSuccess ( EdgeEvent result ) {
systemContext . getClusterService ( ) . onEdgeEventUpdate ( tenantId , edgeId ) ;
systemContext . getClusterService ( ) . onEdgeEventUpdate ( tenantId , edgeId ) ;
}
}
@ -756,8 +818,26 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
. addAllSessions ( sessionsList ) . build ( ) . toByteArray ( ) ) ;
. addAllSessions ( sessionsList ) . build ( ) . toByteArray ( ) ) ;
}
}
void initSessionTimeout ( TbActorCtx ctx ) {
void init ( TbActorCtx ctx ) {
schedulePeriodicMsgWithDelay ( ctx , SessionTimeoutCheckMsg . instance ( ) , systemContext . getSessionReportTimeout ( ) , systemContext . getSessionReportTimeout ( ) ) ;
schedulePeriodicMsgWithDelay ( ctx , SessionTimeoutCheckMsg . instance ( ) , systemContext . getSessionReportTimeout ( ) , systemContext . getSessionReportTimeout ( ) ) ;
PageLink pageLink = new PageLink ( 1024 ) ;
PageData < Rpc > pageData ;
do {
pageData = systemContext . getTbRpcService ( ) . findAllByDeviceIdAndStatus ( tenantId , deviceId , RpcStatus . QUEUED , pageLink ) ;
pageData . getData ( ) . forEach ( rpc - > {
ToDeviceRpcRequest msg = JacksonUtil . convertValue ( rpc . getRequest ( ) , ToDeviceRpcRequest . class ) ;
long timeout = rpc . getExpirationTime ( ) - System . currentTimeMillis ( ) ;
if ( timeout < = 0 ) {
rpc . setStatus ( RpcStatus . TIMEOUT ) ;
systemContext . getTbRpcService ( ) . save ( tenantId , rpc ) ;
} else {
registerPendingRpcRequest ( ctx , new ToDeviceRpcRequestActorMsg ( systemContext . getServiceId ( ) , msg ) , false , creteToDeviceRpcRequestMsg ( msg ) , timeout ) ;
}
} ) ;
if ( pageData . hasNext ( ) ) {
pageLink = pageLink . nextPageLink ( ) ;
}
} while ( pageData . hasNext ( ) ) ;
}
}
void checkSessionsTimeout ( ) {
void checkSessionsTimeout ( ) {