@ -50,9 +50,9 @@ import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.TenantId ;
import org.thingsboard.server.common.data.id.TenantId ;
import org.thingsboard.server.common.data.id.TenantProfileId ;
import org.thingsboard.server.common.data.id.TenantProfileId ;
import org.thingsboard.server.common.data.limit.LimitedApi ;
import org.thingsboard.server.common.data.limit.LimitedApi ;
import org.thingsboard.server.common.data.msg.TbMsgType ;
import org.thingsboard.server.common.data.notification.rule.trigger.RateLimitsTrigger ;
import org.thingsboard.server.common.data.notification.rule.trigger.RateLimitsTrigger ;
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent ;
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent ;
import org.thingsboard.server.common.data.msg.TbMsgType ;
import org.thingsboard.server.common.data.rpc.RpcStatus ;
import org.thingsboard.server.common.data.rpc.RpcStatus ;
import org.thingsboard.server.common.msg.TbMsg ;
import org.thingsboard.server.common.msg.TbMsg ;
import org.thingsboard.server.common.msg.TbMsgMetaData ;
import org.thingsboard.server.common.msg.TbMsgMetaData ;
@ -96,9 +96,9 @@ import org.thingsboard.server.queue.TbQueueProducer;
import org.thingsboard.server.queue.TbQueueRequestTemplate ;
import org.thingsboard.server.queue.TbQueueRequestTemplate ;
import org.thingsboard.server.queue.common.AsyncCallbackTemplate ;
import org.thingsboard.server.queue.common.AsyncCallbackTemplate ;
import org.thingsboard.server.queue.common.TbProtoQueueMsg ;
import org.thingsboard.server.queue.common.TbProtoQueueMsg ;
import org.thingsboard.server.queue.discovery.TopicService ;
import org.thingsboard.server.queue.discovery.PartitionService ;
import org.thingsboard.server.queue.discovery.PartitionService ;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider ;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider ;
import org.thingsboard.server.queue.discovery.TopicService ;
import org.thingsboard.server.queue.provider.TbQueueProducerProvider ;
import org.thingsboard.server.queue.provider.TbQueueProducerProvider ;
import org.thingsboard.server.queue.provider.TbTransportQueueFactory ;
import org.thingsboard.server.queue.provider.TbTransportQueueFactory ;
import org.thingsboard.server.queue.scheduler.SchedulerComponent ;
import org.thingsboard.server.queue.scheduler.SchedulerComponent ;
@ -109,16 +109,13 @@ import org.thingsboard.server.queue.util.TbTransportComponent;
import javax.annotation.PostConstruct ;
import javax.annotation.PostConstruct ;
import javax.annotation.PreDestroy ;
import javax.annotation.PreDestroy ;
import java.util.Collections ;
import java.util.Collections ;
import java.util.HashSet ;
import java.util.LinkedHashMap ;
import java.util.LinkedHashMap ;
import java.util.List ;
import java.util.List ;
import java.util.Map ;
import java.util.Map ;
import java.util.Optional ;
import java.util.Optional ;
import java.util.Random ;
import java.util.Random ;
import java.util.Set ;
import java.util.UUID ;
import java.util.UUID ;
import java.util.concurrent.ConcurrentHashMap ;
import java.util.concurrent.ConcurrentHashMap ;
import java.util.concurrent.ConcurrentMap ;
import java.util.concurrent.ExecutionException ;
import java.util.concurrent.ExecutionException ;
import java.util.concurrent.ExecutorService ;
import java.util.concurrent.ExecutorService ;
import java.util.concurrent.Executors ;
import java.util.concurrent.Executors ;
@ -133,14 +130,12 @@ import java.util.stream.Collectors;
@Slf4j
@Slf4j
@Service
@Service
@TbTransportComponent
@TbTransportComponent
public class DefaultTransportService implements TransportService {
public class DefaultTransportService extends TransportActivityManager implements TransportService {
public static final String OVERWRITE_ACTIVITY_TIME = "overwriteActivityTime" ;
public static final String OVERWRITE_ACTIVITY_TIME = "overwriteActivityTime" ;
public static final String SESSION_EXPIRED_MESSAGE = "Session has expired due to last activity time!" ;
public static final TransportProtos . SessionEventMsg SESSION_EVENT_MSG_OPEN = TransportProtos . SessionEventMsg . newBuilder ( )
public static final TransportProtos . SessionEventMsg SESSION_EVENT_MSG_OPEN = getSessionEventMsg ( TransportProtos . SessionEvent . OPEN ) ;
. setSessionType ( TransportProtos . SessionType . ASYNC )
public static final TransportProtos . SessionEventMsg SESSION_EVENT_MSG_CLOSED = getSessionEventMsg ( TransportProtos . SessionEvent . CLOSED ) ;
. setEvent ( TransportProtos . SessionEvent . OPEN ) . build ( ) ;
public static final TransportProtos . SessionCloseNotificationProto SESSION_CLOSE_NOTIFICATION_PROTO = TransportProtos . SessionCloseNotificationProto . newBuilder ( )
. setMessage ( SESSION_EXPIRED_MESSAGE ) . build ( ) ;
public static final TransportProtos . SubscribeToAttributeUpdatesMsg SUBSCRIBE_TO_ATTRIBUTE_UPDATES_ASYNC_MSG = TransportProtos . SubscribeToAttributeUpdatesMsg . newBuilder ( )
public static final TransportProtos . SubscribeToAttributeUpdatesMsg SUBSCRIBE_TO_ATTRIBUTE_UPDATES_ASYNC_MSG = TransportProtos . SubscribeToAttributeUpdatesMsg . newBuilder ( )
. setSessionType ( TransportProtos . SessionType . ASYNC ) . build ( ) ;
. setSessionType ( TransportProtos . SessionType . ASYNC ) . build ( ) ;
public static final TransportProtos . SubscribeToRPCMsg SUBSCRIBE_TO_RPC_ASYNC_MSG = TransportProtos . SubscribeToRPCMsg . newBuilder ( )
public static final TransportProtos . SubscribeToRPCMsg SUBSCRIBE_TO_RPC_ASYNC_MSG = TransportProtos . SubscribeToRPCMsg . newBuilder ( )
@ -152,10 +147,6 @@ public class DefaultTransportService implements TransportService {
private boolean logEnabled ;
private boolean logEnabled ;
@Value ( "${transport.log.max_length:1024}" )
@Value ( "${transport.log.max_length:1024}" )
private int logMaxLength ;
private int logMaxLength ;
@Value ( "${transport.sessions.inactivity_timeout}" )
private long sessionInactivityTimeout ;
@Value ( "${transport.sessions.report_timeout}" )
private long sessionReportTimeout ;
@Value ( "${transport.client_side_rpc.timeout:60000}" )
@Value ( "${transport.client_side_rpc.timeout:60000}" )
private long clientSideRpcTimeout ;
private long clientSideRpcTimeout ;
@Value ( "${queue.transport.poll_interval}" )
@Value ( "${queue.transport.poll_interval}" )
@ -182,11 +173,9 @@ public class DefaultTransportService implements TransportService {
private final TransportRateLimitService rateLimitService ;
private final TransportRateLimitService rateLimitService ;
private final DataDecodingEncodingService dataDecodingEncodingService ;
private final DataDecodingEncodingService dataDecodingEncodingService ;
private final SchedulerComponent scheduler ;
private final ApplicationEventPublisher eventPublisher ;
private final ApplicationEventPublisher eventPublisher ;
private final TransportResourceCache transportResourceCache ;
private final TransportResourceCache transportResourceCache ;
private final NotificationRuleProcessor notificationRuleProcessor ;
private final NotificationRuleProcessor notificationRuleProcessor ;
private final EntityLimitsCache entityLimitsCache ;
private final EntityLimitsCache entityLimitsCache ;
protected TbQueueRequestTemplate < TbProtoQueueMsg < TransportApiRequestMsg > , TbProtoQueueMsg < TransportApiResponseMsg > > transportApiRequestTemplate ;
protected TbQueueRequestTemplate < TbProtoQueueMsg < TransportApiRequestMsg > , TbProtoQueueMsg < TransportApiResponseMsg > > transportApiRequestTemplate ;
@ -201,8 +190,6 @@ public class DefaultTransportService implements TransportService {
protected ExecutorService transportCallbackExecutor ;
protected ExecutorService transportCallbackExecutor ;
private ExecutorService mainConsumerExecutor ;
private ExecutorService mainConsumerExecutor ;
public final ConcurrentMap < UUID , SessionMetaData > sessions = new ConcurrentHashMap < > ( ) ;
private final ConcurrentMap < UUID , SessionActivityData > sessionsActivity = new ConcurrentHashMap < > ( ) ;
private final Map < String , RpcRequestMetadata > toServerRpcPendingMap = new ConcurrentHashMap < > ( ) ;
private final Map < String , RpcRequestMetadata > toServerRpcPendingMap = new ConcurrentHashMap < > ( ) ;
private volatile boolean stopped = false ;
private volatile boolean stopped = false ;
@ -238,11 +225,11 @@ public class DefaultTransportService implements TransportService {
@PostConstruct
@PostConstruct
public void init ( ) {
public void init ( ) {
super . init ( ) ;
this . ruleEngineProducerStats = statsFactory . createMessagesStats ( StatsType . RULE_ENGINE . getName ( ) + ".producer" ) ;
this . ruleEngineProducerStats = statsFactory . createMessagesStats ( StatsType . RULE_ENGINE . getName ( ) + ".producer" ) ;
this . tbCoreProducerStats = statsFactory . createMessagesStats ( StatsType . CORE . getName ( ) + ".producer" ) ;
this . tbCoreProducerStats = statsFactory . createMessagesStats ( StatsType . CORE . getName ( ) + ".producer" ) ;
this . transportApiStats = statsFactory . createMessagesStats ( StatsType . TRANSPORT . getName ( ) + ".producer" ) ;
this . transportApiStats = statsFactory . createMessagesStats ( StatsType . TRANSPORT . getName ( ) + ".producer" ) ;
this . transportCallbackExecutor = ThingsBoardExecutors . newWorkStealingPool ( 20 , getClass ( ) ) ;
this . transportCallbackExecutor = ThingsBoardExecutors . newWorkStealingPool ( 20 , getClass ( ) ) ;
this . scheduler . scheduleAtFixedRate ( this : : checkInactivityAndReportActivity , new Random ( ) . nextInt ( ( int ) sessionReportTimeout ) , sessionReportTimeout , TimeUnit . MILLISECONDS ) ;
this . scheduler . scheduleAtFixedRate ( this : : invalidateRateLimits , new Random ( ) . nextInt ( ( int ) sessionReportTimeout ) , sessionReportTimeout , TimeUnit . MILLISECONDS ) ;
this . scheduler . scheduleAtFixedRate ( this : : invalidateRateLimits , new Random ( ) . nextInt ( ( int ) sessionReportTimeout ) , sessionReportTimeout , TimeUnit . MILLISECONDS ) ;
transportApiRequestTemplate = queueProvider . createTransportApiRequestTemplate ( ) ;
transportApiRequestTemplate = queueProvider . createTransportApiRequestTemplate ( ) ;
transportApiRequestTemplate . setMessagesStats ( transportApiStats ) ;
transportApiRequestTemplate . setMessagesStats ( transportApiStats ) ;
@ -558,7 +545,7 @@ public class DefaultTransportService implements TransportService {
@Override
@Override
public void process ( TransportProtos . SessionInfoProto sessionInfo , TransportProtos . SessionEventMsg msg , TransportServiceCallback < Void > callback ) {
public void process ( TransportProtos . SessionInfoProto sessionInfo , TransportProtos . SessionEventMsg msg , TransportServiceCallback < Void > callback ) {
if ( checkLimits ( sessionInfo , msg , callback ) ) {
if ( checkLimits ( sessionInfo , msg , callback ) ) {
report ActivityInternal ( sessionInfo ) ;
record ActivityInternal ( sessionInfo ) ;
sendToDeviceActor ( sessionInfo , TransportToDeviceActorMsg . newBuilder ( ) . setSessionInfo ( sessionInfo )
sendToDeviceActor ( sessionInfo , TransportToDeviceActorMsg . newBuilder ( ) . setSessionInfo ( sessionInfo )
. setSessionEvent ( msg ) . build ( ) , callback ) ;
. setSessionEvent ( msg ) . build ( ) , callback ) ;
}
}
@ -578,7 +565,7 @@ public class DefaultTransportService implements TransportService {
}
}
}
}
report ActivityInternal ( sessionInfo ) ;
record ActivityInternal ( sessionInfo ) ;
sendToDeviceActor ( sessionInfo , msg , callback ) ;
sendToDeviceActor ( sessionInfo , msg , callback ) ;
}
}
}
}
@ -595,7 +582,7 @@ public class DefaultTransportService implements TransportService {
dataPoints + = tsKv . getKvCount ( ) ;
dataPoints + = tsKv . getKvCount ( ) ;
}
}
if ( checkLimits ( sessionInfo , msg , callback , dataPoints ) ) {
if ( checkLimits ( sessionInfo , msg , callback , dataPoints ) ) {
report ActivityInternal ( sessionInfo ) ;
record ActivityInternal ( sessionInfo ) ;
TenantId tenantId = getTenantId ( sessionInfo ) ;
TenantId tenantId = getTenantId ( sessionInfo ) ;
DeviceId deviceId = new DeviceId ( new UUID ( sessionInfo . getDeviceIdMSB ( ) , sessionInfo . getDeviceIdLSB ( ) ) ) ;
DeviceId deviceId = new DeviceId ( new UUID ( sessionInfo . getDeviceIdMSB ( ) , sessionInfo . getDeviceIdLSB ( ) ) ) ;
CustomerId customerId = getCustomerId ( sessionInfo ) ;
CustomerId customerId = getCustomerId ( sessionInfo ) ;
@ -619,7 +606,7 @@ public class DefaultTransportService implements TransportService {
@Override
@Override
public void process ( TransportProtos . SessionInfoProto sessionInfo , TransportProtos . PostAttributeMsg msg , TbMsgMetaData md , TransportServiceCallback < Void > callback ) {
public void process ( TransportProtos . SessionInfoProto sessionInfo , TransportProtos . PostAttributeMsg msg , TbMsgMetaData md , TransportServiceCallback < Void > callback ) {
if ( checkLimits ( sessionInfo , msg , callback , msg . getKvCount ( ) ) ) {
if ( checkLimits ( sessionInfo , msg , callback , msg . getKvCount ( ) ) ) {
report ActivityInternal ( sessionInfo ) ;
record ActivityInternal ( sessionInfo ) ;
TenantId tenantId = getTenantId ( sessionInfo ) ;
TenantId tenantId = getTenantId ( sessionInfo ) ;
DeviceId deviceId = new DeviceId ( new UUID ( sessionInfo . getDeviceIdMSB ( ) , sessionInfo . getDeviceIdLSB ( ) ) ) ;
DeviceId deviceId = new DeviceId ( new UUID ( sessionInfo . getDeviceIdMSB ( ) , sessionInfo . getDeviceIdLSB ( ) ) ) ;
JsonObject json = JsonUtils . getJsonObject ( msg . getKvList ( ) ) ;
JsonObject json = JsonUtils . getJsonObject ( msg . getKvList ( ) ) ;
@ -639,7 +626,7 @@ public class DefaultTransportService implements TransportService {
@Override
@Override
public void process ( TransportProtos . SessionInfoProto sessionInfo , TransportProtos . GetAttributeRequestMsg msg , TransportServiceCallback < Void > callback ) {
public void process ( TransportProtos . SessionInfoProto sessionInfo , TransportProtos . GetAttributeRequestMsg msg , TransportServiceCallback < Void > callback ) {
if ( checkLimits ( sessionInfo , msg , callback ) ) {
if ( checkLimits ( sessionInfo , msg , callback ) ) {
report ActivityInternal ( sessionInfo ) ;
record ActivityInternal ( sessionInfo ) ;
sendToDeviceActor ( sessionInfo , TransportToDeviceActorMsg . newBuilder ( ) . setSessionInfo ( sessionInfo )
sendToDeviceActor ( sessionInfo , TransportToDeviceActorMsg . newBuilder ( ) . setSessionInfo ( sessionInfo )
. setGetAttributes ( msg ) . build ( ) , new ApiStatsProxyCallback < > ( getTenantId ( sessionInfo ) , getCustomerId ( sessionInfo ) , 1 , callback ) ) ;
. setGetAttributes ( msg ) . build ( ) , new ApiStatsProxyCallback < > ( getTenantId ( sessionInfo ) , getCustomerId ( sessionInfo ) , 1 , callback ) ) ;
}
}
@ -652,7 +639,7 @@ public class DefaultTransportService implements TransportService {
if ( sessionMetaData ! = null ) {
if ( sessionMetaData ! = null ) {
sessionMetaData . setSubscribedToAttributes ( ! msg . getUnsubscribe ( ) ) ;
sessionMetaData . setSubscribedToAttributes ( ! msg . getUnsubscribe ( ) ) ;
}
}
report ActivityInternal ( sessionInfo ) ;
record ActivityInternal ( sessionInfo ) ;
sendToDeviceActor ( sessionInfo , TransportToDeviceActorMsg . newBuilder ( ) . setSessionInfo ( sessionInfo ) . setSubscribeToAttributes ( msg ) . build ( ) ,
sendToDeviceActor ( sessionInfo , TransportToDeviceActorMsg . newBuilder ( ) . setSessionInfo ( sessionInfo ) . setSubscribeToAttributes ( msg ) . build ( ) ,
new ApiStatsProxyCallback < > ( getTenantId ( sessionInfo ) , getCustomerId ( sessionInfo ) , 1 , callback ) ) ;
new ApiStatsProxyCallback < > ( getTenantId ( sessionInfo ) , getCustomerId ( sessionInfo ) , 1 , callback ) ) ;
}
}
@ -665,7 +652,7 @@ public class DefaultTransportService implements TransportService {
if ( sessionMetaData ! = null ) {
if ( sessionMetaData ! = null ) {
sessionMetaData . setSubscribedToRPC ( ! msg . getUnsubscribe ( ) ) ;
sessionMetaData . setSubscribedToRPC ( ! msg . getUnsubscribe ( ) ) ;
}
}
report ActivityInternal ( sessionInfo ) ;
record ActivityInternal ( sessionInfo ) ;
sendToDeviceActor ( sessionInfo , TransportToDeviceActorMsg . newBuilder ( ) . setSessionInfo ( sessionInfo ) . setSubscribeToRPC ( msg ) . build ( ) ,
sendToDeviceActor ( sessionInfo , TransportToDeviceActorMsg . newBuilder ( ) . setSessionInfo ( sessionInfo ) . setSubscribeToRPC ( msg ) . build ( ) ,
new ApiStatsProxyCallback < > ( getTenantId ( sessionInfo ) , getCustomerId ( sessionInfo ) , 1 , callback ) ) ;
new ApiStatsProxyCallback < > ( getTenantId ( sessionInfo ) , getCustomerId ( sessionInfo ) , 1 , callback ) ) ;
}
}
@ -674,7 +661,7 @@ public class DefaultTransportService implements TransportService {
@Override
@Override
public void process ( TransportProtos . SessionInfoProto sessionInfo , TransportProtos . ToDeviceRpcResponseMsg msg , TransportServiceCallback < Void > callback ) {
public void process ( TransportProtos . SessionInfoProto sessionInfo , TransportProtos . ToDeviceRpcResponseMsg msg , TransportServiceCallback < Void > callback ) {
if ( checkLimits ( sessionInfo , msg , callback ) ) {
if ( checkLimits ( sessionInfo , msg , callback ) ) {
report ActivityInternal ( sessionInfo ) ;
record ActivityInternal ( sessionInfo ) ;
sendToDeviceActor ( sessionInfo , TransportToDeviceActorMsg . newBuilder ( ) . setSessionInfo ( sessionInfo ) . setToDeviceRPCCallResponse ( msg ) . build ( ) ,
sendToDeviceActor ( sessionInfo , TransportToDeviceActorMsg . newBuilder ( ) . setSessionInfo ( sessionInfo ) . setToDeviceRPCCallResponse ( msg ) . build ( ) ,
new ApiStatsProxyCallback < > ( getTenantId ( sessionInfo ) , getCustomerId ( sessionInfo ) , 1 , callback ) ) ;
new ApiStatsProxyCallback < > ( getTenantId ( sessionInfo ) , getCustomerId ( sessionInfo ) , 1 , callback ) ) ;
}
}
@ -683,7 +670,7 @@ public class DefaultTransportService implements TransportService {
@Override
@Override
public void notifyAboutUplink ( TransportProtos . SessionInfoProto sessionInfo , TransportProtos . UplinkNotificationMsg msg , TransportServiceCallback < Void > callback ) {
public void notifyAboutUplink ( TransportProtos . SessionInfoProto sessionInfo , TransportProtos . UplinkNotificationMsg msg , TransportServiceCallback < Void > callback ) {
if ( checkLimits ( sessionInfo , msg , callback ) ) {
if ( checkLimits ( sessionInfo , msg , callback ) ) {
report ActivityInternal ( sessionInfo ) ;
record ActivityInternal ( sessionInfo ) ;
sendToDeviceActor ( sessionInfo , TransportToDeviceActorMsg . newBuilder ( ) . setSessionInfo ( sessionInfo ) . setUplinkNotificationMsg ( msg ) . build ( ) , callback ) ;
sendToDeviceActor ( sessionInfo , TransportToDeviceActorMsg . newBuilder ( ) . setSessionInfo ( sessionInfo ) . setUplinkNotificationMsg ( msg ) . build ( ) , callback ) ;
}
}
}
}
@ -704,7 +691,7 @@ public class DefaultTransportService implements TransportService {
if ( checkLimits ( sessionInfo , responseMsg , callback ) ) {
if ( checkLimits ( sessionInfo , responseMsg , callback ) ) {
if ( reportActivity ) {
if ( reportActivity ) {
report ActivityInternal ( sessionInfo ) ;
record ActivityInternal ( sessionInfo ) ;
}
}
sendToDeviceActor ( sessionInfo , TransportToDeviceActorMsg . newBuilder ( ) . setSessionInfo ( sessionInfo ) . setRpcResponseStatusMsg ( responseMsg ) . build ( ) ,
sendToDeviceActor ( sessionInfo , TransportToDeviceActorMsg . newBuilder ( ) . setSessionInfo ( sessionInfo ) . setRpcResponseStatusMsg ( responseMsg ) . build ( ) ,
new ApiStatsProxyCallback < > ( getTenantId ( sessionInfo ) , getCustomerId ( sessionInfo ) , 1 , TransportServiceCallback . EMPTY ) ) ;
new ApiStatsProxyCallback < > ( getTenantId ( sessionInfo ) , getCustomerId ( sessionInfo ) , 1 , TransportServiceCallback . EMPTY ) ) ;
@ -736,7 +723,7 @@ public class DefaultTransportService implements TransportService {
@Override
@Override
public void process ( TransportProtos . SessionInfoProto sessionInfo , TransportProtos . ToServerRpcRequestMsg msg , TransportServiceCallback < Void > callback ) {
public void process ( TransportProtos . SessionInfoProto sessionInfo , TransportProtos . ToServerRpcRequestMsg msg , TransportServiceCallback < Void > callback ) {
if ( checkLimits ( sessionInfo , msg , callback ) ) {
if ( checkLimits ( sessionInfo , msg , callback ) ) {
report ActivityInternal ( sessionInfo ) ;
record ActivityInternal ( sessionInfo ) ;
UUID sessionId = toSessionId ( sessionInfo ) ;
UUID sessionId = toSessionId ( sessionInfo ) ;
TenantId tenantId = getTenantId ( sessionInfo ) ;
TenantId tenantId = getTenantId ( sessionInfo ) ;
DeviceId deviceId = getDeviceId ( sessionInfo ) ;
DeviceId deviceId = getDeviceId ( sessionInfo ) ;
@ -761,7 +748,7 @@ public class DefaultTransportService implements TransportService {
@Override
@Override
public void process ( TransportProtos . SessionInfoProto sessionInfo , TransportProtos . ClaimDeviceMsg msg , TransportServiceCallback < Void > callback ) {
public void process ( TransportProtos . SessionInfoProto sessionInfo , TransportProtos . ClaimDeviceMsg msg , TransportServiceCallback < Void > callback ) {
if ( checkLimits ( sessionInfo , msg , callback ) ) {
if ( checkLimits ( sessionInfo , msg , callback ) ) {
report ActivityInternal ( sessionInfo ) ;
record ActivityInternal ( sessionInfo ) ;
sendToDeviceActor ( sessionInfo , TransportToDeviceActorMsg . newBuilder ( ) . setSessionInfo ( sessionInfo )
sendToDeviceActor ( sessionInfo , TransportToDeviceActorMsg . newBuilder ( ) . setSessionInfo ( sessionInfo )
. setClaimDevice ( msg ) . build ( ) , callback ) ;
. setClaimDevice ( msg ) . build ( ) , callback ) ;
}
}
@ -780,69 +767,12 @@ public class DefaultTransportService implements TransportService {
}
}
@Override
@Override
public void report Activity ( TransportProtos . SessionInfoProto sessionInfo ) {
public void record Activity ( TransportProtos . SessionInfoProto sessionInfo ) {
report ActivityInternal ( sessionInfo ) ;
record ActivityInternal ( sessionInfo ) ;
}
}
private void reportActivityInternal ( TransportProtos . SessionInfoProto sessionInfo ) {
private void recordActivityInternal ( TransportProtos . SessionInfoProto sessionInfo ) {
UUID sessionId = toSessionId ( sessionInfo ) ;
onActivity ( toSessionId ( sessionInfo ) , sessionInfo , getCurrentTimeMillis ( ) ) ;
SessionActivityData sessionMetaData = sessionsActivity . computeIfAbsent ( sessionId , id - > new SessionActivityData ( sessionInfo ) ) ;
sessionMetaData . updateLastActivityTime ( ) ;
}
private void checkInactivityAndReportActivity ( ) {
long expTime = System . currentTimeMillis ( ) - sessionInactivityTimeout ;
Set < UUID > sessionsToRemove = new HashSet < > ( ) ;
sessionsActivity . forEach ( ( uuid , sessionAD ) - > {
long lastActivityTime = sessionAD . getLastActivityTime ( ) ;
SessionMetaData sessionMD = sessions . get ( uuid ) ;
if ( sessionMD ! = null ) {
sessionAD . setSessionInfo ( sessionMD . getSessionInfo ( ) ) ;
} else {
sessionsToRemove . add ( uuid ) ;
}
TransportProtos . SessionInfoProto sessionInfo = sessionAD . getSessionInfo ( ) ;
if ( sessionInfo . getGwSessionIdMSB ( ) ! = 0 & & sessionInfo . getGwSessionIdLSB ( ) ! = 0 ) {
var gwSessionId = new UUID ( sessionInfo . getGwSessionIdMSB ( ) , sessionInfo . getGwSessionIdLSB ( ) ) ;
SessionMetaData gwMetaData = sessions . get ( gwSessionId ) ;
SessionActivityData gwActivityData = sessionsActivity . get ( gwSessionId ) ;
if ( gwMetaData ! = null & & gwMetaData . isOverwriteActivityTime ( ) ) {
lastActivityTime = Math . max ( gwActivityData . getLastActivityTime ( ) , lastActivityTime ) ;
}
}
if ( lastActivityTime < expTime ) {
if ( sessionMD ! = null ) {
if ( log . isDebugEnabled ( ) ) {
log . debug ( "[{}] Session has expired due to last activity time: {}" , toSessionId ( sessionInfo ) , lastActivityTime ) ;
}
sessions . remove ( uuid ) ;
sessionsToRemove . add ( uuid ) ;
process ( sessionInfo , SESSION_EVENT_MSG_CLOSED , null ) ;
sessionMD . getListener ( ) . onRemoteSessionCloseCommand ( uuid , SESSION_CLOSE_NOTIFICATION_PROTO ) ;
}
} else {
if ( lastActivityTime > sessionAD . getLastReportedActivityTime ( ) ) {
final long lastActivityTimeFinal = lastActivityTime ;
process ( sessionInfo , TransportProtos . SubscriptionInfoProto . newBuilder ( )
. setAttributeSubscription ( sessionMD ! = null & & sessionMD . isSubscribedToAttributes ( ) )
. setRpcSubscription ( sessionMD ! = null & & sessionMD . isSubscribedToRPC ( ) )
. setLastActivityTime ( lastActivityTime ) . build ( ) , new TransportServiceCallback < Void > ( ) {
@Override
public void onSuccess ( Void msg ) {
sessionAD . setLastReportedActivityTime ( lastActivityTimeFinal ) ;
}
@Override
public void onError ( Throwable e ) {
log . warn ( "[{}] Failed to report last activity time" , uuid , e ) ;
}
} ) ;
}
}
} ) ;
// Removes all closed or short-lived sessions.
sessionsToRemove . forEach ( sessionsActivity : : remove ) ;
}
}
@Override
@Override
@ -913,7 +843,7 @@ public class DefaultTransportService implements TransportService {
}
}
TransportProtos . PostTelemetryMsg . Builder request = TransportProtos . PostTelemetryMsg . newBuilder ( ) ;
TransportProtos . PostTelemetryMsg . Builder request = TransportProtos . PostTelemetryMsg . newBuilder ( ) ;
TransportProtos . TsKvListProto . Builder builder = TransportProtos . TsKvListProto . newBuilder ( ) ;
TransportProtos . TsKvListProto . Builder builder = TransportProtos . TsKvListProto . newBuilder ( ) ;
builder . setTs ( TimeUnit . MILLISECONDS . toSeconds ( System . c urrentTimeMillis( ) ) * 1000L + ( atomicTs . getAndIncrement ( ) % 1000 ) ) ;
builder . setTs ( TimeUnit . MILLISECONDS . toSeconds ( getC urrentTimeMillis( ) ) * 1000L + ( atomicTs . getAndIncrement ( ) % 1000 ) ) ;
builder . addKv ( TransportProtos . KeyValueProto . newBuilder ( )
builder . addKv ( TransportProtos . KeyValueProto . newBuilder ( )
. setKey ( "transportLog" )
. setKey ( "transportLog" )
. setType ( TransportProtos . KeyValueType . STRING_V )
. setType ( TransportProtos . KeyValueType . STRING_V )
@ -1001,8 +931,8 @@ public class DefaultTransportService implements TransportService {
Optional < Tenant > profileOpt = dataDecodingEncodingService . decode ( msg . getData ( ) . toByteArray ( ) ) ;
Optional < Tenant > profileOpt = dataDecodingEncodingService . decode ( msg . getData ( ) . toByteArray ( ) ) ;
if ( profileOpt . isPresent ( ) ) {
if ( profileOpt . isPresent ( ) ) {
Tenant tenant = profileOpt . get ( ) ;
Tenant tenant = profileOpt . get ( ) ;
partitionService . removeTenant ( tenant . getId ( ) ) ;
boolean updated = tenantProfileCache . put ( tenant . getId ( ) , tenant . getTenantProfileId ( ) ) ;
boolean updated = tenantProfileCache . put ( tenant . getId ( ) , tenant . getTenantProfileId ( ) ) ;
partitionService . evictTenantInfo ( tenant . getId ( ) ) ;
if ( updated ) {
if ( updated ) {
rateLimitService . update ( tenant . getId ( ) ) ;
rateLimitService . update ( tenant . getId ( ) ) ;
}
}
@ -1027,7 +957,9 @@ public class DefaultTransportService implements TransportService {
} else if ( EntityType . TENANT_PROFILE . equals ( entityType ) ) {
} else if ( EntityType . TENANT_PROFILE . equals ( entityType ) ) {
tenantProfileCache . remove ( new TenantProfileId ( entityUuid ) ) ;
tenantProfileCache . remove ( new TenantProfileId ( entityUuid ) ) ;
} else if ( EntityType . TENANT . equals ( entityType ) ) {
} else if ( EntityType . TENANT . equals ( entityType ) ) {
rateLimitService . remove ( TenantId . fromUUID ( entityUuid ) ) ;
TenantId tenantId = TenantId . fromUUID ( entityUuid ) ;
rateLimitService . remove ( tenantId ) ;
partitionService . removeTenant ( tenantId ) ;
} else if ( EntityType . DEVICE . equals ( entityType ) ) {
} else if ( EntityType . DEVICE . equals ( entityType ) ) {
rateLimitService . remove ( new DeviceId ( entityUuid ) ) ;
rateLimitService . remove ( new DeviceId ( entityUuid ) ) ;
onDeviceDeleted ( new DeviceId ( entityUuid ) ) ;
onDeviceDeleted ( new DeviceId ( entityUuid ) ) ;
@ -1160,12 +1092,6 @@ public class DefaultTransportService implements TransportService {
return new DeviceId ( new UUID ( sessionInfo . getDeviceIdMSB ( ) , sessionInfo . getDeviceIdLSB ( ) ) ) ;
return new DeviceId ( new UUID ( sessionInfo . getDeviceIdMSB ( ) , sessionInfo . getDeviceIdLSB ( ) ) ) ;
}
}
private static TransportProtos . SessionEventMsg getSessionEventMsg ( TransportProtos . SessionEvent event ) {
return TransportProtos . SessionEventMsg . newBuilder ( )
. setSessionType ( TransportProtos . SessionType . ASYNC )
. setEvent ( event ) . build ( ) ;
}
protected void sendToDeviceActor ( TransportProtos . SessionInfoProto sessionInfo , TransportToDeviceActorMsg toDeviceActorMsg , TransportServiceCallback < Void > callback ) {
protected void sendToDeviceActor ( TransportProtos . SessionInfoProto sessionInfo , TransportToDeviceActorMsg toDeviceActorMsg , TransportServiceCallback < Void > callback ) {
ToCoreMsg toCoreMsg = ToCoreMsg . newBuilder ( ) . setToDeviceActorMsg ( toDeviceActorMsg ) . build ( ) ;
ToCoreMsg toCoreMsg = ToCoreMsg . newBuilder ( ) . setToDeviceActorMsg ( toDeviceActorMsg ) . build ( ) ;
sendToCore ( getTenantId ( sessionInfo ) , getDeviceId ( sessionInfo ) , toCoreMsg , getRoutingKey ( sessionInfo ) , callback ) ;
sendToCore ( getTenantId ( sessionInfo ) , getDeviceId ( sessionInfo ) , toCoreMsg , getRoutingKey ( sessionInfo ) , callback ) ;