@ -24,6 +24,7 @@ import org.springframework.context.event.EventListener;
import org.springframework.core.annotation.Order ;
import org.springframework.scheduling.annotation.Scheduled ;
import org.springframework.stereotype.Service ;
import org.thingsboard.common.util.JacksonUtil ;
import org.thingsboard.common.util.ThingsBoardThreadFactory ;
import org.thingsboard.rule.engine.api.RpcError ;
import org.thingsboard.server.actors.ActorSystemContext ;
@ -35,7 +36,7 @@ import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TbCallback ;
import org.thingsboard.server.common.stats.StatsFactory ;
import org.thingsboard.server.common.transport.util.DataDecodingEncodingService ;
import org.thingsboard.common.util.JacksonUtil ;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache ;
import org.thingsboard.server.gen.transport.TransportProtos.DeviceStateServiceMsgProto ;
import org.thingsboard.server.gen.transport.TransportProtos.EdgeNotificationMsgProto ;
import org.thingsboard.server.gen.transport.TransportProtos.FromDeviceRPCResponseProto ;
@ -49,6 +50,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.TbSubscriptionCloseP
import org.thingsboard.server.gen.transport.TransportProtos.TbTimeSeriesUpdateProto ;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.ToFirmwareStateServiceMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg ;
import org.thingsboard.server.gen.transport.TransportProtos.TransportToDeviceActorMsg ;
import org.thingsboard.server.queue.TbQueueConsumer ;
@ -58,8 +60,8 @@ import org.thingsboard.server.queue.provider.TbCoreQueueFactory;
import org.thingsboard.server.queue.util.TbCoreComponent ;
import org.thingsboard.server.service.apiusage.TbApiUsageStateService ;
import org.thingsboard.server.service.edge.EdgeNotificationService ;
import org.thingsboard.server.service.firmware.FirmwareStateService ;
import org.thingsboard.server.service.profile.TbDeviceProfileCache ;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache ;
import org.thingsboard.server.service.queue.processing.AbstractConsumerService ;
import org.thingsboard.server.service.rpc.FromDeviceRpcResponse ;
import org.thingsboard.server.service.rpc.TbCoreDeviceRpcService ;
@ -97,6 +99,11 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
@Value ( "${queue.core.stats.enabled:false}" )
private boolean statsEnabled ;
@Value ( "${queue.core.firmware.pack-interval-ms:60000}" )
private long firmwarePackInterval ;
@Value ( "${queue.core.firmware.pack-size:100}" )
private int firmwarePackSize ;
private final TbQueueConsumer < TbProtoQueueMsg < ToCoreMsg > > mainConsumer ;
private final DeviceStateService stateService ;
private final TbApiUsageStateService statsService ;
@ -104,11 +111,15 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
private final SubscriptionManagerService subscriptionManagerService ;
private final TbCoreDeviceRpcService tbCoreDeviceRpcService ;
private final EdgeNotificationService edgeNotificationService ;
private final FirmwareStateService firmwareStateService ;
private final TbCoreConsumerStats stats ;
protected final TbQueueConsumer < TbProtoQueueMsg < ToUsageStatsServiceMsg > > usageStatsConsumer ;
private final TbQueueConsumer < TbProtoQueueMsg < ToFirmwareStateServiceMsg > > firmwareStatesConsumer ;
protected volatile ExecutorService usageStatsExecutor ;
private volatile ExecutorService firmwareStatesExecutor ;
public DefaultTbCoreConsumerService ( TbCoreQueueFactory tbCoreQueueFactory ,
ActorSystemContext actorContext ,
DeviceStateService stateService ,
@ -121,10 +132,12 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
TbApiUsageStateService statsService ,
TbTenantProfileCache tenantProfileCache ,
TbApiUsageStateService apiUsageStateService ,
EdgeNotificationService edgeNotificationService ) {
EdgeNotificationService edgeNotificationService ,
FirmwareStateService firmwareStateService ) {
super ( actorContext , encodingService , tenantProfileCache , deviceProfileCache , apiUsageStateService , tbCoreQueueFactory . createToCoreNotificationsMsgConsumer ( ) ) ;
this . mainConsumer = tbCoreQueueFactory . createToCoreMsgConsumer ( ) ;
this . usageStatsConsumer = tbCoreQueueFactory . createToUsageStatsServiceMsgConsumer ( ) ;
this . firmwareStatesConsumer = tbCoreQueueFactory . createToFirmwareStateServiceMsgConsumer ( ) ;
this . stateService = stateService ;
this . localSubscriptionService = localSubscriptionService ;
this . subscriptionManagerService = subscriptionManagerService ;
@ -132,12 +145,14 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
this . edgeNotificationService = edgeNotificationService ;
this . stats = new TbCoreConsumerStats ( statsFactory ) ;
this . statsService = statsService ;
this . firmwareStateService = firmwareStateService ;
}
@PostConstruct
public void init ( ) {
super . init ( "tb-core-consumer" , "tb-core-notifications-consumer" ) ;
this . usageStatsExecutor = Executors . newCachedThreadPool ( ThingsBoardThreadFactory . forName ( "tb-core-usage-stats-consumer" ) ) ;
this . firmwareStatesExecutor = Executors . newCachedThreadPool ( ThingsBoardThreadFactory . forName ( "tb-core-firmware-notifications-consumer" ) ) ;
}
@PreDestroy
@ -146,6 +161,9 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
if ( usageStatsExecutor ! = null ) {
usageStatsExecutor . shutdownNow ( ) ;
}
if ( firmwareStatesExecutor ! = null ) {
firmwareStatesExecutor . shutdownNow ( ) ;
}
}
@EventListener ( ApplicationReadyEvent . class )
@ -153,6 +171,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
public void onApplicationEvent ( ApplicationReadyEvent event ) {
super . onApplicationEvent ( event ) ;
launchUsageStatsConsumer ( ) ;
launchFirmwareUpdateNotificationConsumer ( ) ;
}
@Override
@ -167,6 +186,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
. map ( tpi - > tpi . newByTopic ( usageStatsConsumer . getTopic ( ) ) )
. collect ( Collectors . toSet ( ) ) ) ;
}
this . firmwareStatesConsumer . subscribe ( ) ;
}
@Override
@ -336,10 +356,57 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
} ) ;
}
private void launchFirmwareUpdateNotificationConsumer ( ) {
long maxProcessingTimeoutPerRecord = firmwarePackInterval / firmwarePackSize ;
firmwareStatesExecutor . submit ( ( ) - > {
while ( ! stopped ) {
try {
List < TbProtoQueueMsg < ToFirmwareStateServiceMsg > > msgs = firmwareStatesConsumer . poll ( getNotificationPollDuration ( ) ) ;
if ( msgs . isEmpty ( ) ) {
Thread . sleep ( maxProcessingTimeoutPerRecord ) ;
continue ;
}
long timeToSleep = maxProcessingTimeoutPerRecord ;
for ( TbProtoQueueMsg < ToFirmwareStateServiceMsg > msg : msgs ) {
try {
long startTime = System . currentTimeMillis ( ) ;
boolean isSuccessUpdate = handleFirmwareUpdates ( msg ) ;
long endTime = System . currentTimeMillis ( ) ;
long spentTime = endTime - startTime ;
timeToSleep = timeToSleep - spentTime ;
if ( isSuccessUpdate & & timeToSleep > 0 ) {
log . debug ( "Spent time per record is: [{}]!" , spentTime ) ;
Thread . sleep ( timeToSleep ) ;
timeToSleep = maxProcessingTimeoutPerRecord ;
}
} catch ( Throwable e ) {
log . warn ( "Failed to process firmware update msg: {}" , msg , e ) ;
}
}
firmwareStatesConsumer . commit ( ) ;
} catch ( Exception e ) {
if ( ! stopped ) {
log . warn ( "Failed to obtain usage stats from queue." , e ) ;
try {
Thread . sleep ( getNotificationPollDuration ( ) ) ;
} catch ( InterruptedException e2 ) {
log . trace ( "Failed to wait until the server has capacity to handle new firmware updates" , e2 ) ;
}
}
}
}
log . info ( "TB Firmware States Consumer stopped." ) ;
} ) ;
}
private void handleUsageStats ( TbProtoQueueMsg < ToUsageStatsServiceMsg > msg , TbCallback callback ) {
statsService . process ( msg , callback ) ;
}
private boolean handleFirmwareUpdates ( TbProtoQueueMsg < ToFirmwareStateServiceMsg > msg ) {
return firmwareStateService . process ( msg . getValue ( ) ) ;
}
private void forwardToCoreRpcService ( FromDeviceRPCResponseProto proto , TbCallback callback ) {
RpcError error = proto . getError ( ) > 0 ? RpcError . values ( ) [ proto . getError ( ) ] : null ;
FromDeviceRpcResponse response = new FromDeviceRpcResponse ( new UUID ( proto . getRequestIdMSB ( ) , proto . getRequestIdLSB ( ) )
@ -448,6 +515,9 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
if ( usageStatsConsumer ! = null ) {
usageStatsConsumer . unsubscribe ( ) ;
}
if ( firmwareStatesConsumer ! = null ) {
firmwareStatesConsumer . unsubscribe ( ) ;
}
}
}