diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java index 6e8e8c20e4..57b10395fb 100644 --- a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java @@ -16,6 +16,8 @@ package org.thingsboard.server.service.apiusage; import com.google.common.util.concurrent.FutureCallback; +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; import org.checkerframework.checker.nullness.qual.Nullable; @@ -23,9 +25,9 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Service; -import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.rule.engine.api.MailService; +import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.ApiFeature; import org.thingsboard.server.common.data.ApiUsageRecordKey; import org.thingsboard.server.common.data.ApiUsageState; @@ -49,8 +51,8 @@ import org.thingsboard.server.common.data.tenant.profile.TenantProfileConfigurat import org.thingsboard.server.common.data.tenant.profile.TenantProfileData; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TbCallback; +import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.common.msg.tools.SchedulerUtils; -import org.thingsboard.server.dao.customer.CustomerService; import org.thingsboard.server.dao.tenant.TbTenantProfileCache; import org.thingsboard.server.dao.tenant.TenantService; import org.thingsboard.server.dao.timeseries.TimeseriesService; @@ -58,11 +60,9 @@ import org.thingsboard.server.dao.usagerecord.ApiUsageStateService; import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg; import org.thingsboard.server.gen.transport.TransportProtos.UsageStatsKVProto; import org.thingsboard.server.queue.common.TbProtoQueueMsg; -import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; import org.thingsboard.server.queue.discovery.PartitionService; -import org.thingsboard.server.queue.discovery.TbApplicationEventListener; -import org.thingsboard.server.queue.scheduler.SchedulerComponent; -import org.thingsboard.server.cluster.TbClusterService; +import org.thingsboard.server.service.executors.DbCallbackExecutorService; +import org.thingsboard.server.service.partition.AbstractPartitionBasedService; import org.thingsboard.server.service.telemetry.InternalTelemetryService; import javax.annotation.PostConstruct; @@ -79,7 +79,6 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; -import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; @@ -87,7 +86,7 @@ import java.util.stream.Collectors; @Slf4j @Service -public class DefaultTbApiUsageStateService extends TbApplicationEventListener implements TbApiUsageStateService { +public class DefaultTbApiUsageStateService extends AbstractPartitionBasedService implements TbApiUsageStateService { public static final String HOURLY = "Hourly"; public static final FutureCallback VOID_CALLBACK = new FutureCallback() { @@ -102,23 +101,22 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener myUsageStates = new ConcurrentHashMap<>(); + final Map myUsageStates = new ConcurrentHashMap<>(); // Entities that should be processed on other servers - private final Map otherUsageStates = new ConcurrentHashMap<>(); + final Map otherUsageStates = new ConcurrentHashMap<>(); - private final Set deletedEntities = Collections.newSetFromMap(new ConcurrentHashMap<>()); + final Set deletedEntities = Collections.newSetFromMap(new ConcurrentHashMap<>()); @Value("${usage.stats.report.enabled:true}") private boolean enabled; @@ -133,33 +131,37 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener msg, TbCallback callback) { ToUsageStatsServiceMsg statsMsg = msg.getValue(); @@ -216,19 +218,6 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener { - return !partitionService.resolve(ServiceType.TB_CORE, entry.getValue().getTenantId(), entry.getKey()).isMyPartition(); - }); - otherUsageStates.entrySet().removeIf(entry -> { - return partitionService.resolve(ServiceType.TB_CORE, entry.getValue().getTenantId(), entry.getKey()).isMyPartition(); - }); - initStatesFromDataBase(); - } - } - @Override public ApiUsageState getApiUsageState(TenantId tenantId) { TenantApiUsageState tenantState = (TenantApiUsageState) myUsageStates.get(tenantId); @@ -311,6 +300,18 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener entityIds = partitionedEntities.get(tpi); + if (entityIds != null) { + entityIds.add(entityId); + myUsageStates.put(entityId, state); + } else { + log.debug("[{}] belongs to external partition {}", entityId, tpi.getFullTopicName()); + throw new RuntimeException(entityId.getEntityType() + " belongs to external partition " + tpi.getFullTopicName() + "!"); + } + } + private void updateProfileThresholds(TenantId tenantId, ApiUsageStateId id, TenantProfileConfiguration oldData, TenantProfileConfiguration newData) { long ts = System.currentTimeMillis(); @@ -339,6 +340,11 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener result) { log.info("[{}] Detected update of the API state for {}: {}", state.getEntityId(), state.getEntityType(), result); apiUsageStateService.update(state.getApiUsageState()); @@ -420,7 +426,7 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener + partitionService.resolve(ServiceType.TB_CORE, entry.getValue().getTenantId(), entry.getKey()).isMyPartition()); + } + + @Override + protected void onAddedPartitions(Set addedPartitions) { try { log.info("Initializing tenant states."); updateLock.lock(); try { - ExecutorService tmpInitExecutor = ThingsBoardExecutors.newWorkStealingPool(20, "init-tenant-states-from-db"); - try { - PageDataIterable tenantIterator = new PageDataIterable<>(tenantService::findTenants, 1024); - List> futures = new ArrayList<>(); - for (Tenant tenant : tenantIterator) { - if (!myUsageStates.containsKey(tenant.getId()) && partitionService.resolve(ServiceType.TB_CORE, tenant.getId(), tenant.getId()).isMyPartition()) { + PageDataIterable tenantIterator = new PageDataIterable<>(tenantService::findTenants, 1024); + List> futures = new ArrayList<>(); + for (Tenant tenant : tenantIterator) { + TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenant.getId(), tenant.getId()); + if (addedPartitions.contains(tpi)) { + if (!myUsageStates.containsKey(tenant.getId()) && tpi.isMyPartition()) { log.debug("[{}] Initializing tenant state.", tenant.getId()); - futures.add(tmpInitExecutor.submit(() -> { + futures.add(dbExecutor.submit(() -> { try { updateTenantState((TenantApiUsageState) getOrFetchState(tenant.getId(), tenant.getId()), tenantProfileCache.get(tenant.getTenantProfileId())); log.debug("[{}] Initialized tenant state.", tenant.getId()); } catch (Exception e) { log.warn("[{}] Failed to initialize tenant API state", tenant.getId(), e); } + return null; })); } + } else { + log.debug("[{}][{}] Tenant doesn't belong to current partition. tpi [{}]", tenant.getName(), tenant.getId(), tpi); } - for (Future future : futures) { - future.get(); - } - } finally { - tmpInitExecutor.shutdownNow(); } + Futures.whenAllComplete(futures); } finally { updateLock.unlock(); } - log.info("Initialized tenant states."); + log.info("Initialized {} tenant states.", myUsageStates.size()); } catch (Exception e) { log.warn("Unknown failure", e); } @@ -521,6 +538,7 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener extends TbApplicationEventListener { + + protected final ConcurrentMap> partitionedEntities = new ConcurrentHashMap<>(); + final Queue> subscribeQueue = new ConcurrentLinkedQueue<>(); + + protected ListeningScheduledExecutorService scheduledExecutor; + + abstract protected String getSchedulerExecutorName(); + + abstract protected void onAddedPartitions(Set addedPartitions); + + abstract protected void cleanupEntityOnPartitionRemoval(T entityId); + + public Set getPartitionedEntities(TopicPartitionInfo tpi) { + return partitionedEntities.get(tpi); + } + + protected void init() { + // Should be always single threaded due to absence of locks. + scheduledExecutor = MoreExecutors.listeningDecorator(Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("device-state-scheduled"))); + } + + protected ServiceType getServiceType() { + return ServiceType.TB_CORE; + } + + protected void stop() { + if (scheduledExecutor != null) { + scheduledExecutor.shutdownNow(); + } + } + + /** + * DiscoveryService will call this event from the single thread (one-by-one). + * Events order is guaranteed by DiscoveryService. + * The only concurrency is expected from the [main] thread on Application started. + * Async implementation. Locks is not allowed by design. + * Any locks or delays in this module will affect DiscoveryService and entire system + */ + @Override + protected void onTbApplicationEvent(PartitionChangeEvent partitionChangeEvent) { + if (getServiceType().equals(partitionChangeEvent.getServiceType())) { + log.debug("onTbApplicationEvent, processing event: {}", partitionChangeEvent); + subscribeQueue.add(partitionChangeEvent.getPartitions()); + scheduledExecutor.submit(this::pollInitStateFromDB); + } + } + + protected void pollInitStateFromDB() { + final Set partitions = getLatestPartitions(); + if (partitions == null) { + log.debug("Nothing to do. Partitions are empty."); + return; + } + initStateFromDB(partitions); + } + + private void initStateFromDB(Set partitions) { + try { + log.info("CURRENT PARTITIONS: {}", partitionedEntities.keySet()); + log.info("NEW PARTITIONS: {}", partitions); + + Set addedPartitions = new HashSet<>(partitions); + addedPartitions.removeAll(partitionedEntities.keySet()); + + log.info("ADDED PARTITIONS: {}", addedPartitions); + + Set removedPartitions = new HashSet<>(partitionedEntities.keySet()); + removedPartitions.removeAll(partitions); + + log.info("REMOVED PARTITIONS: {}", removedPartitions); + + // We no longer manage current partition of entities; + removedPartitions.forEach(partition -> { + Set entities = partitionedEntities.remove(partition); + entities.forEach(this::cleanupEntityOnPartitionRemoval); + }); + + onRepartitionEvent(); + + addedPartitions.forEach(tpi -> partitionedEntities.computeIfAbsent(tpi, key -> ConcurrentHashMap.newKeySet())); + + if (!addedPartitions.isEmpty()) { + onAddedPartitions(addedPartitions); + } + + log.info("Managing following partitions:"); + partitionedEntities.forEach((tpi, entities) -> { + log.info("[{}]: {} entities", tpi.getFullTopicName(), entities.size()); + }); + } catch (Throwable t) { + log.warn("Failed to init entities state from DB", t); + } + } + + protected void onRepartitionEvent() { + } + + private Set getLatestPartitions() { + 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; + } + +} 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 8a24cdbf1a..8fc84de9e5 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 @@ -20,15 +20,15 @@ import com.google.common.base.Function; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; -import com.google.common.util.concurrent.ListeningScheduledExecutorService; -import com.google.common.util.concurrent.MoreExecutors; import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; import org.springframework.util.StringUtils; +import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.ThingsBoardThreadFactory; +import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Tenant; @@ -52,13 +52,10 @@ import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.tenant.TenantService; import org.thingsboard.server.dao.timeseries.TimeseriesService; -import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.gen.transport.TransportProtos; -import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; 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.cluster.TbClusterService; +import org.thingsboard.server.service.partition.AbstractPartitionBasedService; import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; import javax.annotation.Nonnull; @@ -68,14 +65,11 @@ import javax.annotation.PreDestroy; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; -import java.util.HashSet; import java.util.List; -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; @@ -94,7 +88,7 @@ import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE; @Service @TbCoreComponent @Slf4j -public class DefaultDeviceStateService extends TbApplicationEventListener implements DeviceStateService { +public class DefaultDeviceStateService extends AbstractPartitionBasedService implements DeviceStateService { public static final String ACTIVITY_STATE = "active"; public static final String LAST_CONNECT_TIME = "lastConnectTime"; @@ -131,12 +125,9 @@ public class DefaultDeviceStateService extends TbApplicationEventListener> partitionedDevices = new ConcurrentHashMap<>(); - final ConcurrentMap deviceStates = new ConcurrentHashMap<>(); - final Queue> subscribeQueue = new ConcurrentLinkedQueue<>(); + final ConcurrentMap deviceStates = new ConcurrentHashMap<>(); public DefaultDeviceStateService(TenantService tenantService, DeviceService deviceService, AttributesService attributesService, TimeseriesService tsService, @@ -156,25 +147,30 @@ public class DefaultDeviceStateService extends TbApplicationEventListener 0 && lastReportedActivity > stateData.getState().getLastActivityTime()) { updateActivityState(deviceId, stateData, lastReportedActivity); } - cleanDeviceStateIfBelongsExternalPartition(tenantId, deviceId); } void updateActivityState(DeviceId deviceId, DeviceStateData stateData, long lastReportedActivity) { @@ -208,18 +206,20 @@ public class DefaultDeviceStateService extends TbApplicationEventListener() { + Futures.addCallback(fetchDeviceState(device), new FutureCallback<>() { @Override public void onSuccess(@Nullable DeviceStateData state) { TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, device.getId()); - if (partitionedDevices.containsKey(tpi)) { + if (partitionedEntities.containsKey(tpi)) { addDeviceUsingState(tpi, state); save(deviceId, ACTIVITY_STATE, false); callback.onSuccess(); @@ -286,105 +289,25 @@ public class DefaultDeviceStateService extends TbApplicationEventListener 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: {}", partitions); - - Set addedPartitions = new HashSet<>(partitions); - addedPartitions.removeAll(partitionedDevices.keySet()); - - log.info("ADDED PARTITIONS: {}", addedPartitions); - - Set removedPartitions = new HashSet<>(partitionedDevices.keySet()); - removedPartitions.removeAll(partitions); - - log.info("REMOVED PARTITIONS: {}", removedPartitions); - - // We no longer manage current partition of devices; - removedPartitions.forEach(partition -> { - Set devices = partitionedDevices.remove(partition); - devices.forEach(this::cleanUpDeviceStateMap); - }); - - addedPartitions.forEach(tpi -> partitionedDevices.computeIfAbsent(tpi, key -> ConcurrentHashMap.newKeySet())); - - initPartitions(addedPartitions); - - scheduledExecutor.submit(() -> { - log.info("Managing following partitions:"); - partitionedDevices.forEach((tpi, devices) -> { - log.info("[{}]: {} devices", tpi.getFullTopicName(), devices.size()); - }); - }); - } catch (Throwable t) { - log.warn("Failed to init device states from DB", t); - } - } - - //TODO 3.0: replace this dummy search with new functionality to search by partitions using SQL capabilities. - //Adding only entities that are in new partitions - boolean initPartitions(Set addedPartitions) { - if (addedPartitions.isEmpty()) { - return false; - } - + protected void onAddedPartitions(Set addedPartitions) { List tenants = tenantService.findTenants(new PageLink(Integer.MAX_VALUE)).getData(); for (Tenant tenant : tenants) { log.debug("Finding devices for tenant [{}]", tenant.getName()); final PageLink pageLink = new PageLink(initFetchPackSize); - scheduledExecutor.submit(() -> processPageAndSubmitNextPage(addedPartitions, tenant, pageLink, scheduledExecutor)); - + processPageAndSubmitNextPage(addedPartitions, tenant, pageLink); } - return true; } - private void processPageAndSubmitNextPage(final Set addedPartitions, final Tenant tenant, final PageLink pageLink, final ExecutorService executor) { + private void processPageAndSubmitNextPage(final Set addedPartitions, final Tenant tenant, final PageLink pageLink) { log.trace("[{}] Process page {} from {}", tenant, pageLink.getPage(), pageLink.getPageSize()); List> fetchFutures = new ArrayList<>(); PageData page = deviceService.findDevicesByTenantId(tenant.getId(), pageLink); for (Device device : page.getData()) { TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenant.getId(), device.getId()); - if (addedPartitions.contains(tpi)) { + if (addedPartitions.contains(tpi) && !deviceStates.containsKey(device.getId())) { log.debug("[{}][{}] Device belong to current partition. tpi [{}]. Fetching state from DB", device.getName(), device.getId(), tpi); - ListenableFuture future = Futures.transform(fetchDeviceState(device), new Function() { + ListenableFuture future = Futures.transform(fetchDeviceState(device), new Function<>() { @Nullable @Override public Void apply(@Nullable DeviceStateData state) { @@ -403,7 +326,7 @@ public class DefaultDeviceStateService extends TbApplicationEventListener>() { + Futures.addCallback(Futures.successfulAsList(fetchFutures), new FutureCallback<>() { @Override public void onSuccess(List result) { log.trace("[{}] Success init device state from DB for batch size {}", tenant.getId(), result.size()); @@ -419,7 +342,7 @@ public class DefaultDeviceStateService extends TbApplicationEventListener processPageAndSubmitNextPage(addedPartitions, tenant, nextPageLink, executor)); + processPageAndSubmitNextPage(addedPartitions, tenant, nextPageLink); } } @@ -435,10 +358,10 @@ public class DefaultDeviceStateService extends TbApplicationEventListener deviceIds = partitionedDevices.get(tpi); + Set deviceIds = partitionedEntities.get(tpi); if (deviceIds != null) { deviceIds.add(state.getDeviceId()); - deviceStates.put(state.getDeviceId(), state); + deviceStates.putIfAbsent(state.getDeviceId(), state); } else { log.debug("[{}] Device belongs to external partition {}", state.getDeviceId(), tpi.getFullTopicName()); throw new RuntimeException("Device belongs to external partition " + tpi.getFullTopicName() + "!"); @@ -447,7 +370,7 @@ public class DefaultDeviceStateService extends TbApplicationEventListener { + partitionedEntities.forEach((tpi, deviceIds) -> { log.debug("Calculating state updates. tpi {} for {} devices", tpi.getFullTopicName(), deviceIds.size()); for (DeviceId deviceId : deviceIds) { updateInactivityStateIfExpired(ts, deviceId); @@ -464,16 +387,22 @@ public class DefaultDeviceStateService extends TbApplicationEventListener deviceIdSet = partitionedDevices.get(tpi); - deviceIdSet.remove(deviceId); + Set deviceIdSet = partitionedEntities.get(tpi); + if (deviceIdSet != null) { + deviceIdSet.remove(deviceId); + } + } + + @Override + protected void cleanupEntityOnPartitionRemoval(DeviceId deviceId) { + cleanupEntity(deviceId); } - private void cleanUpDeviceStateMap(DeviceId deviceId) { + private void cleanupEntity(DeviceId deviceId) { deviceStates.remove(deviceId); } + private ListenableFuture fetchDeviceState(Device device) { ListenableFuture future; if (persistToTelemetry) { diff --git a/application/src/test/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateServiceTest.java b/application/src/test/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateServiceTest.java new file mode 100644 index 0000000000..fde595ea3b --- /dev/null +++ b/application/src/test/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateServiceTest.java @@ -0,0 +1,86 @@ +/** + * Copyright © 2016-2022 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.service.apiusage; + +import org.junit.Before; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.Mock; +import org.mockito.Mockito; +import org.mockito.junit.MockitoJUnitRunner; +import org.thingsboard.rule.engine.api.MailService; +import org.thingsboard.server.cluster.TbClusterService; +import org.thingsboard.server.common.data.ApiUsageState; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.msg.queue.ServiceType; +import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; +import org.thingsboard.server.dao.tenant.TbTenantProfileCache; +import org.thingsboard.server.dao.tenant.TenantService; +import org.thingsboard.server.dao.timeseries.TimeseriesService; +import org.thingsboard.server.dao.usagerecord.ApiUsageStateService; +import org.thingsboard.server.queue.discovery.PartitionService; +import org.thingsboard.server.queue.scheduler.SchedulerComponent; +import org.thingsboard.server.service.executors.DbCallbackExecutorService; + +import java.util.UUID; + +import static org.hamcrest.CoreMatchers.is; +import static org.hamcrest.MatcherAssert.assertThat; +import static org.mockito.BDDMockito.willReturn; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.times; + +@RunWith(MockitoJUnitRunner.class) +public class DefaultTbApiUsageStateServiceTest { + + @Mock + TenantService tenantService; + @Mock + TimeseriesService tsService; + @Mock + TbClusterService clusterService; + @Mock + PartitionService partitionService; + @Mock + TenantApiUsageState tenantUsageStateMock; + @Mock + ApiUsageStateService apiUsageStateService; + @Mock + TbTenantProfileCache tenantProfileCache; + @Mock + MailService mailService; + @Mock + DbCallbackExecutorService dbExecutor; + + TenantId tenantId = TenantId.fromUUID(UUID.fromString("00797a3b-7aeb-4b5b-b57a-c2a810d0f112")); + + DefaultTbApiUsageStateService service; + + @Before + public void setUp() { + service = spy(new DefaultTbApiUsageStateService(clusterService, partitionService, tenantService, tsService, apiUsageStateService, tenantProfileCache, mailService, dbExecutor)); + } + + @Test + public void givenTenantIdFromEntityStatesMap_whenGetApiUsageState() { + service.myUsageStates.put(tenantId, tenantUsageStateMock); + ApiUsageState tenantUsageState = service.getApiUsageState(tenantId); + assertThat(tenantUsageState, is(tenantUsageStateMock.getApiUsageState())); + Mockito.verify(service, never()).getOrFetchState(tenantId, tenantId); + } + +} \ No newline at end of file diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java index 7478d21816..96800d0c8a 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java @@ -158,6 +158,8 @@ public class HashPartitionService implements PartitionService { } }); + tpiCache.clear(); + oldPartitions.forEach((serviceQueueKey, partitions) -> { if (!myPartitions.containsKey(serviceQueueKey)) { log.info("[{}] NO MORE PARTITIONS FOR CURRENT KEY", serviceQueueKey); @@ -174,7 +176,6 @@ public class HashPartitionService implements PartitionService { applicationEventPublisher.publishEvent(new PartitionChangeEvent(this, serviceQueueKey, tpiList)); } }); - tpiCache.clear(); if (currentOtherServices == null) { currentOtherServices = new ArrayList<>(otherServices);