From e9592174c2293b842ef6f1946dbd11d2e27056e4 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Tue, 27 Oct 2020 19:39:14 +0200 Subject: [PATCH] Api Tenant States --- .../DefaultTbApiUsageStateService.java | 47 ++++++++++++++----- 1 file changed, 35 insertions(+), 12 deletions(-) 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 1b1a2853c7..17563bcbbd 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 @@ -5,7 +5,7 @@ * 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 + * 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, @@ -23,6 +23,7 @@ import org.springframework.data.util.Pair; import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.ApiUsageRecordKey; import org.thingsboard.server.common.data.ApiUsageState; +import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.TenantProfile; import org.thingsboard.server.common.data.id.ApiUsageStateId; import org.thingsboard.server.common.data.id.TenantId; @@ -31,11 +32,13 @@ import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.BooleanDataEntry; import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.common.data.page.PageDataIterable; import org.thingsboard.server.common.data.tenant.profile.TenantProfileConfiguration; 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.tools.SchedulerUtils; +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.gen.transport.TransportProtos.ToUsageStatsServiceMsg; @@ -69,17 +72,14 @@ public class DefaultTbApiUsageStateService implements TbApiUsageStateService { public static final String HOURLY = "Hourly"; public static final FutureCallback VOID_CALLBACK = new FutureCallback() { @Override - public void onSuccess(@Nullable Void result) { - - } + public void onSuccess(@Nullable Void result) {} @Override - public void onFailure(Throwable t) { - - } + public void onFailure(Throwable t) {} }; private final TbClusterService clusterService; private final PartitionService partitionService; + private final TenantService tenantService; private final ApiUsageStateService apiUsageStateService; private final TimeseriesService tsService; private final TelemetrySubscriptionService tsWsService; @@ -101,12 +101,13 @@ public class DefaultTbApiUsageStateService implements TbApiUsageStateService { public DefaultTbApiUsageStateService(TbClusterService clusterService, PartitionService partitionService, - ApiUsageStateService apiUsageStateService, + TenantService tenantService, ApiUsageStateService apiUsageStateService, TimeseriesService tsService, TelemetrySubscriptionService tsWsService, SchedulerComponent scheduler, TbTenantProfileCache tenantProfileCache) { this.clusterService = clusterService; this.partitionService = partitionService; + this.tenantService = tenantService; this.apiUsageStateService = apiUsageStateService; this.tsService = tsService; this.tsWsService = tsWsService; @@ -117,6 +118,8 @@ public class DefaultTbApiUsageStateService implements TbApiUsageStateService { @PostConstruct public void init() { if (enabled) { + log.info("Starting api usage service."); + initStatesFromDataBase(); scheduler.scheduleAtFixedRate(this::checkStartOfNextCycle, nextCycleCheckInterval, nextCycleCheckInterval, TimeUnit.MILLISECONDS); } } @@ -164,6 +167,7 @@ public class DefaultTbApiUsageStateService implements TbApiUsageStateService { if (partitionChangeEvent.getServiceType().equals(ServiceType.TB_CORE)) { myTenantStates.entrySet().removeIf(entry -> !partitionService.resolve(ServiceType.TB_CORE, entry.getKey(), entry.getKey()).isMyPartition()); otherTenantStates.entrySet().removeIf(entry -> partitionService.resolve(ServiceType.TB_CORE, entry.getKey(), entry.getKey()).isMyPartition()); + initStatesFromDataBase(); } } @@ -299,12 +303,10 @@ public class DefaultTbApiUsageStateService implements TbApiUsageStateService { for (TsKvEntry tsKvEntry : dbValues) { if (tsKvEntry.getKey().equals(key.getApiCountKey())) { cycleEntryFound = true; - tenantState.put(key, tsKvEntry.getLongValue().get()); + tenantState.put(key, tsKvEntry.getTs() == tenantState.getCurrentCycleTs() ? tsKvEntry.getLongValue().get() : 0L); } else if (tsKvEntry.getKey().equals(key.getApiCountKey() + HOURLY)) { hourlyEntryFound = true; - if (tsKvEntry.getTs() == tenantState.getCurrentHourTs()) { - tenantState.putHourly(key, tsKvEntry.getLongValue().get()); - } + tenantState.putHourly(key, tsKvEntry.getTs() == tenantState.getCurrentHourTs() ? tsKvEntry.getLongValue().get() : 0L); } if (cycleEntryFound && hourlyEntryFound) { break; @@ -319,4 +321,25 @@ public class DefaultTbApiUsageStateService implements TbApiUsageStateService { return tenantState; } + private void initStatesFromDataBase() { + try { + PageDataIterable tenantIterator = new PageDataIterable<>(tenantService::findTenants, 1024); + for (Tenant tenant : tenantIterator) { + if (!myTenantStates.containsKey(tenant.getId()) && partitionService.resolve(ServiceType.TB_CORE, tenant.getId(), tenant.getId()).isMyPartition()) { + updateLock.lock(); + try { + updateTenantState(getOrFetchState(tenant.getId()), tenantProfileCache.get(tenant.getTenantProfileId())); + } catch (Exception e) { + log.warn("[{}] Failed to initialize tenant API state", tenant.getId(), e); + } finally { + updateLock.unlock(); + } + } + } + log.info("Api usage service started."); + } catch (Exception e) { + log.warn("Unknown failure", e); + } + } + }