diff --git a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java index d72cb43334..9d928e622a 100644 --- a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java @@ -24,6 +24,7 @@ import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListeningExecutorService; import com.google.common.util.concurrent.MoreExecutors; import lombok.Getter; +import lombok.RequiredArgsConstructor; import lombok.Setter; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.tuple.Pair; @@ -114,6 +115,7 @@ import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE; @Service @TbCoreComponent @Slf4j +@RequiredArgsConstructor public class DefaultDeviceStateService extends AbstractPartitionBasedService implements DeviceStateService { public static final String ACTIVITY_STATE = "active"; @@ -148,13 +150,11 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService deviceStates = new ConcurrentHashMap<>(); - public DefaultDeviceStateService(TenantService tenantService, DeviceService deviceService, - AttributesService attributesService, TimeseriesService tsService, - TbClusterService clusterService, PartitionService partitionService, - TbServiceInfoProvider serviceInfoProvider, - EntityQueryRepository entityQueryRepository, - DbTypeInfoComponent dbTypeInfoComponent, - TbApiUsageReportClient apiUsageReportClient) { - this.tenantService = tenantService; - this.deviceService = deviceService; - this.attributesService = attributesService; - this.tsService = tsService; - this.clusterService = clusterService; - this.partitionService = partitionService; - this.serviceInfoProvider = serviceInfoProvider; - this.entityQueryRepository = entityQueryRepository; - this.dbTypeInfoComponent = dbTypeInfoComponent; - this.apiUsageReportClient = apiUsageReportClient; - } - @PostConstruct public void init() { super.init(); deviceStateExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool( Math.max(4, Runtime.getRuntime().availableProcessors()), "device-state")); - scheduledExecutor.scheduleAtFixedRate(this::checkStates, new Random().nextInt(defaultStateCheckIntervalInSec), defaultStateCheckIntervalInSec, TimeUnit.SECONDS); + scheduledExecutor.scheduleWithFixedDelay(this::checkStates, new Random().nextInt(defaultStateCheckIntervalInSec), defaultStateCheckIntervalInSec, TimeUnit.SECONDS); + scheduledExecutor.scheduleWithFixedDelay(this::reportActivityStats, defaultActivityStatsIntervalInSec, defaultActivityStatsIntervalInSec, TimeUnit.SECONDS); } @PreDestroy @@ -468,29 +454,31 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService> stats = new HashMap<>(); - for (DeviceStateData stateData : deviceStates.values()) { - Pair tenantDevicesActivity = stats.computeIfAbsent(stateData.getTenantId(), - tenantId -> Pair.of(new AtomicInteger(), new AtomicInteger())); - if (stateData.getState().isActive()) { - tenantDevicesActivity.getLeft().incrementAndGet(); - } else { - tenantDevicesActivity.getRight().incrementAndGet(); + void reportActivityStats() { + try{ + Map> stats = new HashMap<>(); + for (DeviceStateData stateData : deviceStates.values()) { + Pair tenantDevicesActivity = stats.computeIfAbsent(stateData.getTenantId(), + tenantId -> Pair.of(new AtomicInteger(), new AtomicInteger())); + if (stateData.getState().isActive()) { + tenantDevicesActivity.getLeft().incrementAndGet(); + } else { + tenantDevicesActivity.getRight().incrementAndGet(); + } } - } - stats.forEach((tenantId, tenantDevicesActivity) -> { - int active = tenantDevicesActivity.getLeft().get(); - int inactive = tenantDevicesActivity.getRight().get(); - apiUsageReportClient.report(tenantId, null, ApiUsageRecordKey.ACTIVE_DEVICES, active); - apiUsageReportClient.report(tenantId, null, ApiUsageRecordKey.INACTIVE_DEVICES, inactive); - if (active > 0) { - log.info("[{}] Active devices: {}, inactive devices: {}", tenantId, active, inactive); - } - }); + stats.forEach((tenantId, tenantDevicesActivity) -> { + int active = tenantDevicesActivity.getLeft().get(); + int inactive = tenantDevicesActivity.getRight().get(); + apiUsageReportClient.report(tenantId, null, ApiUsageRecordKey.ACTIVE_DEVICES, active); + apiUsageReportClient.report(tenantId, null, ApiUsageRecordKey.INACTIVE_DEVICES, inactive); + if (active > 0) { + log.info("[{}] Active devices: {}, inactive devices: {}", tenantId, active, inactive); + } + }); + } catch (Throwable t) { + log.warn("Failed to report activity states", t); + } } void updateInactivityStateIfExpired(long ts, DeviceId deviceId, DeviceStateData stateData) { diff --git a/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java b/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java index a8a68762d8..8245284af6 100644 --- a/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java +++ b/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java @@ -50,8 +50,6 @@ import static org.thingsboard.server.service.state.DefaultDeviceStateService.INA @RunWith(MockitoJUnitRunner.class) public class DefaultDeviceStateServiceTest { - @Mock - TenantService tenantService; @Mock DeviceService deviceService; @Mock @@ -64,8 +62,6 @@ public class DefaultDeviceStateServiceTest { PartitionService partitionService; @Mock DeviceStateData deviceStateDataMock; - @Mock - TbServiceInfoProvider serviceInfoProvider; DeviceId deviceId = DeviceId.fromString("00797a3b-7aeb-4b5b-b57a-c2a810d0f112"); @@ -73,7 +69,7 @@ public class DefaultDeviceStateServiceTest { @Before public void setUp() { - service = spy(new DefaultDeviceStateService(tenantService, deviceService, attributesService, tsService, clusterService, partitionService, serviceInfoProvider, null, null, null)); + service = spy(new DefaultDeviceStateService(deviceService, attributesService, tsService, clusterService, partitionService, null, null, null)); } @Test