From 0e0ab6ca3982d7e20c47267e1967010ad9e607f7 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Thu, 22 Oct 2020 15:56:31 +0300 Subject: [PATCH] Improvements to API State --- .../server/service/apiusage/ApiFeature.java | 5 + .../DefaultTbApiUsageStateService.java | 92 ++++++++++++++----- .../apiusage/TbApiUsageStateService.java | 7 +- .../service/apiusage/TenantApiUsageState.java | 43 ++++++--- .../queue/DefaultTbClusterService.java | 27 ++++-- .../service/queue/TbClusterService.java | 3 + .../processing/AbstractConsumerService.java | 4 +- .../transport/DefaultTransportApiService.java | 22 +++-- .../server/common/data/ApiUsageRecordKey.java | 1 - .../server/common/data/ApiUsageState.java | 30 ++++-- .../MultipleTbQueueCallbackWrapper.java | 45 +++++++++ common/queue/src/main/proto/queue.proto | 1 + .../transport/mqtt/MqttTransportHandler.java | 9 +- .../common/transport/TransportService.java | 4 +- .../DefaultTransportRateLimitService.java | 11 ++- .../limits/TransportRateLimitService.java | 1 + .../limits/TransportRateLimitType.java | 3 +- .../service/DefaultTransportService.java | 17 ++-- .../DefaultTransportTenantProfileCache.java | 26 ++++-- .../server/dao/model/ModelConstants.java | 4 + .../dao/model/sql/ApiUsageStateEntity.java | 19 ++++ ...mpl.java => ApiUsageStateServiceImpl.java} | 4 +- .../resources/sql/schema-entities-hsql.sql | 4 + .../main/resources/sql/schema-entities.sql | 4 + .../service/BaseApiUsageStateServiceTest.java | 13 +++ 25 files changed, 310 insertions(+), 89 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/apiusage/ApiFeature.java create mode 100644 common/queue/src/main/java/org/thingsboard/server/queue/common/MultipleTbQueueCallbackWrapper.java rename dao/src/main/java/org/thingsboard/server/dao/usagerecord/{ApiApiUsageStateServiceImpl.java => ApiUsageStateServiceImpl.java} (96%) diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/ApiFeature.java b/application/src/main/java/org/thingsboard/server/service/apiusage/ApiFeature.java new file mode 100644 index 0000000000..a893a55873 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/apiusage/ApiFeature.java @@ -0,0 +1,5 @@ +package org.thingsboard.server.service.apiusage; + +public enum ApiFeature { + TRANSPORT, DB, RE, JS +} 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 df7e5013d2..c5e6cab16d 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 @@ -17,6 +17,7 @@ package org.thingsboard.server.service.apiusage; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; +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; @@ -33,15 +34,18 @@ import org.thingsboard.server.dao.timeseries.TimeseriesService; 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.TbQueueCallback; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.discovery.PartitionChangeEvent; import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.scheduler.SchedulerComponent; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.profile.TbTenantProfileCache; +import org.thingsboard.server.service.queue.TbClusterService; import javax.annotation.PostConstruct; import java.util.ArrayList; +import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.UUID; @@ -57,12 +61,17 @@ import java.util.concurrent.locks.ReentrantLock; public class DefaultTbApiUsageStateService implements TbApiUsageStateService { public static final String HOURLY = "HOURLY_"; + private final TbClusterService clusterService; private final PartitionService partitionService; private final ApiUsageStateService apiUsageStateService; private final TimeseriesService tsService; private final SchedulerComponent scheduler; private final TbTenantProfileCache tenantProfileCache; - private final Map tenantStates = new ConcurrentHashMap<>(); + + // Tenants that should be processed on this server + private final Map myTenantStates = new ConcurrentHashMap<>(); + // Tenants that should be processed on other servers + private final Map otherTenantStates = new ConcurrentHashMap<>(); @Value("${usage.stats.report.enabled:true}") private boolean enabled; @@ -72,7 +81,13 @@ public class DefaultTbApiUsageStateService implements TbApiUsageStateService { private final Lock updateLock = new ReentrantLock(); - public DefaultTbApiUsageStateService(PartitionService partitionService, ApiUsageStateService apiUsageStateService, TimeseriesService tsService, SchedulerComponent scheduler, TbTenantProfileCache tenantProfileCache) { + public DefaultTbApiUsageStateService(TbClusterService clusterService, + PartitionService partitionService, + ApiUsageStateService apiUsageStateService, + TimeseriesService tsService, + SchedulerComponent scheduler, + TbTenantProfileCache tenantProfileCache) { + this.clusterService = clusterService; this.partitionService = partitionService; this.apiUsageStateService = apiUsageStateService; this.tsService = tsService; @@ -93,7 +108,7 @@ public class DefaultTbApiUsageStateService implements TbApiUsageStateService { TenantId tenantId = new TenantId(new UUID(statsMsg.getTenantIdMSB(), statsMsg.getTenantIdLSB())); TenantApiUsageState tenantState; List updatedEntries; - boolean stateUpdated = false; + Map result = new HashMap<>(); updateLock.lock(); try { tenantState = getOrFetchState(tenantId); @@ -110,32 +125,49 @@ public class DefaultTbApiUsageStateService implements TbApiUsageStateService { updatedEntries.add(new BasicTsKvEntry(ts, new LongDataEntry(recordKey.name(), newValue))); long newHourlyValue = tenantState.addToHourly(recordKey, kvProto.getValue()); updatedEntries.add(new BasicTsKvEntry(hourTs, new LongDataEntry(HOURLY + recordKey.name(), newHourlyValue))); - stateUpdated |= tenantState.checkStateUpdatedDueToThreshold(recordKey); + Pair update = tenantState.checkStateUpdatedDueToThreshold(recordKey); + if (update != null) { + result.put(update.getFirst(), update.getSecond()); + } } } finally { updateLock.unlock(); } tsService.save(tenantId, tenantState.getApiUsageState().getId(), updatedEntries, 0L); - if (stateUpdated) { - // Save new state into the database; - apiUsageStateService.update(tenantState.getApiUsageState()); - //TODO: clear cache on cluster repartition. - //TODO: update profiles on tenant and profile updates. - //TODO: broadcast to everyone notifications about enabled/disabled features. + if (!result.isEmpty()) { + persistAndNotify(tenantState, result); } } @Override public void onApplicationEvent(PartitionChangeEvent partitionChangeEvent) { if (partitionChangeEvent.getServiceType().equals(ServiceType.TB_CORE)) { - tenantStates.entrySet().removeIf(entry -> !partitionService.resolve(ServiceType.TB_CORE, entry.getKey(), entry.getKey()).isMyPartition()); + 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()); } } @Override - public TenantApiUsageState getApiUsageState(TenantId tenantId) { - //We should always get it from the map of from the database; - return null; + public ApiUsageState getApiUsageState(TenantId tenantId) { + if (partitionService.resolve(ServiceType.TB_CORE, tenantId, tenantId).isMyPartition()) { + TenantApiUsageState state = getOrFetchState(tenantId); + return state.getApiUsageState(); + } else { + ApiUsageState state = otherTenantStates.get(tenantId); + if (state == null) { + updateLock.lock(); + try { + state = otherTenantStates.get(tenantId); + if (state == null) { + state = apiUsageStateService.findTenantApiUsageState(tenantId); + otherTenantStates.put(tenantId, state); + } + } finally { + updateLock.unlock(); + } + } + return state; + } } @Override @@ -143,7 +175,7 @@ public class DefaultTbApiUsageStateService implements TbApiUsageStateService { TenantProfile tenantProfile = tenantProfileCache.get(tenantProfileId); updateLock.lock(); try { - tenantStates.values().forEach(state -> { + myTenantStates.values().forEach(state -> { if (tenantProfile.getId().equals(state.getTenantProfileId())) { updateTenantState(state, tenantProfile); } @@ -158,7 +190,7 @@ public class DefaultTbApiUsageStateService implements TbApiUsageStateService { TenantProfile tenantProfile = tenantProfileCache.get(tenantId); updateLock.lock(); try { - TenantApiUsageState state = tenantStates.get(tenantId); + TenantApiUsageState state = myTenantStates.get(tenantId); if (state != null && !state.getTenantProfileId().equals(tenantProfile.getId())) { updateTenantState(state, tenantProfile); } @@ -167,12 +199,28 @@ public class DefaultTbApiUsageStateService implements TbApiUsageStateService { } } + @Override + public void onApiUsageStateUpdate(TenantId tenantId) { + + } private void updateTenantState(TenantApiUsageState state, TenantProfile tenantProfile) { state.setTenantProfileData(tenantProfile.getProfileData()); - if (state.checkStateUpdatedDueToThresholds()) { - apiUsageStateService.update(state.getApiUsageState()); - //TODO: send notification to cluster; + Map result = state.checkStateUpdatedDueToThresholds(); + if (!result.isEmpty()) { + persistAndNotify(state, result); + } + } + + private void persistAndNotify(TenantApiUsageState state, Map result) { + // TODO: + // 1. Broadcast to everyone notifications about enabled/disabled features. + // 2. Report rule engine and js executor metrics + // 4. UI for configuration of the thresholds + // 5. Max rule node executions per message. + apiUsageStateService.update(state.getApiUsageState()); + if (result.containsKey(ApiFeature.TRANSPORT)) { + clusterService.onApiStateChange(state.getApiUsageState(), null); } } @@ -180,7 +228,7 @@ public class DefaultTbApiUsageStateService implements TbApiUsageStateService { updateLock.lock(); try { long now = System.currentTimeMillis(); - tenantStates.values().forEach(state -> { + myTenantStates.values().forEach(state -> { if ((state.getNextCycleTs() > now) && (state.getNextCycleTs() - now < TimeUnit.HOURS.toMillis(1))) { state.setCycles(state.getNextCycleTs(), SchedulerUtils.getStartOfNextNextMonth()); } @@ -191,7 +239,7 @@ public class DefaultTbApiUsageStateService implements TbApiUsageStateService { } private TenantApiUsageState getOrFetchState(TenantId tenantId) { - TenantApiUsageState tenantState = tenantStates.get(tenantId); + TenantApiUsageState tenantState = myTenantStates.get(tenantId); if (tenantState == null) { ApiUsageState dbStateEntity = apiUsageStateService.findTenantApiUsageState(tenantId); if (dbStateEntity == null) { @@ -223,7 +271,7 @@ public class DefaultTbApiUsageStateService implements TbApiUsageStateService { } } } - tenantStates.put(tenantId, tenantState); + myTenantStates.put(tenantId, tenantState); } catch (InterruptedException | ExecutionException e) { log.warn("[{}] Failed to fetch api usage state from db.", tenantId, e); } diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/TbApiUsageStateService.java b/application/src/main/java/org/thingsboard/server/service/apiusage/TbApiUsageStateService.java index 4816dc1037..e19c0a4154 100644 --- a/application/src/main/java/org/thingsboard/server/service/apiusage/TbApiUsageStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/apiusage/TbApiUsageStateService.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, @@ -16,6 +16,7 @@ package org.thingsboard.server.service.apiusage; import org.springframework.context.ApplicationListener; +import org.thingsboard.server.common.data.ApiUsageState; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantProfileId; import org.thingsboard.server.common.msg.queue.TbCallback; @@ -27,9 +28,11 @@ public interface TbApiUsageStateService extends ApplicationListener msg, TbCallback callback); - TenantApiUsageState getApiUsageState(TenantId tenantId); + ApiUsageState getApiUsageState(TenantId tenantId); void onTenantProfileUpdate(TenantProfileId tenantProfileId); void onTenantUpdate(TenantId tenantId); + + void onApiUsageStateUpdate(TenantId tenantId); } diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/TenantApiUsageState.java b/application/src/main/java/org/thingsboard/server/service/apiusage/TenantApiUsageState.java index 0492ba9793..71bc3f1f6d 100644 --- a/application/src/main/java/org/thingsboard/server/service/apiusage/TenantApiUsageState.java +++ b/application/src/main/java/org/thingsboard/server/service/apiusage/TenantApiUsageState.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, @@ -17,14 +17,17 @@ package org.thingsboard.server.service.apiusage; import lombok.Getter; import lombok.Setter; +import org.springframework.data.util.Pair; import org.thingsboard.server.common.data.ApiUsageRecordKey; import org.thingsboard.server.common.data.ApiUsageState; import org.thingsboard.server.common.data.TenantProfile; import org.thingsboard.server.common.data.TenantProfileData; import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantProfileId; import org.thingsboard.server.common.msg.tools.SchedulerUtils; +import java.util.HashMap; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; @@ -103,11 +106,17 @@ public class TenantApiUsageState { return Long.parseLong((String) threshold); } else if (threshold instanceof Long) { return (Long) threshold; + } else if (threshold instanceof Integer) { + return (Integer) threshold; } } return 0L; } + public TenantId getTenantId() { + return apiUsageState.getTenantId(); + } + public EntityId getEntityId() { return apiUsageState.getEntityId(); } @@ -121,7 +130,7 @@ public class TenantApiUsageState { } public boolean isRuleEngineEnabled() { - return apiUsageState.isRuleEngineEnabled(); + return apiUsageState.isReExecEnabled(); } public boolean isJsExecEnabled() { @@ -137,7 +146,7 @@ public class TenantApiUsageState { } public void setRuleEngineEnabled(boolean ruleEngineEnabled) { - apiUsageState.setRuleEngineEnabled(ruleEngineEnabled); + apiUsageState.setReExecEnabled(ruleEngineEnabled); } public void setJsExecEnabled(boolean jsExecEnabled) { @@ -147,7 +156,6 @@ public class TenantApiUsageState { public boolean isFeatureEnabled(ApiUsageRecordKey recordKey) { switch (recordKey) { case MSG_COUNT: - case MSG_BYTES_COUNT: case DP_TRANSPORT_COUNT: return isTransportEnabled(); case RE_EXEC_COUNT: @@ -161,38 +169,47 @@ public class TenantApiUsageState { } } - public boolean setFeatureValue(ApiUsageRecordKey recordKey, boolean value) { + public ApiFeature setFeatureValue(ApiUsageRecordKey recordKey, boolean value) { + ApiFeature feature = null; boolean currentValue = isFeatureEnabled(recordKey); switch (recordKey) { case MSG_COUNT: - case MSG_BYTES_COUNT: case DP_TRANSPORT_COUNT: + feature = ApiFeature.TRANSPORT; setTransportEnabled(value); break; case RE_EXEC_COUNT: + feature = ApiFeature.RE; setRuleEngineEnabled(value); break; case DP_STORAGE_COUNT: + feature = ApiFeature.DB; setDbStorageEnabled(value); break; case JS_EXEC_COUNT: + feature = ApiFeature.JS; setJsExecEnabled(value); break; } - return currentValue == value; + return currentValue == value ? null : feature; } - public boolean checkStateUpdatedDueToThresholds() { - boolean update = false; + public Map checkStateUpdatedDueToThresholds() { + Map result = new HashMap<>(); for (ApiUsageRecordKey key : ApiUsageRecordKey.values()) { - update |= checkStateUpdatedDueToThreshold(key); + Pair featureUpdate = checkStateUpdatedDueToThreshold(key); + if (featureUpdate != null) { + result.put(featureUpdate.getFirst(), featureUpdate.getSecond()); + } } - return update; + return result; } - public boolean checkStateUpdatedDueToThreshold(ApiUsageRecordKey recordKey) { + public Pair checkStateUpdatedDueToThreshold(ApiUsageRecordKey recordKey) { long value = get(recordKey); long threshold = getProfileThreshold(recordKey); - return setFeatureValue(recordKey, threshold == 0 || value < threshold); + boolean featureValue = threshold == 0 || value < threshold; + ApiFeature feature = setFeatureValue(recordKey, featureValue); + return feature != null ? Pair.of(feature, featureValue) : null; } } diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java index 9451b58e9f..b26baf8d2e 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.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, @@ -21,6 +21,7 @@ import org.springframework.beans.factory.annotation.Value; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; import org.thingsboard.rule.engine.api.msg.ToDeviceActorNotificationMsg; +import org.thingsboard.server.common.data.ApiUsageState; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.HasName; @@ -47,6 +48,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotifica import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg; import org.thingsboard.server.queue.TbQueueCallback; import org.thingsboard.server.queue.TbQueueProducer; +import org.thingsboard.server.queue.common.MultipleTbQueueCallbackWrapper; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.provider.TbQueueProducerProvider; @@ -206,6 +208,12 @@ public class DefaultTbClusterService implements TbClusterService { onEntityChange(TenantId.SYS_TENANT_ID, tenant.getId(), tenant, callback); } + @Override + public void onApiStateChange(ApiUsageState apiUsageState, TbQueueCallback callback) { + onEntityChange(apiUsageState.getTenantId(), apiUsageState.getId(), apiUsageState, callback); + broadcast(new ComponentLifecycleMsg(apiUsageState.getTenantId(), apiUsageState.getId(), ComponentLifecycleEvent.UPDATED)); + } + @Override public void onDeviceProfileDelete(DeviceProfile entity, TbQueueCallback callback) { onEntityDelete(entity.getTenantId(), entity.getId(), entity.getName(), callback); @@ -221,13 +229,14 @@ public class DefaultTbClusterService implements TbClusterService { onEntityDelete(TenantId.SYS_TENANT_ID, entity.getId(), entity.getName(), callback); } - public void onEntityChange(TenantId tenantId, EntityId entityid, T entity, TbQueueCallback callback) { - log.trace("[{}][{}][{}] Processing [{}] change event", tenantId, entityid.getEntityType(), entityid.getId(), entity.getName()); + public void onEntityChange(TenantId tenantId, EntityId entityid, T entity, TbQueueCallback callback) { + String entityName = (entity instanceof HasName) ? ((HasName) entity).getName() : entity.getClass().getName(); + log.trace("[{}][{}][{}] Processing [{}] change event", tenantId, entityid.getEntityType(), entityid.getId(), entityName); TransportProtos.EntityUpdateMsg entityUpdateMsg = TransportProtos.EntityUpdateMsg.newBuilder() .setEntityType(entityid.getEntityType().name()) .setData(ByteString.copyFrom(encodingService.encode(entity))).build(); ToTransportMsg transportMsg = ToTransportMsg.newBuilder().setEntityUpdateMsg(entityUpdateMsg).build(); - broadcast(transportMsg); + broadcast(transportMsg, callback); } private void onEntityDelete(TenantId tenantId, EntityId entityId, String name, TbQueueCallback callback) { @@ -238,15 +247,16 @@ public class DefaultTbClusterService implements TbClusterService { .setEntityIdLSB(entityId.getId().getLeastSignificantBits()) .build(); ToTransportMsg transportMsg = ToTransportMsg.newBuilder().setEntityDeleteMsg(entityDeleteMsg).build(); - broadcast(transportMsg); + broadcast(transportMsg, callback); } - private void broadcast(ToTransportMsg transportMsg) { + private void broadcast(ToTransportMsg transportMsg, TbQueueCallback callback) { TbQueueProducer> toTransportNfProducer = producerProvider.getTransportNotificationsMsgProducer(); Set tbTransportServices = partitionService.getAllServiceIds(ServiceType.TB_TRANSPORT); + TbQueueCallback proxyCallback = callback != null ? new MultipleTbQueueCallbackWrapper(tbTransportServices.size(), callback) : null; for (String transportServiceId : tbTransportServices) { TopicPartitionInfo tpi = partitionService.getNotificationsTopic(ServiceType.TB_TRANSPORT, transportServiceId); - toTransportNfProducer.send(tpi, new TbProtoQueueMsg<>(UUID.randomUUID(), transportMsg), null); + toTransportNfProducer.send(tpi, new TbProtoQueueMsg<>(UUID.randomUUID(), transportMsg), proxyCallback); toTransportNfs.incrementAndGet(); } } @@ -256,7 +266,8 @@ public class DefaultTbClusterService implements TbClusterService { TbQueueProducer> toRuleEngineProducer = producerProvider.getRuleEngineNotificationsMsgProducer(); Set tbRuleEngineServices = new HashSet<>(partitionService.getAllServiceIds(ServiceType.TB_RULE_ENGINE)); if (msg.getEntityId().getEntityType().equals(EntityType.TENANT) - || msg.getEntityId().getEntityType().equals(EntityType.DEVICE_PROFILE)) { + || msg.getEntityId().getEntityType().equals(EntityType.DEVICE_PROFILE) + || msg.getEntityId().getEntityType().equals(EntityType.API_USAGE_STATE)) { TbQueueProducer> toCoreNfProducer = producerProvider.getTbCoreNotificationsMsgProducer(); Set tbCoreServices = partitionService.getAllServiceIds(ServiceType.TB_CORE); for (String serviceId : tbCoreServices) { diff --git a/application/src/main/java/org/thingsboard/server/service/queue/TbClusterService.java b/application/src/main/java/org/thingsboard/server/service/queue/TbClusterService.java index 838dbe0eef..b36af16716 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/TbClusterService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/TbClusterService.java @@ -16,6 +16,7 @@ package org.thingsboard.server.service.queue; import org.thingsboard.rule.engine.api.msg.ToDeviceActorNotificationMsg; +import org.thingsboard.server.common.data.ApiUsageState; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.TenantProfile; @@ -63,4 +64,6 @@ public interface TbClusterService { void onTenantChange(Tenant tenant, TbQueueCallback callback); void onTenantDelete(Tenant tenant, TbQueueCallback callback); + + void onApiStateChange(ApiUsageState apiUsageState, TbQueueCallback callback); } diff --git a/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java index 0532d2f9c8..02fe877cd0 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.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, @@ -161,6 +161,8 @@ public abstract class AbstractConsumerService deviceCreationLocks = new ConcurrentHashMap<>(); public DefaultTransportApiService(TbDeviceProfileCache deviceProfileCache, - TbTenantProfileCache tenantProfileCache, DeviceService deviceService, + TbTenantProfileCache tenantProfileCache, TbApiUsageStateService apiUsageStateService, DeviceService deviceService, RelationService relationService, DeviceCredentialsService deviceCredentialsService, DeviceStateService deviceStateService, DbCallbackExecutorService dbCallbackExecutorService, TbClusterService tbClusterService, DataDecodingEncodingService dataDecodingEncodingService, DeviceProvisionService deviceProvisionService) { this.deviceProfileCache = deviceProfileCache; this.tenantProfileCache = tenantProfileCache; + this.apiUsageStateService = apiUsageStateService; this.deviceService = deviceService; this.relationService = relationService; this.deviceCredentialsService = deviceCredentialsService; @@ -316,18 +319,21 @@ public class DefaultTransportApiService implements TransportApiService { private ListenableFuture handle(GetEntityProfileRequestMsg requestMsg) { EntityType entityType = EntityType.valueOf(requestMsg.getEntityType()); UUID entityUuid = new UUID(requestMsg.getEntityIdMSB(), requestMsg.getEntityIdLSB()); - ByteString data; + GetEntityProfileResponseMsg.Builder builder = GetEntityProfileResponseMsg.newBuilder(); if (entityType.equals(EntityType.DEVICE_PROFILE)) { DeviceProfileId deviceProfileId = new DeviceProfileId(entityUuid); DeviceProfile deviceProfile = deviceProfileCache.find(deviceProfileId); - data = ByteString.copyFrom(dataDecodingEncodingService.encode(deviceProfile)); + builder.setData(ByteString.copyFrom(dataDecodingEncodingService.encode(deviceProfile))); } else if (entityType.equals(EntityType.TENANT)) { - TenantProfile tenantProfile = tenantProfileCache.get(new TenantId(entityUuid)); - data = ByteString.copyFrom(dataDecodingEncodingService.encode(tenantProfile)); + TenantId tenantId = new TenantId(entityUuid); + TenantProfile tenantProfile = tenantProfileCache.get(tenantId); + ApiUsageState state = apiUsageStateService.getApiUsageState(tenantId); + builder.setData(ByteString.copyFrom(dataDecodingEncodingService.encode(tenantProfile))); + builder.setApiState(ByteString.copyFrom(dataDecodingEncodingService.encode(state))); } else { throw new RuntimeException("Invalid entity profile request: " + entityType); } - return Futures.immediateFuture(TransportApiResponseMsg.newBuilder().setEntityProfileResponseMsg(GetEntityProfileResponseMsg.newBuilder().setData(data).build()).build()); + return Futures.immediateFuture(TransportApiResponseMsg.newBuilder().setEntityProfileResponseMsg(builder).build()); } private ListenableFuture getDeviceInfo(DeviceId deviceId, DeviceCredentials credentials) { diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageRecordKey.java b/common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageRecordKey.java index 9c8690a88f..bc730e7e5a 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageRecordKey.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageRecordKey.java @@ -18,7 +18,6 @@ package org.thingsboard.server.common.data; public enum ApiUsageRecordKey { MSG_COUNT, - MSG_BYTES_COUNT, DP_TRANSPORT_COUNT, DP_STORAGE_COUNT, RE_EXEC_COUNT, diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageState.java b/common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageState.java index 05186c7599..5144861790 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageState.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageState.java @@ -29,18 +29,24 @@ public class ApiUsageState extends BaseData implements HasTenan private static final long serialVersionUID = 8250339805336035966L; - @Getter @Setter + @Getter + @Setter private TenantId tenantId; - @Getter @Setter + @Getter + @Setter private EntityId entityId; - @Getter @Setter - private boolean transportEnabled; - @Getter @Setter - private boolean dbStorageEnabled; - @Getter @Setter - private boolean ruleEngineEnabled; - @Getter @Setter - private boolean jsExecEnabled; + @Getter + @Setter + private boolean transportEnabled = true; + @Getter + @Setter + private boolean dbStorageEnabled = true; + @Getter + @Setter + private boolean reExecEnabled = true; + @Getter + @Setter + private boolean jsExecEnabled = true; public ApiUsageState() { super(); @@ -54,5 +60,9 @@ public class ApiUsageState extends BaseData implements HasTenan super(ur); this.tenantId = ur.getTenantId(); this.entityId = ur.getEntityId(); + this.transportEnabled = ur.isTransportEnabled(); + this.dbStorageEnabled = ur.isDbStorageEnabled(); + this.reExecEnabled = ur.isReExecEnabled(); + this.jsExecEnabled = ur.isJsExecEnabled(); } } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/MultipleTbQueueCallbackWrapper.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/MultipleTbQueueCallbackWrapper.java new file mode 100644 index 0000000000..ddedb3a606 --- /dev/null +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/MultipleTbQueueCallbackWrapper.java @@ -0,0 +1,45 @@ +/** + * Copyright © 2016-2020 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.queue.common; + +import org.thingsboard.server.common.msg.queue.RuleEngineException; +import org.thingsboard.server.queue.TbQueueCallback; +import org.thingsboard.server.queue.TbQueueMsgMetadata; + +import java.util.concurrent.atomic.AtomicInteger; + +public class MultipleTbQueueCallbackWrapper implements TbQueueCallback { + + private final AtomicInteger tbQueueCallbackCount; + private final TbQueueCallback callback; + + public MultipleTbQueueCallbackWrapper(int tbQueueCallbackCount, TbQueueCallback callback) { + this.tbQueueCallbackCount = new AtomicInteger(tbQueueCallbackCount); + this.callback = callback; + } + + @Override + public void onSuccess(TbQueueMsgMetadata metadata) { + if (tbQueueCallbackCount.decrementAndGet() <= 0) { + callback.onSuccess(metadata); + } + } + + @Override + public void onFailure(Throwable t) { + callback.onFailure(new RuleEngineException(t.getMessage())); + } +} diff --git a/common/queue/src/main/proto/queue.proto b/common/queue/src/main/proto/queue.proto index f716bd90ae..4b5a6028e6 100644 --- a/common/queue/src/main/proto/queue.proto +++ b/common/queue/src/main/proto/queue.proto @@ -186,6 +186,7 @@ message GetEntityProfileRequestMsg { message GetEntityProfileResponseMsg { string entityType = 1; bytes data = 2; + bytes apiState = 3; } message EntityUpdateMsg { diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index 2c4b2db7df..e70da7029e 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.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, @@ -45,6 +45,7 @@ import org.thingsboard.server.common.data.DeviceTransportType; import org.thingsboard.server.common.data.TransportPayloadType; import org.thingsboard.server.common.data.device.profile.MqttTopics; import org.thingsboard.server.common.msg.EncryptionUtil; +import org.thingsboard.server.common.msg.tools.TbRateLimitsException; import org.thingsboard.server.common.transport.SessionMsgListener; import org.thingsboard.server.common.transport.TransportService; import org.thingsboard.server.common.transport.TransportServiceCallback; @@ -630,7 +631,11 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement @Override public void onError(Throwable e) { - log.warn("[{}] Failed to submit session event", sessionId, e); + if (e instanceof TbRateLimitsException) { + log.trace("[{}] Failed to submit session event", sessionId, e); + } else { + log.warn("[{}] Failed to submit session event", sessionId, e); + } ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_SERVER_UNAVAILABLE)); ctx.close(); } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java index 4872d8e477..3b775b9f26 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java @@ -45,7 +45,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceX509Ce */ public interface TransportService { - GetEntityProfileResponseMsg getRoutingInfo(GetEntityProfileRequestMsg msg); + GetEntityProfileResponseMsg getEntityProfile(GetEntityProfileRequestMsg msg); void process(DeviceTransportType transportType, ValidateDeviceTokenRequestMsg msg, TransportServiceCallback callback); @@ -62,8 +62,6 @@ public interface TransportService { void process(ProvisionDeviceRequestMsg msg, TransportServiceCallback callback); - void onProfileUpdate(DeviceProfile deviceProfile); - boolean checkLimits(SessionInfoProto sessionInfo, Object msg, TransportServiceCallback callback); boolean checkLimits(SessionInfoProto sessionInfo, Object msg, TransportServiceCallback callback, int dataPoints, TransportRateLimitType... limits); diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/DefaultTransportRateLimitService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/DefaultTransportRateLimitService.java index 569f5446db..9663ae2e47 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/DefaultTransportRateLimitService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/DefaultTransportRateLimitService.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, @@ -33,6 +33,7 @@ import java.util.concurrent.ConcurrentMap; @Slf4j public class DefaultTransportRateLimitService implements TransportRateLimitService { + private final ConcurrentMap tenantAllowed = new ConcurrentHashMap<>(); private final ConcurrentMap perTenantLimits = new ConcurrentHashMap<>(); private final ConcurrentMap perDeviceLimits = new ConcurrentHashMap<>(); @@ -46,6 +47,9 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi @Override public TransportRateLimitType checkLimits(TenantId tenantId, DeviceId deviceId, int dataPoints, TransportRateLimitType... limits) { + if (!tenantAllowed.getOrDefault(tenantId, Boolean.TRUE)) { + return TransportRateLimitType.TENANT_ADDED_TO_DISABLED_LIST; + } TransportRateLimit[] tenantLimits = getTenantRateLimits(tenantId); TransportRateLimit[] deviceLimits = getDeviceRateLimits(tenantId, deviceId); for (TransportRateLimitType limitType : limits) { @@ -85,6 +89,11 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi perDeviceLimits.remove(deviceId); } + @Override + public void update(TenantId tenantId, boolean allowed) { + tenantAllowed.put(tenantId, allowed); + } + private void mergeLimits(TenantId tenantId, TransportRateLimit[] newRateLimits) { TransportRateLimit[] oldRateLimits = perTenantLimits.get(tenantId); if (oldRateLimits == null) { diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportRateLimitService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportRateLimitService.java index a97fbfc61d..2ab1ec38ac 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportRateLimitService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportRateLimitService.java @@ -31,4 +31,5 @@ public interface TransportRateLimitService { void remove(DeviceId deviceId); + void update(TenantId tenantId, boolean transportEnabled); } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportRateLimitType.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportRateLimitType.java index a3e6da6683..038de86e84 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportRateLimitType.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportRateLimitType.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, @@ -19,6 +19,7 @@ import lombok.Getter; public enum TransportRateLimitType { + TENANT_ADDED_TO_DISABLED_LIST("general.tenant.disabled", true, false), TENANT_MAX_MSGS("transport.tenant.msg", true, true), TENANT_TELEMETRY_MSGS("transport.tenant.telemetry", true, true), TENANT_MAX_DATA_POINTS("transport.tenant.dataPoints", true, false), diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java index 010a3cd899..326f996d14 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.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, @@ -26,6 +26,7 @@ import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.data.ApiUsageRecordKey; +import org.thingsboard.server.common.data.ApiUsageState; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceTransportType; import org.thingsboard.server.common.data.EntityType; @@ -237,7 +238,7 @@ public class DefaultTransportService implements TransportService { } @Override - public TransportProtos.GetEntityProfileResponseMsg getRoutingInfo(TransportProtos.GetEntityProfileRequestMsg msg) { + public TransportProtos.GetEntityProfileResponseMsg getEntityProfile(TransportProtos.GetEntityProfileRequestMsg msg) { TbProtoQueueMsg protoMsg = new TbProtoQueueMsg<>(UUID.randomUUID(), TransportProtos.TransportApiRequestMsg.newBuilder().setEntityProfileRequestMsg(msg).build()); try { @@ -641,8 +642,7 @@ public class DefaultTransportService implements TransportService { onProfileUpdate(deviceProfile); } } else if (EntityType.TENANT_PROFILE.equals(entityType)) { - TenantProfileUpdateResult update = tenantProfileCache.put(msg.getData()); - rateLimitService.update(update); + rateLimitService.update(tenantProfileCache.put(msg.getData())); } else if (EntityType.TENANT.equals(entityType)) { Optional profileOpt = dataDecodingEncodingService.decode(msg.getData().toByteArray()); if (profileOpt.isPresent()) { @@ -652,6 +652,12 @@ public class DefaultTransportService implements TransportService { rateLimitService.update(tenant.getId()); } } + } else if (EntityType.API_USAGE_STATE.equals(entityType)) { + Optional stateOpt = dataDecodingEncodingService.decode(msg.getData().toByteArray()); + if (stateOpt.isPresent()) { + ApiUsageState apiUsageState = stateOpt.get(); + rateLimitService.update(apiUsageState.getTenantId(), apiUsageState.isTransportEnabled()); + } } } else if (toSessionMsg.hasEntityDeleteMsg()) { TransportProtos.EntityDeleteMsg msg = toSessionMsg.getEntityDeleteMsg(); @@ -673,8 +679,7 @@ public class DefaultTransportService implements TransportService { } } - @Override - public void onProfileUpdate(DeviceProfile deviceProfile) { + private void onProfileUpdate(DeviceProfile deviceProfile) { long deviceProfileIdMSB = deviceProfile.getId().getId().getMostSignificantBits(); long deviceProfileIdLSB = deviceProfile.getId().getId().getLeastSignificantBits(); sessions.forEach((id, md) -> { diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportTenantProfileCache.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportTenantProfileCache.java index 784f7b3ae8..18c1512a27 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportTenantProfileCache.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportTenantProfileCache.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, @@ -18,21 +18,19 @@ package org.thingsboard.server.common.transport.service; import com.google.protobuf.ByteString; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Component; -import org.thingsboard.server.common.data.DeviceProfile; +import org.thingsboard.server.common.data.ApiUsageState; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.TenantProfile; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantProfileId; import org.thingsboard.server.common.transport.TransportService; import org.thingsboard.server.common.transport.TransportTenantProfileCache; +import org.thingsboard.server.common.transport.limits.TransportRateLimitService; import org.thingsboard.server.common.transport.profile.TenantProfileUpdateResult; import org.thingsboard.server.common.transport.util.DataDecodingEncodingService; import org.thingsboard.server.gen.transport.TransportProtos; -import org.thingsboard.server.queue.discovery.TenantRoutingInfo; -import org.thingsboard.server.queue.discovery.TenantRoutingInfoService; import org.thingsboard.server.queue.util.TbTransportComponent; import java.util.Collections; @@ -54,8 +52,15 @@ public class DefaultTransportTenantProfileCache implements TransportTenantProfil private final ConcurrentMap> tenantProfileIds = new ConcurrentHashMap<>(); private final DataDecodingEncodingService dataDecodingEncodingService; + private TransportRateLimitService rateLimitService; private TransportService transportService; + @Lazy + @Autowired + public void setRateLimitService(TransportRateLimitService rateLimitService) { + this.rateLimitService = rateLimitService; + } + @Lazy @Autowired public void setTransportService(TransportService transportService) { @@ -77,7 +82,8 @@ public class DefaultTransportTenantProfileCache implements TransportTenantProfil if (profileOpt.isPresent()) { TenantProfile newProfile = profileOpt.get(); log.trace("[{}] put: {}", newProfile.getId(), newProfile); - return new TenantProfileUpdateResult(newProfile, tenantProfileIds.get(newProfile.getId())); + Set affectedTenants = tenantProfileIds.get(newProfile.getId()); + return new TenantProfileUpdateResult(newProfile, affectedTenants != null ? affectedTenants : Collections.emptySet()); } else { log.warn("Failed to decode profile: {}", profileBody.toString()); return new TenantProfileUpdateResult(null, Collections.emptySet()); @@ -127,8 +133,8 @@ public class DefaultTransportTenantProfileCache implements TransportTenantProfil .setEntityIdMSB(tenantId.getId().getMostSignificantBits()) .setEntityIdLSB(tenantId.getId().getLeastSignificantBits()) .build(); - TransportProtos.GetEntityProfileResponseMsg routingInfo = transportService.getRoutingInfo(msg); - Optional profileOpt = dataDecodingEncodingService.decode(routingInfo.getData().toByteArray()); + TransportProtos.GetEntityProfileResponseMsg entityProfileMsg = transportService.getEntityProfile(msg); + Optional profileOpt = dataDecodingEncodingService.decode(entityProfileMsg.getData().toByteArray()); if (profileOpt.isPresent()) { profile = profileOpt.get(); TenantProfile existingProfile = profiles.get(profile.getId()); @@ -140,9 +146,11 @@ public class DefaultTransportTenantProfileCache implements TransportTenantProfil tenantProfileIds.computeIfAbsent(profile.getId(), id -> ConcurrentHashMap.newKeySet()).add(tenantId); tenantIds.put(tenantId, profile.getId()); } else { - log.warn("[{}] Can't decode tenant profile: {}", tenantId, routingInfo.getData()); + log.warn("[{}] Can't decode tenant profile: {}", tenantId, entityProfileMsg.getData()); throw new RuntimeException("Can't decode tenant profile!"); } + Optional apiStateOpt = dataDecodingEncodingService.decode(entityProfileMsg.getApiState().toByteArray()); + apiStateOpt.ifPresent(apiUsageState -> rateLimitService.update(tenantId, apiUsageState.isTransportEnabled())); } } finally { tenantProfileFetchLock.unlock(); diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java index a0f5b45f02..a370fbe7f3 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java @@ -445,6 +445,10 @@ public class ModelConstants { public static final String API_USAGE_STATE_TENANT_ID_COLUMN = TENANT_ID_PROPERTY; public static final String API_USAGE_STATE_ENTITY_TYPE_COLUMN = ENTITY_TYPE_COLUMN; public static final String API_USAGE_STATE_ENTITY_ID_COLUMN = ENTITY_ID_COLUMN; + public static final String API_USAGE_STATE_TRANSPORT_ENABLED_COLUMN = "transport_enabled"; + public static final String API_USAGE_STATE_DB_STORAGE_ENABLED_COLUMN = "db_storage_enabled"; + public static final String API_USAGE_STATE_RE_EXEC_ENABLED_COLUMN = "re_exec_enabled"; + public static final String API_USAGE_STATE_JS_EXEC_ENABLED_COLUMN = "js_exec_enabled"; /** * Cassandra attributes and timeseries constants. diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/ApiUsageStateEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sql/ApiUsageStateEntity.java index a06ec142d3..262252517a 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sql/ApiUsageStateEntity.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sql/ApiUsageStateEntity.java @@ -17,6 +17,8 @@ package org.thingsboard.server.dao.model.sql; import lombok.Data; import lombok.EqualsAndHashCode; +import lombok.Getter; +import lombok.Setter; import org.hibernate.annotations.TypeDef; import org.thingsboard.server.common.data.ApiUsageState; import org.thingsboard.server.common.data.id.EntityIdFactory; @@ -51,6 +53,15 @@ public class ApiUsageStateEntity extends BaseSqlEntity implements @Column(name = ModelConstants.API_USAGE_STATE_ENTITY_ID_COLUMN) private UUID entityId; + @Column(name = ModelConstants.API_USAGE_STATE_TRANSPORT_ENABLED_COLUMN) + private boolean transportEnabled = true; + @Column(name = ModelConstants.API_USAGE_STATE_DB_STORAGE_ENABLED_COLUMN) + private boolean dbStorageEnabled = true; + @Column(name = ModelConstants.API_USAGE_STATE_RE_EXEC_ENABLED_COLUMN) + private boolean reExecEnabled = true; + @Column(name = ModelConstants.API_USAGE_STATE_JS_EXEC_ENABLED_COLUMN) + private boolean jsExecEnabled = true; + public ApiUsageStateEntity() { } @@ -66,6 +77,10 @@ public class ApiUsageStateEntity extends BaseSqlEntity implements this.entityType = ur.getEntityId().getEntityType().name(); this.entityId = ur.getEntityId().getId(); } + this.transportEnabled = ur.isTransportEnabled(); + this.dbStorageEnabled = ur.isDbStorageEnabled(); + this.reExecEnabled = ur.isReExecEnabled(); + this.jsExecEnabled = ur.isJsExecEnabled(); } @Override @@ -78,6 +93,10 @@ public class ApiUsageStateEntity extends BaseSqlEntity implements if (entityId != null) { ur.setEntityId(EntityIdFactory.getByTypeAndUuid(entityType, entityId)); } + ur.setTransportEnabled(transportEnabled); + ur.setDbStorageEnabled(dbStorageEnabled); + ur.setReExecEnabled(reExecEnabled); + ur.setJsExecEnabled(jsExecEnabled); return ur; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/usagerecord/ApiApiUsageStateServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/usagerecord/ApiUsageStateServiceImpl.java similarity index 96% rename from dao/src/main/java/org/thingsboard/server/dao/usagerecord/ApiApiUsageStateServiceImpl.java rename to dao/src/main/java/org/thingsboard/server/dao/usagerecord/ApiUsageStateServiceImpl.java index e59c990346..15be73a0b8 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/usagerecord/ApiApiUsageStateServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/usagerecord/ApiUsageStateServiceImpl.java @@ -31,13 +31,13 @@ import static org.thingsboard.server.dao.service.Validator.validateId; @Service @Slf4j -public class ApiApiUsageStateServiceImpl extends AbstractEntityService implements ApiUsageStateService { +public class ApiUsageStateServiceImpl extends AbstractEntityService implements ApiUsageStateService { public static final String INCORRECT_TENANT_ID = "Incorrect tenantId "; private final ApiUsageStateDao apiUsageStateDao; private final TenantDao tenantDao; - public ApiApiUsageStateServiceImpl(TenantDao tenantDao, ApiUsageStateDao apiUsageStateDao) { + public ApiUsageStateServiceImpl(TenantDao tenantDao, ApiUsageStateDao apiUsageStateDao) { this.tenantDao = tenantDao; this.apiUsageStateDao = apiUsageStateDao; } diff --git a/dao/src/main/resources/sql/schema-entities-hsql.sql b/dao/src/main/resources/sql/schema-entities-hsql.sql index c671cd2e19..1d82e2d493 100644 --- a/dao/src/main/resources/sql/schema-entities-hsql.sql +++ b/dao/src/main/resources/sql/schema-entities-hsql.sql @@ -411,5 +411,9 @@ CREATE TABLE IF NOT EXISTS api_usage_state ( tenant_id uuid, entity_type varchar(32), entity_id uuid, + transport_enabled boolean, + db_storage_enabled boolean, + re_exec_enabled boolean, + js_exec_enabled boolean, CONSTRAINT api_usage_state_unq_key UNIQUE (tenant_id, entity_id) ); diff --git a/dao/src/main/resources/sql/schema-entities.sql b/dao/src/main/resources/sql/schema-entities.sql index d8539d9af9..c26fce6080 100644 --- a/dao/src/main/resources/sql/schema-entities.sql +++ b/dao/src/main/resources/sql/schema-entities.sql @@ -437,6 +437,10 @@ CREATE TABLE IF NOT EXISTS api_usage_state ( tenant_id uuid, entity_type varchar(32), entity_id uuid, + transport_enabled boolean, + db_storage_enabled boolean, + re_exec_enabled boolean, + js_exec_enabled boolean, CONSTRAINT api_usage_state_unq_key UNIQUE (tenant_id, entity_id) ); diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/BaseApiUsageStateServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/BaseApiUsageStateServiceTest.java index 5f56506997..a33e190a03 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/BaseApiUsageStateServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/BaseApiUsageStateServiceTest.java @@ -48,4 +48,17 @@ public abstract class BaseApiUsageStateServiceTest extends AbstractServiceTest { Assert.assertNotNull(apiUsageState); } + @Test + public void testUpdateApiUsageState(){ + ApiUsageState apiUsageState = apiUsageStateService.findTenantApiUsageState(tenantId); + Assert.assertNotNull(apiUsageState); + Assert.assertTrue(apiUsageState.isTransportEnabled()); + apiUsageState.setTransportEnabled(false); + apiUsageState = apiUsageStateService.update(apiUsageState); + Assert.assertNotNull(apiUsageState); + apiUsageState = apiUsageStateService.findTenantApiUsageState(tenantId); + Assert.assertNotNull(apiUsageState); + Assert.assertFalse(apiUsageState.isTransportEnabled()); + } + }