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 9abc32cf6b..f3798ecc5b 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 @@ -59,8 +59,7 @@ import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.discovery.TbApplicationEventListener; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.queue.TbClusterService; -import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; -import org.thingsboard.server.utils.EventDeduplicationExecutor; +import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService import javax.annotation.Nullable; import javax.annotation.PostConstruct; @@ -71,10 +70,12 @@ import java.util.Collections; import java.util.HashSet; import java.util.List; import java.util.Optional; +import java.util.Queue; import java.util.Random; import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; @@ -130,13 +131,13 @@ public class DefaultDeviceStateService extends TbApplicationEventListener> partitionedDevices = new ConcurrentHashMap<>(); private final ConcurrentMap deviceStates = new ConcurrentHashMap<>(); private final ConcurrentMap deviceLastSavedActivity = new ConcurrentHashMap<>(); - private volatile EventDeduplicationExecutor> deduplicationExecutor; + final Queue> subscribeQueue = new ConcurrentLinkedQueue<>(); public DefaultDeviceStateService(TenantService tenantService, DeviceService deviceService, AttributesService attributesService, TimeseriesService tsService, @@ -156,21 +157,20 @@ public class DefaultDeviceStateService extends TbApplicationEventListener(DefaultDeviceStateService.class.getSimpleName(), queueExecutor, this::initStateFromDB); + scheduledExecutor = MoreExecutors.listeningDecorator(Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("device-state-scheduled"))); + scheduledExecutor.scheduleAtFixedRate(this::updateState, new Random().nextInt(defaultStateCheckIntervalInSec), defaultStateCheckIntervalInSec, TimeUnit.SECONDS); } @PreDestroy public void stop() { - if (executorService != null) { - executorService.shutdownNow(); + if (dbCallbackExecutorService != null) { + dbCallbackExecutorService.shutdownNow(); } - if (queueExecutor != null) { - queueExecutor.shutdownNow(); + if (scheduledExecutor != null) { + scheduledExecutor.shutdownNow(); } } @@ -283,7 +283,7 @@ public class DefaultDeviceStateService extends TbApplicationEventListener pendingPartitions) { + void pollInitStateFromDB() { + final Set partitions = getLatestPartitionsFromQueue(); + if (partitions == null) { + log.info("Device state service. Nothing to do. partitions is null"); + return; + } + initStateFromDB(partitions); + } + + // TODO: move to utils + Set getLatestPartitionsFromQueue() { + log.debug("getLatestPartitionsFromQueue, queue size {}", subscribeQueue.size()); + Set partitions = null; + while (!subscribeQueue.isEmpty()) { + partitions = subscribeQueue.poll(); + log.debug("polled from the queue partitions {}", partitions); + } + log.debug("getLatestPartitionsFromQueue, partitions {}", partitions); + return partitions; + } + + private void initStateFromDB(Set partitions) { try { log.info("CURRENT PARTITIONS: {}", partitionedDevices.keySet()); - log.info("NEW PARTITIONS: {}", pendingPartitions); + log.info("NEW PARTITIONS: {}", partitions); - Set addedPartitions = new HashSet<>(pendingPartitions); + Set addedPartitions = new HashSet<>(partitions); addedPartitions.removeAll(partitionedDevices.keySet()); log.info("ADDED PARTITIONS: {}", addedPartitions); Set removedPartitions = new HashSet<>(partitionedDevices.keySet()); - removedPartitions.removeAll(pendingPartitions); + removedPartitions.removeAll(partitions); log.info("REMOVED PARTITIONS: {}", removedPartitions); @@ -363,7 +393,7 @@ public class DefaultDeviceStateService extends TbApplicationEventListener fetchDeviceState(Device device) { if (persistToTelemetry) { ListenableFuture> tsData = tsService.findLatest(TenantId.SYS_TENANT_ID, device.getId(), PERSISTENT_ATTRIBUTES); - return Futures.transform(tsData, extractDeviceStateData(device), executorService); + return Futures.transform(tsData, extractDeviceStateData(device), dbCallbackExecutorService); } else { ListenableFuture> attrData = attributesService.find(TenantId.SYS_TENANT_ID, device.getId(), DataConstants.SERVER_SCOPE, PERSISTENT_ATTRIBUTES); - return Futures.transform(attrData, extractDeviceStateData(device), executorService); + return Futures.transform(attrData, extractDeviceStateData(device), dbCallbackExecutorService); } } private Function, DeviceStateData> extractDeviceStateData(Device device) { - return new Function<>() { + return new Function, DeviceStateData>() { @Nullable @Override public DeviceStateData apply(@Nullable List data) { @@ -470,7 +500,21 @@ public class DefaultDeviceStateService extends TbApplicationEventListener inactivityTimeoutOpt = + attributesService.find(TenantId.SYS_TENANT_ID, device.getId(), SERVER_SCOPE, INACTIVITY_TIMEOUT).get(); + if (inactivityTimeoutOpt.isPresent() && inactivityTimeoutOpt.get().getLongValue().isPresent() + && inactivityTimeoutOpt.get().getLongValue().get() > 0) { + inactivityTimeout = inactivityTimeoutOpt.get().getLongValue().get(); + } + } catch (Exception ignored) { + } + } + // TODO: voba - do we need to calculate it or it's better to fetch from DB directly? + // boolean active = System.currentTimeMillis() < lastActivityTime + inactivityTimeout; + boolean active = getEntryValue(data, ACTIVITY_STATE, false); DeviceState deviceState = DeviceState.builder() .active(active) .lastConnectTime(getEntryValue(data, LAST_CONNECT_TIME, 0L)) @@ -510,6 +554,17 @@ public class DefaultDeviceStateService extends TbApplicationEventListener kvEntries, String attributeName, boolean defaultValue) { + if (kvEntries != null) { + for (KvEntry entry : kvEntries) { + if (entry != null && !StringUtils.isEmpty(entry.getKey()) && entry.getKey().equals(attributeName)) { + return entry.getBooleanValue().orElse(defaultValue); + } + } + } + return defaultValue; + } + private void pushRuleEngineMessage(DeviceStateData stateData, String msgType) { DeviceState state = stateData.getState(); try {