Browse Source

Added Edge API usage state

pull/15325/head
Andrii Landiak 5 months ago
parent
commit
13b68e3dd6
  1. 6
      application/src/main/data/upgrade/basic/schema_update.sql
  2. 65
      application/src/main/java/org/thingsboard/server/service/apiusage/BaseApiUsageState.java
  3. 37
      application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java
  4. 1
      application/src/main/java/org/thingsboard/server/service/apiusage/TbApiUsageStateService.java
  5. 21
      application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java
  6. 89
      application/src/main/java/org/thingsboard/server/service/edge/rpc/service/EdgeGrpcService.java
  7. 108
      application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeGrpcSession.java
  8. 7
      application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeGrpcSessionDelegate.java
  9. 11
      application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeSession.java
  10. 7
      application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeSessionsHolder.java
  11. 5
      application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/AbstractEdgeGrpcSessionManager.java
  12. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/EdgeGrpcSessionManager.java
  13. 3
      application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java
  14. 155
      application/src/test/java/org/thingsboard/server/edge/EdgeApiUsageDisabledTest.java
  15. 8
      application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java
  16. 31
      application/src/test/java/org/thingsboard/server/service/apiusage/ApiUsageTest.java
  17. 143
      application/src/test/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateServiceTest.java
  18. 126
      application/src/test/java/org/thingsboard/server/transport/coap/CoapTransportFeatureDisabledTest.java
  19. 125
      application/src/test/java/org/thingsboard/server/transport/mqtt/MqttTransportFeatureDisabledTest.java
  20. 4
      common/data/src/main/java/org/thingsboard/server/common/data/ApiFeature.java
  21. 33
      common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageRecordKey.java
  22. 23
      common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageState.java
  23. 2
      common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageStateValue.java
  24. 5
      common/data/src/main/java/org/thingsboard/server/common/data/edqs/fields/ApiUsageStateFields.java
  25. 58
      common/data/src/main/java/org/thingsboard/server/common/data/edqs/fields/FieldsUtil.java
  26. 5
      common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java
  27. 8
      common/edge-api/src/main/java/org/thingsboard/edge/exception/EdgeConnectionException.java
  28. 29
      common/edge-api/src/main/java/org/thingsboard/edge/exception/EdgeFeatureDisabledException.java
  29. 7
      common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java
  30. 1
      common/edge-api/src/main/proto/edge.proto
  31. 28
      common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java
  32. 14
      common/proto/src/main/proto/queue.proto
  33. 13
      common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageReportClient.java
  34. 29
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java
  35. 10
      dao/src/main/java/org/thingsboard/server/dao/edge/BaseEdgeEventService.java
  36. 4
      dao/src/main/java/org/thingsboard/server/dao/edge/EdgeDao.java
  37. 24
      dao/src/main/java/org/thingsboard/server/dao/edge/EdgeEventDao.java
  38. 1
      dao/src/main/java/org/thingsboard/server/dao/edge/PostgresEdgeEventService.java
  39. 1
      dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java
  40. 5
      dao/src/main/java/org/thingsboard/server/dao/model/sql/ApiUsageStateEntity.java
  41. 2
      dao/src/main/java/org/thingsboard/server/dao/sql/usagerecord/ApiUsageStateRepository.java
  42. 3
      dao/src/main/java/org/thingsboard/server/dao/usagerecord/ApiUsageStateServiceImpl.java
  43. 1
      dao/src/main/resources/sql/schema-entities.sql
  44. 85
      dao/src/test/java/org/thingsboard/server/dao/service/ApiUsageStateServiceTest.java
  45. 8
      dao/src/test/java/org/thingsboard/server/dao/service/DeviceServiceTest.java
  46. 6
      dao/src/test/java/org/thingsboard/server/dao/service/EdgeEventServiceTest.java
  47. 6
      edqs/src/test/java/org/thingsboard/server/edqs/repo/ApiUsageStateFilterTest.java
  48. 82
      ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.html
  49. 1
      ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.ts
  50. 2
      ui-ngx/src/app/shared/models/tenant.model.ts
  51. 6
      ui-ngx/src/assets/locale/locale.constant-en_US.json

6
application/src/main/data/upgrade/basic/schema_update.sql

@ -25,3 +25,9 @@ ALTER TABLE calculated_field ADD COLUMN IF NOT EXISTS additional_info varchar;
ALTER TABLE rule_chain ADD COLUMN IF NOT EXISTS notes varchar(1000000); ALTER TABLE rule_chain ADD COLUMN IF NOT EXISTS notes varchar(1000000);
-- RULE CHAIN NOTES MIGRATION END -- RULE CHAIN NOTES MIGRATION END
-- EDGE API USAGE STATE ADDITION START
ALTER TABLE api_usage_state ADD COLUMN IF NOT EXISTS edge varchar(32) DEFAULT 'ENABLED';
-- EDGE API USAGE STATE ADDITION END

65
application/src/main/java/org/thingsboard/server/service/apiusage/BaseApiUsageState.java

@ -33,6 +33,7 @@ import java.util.Map;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
public abstract class BaseApiUsageState { public abstract class BaseApiUsageState {
private final Map<ApiUsageRecordKey, Long> currentCycleValues = new ConcurrentHashMap<>(); private final Map<ApiUsageRecordKey, Long> currentCycleValues = new ConcurrentHashMap<>();
private final Map<ApiUsageRecordKey, Long> currentHourValues = new ConcurrentHashMap<>(); private final Map<ApiUsageRecordKey, Long> currentHourValues = new ConcurrentHashMap<>();
@ -137,55 +138,31 @@ public abstract class BaseApiUsageState {
} }
public ApiUsageStateValue getFeatureValue(ApiFeature feature) { public ApiUsageStateValue getFeatureValue(ApiFeature feature) {
switch (feature) { return switch (feature) {
case TRANSPORT: case TRANSPORT -> apiUsageState.getTransportState();
return apiUsageState.getTransportState(); case RE -> apiUsageState.getReExecState();
case RE: case DB -> apiUsageState.getDbStorageState();
return apiUsageState.getReExecState(); case JS -> apiUsageState.getJsExecState();
case DB: case TBEL -> apiUsageState.getTbelExecState();
return apiUsageState.getDbStorageState(); case EMAIL -> apiUsageState.getEmailExecState();
case JS: case SMS -> apiUsageState.getSmsExecState();
return apiUsageState.getJsExecState(); case ALARM -> apiUsageState.getAlarmExecState();
case TBEL: case EDGE -> apiUsageState.getEdgeState();
return apiUsageState.getTbelExecState(); };
case EMAIL:
return apiUsageState.getEmailExecState();
case SMS:
return apiUsageState.getSmsExecState();
case ALARM:
return apiUsageState.getAlarmExecState();
default:
return ApiUsageStateValue.ENABLED;
}
} }
public boolean setFeatureValue(ApiFeature feature, ApiUsageStateValue value) { public boolean setFeatureValue(ApiFeature feature, ApiUsageStateValue value) {
ApiUsageStateValue currentValue = getFeatureValue(feature); ApiUsageStateValue currentValue = getFeatureValue(feature);
switch (feature) { switch (feature) {
case TRANSPORT: case TRANSPORT -> apiUsageState.setTransportState(value);
apiUsageState.setTransportState(value); case RE -> apiUsageState.setReExecState(value);
break; case DB -> apiUsageState.setDbStorageState(value);
case RE: case JS -> apiUsageState.setJsExecState(value);
apiUsageState.setReExecState(value); case TBEL -> apiUsageState.setTbelExecState(value);
break; case EMAIL -> apiUsageState.setEmailExecState(value);
case DB: case SMS -> apiUsageState.setSmsExecState(value);
apiUsageState.setDbStorageState(value); case ALARM -> apiUsageState.setAlarmExecState(value);
break; case EDGE -> apiUsageState.setEdgeState(value);
case JS:
apiUsageState.setJsExecState(value);
break;
case TBEL:
apiUsageState.setTbelExecState(value);
break;
case EMAIL:
apiUsageState.setEmailExecState(value);
break;
case SMS:
apiUsageState.setSmsExecState(value);
break;
case ALARM:
apiUsageState.setAlarmExecState(value);
break;
} }
return !currentValue.equals(value); return !currentValue.equals(value);
} }

37
application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java

@ -59,7 +59,6 @@ import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.dao.tenant.TenantService; import org.thingsboard.server.dao.tenant.TenantService;
import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.dao.timeseries.TimeseriesService;
import org.thingsboard.server.dao.usagerecord.ApiUsageStateService; import org.thingsboard.server.dao.usagerecord.ApiUsageStateService;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg;
import org.thingsboard.server.gen.transport.TransportProtos.UsageStatsKVProto; import org.thingsboard.server.gen.transport.TransportProtos.UsageStatsKVProto;
import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg;
@ -149,27 +148,7 @@ public class DefaultTbApiUsageStateService extends AbstractPartitionBasedService
public void process(TbProtoQueueMsg<ToUsageStatsServiceMsg> msgPack, TbCallback callback) { public void process(TbProtoQueueMsg<ToUsageStatsServiceMsg> msgPack, TbCallback callback) {
ToUsageStatsServiceMsg serviceMsg = msgPack.getValue(); ToUsageStatsServiceMsg serviceMsg = msgPack.getValue();
String serviceId = serviceMsg.getServiceId(); String serviceId = serviceMsg.getServiceId();
serviceMsg.getMsgsList().forEach(msg -> {
List<TransportProtos.UsageStatsServiceMsg> msgs;
//For backward compatibility, remove after release
if (serviceMsg.getMsgsList().isEmpty()) {
TransportProtos.UsageStatsServiceMsg oldMsg = TransportProtos.UsageStatsServiceMsg.newBuilder()
.setTenantIdMSB(serviceMsg.getTenantIdMSB())
.setTenantIdLSB(serviceMsg.getTenantIdLSB())
.setCustomerIdMSB(serviceMsg.getCustomerIdMSB())
.setCustomerIdLSB(serviceMsg.getCustomerIdLSB())
.setEntityIdMSB(serviceMsg.getEntityIdMSB())
.setEntityIdLSB(serviceMsg.getEntityIdLSB())
.addAllValues(serviceMsg.getValuesList())
.build();
msgs = List.of(oldMsg);
} else {
msgs = serviceMsg.getMsgsList();
}
msgs.forEach(msg -> {
TenantId tenantId = TenantId.fromUUID(new UUID(msg.getTenantIdMSB(), msg.getTenantIdLSB())); TenantId tenantId = TenantId.fromUUID(new UUID(msg.getTenantIdMSB(), msg.getTenantIdLSB()));
EntityId ownerId; EntityId ownerId;
if (msg.getCustomerIdMSB() != 0 && msg.getCustomerIdLSB() != 0) { if (msg.getCustomerIdMSB() != 0 && msg.getCustomerIdLSB() != 0) {
@ -184,7 +163,9 @@ public class DefaultTbApiUsageStateService extends AbstractPartitionBasedService
} }
private void processEntityUsageStats(TenantId tenantId, EntityId ownerId, List<UsageStatsKVProto> values, String serviceId) { private void processEntityUsageStats(TenantId tenantId, EntityId ownerId, List<UsageStatsKVProto> values, String serviceId) {
if (deletedEntities.contains(ownerId)) return; if (deletedEntities.contains(ownerId)) {
return;
}
BaseApiUsageState usageState; BaseApiUsageState usageState;
List<TsKvEntry> updatedEntries; List<TsKvEntry> updatedEntries;
@ -205,14 +186,7 @@ public class DefaultTbApiUsageStateService extends AbstractPartitionBasedService
updatedEntries = new ArrayList<>(ApiUsageRecordKey.values().length); updatedEntries = new ArrayList<>(ApiUsageRecordKey.values().length);
Set<ApiFeature> apiFeatures = new HashSet<>(); Set<ApiFeature> apiFeatures = new HashSet<>();
for (UsageStatsKVProto statsItem : values) { for (UsageStatsKVProto statsItem : values) {
ApiUsageRecordKey recordKey; ApiUsageRecordKey recordKey = ProtoUtils.fromProto(statsItem.getRecordKey());
//For backward compatibility, remove after release
if (StringUtils.isNotEmpty(statsItem.getKey())) {
recordKey = ApiUsageRecordKey.valueOf(statsItem.getKey());
} else {
recordKey = ProtoUtils.fromProto(statsItem.getRecordKey());
}
StatsCalculationResult calculationResult = usageState.calculate(recordKey, statsItem.getValue(), serviceId); StatsCalculationResult calculationResult = usageState.calculate(recordKey, statsItem.getValue(), serviceId);
if (calculationResult.isValueChanged()) { if (calculationResult.isValueChanged()) {
@ -598,4 +572,5 @@ public class DefaultTbApiUsageStateService extends AbstractPartitionBasedService
private void destroy() { private void destroy() {
super.stop(); super.stop();
} }
} }

1
application/src/main/java/org/thingsboard/server/service/apiusage/TbApiUsageStateService.java

@ -39,4 +39,5 @@ public interface TbApiUsageStateService extends TbApiUsageStateClient, RuleEngin
void onCustomerDelete(CustomerId customerId); void onCustomerDelete(CustomerId customerId);
void onApiUsageStateUpdate(TenantId tenantId); void onApiUsageStateUpdate(TenantId tenantId);
} }

21
application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java

@ -15,8 +15,8 @@
*/ */
package org.thingsboard.server.service.edge.rpc; package org.thingsboard.server.service.edge.rpc;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.SettableFuture;
import lombok.RequiredArgsConstructor; import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
@ -28,6 +28,8 @@ import org.thingsboard.server.dao.edge.BaseEdgeEventService;
import org.thingsboard.server.dao.edge.stats.EdgeStatsCounterService; import org.thingsboard.server.dao.edge.stats.EdgeStatsCounterService;
import org.thingsboard.server.dao.edge.stats.EdgeStatsKey; import org.thingsboard.server.dao.edge.stats.EdgeStatsKey;
import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeEventNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeEventNotificationMsg;
import org.thingsboard.server.queue.TbQueueCallback;
import org.thingsboard.server.queue.TbQueueMsgMetadata;
import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.discovery.TopicService;
import org.thingsboard.server.queue.provider.TbQueueProducerProvider; import org.thingsboard.server.queue.provider.TbQueueProducerProvider;
@ -51,10 +53,23 @@ public class KafkaEdgeEventService extends BaseEdgeEventService {
TopicPartitionInfo tpi = topicService.getEdgeEventNotificationsTopic(edgeEvent.getTenantId(), edgeEvent.getEdgeId()); TopicPartitionInfo tpi = topicService.getEdgeEventNotificationsTopic(edgeEvent.getTenantId(), edgeEvent.getEdgeId());
ToEdgeEventNotificationMsg msg = ToEdgeEventNotificationMsg.newBuilder().setEdgeEventMsg(ProtoUtils.toProto(edgeEvent)).build(); ToEdgeEventNotificationMsg msg = ToEdgeEventNotificationMsg.newBuilder().setEdgeEventMsg(ProtoUtils.toProto(edgeEvent)).build();
producerProvider.getTbEdgeEventsMsgProducer().send(tpi, new TbProtoQueueMsg<>(UUID.randomUUID(), msg), null); SettableFuture<Void> result = SettableFuture.create();
producerProvider.getTbEdgeEventsMsgProducer().send(tpi, new TbProtoQueueMsg<>(UUID.randomUUID(), msg), new TbQueueCallback() {
@Override
public void onSuccess(TbQueueMsgMetadata metadata) {
reportEdgeEventUsage(edgeEvent);
result.set(null);
}
@Override
public void onFailure(Throwable t) {
log.warn("[{}][{}] Failed to send edge event to queue", edgeEvent.getTenantId(), edgeEvent.getEdgeId(), t);
result.setException(t);
}
});
statsCounterService.ifPresent(statsCounterService -> statsCounterService.ifPresent(statsCounterService ->
statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_ADDED, edgeEvent.getTenantId(), edgeEvent.getEdgeId(), 1)); statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_ADDED, edgeEvent.getTenantId(), edgeEvent.getEdgeId(), 1));
return Futures.immediateFuture(null); return result;
} }
} }

89
application/src/main/java/org/thingsboard/server/service/edge/rpc/service/EdgeGrpcService.java

@ -16,6 +16,8 @@
package org.thingsboard.server.service.edge.rpc.service; package org.thingsboard.server.service.edge.rpc.service;
import com.fasterxml.jackson.databind.node.ObjectNode; import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.Striped;
import io.grpc.Status;
import io.grpc.stub.StreamObserver; import io.grpc.stub.StreamObserver;
import jakarta.annotation.PostConstruct; import jakarta.annotation.PostConstruct;
import jakarta.annotation.PreDestroy; import jakarta.annotation.PreDestroy;
@ -26,15 +28,19 @@ import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationContext;
import org.springframework.context.annotation.Lazy; import org.springframework.context.annotation.Lazy;
import org.springframework.context.event.EventListener;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.edge.exception.EdgeFeatureDisabledException;
import org.thingsboard.rule.engine.api.AttributesSaveRequest; import org.thingsboard.rule.engine.api.AttributesSaveRequest;
import org.thingsboard.rule.engine.api.TimeseriesSaveRequest; import org.thingsboard.rule.engine.api.TimeseriesSaveRequest;
import org.thingsboard.server.cache.TbTransactionalCache; import org.thingsboard.server.cache.TbTransactionalCache;
import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.common.data.ApiUsageState;
import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EdgeId;
@ -43,6 +49,7 @@ import org.thingsboard.server.common.data.kv.BooleanDataEntry;
import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.kv.LongDataEntry;
import org.thingsboard.server.common.data.msg.TbMsgType; import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.server.common.data.notification.rule.trigger.EdgeConnectionTrigger; import org.thingsboard.server.common.data.notification.rule.trigger.EdgeConnectionTrigger;
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgDataType; import org.thingsboard.server.common.msg.TbMsgDataType;
import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.TbMsgMetaData;
@ -51,6 +58,8 @@ import org.thingsboard.server.common.msg.edge.EdgeHighPriorityMsg;
import org.thingsboard.server.common.msg.edge.EdgeSessionMsg; import org.thingsboard.server.common.msg.edge.EdgeSessionMsg;
import org.thingsboard.server.common.msg.edge.FromEdgeSyncResponse; import org.thingsboard.server.common.msg.edge.FromEdgeSyncResponse;
import org.thingsboard.server.common.msg.edge.ToEdgeSyncRequest; import org.thingsboard.server.common.msg.edge.ToEdgeSyncRequest;
import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg;
import org.thingsboard.server.common.stats.TbApiUsageStateClient;
import org.thingsboard.server.gen.edge.v1.EdgeRpcServiceGrpc; import org.thingsboard.server.gen.edge.v1.EdgeRpcServiceGrpc;
import org.thingsboard.server.gen.edge.v1.RequestMsg; import org.thingsboard.server.gen.edge.v1.RequestMsg;
import org.thingsboard.server.gen.edge.v1.ResponseMsg; import org.thingsboard.server.gen.edge.v1.ResponseMsg;
@ -64,13 +73,17 @@ import org.thingsboard.server.service.edge.rpc.session.EdgeSessionsHolder;
import org.thingsboard.server.service.edge.rpc.session.manager.EdgeGrpcSessionManager; import org.thingsboard.server.service.edge.rpc.session.manager.EdgeGrpcSessionManager;
import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService;
import java.util.List;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ExecutorService; import java.util.concurrent.ExecutorService;
import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.Lock;
import java.util.function.BooleanSupplier;
import java.util.function.Consumer; import java.util.function.Consumer;
import java.util.function.IntConsumer;
import static org.thingsboard.server.service.state.DefaultDeviceStateService.ACTIVITY_STATE; import static org.thingsboard.server.service.state.DefaultDeviceStateService.ACTIVITY_STATE;
import static org.thingsboard.server.service.state.DefaultDeviceStateService.LAST_CONNECT_TIME; import static org.thingsboard.server.service.state.DefaultDeviceStateService.LAST_CONNECT_TIME;
@ -99,7 +112,9 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
private final TbServiceInfoProvider serviceInfoProvider; private final TbServiceInfoProvider serviceInfoProvider;
private final TelemetrySubscriptionService tsSubService; private final TelemetrySubscriptionService tsSubService;
private final TbTransactionalCache<EdgeId, String> edgeIdServiceIdCache; private final TbTransactionalCache<EdgeId, String> edgeIdServiceIdCache;
private final TbApiUsageStateClient apiUsageStateClient;
private final Striped<Lock> tenantLocks = Striped.lock(64);
private final ConcurrentMap<UUID, Consumer<FromEdgeSyncResponse>> localSyncEdgeRequests = new ConcurrentHashMap<>(); private final ConcurrentMap<UUID, Consumer<FromEdgeSyncResponse>> localSyncEdgeRequests = new ConcurrentHashMap<>();
private ScheduledExecutorService executorService; private ScheduledExecutorService executorService;
private ScheduledExecutorService sendDownlinkExecutorService; private ScheduledExecutorService sendDownlinkExecutorService;
@ -117,6 +132,46 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
sessions.forEach(EdgeGrpcSessionManager::onEdgeDisconnect); sessions.forEach(EdgeGrpcSessionManager::onEdgeDisconnect);
} }
@EventListener
public void onComponentLifecycleMsg(ComponentLifecycleMsg msg) {
EntityType entityType = msg.getEntityId().getEntityType();
TenantId tenantId = msg.getTenantId();
if (entityType == EntityType.TENANT && msg.getEvent() == ComponentLifecycleEvent.DELETED) {
closeTenantSessions(tenantId, Status.NOT_FOUND, "Tenant deleted", null,
count -> log.warn("[{}] Tenant deleted but {} edge session(s) still linger - force-closing.", tenantId, count));
} else if (entityType == EntityType.API_USAGE_STATE && msg.getEvent() == ComponentLifecycleEvent.UPDATED) {
// Optimistic check to avoid acquiring the lock when edge is still enabled.
if (isEdgeEnabled(tenantId)) {
return;
}
closeTenantSessions(tenantId, Status.RESOURCE_EXHAUSTED, "Edge feature disabled due to API limits",
// Re-check under the same lock used by onEdgeConnect: the state may have flipped
// back to ENABLED between the optimistic check and lock acquisition.
() -> isEdgeEnabled(tenantId),
count -> log.info("[{}] Edge feature disabled due to API limits. Disconnecting {} edge sessions.", tenantId, count));
}
}
private void closeTenantSessions(TenantId tenantId, Status status, String reason,
BooleanSupplier skipUnderLock, IntConsumer onClose) {
List<EdgeGrpcSessionManager> toClose;
Lock lock = tenantLocks.get(tenantId);
lock.lock();
try {
if (skipUnderLock != null && skipUnderLock.getAsBoolean()) {
return;
}
toClose = sessions.getByTenantId(tenantId);
} finally {
lock.unlock();
}
if (toClose.isEmpty()) {
return;
}
onClose.accept(toClose.size());
toClose.forEach(s -> s.closeWithError(status, reason));
}
@Override @Override
public StreamObserver<RequestMsg> handleMsgs(StreamObserver<ResponseMsg> outputStream) { public StreamObserver<RequestMsg> handleMsgs(StreamObserver<ResponseMsg> outputStream) {
EdgeGrpcSessionManager sessionManager = applicationContext.getBean(EdgeGrpcSessionManager.class); EdgeGrpcSessionManager sessionManager = applicationContext.getBean(EdgeGrpcSessionManager.class);
@ -195,17 +250,27 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
EdgeSessionState state = edgeSession.getState(); EdgeSessionState state = edgeSession.getState();
Edge edge = state.getEdge(); Edge edge = state.getEdge();
TenantId tenantId = state.getTenantId(); TenantId tenantId = state.getTenantId();
log.info("[{}][{}] edge [{}] connected successfully.", tenantId, state.getSessionId(), edgeId); EdgeGrpcSessionManager replaced;
if (sessions.hasByEdgeId(edgeId)) { Lock lock = tenantLocks.get(tenantId);
EdgeGrpcSessionManager existingSession = sessions.getByEdgeId(edgeId); lock.lock();
if (existingSession != null) { try {
UUID sessionId = existingSession.getState().getSessionId(); if (!isEdgeEnabled(tenantId)) {
log.info("[{}][{}] Replacing existing session [{}] for edge [{}]", tenantId, state.getSessionId(), sessionId, edgeId); throw new EdgeFeatureDisabledException("Edge feature disabled due to API limits");
existingSession.destroyAndMarkAsZombieIfFailed(); }
sessions.removeBySessionId(sessionId); log.info("[{}][{}] edge [{}] connected successfully.", tenantId, state.getSessionId(), edgeId);
replaced = sessions.getByEdgeId(edgeId);
if (replaced != null) {
sessions.removeBySessionId(replaced.getState().getSessionId());
} }
sessions.put(edgeSession);
} finally {
lock.unlock();
}
if (replaced != null) {
UUID replacedSessionId = replaced.getState().getSessionId();
log.info("[{}][{}] Replacing existing session [{}] for edge [{}]", tenantId, state.getSessionId(), replacedSessionId, edgeId);
replaced.destroyAndMarkAsZombieIfFailed();
} }
sessions.put(edgeSession);
save(tenantId, edgeId, ACTIVITY_STATE, true); save(tenantId, edgeId, ACTIVITY_STATE, true);
long lastConnectTs = System.currentTimeMillis(); long lastConnectTs = System.currentTimeMillis();
save(tenantId, edgeId, LAST_CONNECT_TIME, lastConnectTs); save(tenantId, edgeId, LAST_CONNECT_TIME, lastConnectTs);
@ -389,4 +454,10 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
e.shutdown(); e.shutdown();
} }
} }
private boolean isEdgeEnabled(TenantId tenantId) {
ApiUsageState state = apiUsageStateClient.getApiUsageState(tenantId);
return state == null || state.isEdgeEnabled();
}
} }

108
application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeGrpcSession.java

@ -20,11 +20,13 @@ import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.SettableFuture; import com.google.common.util.concurrent.SettableFuture;
import io.grpc.Status;
import io.grpc.stub.StreamObserver; import io.grpc.stub.StreamObserver;
import lombok.RequiredArgsConstructor; import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.checkerframework.checker.nullness.qual.Nullable; import org.checkerframework.checker.nullness.qual.Nullable;
import org.springframework.data.util.Pair; import org.springframework.data.util.Pair;
import org.thingsboard.edge.exception.EdgeFeatureDisabledException;
import org.thingsboard.rule.engine.api.AttributesSaveRequest; import org.thingsboard.rule.engine.api.AttributesSaveRequest;
import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.DataConstants;
@ -74,6 +76,7 @@ import java.util.UUID;
import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock; import java.util.concurrent.locks.ReentrantLock;
import java.util.function.BiConsumer; import java.util.function.BiConsumer;
@ -99,6 +102,7 @@ public class EdgeGrpcSession implements EdgeSession {
private final EdgeSessionState state = new EdgeSessionState(); private final EdgeSessionState state = new EdgeSessionState();
private final Lock downlinkMsgLock = new ReentrantLock(); private final Lock downlinkMsgLock = new ReentrantLock();
private final ConcurrentLinkedQueue<EdgeEvent> highPriorityQueue = new ConcurrentLinkedQueue<>(); private final ConcurrentLinkedQueue<EdgeEvent> highPriorityQueue = new ConcurrentLinkedQueue<>();
private final AtomicBoolean closed = new AtomicBoolean(false);
private int clientMaxInboundMessageSize; private int clientMaxInboundMessageSize;
@ -273,8 +277,26 @@ public class EdgeGrpcSession implements EdgeSession {
return state.getSendDownlinkMsgsFuture(); return state.getSendDownlinkMsgsFuture();
} }
public void closeWithError(Status status, String errorMsg) {
if (!closed.compareAndSet(false, true)) {
log.debug("[{}][{}] closeWithError skipped — session already closed", getTenantId(), getSessionId());
return;
}
log.debug("[{}][{}] Closing session with error: {}", getTenantId(), getSessionId(), errorMsg);
state.setConnected(false);
try {
outputStream.onError(status.withDescription(errorMsg).asRuntimeException());
} catch (Exception e) {
log.debug("[{}][{}] Failed to close output stream with error: {}", getTenantId(), getSessionId(), e.getMessage());
}
}
@Override @Override
public void close() { public void close() {
if (!closed.compareAndSet(false, true)) {
log.debug("[{}][{}] close skipped — session already closed", getTenantId(), getSessionId());
return;
}
log.debug("[{}][{}] Closing session", getTenantId(), getSessionId()); log.debug("[{}][{}] Closing session", getTenantId(), getSessionId());
state.setConnected(false); state.setConnected(false);
try { try {
@ -421,7 +443,7 @@ public class EdgeGrpcSession implements EdgeSession {
if (state.isConnected() && pageData.hasNext()) { if (state.isConnected() && pageData.hasNext()) {
fetchAndSendEdgeEvents(fetcher, pageLink.nextPageLink(), result); fetchAndSendEdgeEvents(fetcher, pageLink.nextPageLink(), result);
} else { } else {
EdgeEvent latestEdgeEvent = pageData.getData().get(pageData.getData().size() - 1); EdgeEvent latestEdgeEvent = pageData.getData().getLast();
UUID idOffset = latestEdgeEvent.getUuidId(); UUID idOffset = latestEdgeEvent.getUuidId();
if (idOffset != null) { if (idOffset != null) {
Long newStartTs = Uuids.unixTimestamp(idOffset); Long newStartTs = Uuids.unixTimestamp(idOffset);
@ -457,7 +479,7 @@ public class EdgeGrpcSession implements EdgeSession {
stopCurrentSendDownlinkMsgsTask(true); stopCurrentSendDownlinkMsgsTask(true);
return; return;
} }
if (!state.getPendingMsgsMap().values().isEmpty()) { if (!state.getPendingMsgsMap().isEmpty()) {
Edge edge = state.getEdge(); Edge edge = state.getEdge();
List<DownlinkMsg> copy = new ArrayList<>(state.getPendingMsgsMap().values()); List<DownlinkMsg> copy = new ArrayList<>(state.getPendingMsgsMap().values());
if (attempt > 1) { if (attempt > 1) {
@ -551,45 +573,54 @@ public class EdgeGrpcSession implements EdgeSession {
private ConnectResponseMsg processConnect(ConnectRequestMsg request) { private ConnectResponseMsg processConnect(ConnectRequestMsg request) {
log.trace("[{}] processConnect [{}]", getSessionId(), request); log.trace("[{}] processConnect [{}]", getSessionId(), request);
Optional<Edge> optional = ctx.getEdgeService().findEdgeByRoutingKey(TenantId.SYS_TENANT_ID, request.getEdgeRoutingKey()); Optional<Edge> optional = ctx.getEdgeService().findEdgeByRoutingKey(TenantId.SYS_TENANT_ID, request.getEdgeRoutingKey());
if (optional.isPresent()) { if (optional.isEmpty()) {
Edge edge = optional.get(); return buildErrorResponse(ConnectResponseCode.BAD_CREDENTIALS, "Failed to find the edge! Routing key: " + request.getEdgeRoutingKey());
TenantId tenantId = edge.getTenantId(); }
state.setEdge(edge); Edge edge = optional.get();
try { TenantId tenantId = edge.getTenantId();
if (edge.getSecret().equals(request.getEdgeSecret())) { try {
sessionOpenListener.accept(edge.getId(), parentManagerRef); if (!edge.getSecret().equals(request.getEdgeSecret())) {
state.setEdgeVersion(request.getEdgeVersion()); String failureMsg = "Failed to validate the edge! Provided request secret: " + request.getEdgeSecret();
processSaveEdgeVersionAsAttribute(request.getEdgeVersion().name()); processConnectionFailure(edge, failureMsg, "Failed to validate the edge!");
return ConnectResponseMsg.newBuilder() return buildErrorResponse(ConnectResponseCode.BAD_CREDENTIALS, failureMsg);
.setResponseCode(ConnectResponseCode.ACCEPTED)
.setErrorMsg("")
.setConfiguration(EdgeMsgConstructorUtils.constructEdgeConfiguration(edge))
.setMaxInboundMessageSize(maxInboundMessageSize)
.build();
}
String error = "Failed to validate the edge!";
String failureMsg = String.format("%s Provided request secret: %s", error, request.getEdgeSecret());
ctx.getRuleProcessor().process(EdgeCommunicationFailureTrigger.builder().tenantId(tenantId).edgeId(edge.getId())
.customerId(edge.getCustomerId()).edgeName(edge.getName()).failureMsg(failureMsg).error(error).build());
return ConnectResponseMsg.newBuilder()
.setResponseCode(ConnectResponseCode.BAD_CREDENTIALS)
.setErrorMsg(failureMsg)
.setConfiguration(EdgeConfiguration.getDefaultInstance()).build();
} catch (Exception e) {
String failureMsg = "Failed to process edge connection!";
ctx.getRuleProcessor().process(EdgeCommunicationFailureTrigger.builder().tenantId(tenantId).edgeId(edge.getId())
.customerId(edge.getCustomerId()).edgeName(edge.getName()).failureMsg(failureMsg).error(e.getMessage()).build());
log.error(failureMsg, e);
return ConnectResponseMsg.newBuilder()
.setResponseCode(ConnectResponseCode.SERVER_UNAVAILABLE)
.setErrorMsg(failureMsg)
.setConfiguration(EdgeConfiguration.getDefaultInstance()).build();
} }
state.setEdge(edge);
sessionOpenListener.accept(edge.getId(), parentManagerRef);
state.setEdgeVersion(request.getEdgeVersion());
processSaveEdgeVersionAsAttribute(request.getEdgeVersion().name());
return ConnectResponseMsg.newBuilder()
.setResponseCode(ConnectResponseCode.ACCEPTED)
.setErrorMsg("")
.setConfiguration(EdgeMsgConstructorUtils.constructEdgeConfiguration(edge))
.setMaxInboundMessageSize(maxInboundMessageSize)
.build();
} catch (EdgeFeatureDisabledException e) {
log.trace("[{}][{}] {}", tenantId, edge.getId(), e.getMessage());
// Intentionally skip processConnectionFailure: edges retry aggressively when disabled,
// which would flood the rule engine with EdgeCommunicationFailure events. The
// FEATURE_DISABLED response code already conveys the cause to the operator.
return buildErrorResponse(ConnectResponseCode.FEATURE_DISABLED, e.getMessage());
} catch (Exception e) {
String failureMsg = "Failed to process edge connection!";
processConnectionFailure(edge, failureMsg, e.getMessage());
log.error(failureMsg, e);
return buildErrorResponse(ConnectResponseCode.SERVER_UNAVAILABLE, failureMsg);
} }
}
private void processConnectionFailure(Edge edge, String failureMsg, String error) {
ctx.getRuleProcessor().process(EdgeCommunicationFailureTrigger.builder()
.tenantId(edge.getTenantId()).edgeId(edge.getId())
.customerId(edge.getCustomerId()).edgeName(edge.getName())
.failureMsg(failureMsg).error(error).build());
}
private ConnectResponseMsg buildErrorResponse(ConnectResponseCode code, String errorMsg) {
return ConnectResponseMsg.newBuilder() return ConnectResponseMsg.newBuilder()
.setResponseCode(ConnectResponseCode.BAD_CREDENTIALS) .setResponseCode(code)
.setErrorMsg("Failed to find the edge! Routing key: " + request.getEdgeRoutingKey()) .setErrorMsg(errorMsg)
.setConfiguration(EdgeConfiguration.getDefaultInstance()).build(); .setConfiguration(EdgeConfiguration.getDefaultInstance())
.build();
} }
private void processSaveEdgeVersionAsAttribute(String edgeVersion) { private void processSaveEdgeVersionAsAttribute(String edgeVersion) {
@ -618,4 +649,5 @@ public class EdgeGrpcSession implements EdgeSession {
private UUID getSessionId() { private UUID getSessionId() {
return state.getSessionId(); return state.getSessionId();
} }
} }

7
application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeGrpcSessionDelegate.java

@ -15,6 +15,7 @@
*/ */
package org.thingsboard.server.service.edge.rpc.session; package org.thingsboard.server.service.edge.rpc.session;
import io.grpc.Status;
import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.service.edge.rpc.session.manager.EdgeGrpcSessionManager; import org.thingsboard.server.service.edge.rpc.session.manager.EdgeGrpcSessionManager;
@ -31,4 +32,10 @@ public abstract class EdgeGrpcSessionDelegate implements EdgeGrpcSessionManager
public void startSyncProcess(boolean fullSync) { public void startSyncProcess(boolean fullSync) {
getSession().startSyncProcess(fullSync); getSession().startSyncProcess(fullSync);
} }
@Override
public void closeWithError(Status status, String errorMsg) {
getSession().closeWithError(status, errorMsg);
}
} }

11
application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeSession.java

@ -16,6 +16,7 @@
package org.thingsboard.server.service.edge.rpc.session; package org.thingsboard.server.service.edge.rpc.session;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import io.grpc.Status;
import io.grpc.stub.StreamObserver; import io.grpc.stub.StreamObserver;
import org.springframework.data.util.Pair; import org.springframework.data.util.Pair;
import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.edge.EdgeEvent;
@ -31,13 +32,23 @@ import java.util.List;
public interface EdgeSession extends Closeable { public interface EdgeSession extends Closeable {
StreamObserver<RequestMsg> initInputStream(); StreamObserver<RequestMsg> initInputStream();
EdgeSessionState getState(); EdgeSessionState getState();
void startSyncProcess(boolean fullSync); void startSyncProcess(boolean fullSync);
void sendDownlinkMsg(ResponseMsg responseMsg); void sendDownlinkMsg(ResponseMsg responseMsg);
void addHighPriorityEvent(EdgeEvent edgeEvent); void addHighPriorityEvent(EdgeEvent edgeEvent);
void processHighPriorityEvents(); void processHighPriorityEvents();
boolean hasHighPriorityEvents(); boolean hasHighPriorityEvents();
ListenableFuture<Pair<Long, Long>> fetchAndSendEdgeEvents(EdgeEventFetcher fetcher); ListenableFuture<Pair<Long, Long>> fetchAndSendEdgeEvents(EdgeEventFetcher fetcher);
ListenableFuture<Boolean> sendDownlinkMsgsPack(List<DownlinkMsg> downlinkMsgsPack); ListenableFuture<Boolean> sendDownlinkMsgsPack(List<DownlinkMsg> downlinkMsgsPack);
void closeWithError(Status status, String errorMsg);
} }

7
application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeSessionsHolder.java

@ -20,11 +20,13 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.edge.rpc.EdgeSessionState; import org.thingsboard.server.service.edge.rpc.EdgeSessionState;
import org.thingsboard.server.service.edge.rpc.session.manager.EdgeGrpcSessionManager; import org.thingsboard.server.service.edge.rpc.session.manager.EdgeGrpcSessionManager;
import java.util.HashSet; import java.util.HashSet;
import java.util.List;
import java.util.Set; import java.util.Set;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
@ -71,6 +73,10 @@ public class EdgeSessionsHolder {
return sessionsById.remove(sessionId); return sessionsById.remove(sessionId);
} }
public List<EdgeGrpcSessionManager> getByTenantId(TenantId tenantId) {
return sessions.values().stream().filter(s -> tenantId.equals(s.getState().getTenantId())).toList();
}
public void remove(EdgeGrpcSessionManager session) { public void remove(EdgeGrpcSessionManager session) {
if (session == null) { if (session == null) {
log.warn("Can't remove session from holder because it's null"); log.warn("Can't remove session from holder because it's null");
@ -80,4 +86,5 @@ public class EdgeSessionsHolder {
removeByEdgeId(sessionState.getEdgeId()); removeByEdgeId(sessionState.getEdgeId());
removeBySessionId(sessionState.getSessionId()); removeBySessionId(sessionState.getSessionId());
} }
} }

5
application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/AbstractEdgeGrpcSessionManager.java

@ -108,8 +108,8 @@ public abstract class AbstractEdgeGrpcSessionManager extends EdgeGrpcSessionDele
CustomerId stateCustomerId = state.getEdge().getCustomerId(); CustomerId stateCustomerId = state.getEdge().getCustomerId();
state.setEdge(edge); state.setEdge(edge);
if (stateCustomerId != null && !stateCustomerId.equals(edge.getCustomerId())) { if (stateCustomerId != null && !stateCustomerId.equals(edge.getCustomerId())) {
// do not send edge configuration message on customer update // do not send an edge configuration message on a customer update
// message send by separate flow from assign_to or unassing_from customer // message send by separate flow from assign_to or unassign_from customer
return; return;
} }
EdgeUpdateMsg edgeConfig = EdgeUpdateMsg.newBuilder() EdgeUpdateMsg edgeConfig = EdgeUpdateMsg.newBuilder()
@ -146,4 +146,5 @@ public abstract class AbstractEdgeGrpcSessionManager extends EdgeGrpcSessionDele
} }
} }
} }
} }

2
application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/EdgeGrpcSessionManager.java

@ -15,6 +15,7 @@
*/ */
package org.thingsboard.server.service.edge.rpc.session.manager; package org.thingsboard.server.service.edge.rpc.session.manager;
import io.grpc.Status;
import io.grpc.stub.StreamObserver; import io.grpc.stub.StreamObserver;
import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.edge.EdgeEvent;
@ -41,6 +42,7 @@ public interface EdgeGrpcSessionManager {
void onConfigurationUpdate(Edge edge); void onConfigurationUpdate(Edge edge);
void onEdgeDisconnect(); void onEdgeDisconnect();
void onEdgeRemoval(); void onEdgeRemoval();
void closeWithError(Status status, String errorMsg);
void destroyAndMarkAsZombieIfFailed(); void destroyAndMarkAsZombieIfFailed();
boolean destroy(); boolean destroy();

3
application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java

@ -115,9 +115,6 @@ import java.util.stream.Collectors;
import static org.thingsboard.server.service.transport.BasicCredentialsValidationResult.PASSWORD_MISMATCH; import static org.thingsboard.server.service.transport.BasicCredentialsValidationResult.PASSWORD_MISMATCH;
import static org.thingsboard.server.service.transport.BasicCredentialsValidationResult.VALID; import static org.thingsboard.server.service.transport.BasicCredentialsValidationResult.VALID;
/**
* Created by ashvayka on 05.10.18.
*/
@Slf4j @Slf4j
@Service @Service
@TbCoreComponent @TbCoreComponent

155
application/src/test/java/org/thingsboard/server/edge/EdgeApiUsageDisabledTest.java

@ -0,0 +1,155 @@
/**
* Copyright © 2016-2026 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.edge;
import org.junit.After;
import org.junit.Assert;
import org.junit.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.test.context.TestPropertySource;
import org.thingsboard.edge.exception.EdgeFeatureDisabledException;
import org.thingsboard.server.common.data.ApiUsageRecordKey;
import org.thingsboard.server.common.data.ApiUsageStateValue;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.stats.TbApiUsageReportClient;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.dao.usagerecord.ApiUsageStateService;
import org.thingsboard.server.edge.imitator.EdgeImitator;
import org.thingsboard.server.service.edge.rpc.session.EdgeSessionsHolder;
import java.util.concurrent.TimeUnit;
import static org.awaitility.Awaitility.await;
@DaoSqlTest
@TestPropertySource(properties = {
"usage.stats.report.enabled=true",
"usage.stats.report.interval=2",
"usage.stats.report.urgent_interval=1"
})
public class EdgeApiUsageDisabledTest extends AbstractEdgeTest {
private static final int MAX_EDGE_EVENTS = 1;
@Autowired
private ApiUsageStateService apiUsageStateService;
@Autowired
private TbApiUsageReportClient apiUsageReportClient;
@Autowired
private EdgeSessionsHolder edgeSessionsHolder;
@After
public void restoreEdgeLimit() {
try {
loginSysAdmin();
updateDefaultTenantProfileConfig(cfg -> cfg.setMaxEdgeEvents(0));
} catch (Exception ignored) {}
}
@Test
public void testLiveSessionForceClosedWhenEdgeStateDisabled() throws Exception {
await().atMost(10, TimeUnit.SECONDS).until(() -> sessionConnected(tenantId));
loginSysAdmin();
updateDefaultTenantProfileConfig(cfg -> cfg.setMaxEdgeEvents(MAX_EDGE_EVENTS));
loginTenantAdmin();
for (int i = 0; i < MAX_EDGE_EVENTS + 5; i++) {
apiUsageReportClient.report(tenantId, null, ApiUsageRecordKey.EDGE_EVENT_COUNT);
}
await().atMost(15, TimeUnit.SECONDS).until(() -> apiUsageStateService.findTenantApiUsageState(tenantId).getEdgeState() == ApiUsageStateValue.DISABLED);
await().atMost(10, TimeUnit.SECONDS).until(() -> !sessionConnected(tenantId));
}
@Test
public void testReconnectRejectedWhenEdgeStateDisabled() throws Exception {
await().atMost(10, TimeUnit.SECONDS).until(() -> sessionConnected(tenantId));
loginSysAdmin();
updateDefaultTenantProfileConfig(cfg -> cfg.setMaxEdgeEvents(MAX_EDGE_EVENTS));
loginTenantAdmin();
for (int i = 0; i < MAX_EDGE_EVENTS + 5; i++) {
apiUsageReportClient.report(tenantId, null, ApiUsageRecordKey.EDGE_EVENT_COUNT);
}
await().atMost(15, TimeUnit.SECONDS).until(() -> !sessionConnected(tenantId));
EdgeImitator rejectedImitator = new EdgeImitator(EDGE_HOST, EDGE_PORT, edge.getRoutingKey(), edge.getSecret());
rejectedImitator.connect();
await().atMost(10, TimeUnit.SECONDS)
.untilAsserted(() -> {
Exception closeException = rejectedImitator.getCloseException();
Assert.assertNotNull("Imitator must receive a connection-rejected callback", closeException);
Assert.assertTrue("Expected EdgeFeatureDisabledException, got: " + closeException, closeException instanceof EdgeFeatureDisabledException);
});
Assert.assertFalse("Edge session must not be admitted when edge feature is disabled",
sessionConnected(tenantId));
try {
rejectedImitator.disconnect();
} catch (Exception ignored) {}
}
@Test
public void testReconnectAdmittedAfterEdgeStateReEnabled() throws Exception {
await().atMost(10, TimeUnit.SECONDS).until(() -> sessionConnected(tenantId));
loginSysAdmin();
updateDefaultTenantProfileConfig(cfg -> cfg.setMaxEdgeEvents(MAX_EDGE_EVENTS));
loginTenantAdmin();
for (int i = 0; i < MAX_EDGE_EVENTS + 5; i++) {
apiUsageReportClient.report(tenantId, null, ApiUsageRecordKey.EDGE_EVENT_COUNT);
}
await().atMost(15, TimeUnit.SECONDS)
.until(() -> apiUsageStateService.findTenantApiUsageState(tenantId).getEdgeState() == ApiUsageStateValue.DISABLED);
await().atMost(10, TimeUnit.SECONDS)
.until(() -> !sessionConnected(tenantId));
loginSysAdmin();
updateDefaultTenantProfileConfig(cfg -> cfg.setMaxEdgeEvents(0));
loginTenantAdmin();
await().atMost(15, TimeUnit.SECONDS)
.until(() -> apiUsageStateService.findTenantApiUsageState(tenantId).getEdgeState() == ApiUsageStateValue.ENABLED);
EdgeImitator reconnectImitator = new EdgeImitator(EDGE_HOST, EDGE_PORT, edge.getRoutingKey(), edge.getSecret());
reconnectImitator.connect();
try {
await().atMost(15, TimeUnit.SECONDS)
.until(() -> sessionConnected(tenantId));
Assert.assertNull("Imitator must not receive a close callback after re-enable",
reconnectImitator.getCloseException());
} finally {
try {
reconnectImitator.disconnect();
} catch (Exception ignored) {}
}
}
private boolean sessionConnected(TenantId tenantId) {
return edgeSessionsHolder.getByTenantId(tenantId).stream().anyMatch(s -> s.getState().isConnected());
}
}

8
application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java

@ -30,9 +30,9 @@ import org.thingsboard.edge.rpc.EdgeRpcClient;
import org.thingsboard.server.controller.AbstractWebTest; import org.thingsboard.server.controller.AbstractWebTest;
import org.thingsboard.server.gen.edge.v1.AdminSettingsUpdateMsg; import org.thingsboard.server.gen.edge.v1.AdminSettingsUpdateMsg;
import org.thingsboard.server.gen.edge.v1.AiModelUpdateMsg; import org.thingsboard.server.gen.edge.v1.AiModelUpdateMsg;
import org.thingsboard.server.gen.edge.v1.ApiKeyUpdateMsg;
import org.thingsboard.server.gen.edge.v1.AlarmCommentUpdateMsg; import org.thingsboard.server.gen.edge.v1.AlarmCommentUpdateMsg;
import org.thingsboard.server.gen.edge.v1.AlarmUpdateMsg; import org.thingsboard.server.gen.edge.v1.AlarmUpdateMsg;
import org.thingsboard.server.gen.edge.v1.ApiKeyUpdateMsg;
import org.thingsboard.server.gen.edge.v1.AssetProfileUpdateMsg; import org.thingsboard.server.gen.edge.v1.AssetProfileUpdateMsg;
import org.thingsboard.server.gen.edge.v1.AssetUpdateMsg; import org.thingsboard.server.gen.edge.v1.AssetUpdateMsg;
import org.thingsboard.server.gen.edge.v1.CalculatedFieldUpdateMsg; import org.thingsboard.server.gen.edge.v1.CalculatedFieldUpdateMsg;
@ -87,6 +87,7 @@ import java.util.stream.Collectors;
public class EdgeImitator { public class EdgeImitator {
private static final int MAX_DOWNLINK_FAILS = 2; private static final int MAX_DOWNLINK_FAILS = 2;
private final String routingKey; private final String routingKey;
private final String routingSecret; private final String routingSecret;
@ -106,6 +107,8 @@ public class EdgeImitator {
@Getter @Getter
private EdgeConfiguration configuration; private EdgeConfiguration configuration;
@Getter
private volatile Exception closeException;
private final ConcurrentLinkedDeque<AbstractMessage> downlinkMsgs; private final ConcurrentLinkedDeque<AbstractMessage> downlinkMsgs;
//Returns collection copy as Unmodifiable list //Returns collection copy as Unmodifiable list
@ -195,6 +198,9 @@ public class EdgeImitator {
private void onClose(Exception e) { private void onClose(Exception e) {
log.info("onClose: {}", e.getMessage()); log.info("onClose: {}", e.getMessage());
if (this.closeException == null) {
this.closeException = e;
}
} }
private ListenableFuture<List<Void>> processDownlinkMsg(DownlinkMsg downlinkMsg) { private ListenableFuture<List<Void>> processDownlinkMsg(DownlinkMsg downlinkMsg) {

31
application/src/test/java/org/thingsboard/server/service/apiusage/ApiUsageTest.java

@ -54,12 +54,11 @@ import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.
}) })
public class ApiUsageTest extends AbstractControllerTest { public class ApiUsageTest extends AbstractControllerTest {
private Tenant savedTenant;
private User tenantAdmin;
private static final int MAX_DP_ENABLE_VALUE = 12; private static final int MAX_DP_ENABLE_VALUE = 12;
private static final int MAX_SMS_ENABLE_VALUE = 10; private static final int MAX_SMS_ENABLE_VALUE = 10;
private static final int MAX_EDGE_ENABLE_VALUE = 10;
private static final double WARN_THRESHOLD_VALUE = 0.5; private static final double WARN_THRESHOLD_VALUE = 0.5;
@Autowired @Autowired
private ApiUsageStateService apiUsageStateService; private ApiUsageStateService apiUsageStateService;
@Autowired @Autowired
@ -76,16 +75,16 @@ public class ApiUsageTest extends AbstractControllerTest {
Tenant tenant = new Tenant(); Tenant tenant = new Tenant();
tenant.setTitle("My tenant"); tenant.setTitle("My tenant");
tenant.setTenantProfileId(savedTenantProfile.getId()); tenant.setTenantProfileId(savedTenantProfile.getId());
savedTenant = saveTenant(tenant); Tenant savedTenant = saveTenant(tenant);
tenantId = savedTenant.getId(); tenantId = savedTenant.getId();
assertNotNull(savedTenant); assertNotNull(savedTenant);
tenantAdmin = new User(); User tenantAdmin = new User();
tenantAdmin.setAuthority(Authority.TENANT_ADMIN); tenantAdmin.setAuthority(Authority.TENANT_ADMIN);
tenantAdmin.setTenantId(savedTenant.getId()); tenantAdmin.setTenantId(savedTenant.getId());
tenantAdmin.setEmail("tenant2@thingsboard.org"); tenantAdmin.setEmail("tenant2@thingsboard.org");
tenantAdmin = createUserAndLogin(tenantAdmin, "testPassword1"); createUserAndLogin(tenantAdmin, "testPassword1");
} }
@Test @Test
@ -137,6 +136,25 @@ public class ApiUsageTest extends AbstractControllerTest {
assertEquals(ApiUsageStateValue.DISABLED, getUsageState().getSmsExecState())); assertEquals(ApiUsageStateValue.DISABLED, getUsageState().getSmsExecState()));
} }
@Test
public void testEdgeApiUsage() {
long edgeWarnThreshold = (long) (MAX_EDGE_ENABLE_VALUE * WARN_THRESHOLD_VALUE);
for (int i = 0; i < edgeWarnThreshold; i++) {
apiUsageReportClient.report(tenantId, null, ApiUsageRecordKey.EDGE_EVENT_COUNT);
}
await().atMost(TIMEOUT, TimeUnit.SECONDS).untilAsserted(() -> assertEquals(ApiUsageStateValue.WARNING, getUsageState().getEdgeState()));
long edgeDisableCount = MAX_EDGE_ENABLE_VALUE - edgeWarnThreshold;
for (int i = 0; i < edgeDisableCount; i++) {
apiUsageReportClient.report(tenantId, null, ApiUsageRecordKey.EDGE_EVENT_COUNT);
}
await().atMost(TIMEOUT, TimeUnit.SECONDS).untilAsserted(() -> assertEquals(ApiUsageStateValue.DISABLED, getUsageState().getEdgeState()));
}
private ApiUsageState getUsageState() { private ApiUsageState getUsageState() {
return apiUsageStateService.findTenantApiUsageState(tenantId); return apiUsageStateService.findTenantApiUsageState(tenantId);
} }
@ -151,6 +169,7 @@ public class ApiUsageTest extends AbstractControllerTest {
.maxDPStorageDays(MAX_DP_ENABLE_VALUE) .maxDPStorageDays(MAX_DP_ENABLE_VALUE)
.maxSms(MAX_SMS_ENABLE_VALUE) .maxSms(MAX_SMS_ENABLE_VALUE)
.smsEnabled(true) .smsEnabled(true)
.maxEdgeEvents(MAX_EDGE_ENABLE_VALUE)
.warnThreshold(WARN_THRESHOLD_VALUE) .warnThreshold(WARN_THRESHOLD_VALUE)
.build(); .build();

143
application/src/test/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateServiceTest.java

@ -31,6 +31,7 @@ import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
import org.thingsboard.server.common.data.tenant.profile.TenantProfileData; import org.thingsboard.server.common.data.tenant.profile.TenantProfileData;
import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.common.util.ProtoUtils;
import org.thingsboard.server.controller.AbstractControllerTest; import org.thingsboard.server.controller.AbstractControllerTest;
import org.thingsboard.server.dao.service.DaoSqlTest; import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.dao.usagerecord.ApiUsageStateService; import org.thingsboard.server.dao.usagerecord.ApiUsageStateService;
@ -64,7 +65,6 @@ public class DefaultTbApiUsageStateServiceTest extends AbstractControllerTest {
private ApiUsageStateService apiUsageStateService; private ApiUsageStateService apiUsageStateService;
private TenantId tenantId; private TenantId tenantId;
private Tenant savedTenant;
private TenantProfile savedTenantProfile; private TenantProfile savedTenantProfile;
private static final int MAX_ENABLE_VALUE = 5000; private static final int MAX_ENABLE_VALUE = 5000;
@ -83,48 +83,20 @@ public class DefaultTbApiUsageStateServiceTest extends AbstractControllerTest {
Tenant tenant = new Tenant(); Tenant tenant = new Tenant();
tenant.setTitle("My tenant"); tenant.setTitle("My tenant");
tenant.setTenantProfileId(savedTenantProfile.getId()); tenant.setTenantProfileId(savedTenantProfile.getId());
savedTenant = saveTenant(tenant); Tenant savedTenant = saveTenant(tenant);
tenantId = savedTenant.getId(); tenantId = savedTenant.getId();
Assert.assertNotNull(savedTenant); Assert.assertNotNull(savedTenant);
} }
@Test @Test
public void testProcess_transitionFromWarningToDisabled() { public void testProcess_transitionFromWarningToDisabled() {
TransportProtos.ToUsageStatsServiceMsg.Builder warningMsgBuilder = TransportProtos.ToUsageStatsServiceMsg.newBuilder() sendUsageStats(tenantId, ApiUsageRecordKey.STORAGE_DP_COUNT, VALUE_WARNING);
.setTenantIdMSB(tenantId.getId().getMostSignificantBits()) await().atMost(5, TimeUnit.SECONDS).until(() ->
.setTenantIdLSB(tenantId.getId().getLeastSignificantBits()) apiUsageStateService.findTenantApiUsageState(tenantId).getDbStorageState() == ApiUsageStateValue.WARNING);
.setCustomerIdMSB(0)
.setCustomerIdLSB(0)
.setServiceId("testService");
warningMsgBuilder.addValues(TransportProtos.UsageStatsKVProto.newBuilder()
.setKey(ApiUsageRecordKey.STORAGE_DP_COUNT.name())
.setValue(VALUE_WARNING)
.build());
TransportProtos.ToUsageStatsServiceMsg warningStatsMsg = warningMsgBuilder.build();
TbProtoQueueMsg<TransportProtos.ToUsageStatsServiceMsg> warningMsg = new TbProtoQueueMsg<>(UUID.randomUUID(), warningStatsMsg);
service.process(warningMsg, TbCallback.EMPTY);
assertEquals(ApiUsageStateValue.WARNING, apiUsageStateService.findTenantApiUsageState(tenantId).getDbStorageState());
TransportProtos.ToUsageStatsServiceMsg.Builder disableMsgBuilder = TransportProtos.ToUsageStatsServiceMsg.newBuilder()
.setTenantIdMSB(tenantId.getId().getMostSignificantBits())
.setTenantIdLSB(tenantId.getId().getLeastSignificantBits())
.setCustomerIdMSB(0)
.setCustomerIdLSB(0)
.setServiceId("testService");
disableMsgBuilder.addValues(TransportProtos.UsageStatsKVProto.newBuilder()
.setKey(ApiUsageRecordKey.STORAGE_DP_COUNT.name())
.setValue(VALUE_DISABLE)
.build());
TransportProtos.ToUsageStatsServiceMsg disableStatsMsg = disableMsgBuilder.build();
TbProtoQueueMsg<TransportProtos.ToUsageStatsServiceMsg> disableMsg = new TbProtoQueueMsg<>(UUID.randomUUID(), disableStatsMsg);
service.process(disableMsg, TbCallback.EMPTY); sendUsageStats(tenantId, ApiUsageRecordKey.STORAGE_DP_COUNT, VALUE_DISABLE);
assertEquals(ApiUsageStateValue.DISABLED, apiUsageStateService.findTenantApiUsageState(tenantId).getDbStorageState()); await().atMost(5, TimeUnit.SECONDS).until(() ->
apiUsageStateService.findTenantApiUsageState(tenantId).getDbStorageState() == ApiUsageStateValue.DISABLED);
} }
@Test @Test
@ -138,6 +110,7 @@ public class DefaultTbApiUsageStateServiceTest extends AbstractControllerTest {
apiUsageState.setTransportState(ApiUsageStateValue.ENABLED); apiUsageState.setTransportState(ApiUsageStateValue.ENABLED);
apiUsageState.setEmailExecState(ApiUsageStateValue.ENABLED); apiUsageState.setEmailExecState(ApiUsageStateValue.ENABLED);
apiUsageState.setJsExecState(ApiUsageStateValue.ENABLED); apiUsageState.setJsExecState(ApiUsageStateValue.ENABLED);
apiUsageState.setEdgeState(ApiUsageStateValue.ENABLED);
apiUsageState.setTenantId(tenantId); apiUsageState.setTenantId(tenantId);
apiUsageState.setEntityId(tenantId); apiUsageState.setEntityId(tenantId);
@ -194,7 +167,7 @@ public class DefaultTbApiUsageStateServiceTest extends AbstractControllerTest {
await().atMost(5, TimeUnit.SECONDS).until(() -> { await().atMost(5, TimeUnit.SECONDS).until(() -> {
Optional<TsKvEntry> smsApiState = tsService.findLatest(finalTenantId, finalApiUsageStateId, SMS_EXEC_COUNT.getApiLimitKey()).get(); Optional<TsKvEntry> smsApiState = tsService.findLatest(finalTenantId, finalApiUsageStateId, SMS_EXEC_COUNT.getApiLimitKey()).get();
return smsApiState.isPresent() && smsApiState.get().getLongValue().get().equals(0L); return smsApiState.isPresent() && smsApiState.get().getLongValue().isPresent() && smsApiState.get().getLongValue().get().equals(0L);
}); });
// enable SMS and check that the ApiUsageState is updated accordingly // enable SMS and check that the ApiUsageState is updated accordingly
@ -216,10 +189,10 @@ public class DefaultTbApiUsageStateServiceTest extends AbstractControllerTest {
await().atMost(5, TimeUnit.SECONDS).until(() -> { await().atMost(5, TimeUnit.SECONDS).until(() -> {
Optional<TsKvEntry> smsApiState = tsService.findLatest(finalTenantId, finalApiUsageStateId, SMS_EXEC_COUNT.getApiLimitKey()).get(); Optional<TsKvEntry> smsApiState = tsService.findLatest(finalTenantId, finalApiUsageStateId, SMS_EXEC_COUNT.getApiLimitKey()).get();
return smsApiState.isPresent() && smsApiState.get().getLongValue().get().equals(10L); return smsApiState.isPresent() && smsApiState.get().getLongValue().isPresent() && smsApiState.get().getLongValue().get().equals(10L);
}); });
//disable SMS and check that the ApiUsageState is updated accordingly // disable SMS and check that the ApiUsageState is updated accordingly
config = DefaultTenantProfileConfiguration.builder() config = DefaultTenantProfileConfiguration.builder()
.smsEnabled(false) .smsEnabled(false)
.build(); .build();
@ -237,10 +210,98 @@ public class DefaultTbApiUsageStateServiceTest extends AbstractControllerTest {
await().atMost(5, TimeUnit.SECONDS).until(() -> { await().atMost(5, TimeUnit.SECONDS).until(() -> {
Optional<TsKvEntry> smsApiState = tsService.findLatest(finalTenantId, finalApiUsageStateId, SMS_EXEC_COUNT.getApiLimitKey()).get(); Optional<TsKvEntry> smsApiState = tsService.findLatest(finalTenantId, finalApiUsageStateId, SMS_EXEC_COUNT.getApiLimitKey()).get();
return smsApiState.isPresent() && smsApiState.get().getLongValue().get().equals(0L); return smsApiState.isPresent() && smsApiState.get().getLongValue().isPresent() && smsApiState.get().getLongValue().get().equals(0L);
}); });
} }
@Test
public void testEdgeStateTransitions() throws Exception {
TenantProfile edgeProfile = createProfileWithThreshold("Edge Test Profile", DefaultTenantProfileConfiguration.builder()
.maxEdgeEvents(100)
.warnThreshold(0.8)
.build());
Tenant edgeTenant = createTenantWithProfile("Edge tenant", edgeProfile);
TenantId edgeTenantId = edgeTenant.getId();
sendUsageStats(edgeTenantId, ApiUsageRecordKey.EDGE_EVENT_COUNT, 85);
await().atMost(5, TimeUnit.SECONDS).until(() ->
apiUsageStateService.findTenantApiUsageState(edgeTenantId).getEdgeState() == ApiUsageStateValue.WARNING);
sendUsageStats(edgeTenantId, ApiUsageRecordKey.EDGE_EVENT_COUNT, 20);
await().atMost(5, TimeUnit.SECONDS).until(() ->
apiUsageStateService.findTenantApiUsageState(edgeTenantId).getEdgeState() == ApiUsageStateValue.DISABLED);
}
@Test
public void testEdgeDisableDoesNotAffectOtherFeatures() throws Exception {
TenantProfile profile = createProfileWithThreshold("Edge Isolation Profile", DefaultTenantProfileConfiguration.builder()
.maxEdgeEvents(50)
.maxDPStorageDays(MAX_ENABLE_VALUE)
.warnThreshold(0.8)
.build());
Tenant tenant = createTenantWithProfile("Edge isolation tenant", profile);
TenantId tid = tenant.getId();
sendUsageStats(tid, ApiUsageRecordKey.EDGE_EVENT_COUNT, 100);
await().atMost(5, TimeUnit.SECONDS).until(() ->
apiUsageStateService.findTenantApiUsageState(tid).getEdgeState() == ApiUsageStateValue.DISABLED);
ApiUsageState state = apiUsageStateService.findTenantApiUsageState(tid);
assertEquals(ApiUsageStateValue.ENABLED, state.getTransportState());
assertEquals(ApiUsageStateValue.ENABLED, state.getReExecState());
assertEquals(ApiUsageStateValue.ENABLED, state.getDbStorageState());
}
@Test
public void testZeroThresholdMeansUnlimited() throws Exception {
TenantProfile profile = createProfileWithThreshold("Unlimited Profile", DefaultTenantProfileConfiguration.builder()
.maxEdgeEvents(0)
.warnThreshold(0.8)
.build());
Tenant tenant = createTenantWithProfile("Unlimited tenant", profile);
TenantId tid = tenant.getId();
sendUsageStats(tid, ApiUsageRecordKey.EDGE_EVENT_COUNT, 1_000_000);
await().atMost(5, TimeUnit.SECONDS).until(() ->
apiUsageStateService.findTenantApiUsageState(tid).getEdgeState() == ApiUsageStateValue.ENABLED);
}
private void sendUsageStats(TenantId tenantId, ApiUsageRecordKey recordKey, long value) {
TransportProtos.UsageStatsServiceMsg statsMsg = TransportProtos.UsageStatsServiceMsg.newBuilder()
.setTenantIdMSB(tenantId.getId().getMostSignificantBits())
.setTenantIdLSB(tenantId.getId().getLeastSignificantBits())
.setCustomerIdMSB(0)
.setCustomerIdLSB(0)
.addValues(TransportProtos.UsageStatsKVProto.newBuilder()
.setRecordKey(ProtoUtils.toProto(recordKey))
.setValue(value)
.build())
.build();
TransportProtos.ToUsageStatsServiceMsg msg = TransportProtos.ToUsageStatsServiceMsg.newBuilder()
.setServiceId("testService")
.addMsgs(statsMsg)
.build();
service.process(new TbProtoQueueMsg<>(UUID.randomUUID(), msg), TbCallback.EMPTY);
}
private TenantProfile createProfileWithThreshold(String name, DefaultTenantProfileConfiguration config) {
TenantProfile profile = new TenantProfile();
profile.setName(name);
TenantProfileData profileData = new TenantProfileData();
profileData.setConfiguration(config);
profile.setProfileData(profileData);
return doPost("/api/tenantProfile", profile, TenantProfile.class);
}
private Tenant createTenantWithProfile(String title, TenantProfile profile) throws Exception {
Tenant tenant = new Tenant();
tenant.setTitle(title);
tenant.setTenantProfileId(profile.getId());
return saveTenant(tenant);
}
private TenantProfile createTenantProfile() { private TenantProfile createTenantProfile() {
TenantProfile tenantProfile = new TenantProfile(); TenantProfile tenantProfile = new TenantProfile();
tenantProfile.setName("Tenant Profile"); tenantProfile.setName("Tenant Profile");
@ -257,4 +318,4 @@ public class DefaultTbApiUsageStateServiceTest extends AbstractControllerTest {
return tenantProfile; return tenantProfile;
} }
} }

126
application/src/test/java/org/thingsboard/server/transport/coap/CoapTransportFeatureDisabledTest.java

@ -0,0 +1,126 @@
/**
* Copyright © 2016-2026 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.transport.coap;
import lombok.extern.slf4j.Slf4j;
import org.awaitility.Awaitility;
import org.eclipse.californium.core.CoapObserveRelation;
import org.eclipse.californium.core.coap.CoAP;
import org.junit.After;
import org.junit.Assert;
import org.junit.Before;
import org.junit.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.test.context.TestPropertySource;
import org.thingsboard.server.common.data.ApiUsageRecordKey;
import org.thingsboard.server.common.data.ApiUsageStateValue;
import org.thingsboard.server.common.data.CoapDeviceType;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.TransportPayloadType;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
import org.thingsboard.server.common.msg.session.FeatureType;
import org.thingsboard.server.common.stats.TbApiUsageReportClient;
import org.thingsboard.server.common.transport.service.DefaultTransportService;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.dao.usagerecord.ApiUsageStateService;
import java.util.concurrent.TimeUnit;
@DaoSqlTest
@TestPropertySource(properties = {
"usage.stats.report.enabled=true",
"usage.stats.report.interval=2",
"usage.stats.report.urgent_interval=1",
})
@Slf4j
public class CoapTransportFeatureDisabledTest extends AbstractCoapIntegrationTest {
private static final int MAX_TRANSPORT_MESSAGES = 10;
private static final double WARN_THRESHOLD = 0.5;
@Autowired
private ApiUsageStateService apiUsageStateService;
@Autowired
private TbApiUsageReportClient apiUsageReportClient;
@Autowired
private DefaultTransportService defaultTransportService;
@Before
public void beforeTest() throws Exception {
loginSysAdmin();
TenantProfile tenantProfile = doGet("/api/tenantProfile/" + tenantProfileId, TenantProfile.class);
DefaultTenantProfileConfiguration config =
(DefaultTenantProfileConfiguration) tenantProfile.getProfileData().getConfiguration();
config.setMaxTransportMessages(MAX_TRANSPORT_MESSAGES);
config.setWarnThreshold(WARN_THRESHOLD);
doPost("/api/tenantProfile", tenantProfile);
CoapTestConfigProperties configProperties = CoapTestConfigProperties.builder()
.deviceName("Coap transport disable test device")
.coapDeviceType(CoapDeviceType.DEFAULT)
.transportPayloadType(TransportPayloadType.JSON)
.build();
processBeforeTest(configProperties);
}
@After
public void afterTest() throws Exception {
try {
processAfterTest();
} finally {
try {
loginSysAdmin();
TenantProfile tenantProfile = doGet("/api/tenantProfile/" + tenantProfileId, TenantProfile.class);
DefaultTenantProfileConfiguration config =
(DefaultTenantProfileConfiguration) tenantProfile.getProfileData().getConfiguration();
config.setMaxTransportMessages(0);
doPost("/api/tenantProfile", tenantProfile);
} catch (Exception ignored) {}
}
}
@Test
public void testCoapObserveSessionClosedWhenTransportDisabled() {
client = new CoapTestClient(accessToken, FeatureType.ATTRIBUTES);
CoapTestCallback callback = new CoapTestCallback();
CoapObserveRelation observeRelation = client.getObserveRelation(callback);
Awaitility.await("await initial observe response")
.atMost(10, TimeUnit.SECONDS)
.until(() -> CoAP.ResponseCode.CONTENT.equals(callback.getResponseCode())
&& callback.getObserve() != null);
Assert.assertFalse("CoAP transport must hold at least one registered session after observe",
defaultTransportService.sessions.isEmpty());
for (int i = 0; i < MAX_TRANSPORT_MESSAGES + 5; i++) {
apiUsageReportClient.report(tenantId, null, ApiUsageRecordKey.TRANSPORT_MSG_COUNT);
}
Awaitility.await("transport state flips to DISABLED")
.atMost(15, TimeUnit.SECONDS)
.until(() -> apiUsageStateService.findTenantApiUsageState(tenantId).getTransportState() == ApiUsageStateValue.DISABLED);
Awaitility.await("CoAP session is removed from DefaultTransportService")
.atMost(10, TimeUnit.SECONDS)
.until(() -> defaultTransportService.sessions.isEmpty());
observeRelation.proactiveCancel();
}
}

125
application/src/test/java/org/thingsboard/server/transport/mqtt/MqttTransportFeatureDisabledTest.java

@ -0,0 +1,125 @@
/**
* Copyright © 2016-2026 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.transport.mqtt;
import lombok.extern.slf4j.Slf4j;
import org.awaitility.Awaitility;
import org.junit.After;
import org.junit.Assert;
import org.junit.Before;
import org.junit.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.test.context.DynamicPropertyRegistry;
import org.springframework.test.context.DynamicPropertySource;
import org.springframework.test.context.TestPropertySource;
import org.thingsboard.server.common.data.ApiUsageRecordKey;
import org.thingsboard.server.common.data.ApiUsageStateValue;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.security.DeviceCredentials;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
import org.thingsboard.server.common.stats.TbApiUsageReportClient;
import org.thingsboard.server.controller.AbstractControllerTest;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.dao.usagerecord.ApiUsageStateService;
import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient;
import java.util.concurrent.TimeUnit;
import static org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest.MQTT_PORT;
@DaoSqlTest
@TestPropertySource(properties = {
"service.integrations.supported=ALL",
"transport.mqtt.enabled=true",
"usage.stats.report.enabled=true",
"usage.stats.report.interval=2",
"usage.stats.report.urgent_interval=1",
})
@Slf4j
public class MqttTransportFeatureDisabledTest extends AbstractControllerTest {
private static final int MAX_TRANSPORT_MESSAGES = 10;
private static final double WARN_THRESHOLD = 0.5;
@DynamicPropertySource
static void props(DynamicPropertyRegistry registry) {
log.warn("transport.mqtt.bind_port = {}", MQTT_PORT);
registry.add("transport.mqtt.bind_port", () -> MQTT_PORT);
}
@Autowired
private ApiUsageStateService apiUsageStateService;
@Autowired
private TbApiUsageReportClient apiUsageReportClient;
private String deviceAccessToken;
@Before
public void beforeTest() throws Exception {
loginSysAdmin();
TenantProfile tenantProfile = doGet("/api/tenantProfile/" + tenantProfileId, TenantProfile.class);
DefaultTenantProfileConfiguration config =
(DefaultTenantProfileConfiguration) tenantProfile.getProfileData().getConfiguration();
config.setMaxTransportMessages(MAX_TRANSPORT_MESSAGES);
config.setWarnThreshold(WARN_THRESHOLD);
doPost("/api/tenantProfile", tenantProfile);
loginTenantAdmin();
Device device = new Device();
device.setName("Transport disable test device");
device.setType("default");
device = doPost("/api/device", device, Device.class);
DeviceCredentials credentials = doGet("/api/device/" + device.getId().getId() + "/credentials", DeviceCredentials.class);
deviceAccessToken = credentials.getCredentialsId();
}
@After
public void afterTest() throws Exception {
try {
loginSysAdmin();
TenantProfile tenantProfile = doGet("/api/tenantProfile/" + tenantProfileId, TenantProfile.class);
DefaultTenantProfileConfiguration config = (DefaultTenantProfileConfiguration) tenantProfile.getProfileData().getConfiguration();
config.setMaxTransportMessages(0);
doPost("/api/tenantProfile", tenantProfile);
} catch (Exception ignored) {}
}
@Test
public void testLiveMqttSessionClosedWhenTransportDisabled() throws Exception {
MqttTestClient client = new MqttTestClient();
client.connectAndWait(deviceAccessToken);
Assert.assertTrue("MQTT client must be connected before flipping transport state", client.isConnected());
for (int i = 0; i < MAX_TRANSPORT_MESSAGES + 5; i++) {
apiUsageReportClient.report(tenantId, null, ApiUsageRecordKey.TRANSPORT_MSG_COUNT);
}
Awaitility.await("transport state flips to DISABLED")
.atMost(15, TimeUnit.SECONDS)
.until(() -> apiUsageStateService.findTenantApiUsageState(tenantId).getTransportState() == ApiUsageStateValue.DISABLED);
Awaitility.await("MQTT client receives server-side disconnect")
.atMost(10, TimeUnit.SECONDS)
.until(() -> !client.isConnected());
try {
client.disconnectForcibly();
} catch (Exception ignored) {}
}
}

4
common/data/src/main/java/org/thingsboard/server/common/data/ApiFeature.java

@ -18,6 +18,7 @@ package org.thingsboard.server.common.data;
import lombok.Getter; import lombok.Getter;
public enum ApiFeature { public enum ApiFeature {
TRANSPORT("transportApiState", "Device API"), TRANSPORT("transportApiState", "Device API"),
DB("dbApiState", "Telemetry persistence"), DB("dbApiState", "Telemetry persistence"),
RE("ruleEngineApiState", "Rule Engine execution"), RE("ruleEngineApiState", "Rule Engine execution"),
@ -25,7 +26,8 @@ public enum ApiFeature {
TBEL("tbelExecutionApiState", "Tbel functions execution"), TBEL("tbelExecutionApiState", "Tbel functions execution"),
EMAIL("emailApiState", "Email messages"), EMAIL("emailApiState", "Email messages"),
SMS("smsApiState", "SMS messages"), SMS("smsApiState", "SMS messages"),
ALARM("alarmApiState", "Alarms"); ALARM("alarmApiState", "Alarms"),
EDGE("edgeApiState", "Edge");
@Getter @Getter
private final String apiStateKey; private final String apiStateKey;

33
common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageRecordKey.java

@ -28,6 +28,7 @@ public enum ApiUsageRecordKey {
EMAIL_EXEC_COUNT(ApiFeature.EMAIL, "emailCount", "emailLimit", "email message", true, true), EMAIL_EXEC_COUNT(ApiFeature.EMAIL, "emailCount", "emailLimit", "email message", true, true),
SMS_EXEC_COUNT(ApiFeature.SMS, "smsCount", "smsLimit", "SMS message", true, true), SMS_EXEC_COUNT(ApiFeature.SMS, "smsCount", "smsLimit", "SMS message", true, true),
CREATED_ALARMS_COUNT(ApiFeature.ALARM, "createdAlarmsCount", "createdAlarmsLimit", "alarm"), CREATED_ALARMS_COUNT(ApiFeature.ALARM, "createdAlarmsCount", "createdAlarmsLimit", "alarm"),
EDGE_EVENT_COUNT(ApiFeature.EDGE, "edgeEventCount", "edgeEventLimit", "edge event"),
ACTIVE_DEVICES("activeDevicesCount"), ACTIVE_DEVICES("activeDevicesCount"),
INACTIVE_DEVICES("inactiveDevicesCount"); INACTIVE_DEVICES("inactiveDevicesCount");
@ -39,6 +40,7 @@ public enum ApiUsageRecordKey {
private static final ApiUsageRecordKey[] EMAIL_RECORD_KEYS = {EMAIL_EXEC_COUNT}; private static final ApiUsageRecordKey[] EMAIL_RECORD_KEYS = {EMAIL_EXEC_COUNT};
private static final ApiUsageRecordKey[] SMS_RECORD_KEYS = {SMS_EXEC_COUNT}; private static final ApiUsageRecordKey[] SMS_RECORD_KEYS = {SMS_EXEC_COUNT};
private static final ApiUsageRecordKey[] ALARM_RECORD_KEYS = {CREATED_ALARMS_COUNT}; private static final ApiUsageRecordKey[] ALARM_RECORD_KEYS = {CREATED_ALARMS_COUNT};
private static final ApiUsageRecordKey[] EDGE_RECORD_KEYS = {EDGE_EVENT_COUNT};
@Getter @Getter
private final ApiFeature apiFeature; private final ApiFeature apiFeature;
@ -71,26 +73,17 @@ public enum ApiUsageRecordKey {
} }
public static ApiUsageRecordKey[] getKeys(ApiFeature feature) { public static ApiUsageRecordKey[] getKeys(ApiFeature feature) {
switch (feature) { return switch (feature) {
case TRANSPORT: case TRANSPORT -> TRANSPORT_RECORD_KEYS;
return TRANSPORT_RECORD_KEYS; case DB -> DB_RECORD_KEYS;
case DB: case RE -> RE_RECORD_KEYS;
return DB_RECORD_KEYS; case JS -> JS_RECORD_KEYS;
case RE: case TBEL -> TBEL_RECORD_KEYS;
return RE_RECORD_KEYS; case EMAIL -> EMAIL_RECORD_KEYS;
case JS: case SMS -> SMS_RECORD_KEYS;
return JS_RECORD_KEYS; case ALARM -> ALARM_RECORD_KEYS;
case TBEL: case EDGE -> EDGE_RECORD_KEYS;
return TBEL_RECORD_KEYS; };
case EMAIL:
return EMAIL_RECORD_KEYS;
case SMS:
return SMS_RECORD_KEYS;
case ALARM:
return ALARM_RECORD_KEYS;
default:
return new ApiUsageRecordKey[]{};
}
} }
} }

23
common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageState.java

@ -23,12 +23,15 @@ import org.thingsboard.server.common.data.id.ApiUsageStateId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import java.io.Serial;
@ToString @ToString
@EqualsAndHashCode(callSuper = true) @EqualsAndHashCode(callSuper = true)
@Getter @Getter
@Setter @Setter
public class ApiUsageState extends BaseData<ApiUsageStateId> implements HasTenantId, HasVersion { public class ApiUsageState extends BaseData<ApiUsageStateId> implements HasTenantId, HasVersion {
@Serial
private static final long serialVersionUID = 8250339805336035966L; private static final long serialVersionUID = 8250339805336035966L;
private TenantId tenantId; private TenantId tenantId;
@ -41,6 +44,7 @@ public class ApiUsageState extends BaseData<ApiUsageStateId> implements HasTenan
private ApiUsageStateValue emailExecState; private ApiUsageStateValue emailExecState;
private ApiUsageStateValue smsExecState; private ApiUsageStateValue smsExecState;
private ApiUsageStateValue alarmExecState; private ApiUsageStateValue alarmExecState;
private ApiUsageStateValue edgeState;
private Long version; private Long version;
public ApiUsageState() { public ApiUsageState() {
@ -63,39 +67,44 @@ public class ApiUsageState extends BaseData<ApiUsageStateId> implements HasTenan
this.emailExecState = ur.getEmailExecState(); this.emailExecState = ur.getEmailExecState();
this.smsExecState = ur.getSmsExecState(); this.smsExecState = ur.getSmsExecState();
this.alarmExecState = ur.getAlarmExecState(); this.alarmExecState = ur.getAlarmExecState();
this.edgeState = ur.getEdgeState();
this.version = ur.getVersion(); this.version = ur.getVersion();
} }
public boolean isTransportEnabled() { public boolean isTransportEnabled() {
return !ApiUsageStateValue.DISABLED.equals(transportState); return transportState != ApiUsageStateValue.DISABLED;
} }
public boolean isReExecEnabled() { public boolean isReExecEnabled() {
return !ApiUsageStateValue.DISABLED.equals(reExecState); return reExecState != ApiUsageStateValue.DISABLED;
} }
public boolean isDbStorageEnabled() { public boolean isDbStorageEnabled() {
return !ApiUsageStateValue.DISABLED.equals(dbStorageState); return dbStorageState != ApiUsageStateValue.DISABLED;
} }
public boolean isJsExecEnabled() { public boolean isJsExecEnabled() {
return !ApiUsageStateValue.DISABLED.equals(jsExecState); return jsExecState != ApiUsageStateValue.DISABLED;
} }
public boolean isTbelExecEnabled() { public boolean isTbelExecEnabled() {
return !ApiUsageStateValue.DISABLED.equals(tbelExecState); return tbelExecState != ApiUsageStateValue.DISABLED;
} }
public boolean isEmailSendEnabled() { public boolean isEmailSendEnabled() {
return !ApiUsageStateValue.DISABLED.equals(emailExecState); return emailExecState != ApiUsageStateValue.DISABLED;
} }
public boolean isSmsSendEnabled() { public boolean isSmsSendEnabled() {
return !ApiUsageStateValue.DISABLED.equals(smsExecState); return smsExecState != ApiUsageStateValue.DISABLED;
} }
public boolean isAlarmCreationEnabled() { public boolean isAlarmCreationEnabled() {
return alarmExecState != ApiUsageStateValue.DISABLED; return alarmExecState != ApiUsageStateValue.DISABLED;
} }
public boolean isEdgeEnabled() {
return edgeState != ApiUsageStateValue.DISABLED;
}
} }

2
common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageStateValue.java

@ -19,8 +19,8 @@ public enum ApiUsageStateValue {
ENABLED, WARNING, DISABLED; ENABLED, WARNING, DISABLED;
public static ApiUsageStateValue toMoreRestricted(ApiUsageStateValue a, ApiUsageStateValue b) { public static ApiUsageStateValue toMoreRestricted(ApiUsageStateValue a, ApiUsageStateValue b) {
return a.ordinal() > b.ordinal() ? a : b; return a.ordinal() > b.ordinal() ? a : b;
} }
} }

5
common/data/src/main/java/org/thingsboard/server/common/data/edqs/fields/ApiUsageStateFields.java

@ -40,11 +40,12 @@ public class ApiUsageStateFields extends AbstractEntityFields {
private ApiUsageStateValue emailExecState; private ApiUsageStateValue emailExecState;
private ApiUsageStateValue smsExecState; private ApiUsageStateValue smsExecState;
private ApiUsageStateValue alarmExecState; private ApiUsageStateValue alarmExecState;
private ApiUsageStateValue edgeState;
public ApiUsageStateFields(UUID id, long createdTime, UUID tenantId, UUID entityId, String entityType, ApiUsageStateValue transportState, ApiUsageStateValue dbStorageState, public ApiUsageStateFields(UUID id, long createdTime, UUID tenantId, UUID entityId, String entityType, ApiUsageStateValue transportState, ApiUsageStateValue dbStorageState,
ApiUsageStateValue reExecState, ApiUsageStateValue jsExecState, ApiUsageStateValue tbelExecState, ApiUsageStateValue reExecState, ApiUsageStateValue jsExecState, ApiUsageStateValue tbelExecState,
ApiUsageStateValue emailExecState, ApiUsageStateValue smsExecState, ApiUsageStateValue alarmExecState, ApiUsageStateValue emailExecState, ApiUsageStateValue smsExecState, ApiUsageStateValue alarmExecState,
Long version) { ApiUsageStateValue edgeState, Long version) {
super(id, createdTime, tenantId, null, null, version); super(id, createdTime, tenantId, null, null, version);
this.entityId = (entityType != null && entityId != null) ? EntityIdFactory.getByTypeAndUuid(entityType, entityId) : null; this.entityId = (entityType != null && entityId != null) ? EntityIdFactory.getByTypeAndUuid(entityType, entityId) : null;
this.transportState = transportState; this.transportState = transportState;
@ -55,5 +56,7 @@ public class ApiUsageStateFields extends AbstractEntityFields {
this.emailExecState = emailExecState; this.emailExecState = emailExecState;
this.smsExecState = smsExecState; this.smsExecState = smsExecState;
this.alarmExecState = alarmExecState; this.alarmExecState = alarmExecState;
this.edgeState = edgeState;
} }
} }

58
common/data/src/main/java/org/thingsboard/server/common/data/edqs/fields/FieldsUtil.java

@ -42,43 +42,26 @@ import java.util.UUID;
public class FieldsUtil { public class FieldsUtil {
public static EntityFields toFields(Object entity) { public static EntityFields toFields(Object entity) {
if (entity instanceof Customer customer) { return switch (entity) {
return toFields(customer); case Customer customer -> toFields(customer);
} else if (entity instanceof Tenant tenant) { case Tenant tenant -> toFields(tenant);
return toFields(tenant); case TenantProfile tenantProfile -> toFields(tenantProfile);
} else if (entity instanceof TenantProfile tenantProfile) { case Device device -> toFields(device);
return toFields(tenantProfile); case Asset asset -> toFields(asset);
} else if (entity instanceof Device device) { case Edge edge -> toFields(edge);
return toFields(device); case EntityView entityView -> toFields(entityView);
} else if (entity instanceof Asset asset) { case User user -> toFields(user);
return toFields(asset); case Dashboard dashboard -> toFields(dashboard);
} else if (entity instanceof Edge edge) { case RuleChain ruleChain -> toFields(ruleChain);
return toFields(edge); case RuleNode ruleNode -> toFields(ruleNode);
} else if (entity instanceof EntityView entityView) { case WidgetType widgetType -> toFields(widgetType);
return toFields(entityView); case WidgetsBundle widgetsBundle -> toFields(widgetsBundle);
} else if (entity instanceof User user) { case DeviceProfile deviceProfile -> toFields(deviceProfile);
return toFields(user); case AssetProfile assetProfile -> toFields(assetProfile);
} else if (entity instanceof Dashboard dashboard) { case QueueStats queueStats -> toFields(queueStats);
return toFields(dashboard); case ApiUsageState apiUsageState -> toFields(apiUsageState);
} else if (entity instanceof RuleChain ruleChain) { default -> throw new IllegalArgumentException("Unsupported entity type: " + entity.getClass().getName());
return toFields(ruleChain); };
} else if (entity instanceof RuleNode ruleNode) {
return toFields(ruleNode);
} else if (entity instanceof WidgetType widgetType) {
return toFields(widgetType);
} else if (entity instanceof WidgetsBundle widgetsBundle) {
return toFields(widgetsBundle);
} else if (entity instanceof DeviceProfile deviceProfile) {
return toFields(deviceProfile);
} else if (entity instanceof AssetProfile assetProfile) {
return toFields(assetProfile);
} else if (entity instanceof QueueStats queueStats) {
return toFields(queueStats);
} else if (entity instanceof ApiUsageState apiUsageState) {
return toFields(apiUsageState);
} else {
throw new IllegalArgumentException("Unsupported entity type: " + entity.getClass().getName());
}
} }
private static CustomerFields toFields(Customer entity) { private static CustomerFields toFields(Customer entity) {
@ -284,6 +267,7 @@ public class FieldsUtil {
.emailExecState(entity.getEmailExecState()) .emailExecState(entity.getEmailExecState())
.smsExecState(entity.getSmsExecState()) .smsExecState(entity.getSmsExecState())
.alarmExecState(entity.getAlarmExecState()) .alarmExecState(entity.getAlarmExecState())
.edgeState(entity.getEdgeState())
.version(entity.getVersion()) .version(entity.getVersion())
.build(); .build();
} }

5
common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java

@ -124,6 +124,8 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura
private long maxSms; private long maxSms;
@Schema(example = "1000") @Schema(example = "1000")
private long maxCreatedAlarms; private long maxCreatedAlarms;
@Schema(example = "10000000")
private long maxEdgeEvents;
@RateLimit(fieldName = "REST requests for tenant") @RateLimit(fieldName = "REST requests for tenant")
private String tenantServerRestLimitsConfiguration; private String tenantServerRestLimitsConfiguration;
@ -218,6 +220,7 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura
case EMAIL_EXEC_COUNT -> maxEmails; case EMAIL_EXEC_COUNT -> maxEmails;
case SMS_EXEC_COUNT -> maxSms; case SMS_EXEC_COUNT -> maxSms;
case CREATED_ALARMS_COUNT -> maxCreatedAlarms; case CREATED_ALARMS_COUNT -> maxCreatedAlarms;
case EDGE_EVENT_COUNT -> maxEdgeEvents;
default -> 0L; default -> 0L;
}; };
} }
@ -226,7 +229,7 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura
public boolean getProfileFeatureEnabled(ApiUsageRecordKey key) { public boolean getProfileFeatureEnabled(ApiUsageRecordKey key) {
switch (key) { switch (key) {
case SMS_EXEC_COUNT: case SMS_EXEC_COUNT:
return smsEnabled == null || Boolean.TRUE.equals(smsEnabled); return smsEnabled == null || smsEnabled;
default: default:
return true; return true;
} }

8
common/edge-api/src/main/java/org/thingsboard/edge/exception/EdgeConnectionException.java

@ -15,15 +15,15 @@
*/ */
package org.thingsboard.edge.exception; package org.thingsboard.edge.exception;
import java.io.Serial;
public class EdgeConnectionException extends RuntimeException { public class EdgeConnectionException extends RuntimeException {
@Serial
private static final long serialVersionUID = -4372754681230555723L; private static final long serialVersionUID = -4372754681230555723L;
public EdgeConnectionException(String message) { public EdgeConnectionException(String message) {
super(message); super(message);
} }
public EdgeConnectionException(String message, Throwable cause) { }
super(message, cause);
}
}

29
common/edge-api/src/main/java/org/thingsboard/edge/exception/EdgeFeatureDisabledException.java

@ -0,0 +1,29 @@
/**
* Copyright © 2016-2026 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.edge.exception;
import java.io.Serial;
public class EdgeFeatureDisabledException extends RuntimeException {
@Serial
private static final long serialVersionUID = -4719918663724404639L;
public EdgeFeatureDisabledException(String message) {
super(message);
}
}

7
common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java

@ -26,6 +26,7 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.edge.exception.EdgeConnectionException; import org.thingsboard.edge.exception.EdgeConnectionException;
import org.thingsboard.edge.exception.EdgeFeatureDisabledException;
import org.thingsboard.server.common.data.ResourceUtils; import org.thingsboard.server.common.data.ResourceUtils;
import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.gen.edge.v1.ConnectRequestMsg; import org.thingsboard.server.gen.edge.v1.ConnectRequestMsg;
@ -170,7 +171,11 @@ public class EdgeGrpcClient implements EdgeRpcClient {
} catch (InterruptedException e) { } catch (InterruptedException e) {
log.error("[{}] Got interruption during disconnect!", edgeKey, e); log.error("[{}] Got interruption during disconnect!", edgeKey, e);
} }
onError.accept(new EdgeConnectionException("Failed to establish the connection! Response code: " + connectResponseMsg.getResponseCode().name())); if (ConnectResponseCode.FEATURE_DISABLED.equals(connectResponseMsg.getResponseCode())) {
onError.accept(new EdgeFeatureDisabledException(connectResponseMsg.getErrorMsg()));
} else {
onError.accept(new EdgeConnectionException("Failed to establish the connection! Response code: " + connectResponseMsg.getResponseCode().name()));
}
} }
} else if (responseMsg.hasEdgeUpdateMsg()) { } else if (responseMsg.hasEdgeUpdateMsg()) {
log.debug("[{}] Edge update message received {}", edgeKey, responseMsg.getEdgeUpdateMsg()); log.debug("[{}] Edge update message received {}", edgeKey, responseMsg.getEdgeUpdateMsg());

1
common/edge-api/src/main/proto/edge.proto

@ -97,6 +97,7 @@ enum ConnectResponseCode {
ACCEPTED = 0; ACCEPTED = 0;
BAD_CREDENTIALS = 1; BAD_CREDENTIALS = 1;
SERVER_UNAVAILABLE = 2; SERVER_UNAVAILABLE = 2;
FEATURE_DISABLED = 3;
} }
message ConnectResponseMsg { message ConnectResponseMsg {

28
common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java

@ -459,6 +459,7 @@ public class ProtoUtils {
case CREATED_ALARMS_COUNT -> ApiUsageRecordKeyProto.CREATED_ALARMS_COUNT; case CREATED_ALARMS_COUNT -> ApiUsageRecordKeyProto.CREATED_ALARMS_COUNT;
case ACTIVE_DEVICES -> ApiUsageRecordKeyProto.ACTIVE_DEVICES; case ACTIVE_DEVICES -> ApiUsageRecordKeyProto.ACTIVE_DEVICES;
case INACTIVE_DEVICES -> ApiUsageRecordKeyProto.INACTIVE_DEVICES; case INACTIVE_DEVICES -> ApiUsageRecordKeyProto.INACTIVE_DEVICES;
case EDGE_EVENT_COUNT -> ApiUsageRecordKeyProto.EDGE_EVENT_COUNT;
}; };
} }
@ -476,6 +477,7 @@ public class ProtoUtils {
case CREATED_ALARMS_COUNT -> ApiUsageRecordKey.CREATED_ALARMS_COUNT; case CREATED_ALARMS_COUNT -> ApiUsageRecordKey.CREATED_ALARMS_COUNT;
case ACTIVE_DEVICES -> ApiUsageRecordKey.ACTIVE_DEVICES; case ACTIVE_DEVICES -> ApiUsageRecordKey.ACTIVE_DEVICES;
case INACTIVE_DEVICES -> ApiUsageRecordKey.INACTIVE_DEVICES; case INACTIVE_DEVICES -> ApiUsageRecordKey.INACTIVE_DEVICES;
case EDGE_EVENT_COUNT -> ApiUsageRecordKey.EDGE_EVENT_COUNT;
}; };
} }
@ -1198,6 +1200,7 @@ public class ProtoUtils {
.setEmailExecState(apiUsageState.getEmailExecState().name()) .setEmailExecState(apiUsageState.getEmailExecState().name())
.setSmsExecState(apiUsageState.getSmsExecState().name()) .setSmsExecState(apiUsageState.getSmsExecState().name())
.setAlarmExecState(apiUsageState.getAlarmExecState().name()) .setAlarmExecState(apiUsageState.getAlarmExecState().name())
.setEdgeState(apiUsageState.getEdgeState().name())
.setVersion(apiUsageState.getVersion()) .setVersion(apiUsageState.getVersion())
.build(); .build();
} }
@ -1215,6 +1218,12 @@ public class ProtoUtils {
apiUsageState.setEmailExecState(ApiUsageStateValue.valueOf(proto.getEmailExecState())); apiUsageState.setEmailExecState(ApiUsageStateValue.valueOf(proto.getEmailExecState()));
apiUsageState.setSmsExecState(ApiUsageStateValue.valueOf(proto.getSmsExecState())); apiUsageState.setSmsExecState(ApiUsageStateValue.valueOf(proto.getSmsExecState()));
apiUsageState.setAlarmExecState(ApiUsageStateValue.valueOf(proto.getAlarmExecState())); apiUsageState.setAlarmExecState(ApiUsageStateValue.valueOf(proto.getAlarmExecState()));
// for backward compatibility, if any message is already in kafka, can simplify after the next release;
if (!proto.getEdgeState().isEmpty()) {
apiUsageState.setEdgeState(ApiUsageStateValue.valueOf(proto.getEdgeState()));
} else {
apiUsageState.setEdgeState(ApiUsageStateValue.ENABLED);
}
apiUsageState.setVersion(proto.getVersion()); apiUsageState.setVersion(proto.getVersion());
return apiUsageState; return apiUsageState;
} }
@ -1309,18 +1318,13 @@ public class ProtoUtils {
public static <T> TransportProtos.EntityUpdateMsg toEntityUpdateProto(T entity) { public static <T> TransportProtos.EntityUpdateMsg toEntityUpdateProto(T entity) {
var builder = TransportProtos.EntityUpdateMsg.newBuilder(); var builder = TransportProtos.EntityUpdateMsg.newBuilder();
if (entity instanceof Device) { switch (entity) {
builder.setDevice(toProto((Device) entity)); case Device device -> builder.setDevice(toProto(device));
} else if (entity instanceof DeviceProfile) { case DeviceProfile deviceProfile -> builder.setDeviceProfile(toProto(deviceProfile));
builder.setDeviceProfile(toProto((DeviceProfile) entity)); case Tenant tenant -> builder.setTenant(toProto(tenant));
} else if (entity instanceof Tenant) { case TenantProfile tenantProfile -> builder.setTenantProfile(toProto(tenantProfile));
builder.setTenant(toProto((Tenant) entity)); case ApiUsageState apiUsageState -> builder.setApiUsageState(toProto(apiUsageState));
} else if (entity instanceof TenantProfile) { default -> log.warn("[{}] entity does not support toProto serialization .", entity.getClass().getSimpleName());
builder.setTenantProfile(toProto((TenantProfile) entity));
} else if (entity instanceof ApiUsageState) {
builder.setApiUsageState(toProto((ApiUsageState) entity));
} else {
log.warn("[{}] entity does not support toProto serialization .", entity.getClass().getSimpleName());
} }
return builder.build(); return builder.build();
} }

14
common/proto/src/main/proto/queue.proto

@ -81,6 +81,7 @@ enum ApiUsageRecordKeyProto {
CREATED_ALARMS_COUNT = 8; CREATED_ALARMS_COUNT = 8;
ACTIVE_DEVICES = 9; ACTIVE_DEVICES = 9;
INACTIVE_DEVICES = 10; INACTIVE_DEVICES = 10;
EDGE_EVENT_COUNT = 11;
} }
message EntityIdProto { message EntityIdProto {
@ -387,6 +388,7 @@ message ApiUsageStateProto {
string smsExecState = 15; string smsExecState = 15;
string alarmExecState = 16; string alarmExecState = 16;
int64 version = 17; int64 version = 17;
string edgeState = 18;
} }
message RepositorySettingsProto { message RepositorySettingsProto {
@ -1783,19 +1785,15 @@ message ToTransportMsg {
} }
message UsageStatsKVProto { message UsageStatsKVProto {
string key = 1 [deprecated=true]; reserved 1;
reserved "key";
int64 value = 2; int64 value = 2;
ApiUsageRecordKeyProto recordKey = 3; ApiUsageRecordKeyProto recordKey = 3;
} }
message ToUsageStatsServiceMsg { message ToUsageStatsServiceMsg {
int64 tenantIdMSB = 1 [deprecated=true]; reserved 1 to 7;
int64 tenantIdLSB = 2 [deprecated=true]; reserved "tenantIdMSB", "tenantIdLSB", "entityIdMSB", "entityIdLSB", "values", "customerIdMSB", "customerIdLSB";
int64 entityIdMSB = 3 [deprecated=true];
int64 entityIdLSB = 4 [deprecated=true];
repeated UsageStatsKVProto values = 5 [deprecated=true];
int64 customerIdMSB = 6 [deprecated=true];
int64 customerIdLSB = 7 [deprecated=true];
string serviceId = 8; string serviceId = 8;
repeated UsageStatsServiceMsg msgs = 9; repeated UsageStatsServiceMsg msgs = 9;
} }

13
common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageReportClient.java

@ -32,9 +32,9 @@ import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.common.stats.TbApiUsageReportClient; import org.thingsboard.server.common.stats.TbApiUsageReportClient;
import org.thingsboard.server.common.util.ProtoUtils; import org.thingsboard.server.common.util.ProtoUtils;
import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.UsageStatsServiceMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg;
import org.thingsboard.server.gen.transport.TransportProtos.UsageStatsKVProto; import org.thingsboard.server.gen.transport.TransportProtos.UsageStatsKVProto;
import org.thingsboard.server.gen.transport.TransportProtos.UsageStatsServiceMsg;
import org.thingsboard.server.queue.TbQueueProducer; import org.thingsboard.server.queue.TbQueueProducer;
import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.discovery.PartitionService;
@ -178,7 +178,9 @@ public class DefaultTbApiUsageReportClient implements TbApiUsageReportClient {
@Override @Override
public void report(TenantId tenantId, CustomerId customerId, ApiUsageRecordKey key, long value) { public void report(TenantId tenantId, CustomerId customerId, ApiUsageRecordKey key, long value) {
if (!enabled) return; if (!enabled) {
return;
}
ReportLevel[] reportLevels = new ReportLevel[3]; ReportLevel[] reportLevels = new ReportLevel[3];
reportLevels[0] = ReportLevel.of(tenantId); reportLevels[0] = ReportLevel.of(tenantId);
@ -199,7 +201,9 @@ public class DefaultTbApiUsageReportClient implements TbApiUsageReportClient {
private void report(ApiUsageRecordKey key, long value, ReportLevel... levels) { private void report(ApiUsageRecordKey key, long value, ReportLevel... levels) {
ConcurrentMap<ReportLevel, AtomicLong> statsForKey = stats.get(key); ConcurrentMap<ReportLevel, AtomicLong> statsForKey = stats.get(key);
for (ReportLevel level : levels) { for (ReportLevel level : levels) {
if (level == null) continue; if (level == null) {
continue;
}
AtomicLong n = statsForKey.computeIfAbsent(level, k -> new AtomicLong()); AtomicLong n = statsForKey.computeIfAbsent(level, k -> new AtomicLong());
if (key.isCounter()) { if (key.isCounter()) {
@ -212,6 +216,7 @@ public class DefaultTbApiUsageReportClient implements TbApiUsageReportClient {
@Data @Data
private static class ReportLevel { private static class ReportLevel {
private final TenantId tenantId; private final TenantId tenantId;
private final CustomerId customerId; private final CustomerId customerId;
@ -231,12 +236,14 @@ public class DefaultTbApiUsageReportClient implements TbApiUsageReportClient {
@Data @Data
private static class ParentEntity { private static class ParentEntity {
private final TenantId tenantId; private final TenantId tenantId;
private final CustomerId customerId; private final CustomerId customerId;
public EntityId getId() { public EntityId getId() {
return customerId != null ? customerId : tenantId; return customerId != null ? customerId : tenantId;
} }
} }
} }

29
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java

@ -807,6 +807,31 @@ public class DefaultTransportService extends TransportActivityManager implements
sessions.remove(toSessionId(sessionInfo)); sessions.remove(toSessionId(sessionInfo));
} }
private void closeTenantSessions(TenantId tenantId, String reason) {
long msb = tenantId.getId().getMostSignificantBits();
long lsb = tenantId.getId().getLeastSignificantBits();
TransportProtos.SessionCloseNotificationProto notification = TransportProtos.SessionCloseNotificationProto.newBuilder()
.setMessage(reason).build();
int closed = 0;
for (Map.Entry<UUID, SessionMetaData> entry : sessions.entrySet()) {
UUID sessionId = entry.getKey();
SessionMetaData md = entry.getValue();
TransportProtos.SessionInfoProto sessionInfo = md.getSessionInfo();
if (sessionInfo.getTenantIdMSB() == msb && sessionInfo.getTenantIdLSB() == lsb) {
transportCallbackExecutor.submit(() -> {
md.getListener().onRemoteSessionCloseCommand(sessionId, notification);
if (md.getSessionType() == TransportProtos.SessionType.SYNC) {
deregisterSession(sessionInfo);
}
});
closed++;
}
}
if (closed > 0) {
log.info("[{}] Transport feature disabled due to API limits. Closing {} sessions.", tenantId, closed);
}
}
@Override @Override
public void log(TransportProtos.SessionInfoProto sessionInfo, String msg) { public void log(TransportProtos.SessionInfoProto sessionInfo, String msg) {
if (!logEnabled || sessionInfo == null || StringUtils.isEmpty(msg)) { if (!logEnabled || sessionInfo == null || StringUtils.isEmpty(msg)) {
@ -1008,7 +1033,9 @@ public class DefaultTransportService extends TransportActivityManager implements
case APIUSAGESTATE: case APIUSAGESTATE:
ApiUsageState apiUsageState = ProtoUtils.fromProto(msg.getApiUsageState()); ApiUsageState apiUsageState = ProtoUtils.fromProto(msg.getApiUsageState());
rateLimitService.update(apiUsageState.getTenantId(), apiUsageState.isTransportEnabled()); rateLimitService.update(apiUsageState.getTenantId(), apiUsageState.isTransportEnabled());
//TODO: if transport is disabled, we should close all sessions and not to check credentials. if (!apiUsageState.isTransportEnabled()) {
closeTenantSessions(apiUsageState.getTenantId(), "Transport feature disabled due to API limits!");
}
break; break;
case DEVICE: case DEVICE:
onDeviceUpdate(ProtoUtils.fromProto(msg.getDevice())); onDeviceUpdate(ProtoUtils.fromProto(msg.getDevice()));

10
dao/src/main/java/org/thingsboard/server/dao/edge/BaseEdgeEventService.java

@ -17,6 +17,7 @@ package org.thingsboard.server.dao.edge;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.thingsboard.server.cache.limits.RateLimitService; import org.thingsboard.server.cache.limits.RateLimitService;
import org.thingsboard.server.common.data.ApiUsageRecordKey;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EdgeId;
@ -25,6 +26,7 @@ import org.thingsboard.server.common.data.limit.LimitedApi;
import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.TimePageLink; import org.thingsboard.server.common.data.page.TimePageLink;
import org.thingsboard.server.common.msg.tools.TbRateLimitsException; import org.thingsboard.server.common.msg.tools.TbRateLimitsException;
import org.thingsboard.server.common.stats.TbApiUsageReportClient;
import org.thingsboard.server.dao.service.DataValidator; import org.thingsboard.server.dao.service.DataValidator;
public abstract class BaseEdgeEventService implements EdgeEventService { public abstract class BaseEdgeEventService implements EdgeEventService {
@ -35,6 +37,8 @@ public abstract class BaseEdgeEventService implements EdgeEventService {
private RateLimitService rateLimitService; private RateLimitService rateLimitService;
@Autowired @Autowired
private DataValidator<EdgeEvent> edgeEventValidator; private DataValidator<EdgeEvent> edgeEventValidator;
@Autowired(required = false)
private TbApiUsageReportClient apiUsageReportClient;
@Override @Override
public PageData<EdgeEvent> findEdgeEvents(TenantId tenantId, EdgeId edgeId, Long seqIdStart, Long seqIdEnd, TimePageLink pageLink) { public PageData<EdgeEvent> findEdgeEvents(TenantId tenantId, EdgeId edgeId, Long seqIdStart, Long seqIdEnd, TimePageLink pageLink) {
@ -56,4 +60,10 @@ public abstract class BaseEdgeEventService implements EdgeEventService {
edgeEventValidator.validate(edgeEvent, EdgeEvent::getTenantId); edgeEventValidator.validate(edgeEvent, EdgeEvent::getTenantId);
} }
protected void reportEdgeEventUsage(EdgeEvent edgeEvent) {
if (apiUsageReportClient != null) {
apiUsageReportClient.report(edgeEvent.getTenantId(), null, ApiUsageRecordKey.EDGE_EVENT_COUNT, 1);
}
}
} }

4
dao/src/main/java/org/thingsboard/server/dao/edge/EdgeDao.java

@ -31,10 +31,6 @@ import java.util.List;
import java.util.Optional; import java.util.Optional;
import java.util.UUID; import java.util.UUID;
/**
* The Interface EdgeDao.
*
*/
public interface EdgeDao extends Dao<Edge>, TenantEntityDao<Edge> { public interface EdgeDao extends Dao<Edge>, TenantEntityDao<Edge> {
Edge save(TenantId tenantId, Edge edge); Edge save(TenantId tenantId, Edge edge);

24
dao/src/main/java/org/thingsboard/server/dao/edge/EdgeEventDao.java

@ -24,36 +24,12 @@ import org.thingsboard.server.dao.Dao;
import java.util.UUID; import java.util.UUID;
/**
* The Interface EdgeEventDao.
*/
public interface EdgeEventDao extends Dao<EdgeEvent> { public interface EdgeEventDao extends Dao<EdgeEvent> {
/**
* Save or update edge event object
*
* @param edgeEvent the event object
* @return saved edge event object future
*/
ListenableFuture<Void> saveAsync(EdgeEvent edgeEvent); ListenableFuture<Void> saveAsync(EdgeEvent edgeEvent);
/**
* Find edge events by tenantId, edgeId and pageLink.
*
* @param tenantId the tenantId
* @param edgeId the edgeId
* @param seqIdStart the seq id start
* @param seqIdEnd the seq id end
* @param pageLink the pageLink
* @return the event list
*/
PageData<EdgeEvent> findEdgeEvents(UUID tenantId, EdgeId edgeId, Long seqIdStart, Long seqIdEnd, TimePageLink pageLink); PageData<EdgeEvent> findEdgeEvents(UUID tenantId, EdgeId edgeId, Long seqIdStart, Long seqIdEnd, TimePageLink pageLink);
/**
* Executes stored procedure to cleanup old edge events.
* @param ttl the ttl for edge events in seconds
*/
void cleanupEvents(long ttl); void cleanupEvents(long ttl);
} }

1
dao/src/main/java/org/thingsboard/server/dao/edge/PostgresEdgeEventService.java

@ -70,6 +70,7 @@ public class PostgresEdgeEventService extends BaseEdgeEventService {
public void onSuccess(Void result) { public void onSuccess(Void result) {
statsCounterService.ifPresent(statsCounterService -> statsCounterService.ifPresent(statsCounterService ->
statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_ADDED, edgeEvent.getTenantId(), edgeEvent.getEdgeId(), 1)); statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_ADDED, edgeEvent.getTenantId(), edgeEvent.getEdgeId(), 1));
reportEdgeEventUsage(edgeEvent);
eventPublisher.publishEvent(SaveEntityEvent.builder() eventPublisher.publishEvent(SaveEntityEvent.builder()
.tenantId(edgeEvent.getTenantId()) .tenantId(edgeEvent.getTenantId())
.entityId(edgeEvent.getEdgeId()) .entityId(edgeEvent.getEdgeId())

1
dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java

@ -539,6 +539,7 @@ public class ModelConstants {
public static final String API_USAGE_STATE_EMAIL_EXEC_COLUMN = "email_exec"; public static final String API_USAGE_STATE_EMAIL_EXEC_COLUMN = "email_exec";
public static final String API_USAGE_STATE_SMS_EXEC_COLUMN = "sms_exec"; public static final String API_USAGE_STATE_SMS_EXEC_COLUMN = "sms_exec";
public static final String API_USAGE_STATE_ALARM_EXEC_COLUMN = "alarm_exec"; public static final String API_USAGE_STATE_ALARM_EXEC_COLUMN = "alarm_exec";
public static final String API_USAGE_STATE_EDGE_COLUMN = "edge";
/** /**
* Resource constants. * Resource constants.

5
dao/src/main/java/org/thingsboard/server/dao/model/sql/ApiUsageStateEntity.java

@ -69,6 +69,9 @@ public class ApiUsageStateEntity extends BaseVersionedEntity<ApiUsageState> impl
@Enumerated(EnumType.STRING) @Enumerated(EnumType.STRING)
@Column(name = ModelConstants.API_USAGE_STATE_ALARM_EXEC_COLUMN) @Column(name = ModelConstants.API_USAGE_STATE_ALARM_EXEC_COLUMN)
private ApiUsageStateValue alarmExecState = ApiUsageStateValue.ENABLED; private ApiUsageStateValue alarmExecState = ApiUsageStateValue.ENABLED;
@Enumerated(EnumType.STRING)
@Column(name = ModelConstants.API_USAGE_STATE_EDGE_COLUMN)
private ApiUsageStateValue edgeState = ApiUsageStateValue.ENABLED;
public ApiUsageStateEntity() { public ApiUsageStateEntity() {
} }
@ -90,6 +93,7 @@ public class ApiUsageStateEntity extends BaseVersionedEntity<ApiUsageState> impl
this.emailExecState = ur.getEmailExecState(); this.emailExecState = ur.getEmailExecState();
this.smsExecState = ur.getSmsExecState(); this.smsExecState = ur.getSmsExecState();
this.alarmExecState = ur.getAlarmExecState(); this.alarmExecState = ur.getAlarmExecState();
this.edgeState = ur.getEdgeState();
} }
@Override @Override
@ -110,6 +114,7 @@ public class ApiUsageStateEntity extends BaseVersionedEntity<ApiUsageState> impl
ur.setEmailExecState(emailExecState); ur.setEmailExecState(emailExecState);
ur.setSmsExecState(smsExecState); ur.setSmsExecState(smsExecState);
ur.setAlarmExecState(alarmExecState); ur.setAlarmExecState(alarmExecState);
ur.setEdgeState(edgeState);
ur.setVersion(version); ur.setVersion(version);
return ur; return ur;
} }

2
dao/src/main/java/org/thingsboard/server/dao/sql/usagerecord/ApiUsageStateRepository.java

@ -51,7 +51,7 @@ public interface ApiUsageStateRepository extends JpaRepository<ApiUsageStateEnti
@Query("SELECT new org.thingsboard.server.common.data.edqs.fields.ApiUsageStateFields(a.id, a.createdTime, a.tenantId," + @Query("SELECT new org.thingsboard.server.common.data.edqs.fields.ApiUsageStateFields(a.id, a.createdTime, a.tenantId," +
"a.entityId, a.entityType, a.transportState, a.dbStorageState, a.reExecState, a.jsExecState, a.tbelExecState, " + "a.entityId, a.entityType, a.transportState, a.dbStorageState, a.reExecState, a.jsExecState, a.tbelExecState, " +
"a.emailExecState, a.smsExecState, a.alarmExecState, a.version) FROM ApiUsageStateEntity a WHERE a.id > :id ORDER BY a.id") "a.emailExecState, a.smsExecState, a.alarmExecState, a.edgeState, a.version) FROM ApiUsageStateEntity a WHERE a.id > :id ORDER BY a.id")
List<ApiUsageStateFields> findNextBatch(@Param("id") UUID id, Limit limit); List<ApiUsageStateFields> findNextBatch(@Param("id") UUID id, Limit limit);
} }

3
dao/src/main/java/org/thingsboard/server/dao/usagerecord/ApiUsageStateServiceImpl.java

@ -127,6 +127,7 @@ public class ApiUsageStateServiceImpl extends AbstractEntityService implements A
apiUsageState.setSmsExecState(smsApiUsageState); apiUsageState.setSmsExecState(smsApiUsageState);
apiUsageState.setEmailExecState(ApiUsageStateValue.ENABLED); apiUsageState.setEmailExecState(ApiUsageStateValue.ENABLED);
apiUsageState.setAlarmExecState(ApiUsageStateValue.ENABLED); apiUsageState.setAlarmExecState(ApiUsageStateValue.ENABLED);
apiUsageState.setEdgeState(ApiUsageStateValue.ENABLED);
apiUsageStateValidator.validate(apiUsageState, ApiUsageState::getTenantId); apiUsageStateValidator.validate(apiUsageState, ApiUsageState::getTenantId);
ApiUsageState saved = apiUsageStateDao.save(apiUsageState.getTenantId(), apiUsageState); ApiUsageState saved = apiUsageStateDao.save(apiUsageState.getTenantId(), apiUsageState);
@ -156,6 +157,8 @@ public class ApiUsageStateServiceImpl extends AbstractEntityService implements A
new StringDataEntry(ApiFeature.SMS.getApiStateKey(), smsApiUsageState.name()))); new StringDataEntry(ApiFeature.SMS.getApiStateKey(), smsApiUsageState.name())));
apiUsageStates.add(new BasicTsKvEntry(saved.getCreatedTime(), apiUsageStates.add(new BasicTsKvEntry(saved.getCreatedTime(),
new StringDataEntry(ApiFeature.ALARM.getApiStateKey(), ApiUsageStateValue.ENABLED.name()))); new StringDataEntry(ApiFeature.ALARM.getApiStateKey(), ApiUsageStateValue.ENABLED.name())));
apiUsageStates.add(new BasicTsKvEntry(saved.getCreatedTime(),
new StringDataEntry(ApiFeature.EDGE.getApiStateKey(), ApiUsageStateValue.ENABLED.name())));
tsService.save(tenantId, saved.getId(), apiUsageStates, 0L); tsService.save(tenantId, saved.getId(), apiUsageStates, 0L);
if (configuration != null) { if (configuration != null) {

1
dao/src/main/resources/sql/schema-entities.sql

@ -706,6 +706,7 @@ CREATE TABLE IF NOT EXISTS api_usage_state (
email_exec varchar(32), email_exec varchar(32),
sms_exec varchar(32), sms_exec varchar(32),
alarm_exec varchar(32), alarm_exec varchar(32),
edge varchar(32),
version BIGINT DEFAULT 1, version BIGINT DEFAULT 1,
CONSTRAINT api_usage_state_unq_key UNIQUE (tenant_id, entity_id) CONSTRAINT api_usage_state_unq_key UNIQUE (tenant_id, entity_id)
); );

85
dao/src/test/java/org/thingsboard/server/dao/service/ApiUsageStateServiceTest.java

@ -23,6 +23,10 @@ import org.thingsboard.server.common.data.ApiUsageStateValue;
import org.thingsboard.server.dao.exception.IncorrectParameterException; import org.thingsboard.server.dao.exception.IncorrectParameterException;
import org.thingsboard.server.dao.usagerecord.ApiUsageStateService; import org.thingsboard.server.dao.usagerecord.ApiUsageStateService;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertTrue;
@DaoSqlTest @DaoSqlTest
public class ApiUsageStateServiceTest extends AbstractServiceTest { public class ApiUsageStateServiceTest extends AbstractServiceTest {
@ -33,7 +37,22 @@ public class ApiUsageStateServiceTest extends AbstractServiceTest {
@Test @Test
public void testFindTenantApiUsageState() { public void testFindTenantApiUsageState() {
ApiUsageState state = apiUsageStateService.findTenantApiUsageState(tenantId); ApiUsageState state = apiUsageStateService.findTenantApiUsageState(tenantId);
Assert.assertNotNull(state); assertNotNull(state);
}
@Test
public void testDefaultStateIsEnabled() {
ApiUsageState state = apiUsageStateService.findTenantApiUsageState(tenantId);
assertNotNull(state);
assertTrue(state.isTransportEnabled());
assertTrue(state.isReExecEnabled());
assertTrue(state.isDbStorageEnabled());
assertTrue(state.isJsExecEnabled());
assertTrue(state.isTbelExecEnabled());
assertTrue(state.isEmailSendEnabled());
assertTrue(state.isSmsSendEnabled());
assertTrue(state.isAlarmCreationEnabled());
assertTrue(state.isEdgeEnabled());
} }
@Test @Test
@ -42,7 +61,7 @@ public class ApiUsageStateServiceTest extends AbstractServiceTest {
state.setTransportState(ApiUsageStateValue.DISABLED); state.setTransportState(ApiUsageStateValue.DISABLED);
ApiUsageState updated = apiUsageStateService.update(state); ApiUsageState updated = apiUsageStateService.update(state);
Assert.assertEquals(ApiUsageStateValue.DISABLED, updated.getTransportState()); assertEquals(ApiUsageStateValue.DISABLED, updated.getTransportState());
} }
@Test @Test
@ -53,20 +72,76 @@ public class ApiUsageStateServiceTest extends AbstractServiceTest {
Assert.assertThrows(IncorrectParameterException.class, () -> apiUsageStateService.update(newState)); Assert.assertThrows(IncorrectParameterException.class, () -> apiUsageStateService.update(newState));
} }
@Test
public void testTransportStateUpdate() {
ApiUsageState state = apiUsageStateService.findTenantApiUsageState(tenantId);
state.setTransportState(ApiUsageStateValue.WARNING);
ApiUsageState updated = apiUsageStateService.update(state);
assertEquals(ApiUsageStateValue.WARNING, updated.getTransportState());
assertTrue(updated.isTransportEnabled());
updated.setTransportState(ApiUsageStateValue.DISABLED);
updated = apiUsageStateService.update(updated);
assertEquals(ApiUsageStateValue.DISABLED, updated.getTransportState());
Assert.assertFalse(updated.isTransportEnabled());
updated.setTransportState(ApiUsageStateValue.ENABLED);
updated = apiUsageStateService.update(updated);
assertEquals(ApiUsageStateValue.ENABLED, updated.getTransportState());
assertTrue(updated.isTransportEnabled());
}
@Test
public void testEdgeStateUpdate() {
ApiUsageState state = apiUsageStateService.findTenantApiUsageState(tenantId);
state.setEdgeState(ApiUsageStateValue.WARNING);
ApiUsageState updated = apiUsageStateService.update(state);
assertEquals(ApiUsageStateValue.WARNING, updated.getEdgeState());
assertTrue(updated.isEdgeEnabled());
updated.setEdgeState(ApiUsageStateValue.DISABLED);
updated = apiUsageStateService.update(updated);
assertEquals(ApiUsageStateValue.DISABLED, updated.getEdgeState());
Assert.assertFalse(updated.isEdgeEnabled());
ApiUsageState fetched = apiUsageStateService.findTenantApiUsageState(tenantId);
assertEquals(ApiUsageStateValue.DISABLED, fetched.getEdgeState());
}
@Test
public void testMultipleStatesIndependent() {
ApiUsageState state = apiUsageStateService.findTenantApiUsageState(tenantId);
state.setEdgeState(ApiUsageStateValue.DISABLED);
state.setTransportState(ApiUsageStateValue.WARNING);
state.setReExecState(ApiUsageStateValue.ENABLED);
ApiUsageState updated = apiUsageStateService.update(state);
assertEquals(ApiUsageStateValue.DISABLED, updated.getEdgeState());
assertEquals(ApiUsageStateValue.WARNING, updated.getTransportState());
assertEquals(ApiUsageStateValue.ENABLED, updated.getReExecState());
Assert.assertFalse(updated.isEdgeEnabled());
assertTrue(updated.isTransportEnabled());
assertTrue(updated.isReExecEnabled());
}
@Test @Test
public void testFindApiUsageStateByEntityId() { public void testFindApiUsageStateByEntityId() {
ApiUsageState state = apiUsageStateService.findApiUsageStateByEntityId(tenantId); ApiUsageState state = apiUsageStateService.findApiUsageStateByEntityId(tenantId);
Assert.assertNotNull(state); assertNotNull(state);
} }
@Test @Test
public void testDeleteByTenantId() { public void testDeleteByTenantId() {
ApiUsageState state = apiUsageStateService.findTenantApiUsageState(tenantId); ApiUsageState state = apiUsageStateService.findTenantApiUsageState(tenantId);
Assert.assertNotNull(state); assertNotNull(state);
apiUsageStateService.deleteByTenantId(tenantId); apiUsageStateService.deleteByTenantId(tenantId);
state = apiUsageStateService.findTenantApiUsageState(tenantId); state = apiUsageStateService.findTenantApiUsageState(tenantId);
Assert.assertNull(state); assertNull(state);
} }
} }

8
dao/src/test/java/org/thingsboard/server/dao/service/DeviceServiceTest.java

@ -27,7 +27,7 @@ import org.junit.Test;
import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Assertions;
import org.mockito.Mockito; import org.mockito.Mockito;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.mock.mockito.SpyBean; import org.springframework.test.context.bean.override.mockito.MockitoSpyBean;
import org.springframework.transaction.PlatformTransactionManager; import org.springframework.transaction.PlatformTransactionManager;
import org.springframework.transaction.TransactionStatus; import org.springframework.transaction.TransactionStatus;
import org.springframework.transaction.support.DefaultTransactionDefinition; import org.springframework.transaction.support.DefaultTransactionDefinition;
@ -68,12 +68,12 @@ import org.thingsboard.server.dao.customer.CustomerService;
import org.thingsboard.server.dao.device.DeviceCredentialsService; import org.thingsboard.server.dao.device.DeviceCredentialsService;
import org.thingsboard.server.dao.device.DeviceProfileService; import org.thingsboard.server.dao.device.DeviceProfileService;
import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.exception.DataValidationException;
import org.thingsboard.server.dao.exception.DeviceCredentialsValidationException; import org.thingsboard.server.dao.exception.DeviceCredentialsValidationException;
import org.thingsboard.server.dao.ota.OtaPackageService; import org.thingsboard.server.dao.ota.OtaPackageService;
import org.thingsboard.server.dao.service.validator.DeviceCredentialsDataValidator; import org.thingsboard.server.dao.service.validator.DeviceCredentialsDataValidator;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache; import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.dao.tenant.TenantProfileService; import org.thingsboard.server.dao.tenant.TenantProfileService;
import org.thingsboard.server.exception.DataValidationException;
import java.nio.ByteBuffer; import java.nio.ByteBuffer;
import java.util.ArrayList; import java.util.ArrayList;
@ -111,10 +111,10 @@ public class DeviceServiceTest extends AbstractServiceTest {
private CalculatedFieldService calculatedFieldService; private CalculatedFieldService calculatedFieldService;
@Autowired @Autowired
private PlatformTransactionManager platformTransactionManager; private PlatformTransactionManager platformTransactionManager;
@SpyBean @MockitoSpyBean
private DeviceCredentialsDataValidator validator; private DeviceCredentialsDataValidator validator;
private IdComparator<Device> idComparator = new IdComparator<>(); private final IdComparator<Device> idComparator = new IdComparator<>();
private TenantId anotherTenantId; private TenantId anotherTenantId;
private static ListeningExecutorService executor; private static ListeningExecutorService executor;

6
dao/src/test/java/org/thingsboard/server/dao/service/EdgeEventServiceTest.java

@ -74,7 +74,7 @@ public class EdgeEventServiceTest extends AbstractServiceTest {
PageData<EdgeEvent> edgeEvents = edgeEventService.findEdgeEvents(tenantId, edgeId, 0L, null, new TimePageLink(1)); PageData<EdgeEvent> edgeEvents = edgeEventService.findEdgeEvents(tenantId, edgeId, 0L, null, new TimePageLink(1));
Assert.assertFalse(edgeEvents.getData().isEmpty()); Assert.assertFalse(edgeEvents.getData().isEmpty());
EdgeEvent saved = edgeEvents.getData().get(0); EdgeEvent saved = edgeEvents.getData().getFirst();
Assert.assertEquals(saved.getTenantId(), edgeEvent.getTenantId()); Assert.assertEquals(saved.getTenantId(), edgeEvent.getTenantId());
Assert.assertEquals(saved.getEdgeId(), edgeEvent.getEdgeId()); Assert.assertEquals(saved.getEdgeId(), edgeEvent.getEdgeId());
Assert.assertEquals(saved.getEntityId(), edgeEvent.getEntityId()); Assert.assertEquals(saved.getEntityId(), edgeEvent.getEntityId());
@ -124,7 +124,7 @@ public class EdgeEventServiceTest extends AbstractServiceTest {
Assert.assertNotNull(edgeEvents.getData()); Assert.assertNotNull(edgeEvents.getData());
Assert.assertEquals(1, edgeEvents.getData().size()); Assert.assertEquals(1, edgeEvents.getData().size());
Assert.assertEquals(Uuids.startOf(eventTime + 1), edgeEvents.getData().get(0).getUuidId()); Assert.assertEquals(Uuids.startOf(eventTime + 1), edgeEvents.getData().getFirst().getUuidId());
Assert.assertFalse(edgeEvents.hasNext()); Assert.assertFalse(edgeEvents.hasNext());
edgeEventDao.cleanupEvents(1); edgeEventDao.cleanupEvents(1);
@ -136,4 +136,4 @@ public class EdgeEventServiceTest extends AbstractServiceTest {
return edgeEventService.saveAsync(edgeEvent); return edgeEventService.saveAsync(edgeEvent);
} }
} }

6
edqs/src/test/java/org/thingsboard/server/edqs/repo/ApiUsageStateFilterTest.java

@ -37,6 +37,7 @@ import org.thingsboard.server.common.data.query.KeyFilter;
import org.thingsboard.server.common.data.query.StringFilterPredicate; import org.thingsboard.server.common.data.query.StringFilterPredicate;
import java.util.Arrays; import java.util.Arrays;
import java.util.List;
import java.util.UUID; import java.util.UUID;
public class ApiUsageStateFilterTest extends AbstractEDQTest { public class ApiUsageStateFilterTest extends AbstractEDQTest {
@ -64,7 +65,7 @@ public class ApiUsageStateFilterTest extends AbstractEDQTest {
var result = repository.findEntityDataByQuery(tenantId, null, getEntityDataQuery(new CustomerId(customerId)), false); var result = repository.findEntityDataByQuery(tenantId, null, getEntityDataQuery(new CustomerId(customerId)), false);
Assert.assertEquals(1, result.getTotalElements()); Assert.assertEquals(1, result.getTotalElements());
var customer = result.getData().get(0); var customer = result.getData().getFirst();
Assert.assertEquals("Customer A", customer.getLatest().get(EntityKeyType.ENTITY_FIELD).get("name").getValue()); Assert.assertEquals("Customer A", customer.getLatest().get(EntityKeyType.ENTITY_FIELD).get("name").getValue());
} }
@ -81,6 +82,7 @@ public class ApiUsageStateFilterTest extends AbstractEDQTest {
apiUsageState.setSmsExecState(ApiUsageStateValue.ENABLED); apiUsageState.setSmsExecState(ApiUsageStateValue.ENABLED);
apiUsageState.setEmailExecState(ApiUsageStateValue.ENABLED); apiUsageState.setEmailExecState(ApiUsageStateValue.ENABLED);
apiUsageState.setAlarmExecState(ApiUsageStateValue.ENABLED); apiUsageState.setAlarmExecState(ApiUsageStateValue.ENABLED);
apiUsageState.setEdgeState(ApiUsageStateValue.ENABLED);
return apiUsageState; return apiUsageState;
} }
@ -99,7 +101,7 @@ public class ApiUsageStateFilterTest extends AbstractEDQTest {
nameFilter.setPredicate(predicate); nameFilter.setPredicate(predicate);
nameFilter.setValueType(EntityKeyValueType.STRING); nameFilter.setValueType(EntityKeyValueType.STRING);
return new EntityDataQuery(filter, pageLink, entityFields, null, Arrays.asList(nameFilter)); return new EntityDataQuery(filter, pageLink, entityFields, null, List.of(nameFilter));
} }
} }

82
ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.html

@ -182,18 +182,18 @@
<mat-hint></mat-hint> <mat-hint></mat-hint>
</mat-form-field> </mat-form-field>
<mat-form-field class="mat-block flex-1" appearance="outline" subscriptSizing="dynamic"> <mat-form-field class="mat-block flex-1" appearance="outline" subscriptSizing="dynamic">
<mat-label translate>tenant-profile.max-transport-messages</mat-label> <mat-label translate>tenant-profile.max-rule-node-executions-per-message</mat-label>
<input matInput required min="0" step="1" <input matInput required min="0" step="1"
formControlName="maxTransportMessages" formControlName="maxRuleNodeExecutionsPerMessage"
type="number"> type="number">
@if (tenantProfileConfigurationForm.get('maxTransportMessages').hasError('required')) { @if (tenantProfileConfigurationForm.get('maxRuleNodeExecutionsPerMessage').hasError('required')) {
<mat-error> <mat-error>
{{ 'tenant-profile.max-transport-messages-required' | translate}} {{ 'tenant-profile.max-rule-node-executions-per-message-required' | translate}}
</mat-error> </mat-error>
} }
@if (tenantProfileConfigurationForm.get('maxTransportMessages').hasError('min')) { @if (tenantProfileConfigurationForm.get('maxRuleNodeExecutionsPerMessage').hasError('min')) {
<mat-error> <mat-error>
{{ 'tenant-profile.max-transport-messages-range' | translate}} {{ 'tenant-profile.max-rule-node-executions-per-message-range' | translate}}
</mat-error> </mat-error>
} }
<mat-hint></mat-hint> <mat-hint></mat-hint>
@ -242,24 +242,58 @@
<mat-hint></mat-hint> <mat-hint></mat-hint>
</mat-form-field> </mat-form-field>
</div> </div>
</ng-template>
</mat-expansion-panel>
</fieldset>
<fieldset class="fields-group">
<legend class="group-title">
{{ 'tenant-profile.api-usage' | translate }} <span translate>tenant-profile.unlimited</span>
</legend>
<div class="fields-element flex flex-1 flex-row xs:flex-col gt-xs:gap-4">
<mat-form-field class="mat-block flex-1" appearance="outline" subscriptSizing="dynamic">
<mat-label translate>tenant-profile.max-transport-messages</mat-label>
<input matInput required min="0" step="1"
formControlName="maxTransportMessages"
type="number">
@if (tenantProfileConfigurationForm.get('maxTransportMessages').hasError('required')) {
<mat-error>
{{ 'tenant-profile.max-transport-messages-required' | translate}}
</mat-error>
}
@if (tenantProfileConfigurationForm.get('maxTransportMessages').hasError('min')) {
<mat-error>
{{ 'tenant-profile.max-transport-messages-range' | translate}}
</mat-error>
}
<mat-hint></mat-hint>
</mat-form-field>
<mat-form-field class="mat-block flex-1" appearance="outline" subscriptSizing="dynamic">
<mat-label translate>tenant-profile.max-edge-events</mat-label>
<input matInput required min="0" step="1"
formControlName="maxEdgeEvents"
type="number">
@if (tenantProfileConfigurationForm.get('maxEdgeEvents').hasError('required')) {
<mat-error>
{{ 'tenant-profile.max-edge-events-required' | translate}}
</mat-error>
}
@if (tenantProfileConfigurationForm.get('maxEdgeEvents').hasError('min')) {
<mat-error>
{{ 'tenant-profile.max-edge-events-range' | translate}}
</mat-error>
}
<mat-hint></mat-hint>
</mat-form-field>
</div>
<mat-expansion-panel class="configuration-panel">
<mat-expansion-panel-header>
<mat-panel-description class="flex items-stretch justify-end" translate>
tenant-profile.advanced-settings
</mat-panel-description>
</mat-expansion-panel-header>
<ng-template matExpansionPanelContent>
<div class="flex flex-1 flex-row xs:flex-col gt-xs:gap-4"> <div class="flex flex-1 flex-row xs:flex-col gt-xs:gap-4">
<mat-form-field class="mat-block flex-1" appearance="outline" subscriptSizing="dynamic">
<mat-label translate>tenant-profile.max-rule-node-executions-per-message</mat-label>
<input matInput required min="0" step="1"
formControlName="maxRuleNodeExecutionsPerMessage"
type="number">
@if (tenantProfileConfigurationForm.get('maxRuleNodeExecutionsPerMessage').hasError('required')) {
<mat-error>
{{ 'tenant-profile.max-rule-node-executions-per-message-required' | translate}}
</mat-error>
}
@if (tenantProfileConfigurationForm.get('maxRuleNodeExecutionsPerMessage').hasError('min')) {
<mat-error>
{{ 'tenant-profile.max-rule-node-executions-per-message-range' | translate}}
</mat-error>
}
<mat-hint></mat-hint>
</mat-form-field>
<mat-form-field class="mat-block flex-1" appearance="outline" subscriptSizing="dynamic"> <mat-form-field class="mat-block flex-1" appearance="outline" subscriptSizing="dynamic">
<mat-label translate>tenant-profile.max-transport-data-points</mat-label> <mat-label translate>tenant-profile.max-transport-data-points</mat-label>
<input matInput required min="0" step="1" <input matInput required min="0" step="1"
@ -277,10 +311,12 @@
} }
<mat-hint></mat-hint> <mat-hint></mat-hint>
</mat-form-field> </mat-form-field>
<div class="flex-1"></div>
</div> </div>
</ng-template> </ng-template>
</mat-expansion-panel> </mat-expansion-panel>
</fieldset> </fieldset>
<fieldset class="fields-group"> <fieldset class="fields-group">
<legend class="group-title"> <legend class="group-title">
{{ 'tenant-profile.calculated-fields' | translate }} <span translate>tenant-profile.unlimited</span> {{ 'tenant-profile.calculated-fields' | translate }} <span translate>tenant-profile.unlimited</span>

1
ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.ts

@ -88,6 +88,7 @@ export class DefaultTenantProfileConfigurationComponent implements ControlValueA
maxSms: [0], maxSms: [0],
smsEnabled: [false], smsEnabled: [false],
maxCreatedAlarms: [0, [Validators.required, Validators.min(0)]], maxCreatedAlarms: [0, [Validators.required, Validators.min(0)]],
maxEdgeEvents: [0, [Validators.required, Validators.min(0)]],
maxDebugModeDurationMinutes: [0, [Validators.min(0)]], maxDebugModeDurationMinutes: [0, [Validators.min(0)]],
defaultStorageTtlDays: [0, [Validators.required, Validators.min(0)]], defaultStorageTtlDays: [0, [Validators.required, Validators.min(0)]],
alarmsTtlDays: [0, [Validators.required, Validators.min(0)]], alarmsTtlDays: [0, [Validators.required, Validators.min(0)]],

2
ui-ngx/src/app/shared/models/tenant.model.ts

@ -71,6 +71,7 @@ export interface DefaultTenantProfileConfiguration {
maxSms: number; maxSms: number;
smsEnabled: boolean; smsEnabled: boolean;
maxCreatedAlarms: number; maxCreatedAlarms: number;
maxEdgeEvents: number;
maxDebugModeDurationMinutes: number; maxDebugModeDurationMinutes: number;
@ -154,6 +155,7 @@ export function createTenantProfileConfiguration(type: TenantProfileType): Tenan
maxSms: 0, maxSms: 0,
smsEnabled: true, smsEnabled: true,
maxCreatedAlarms: 0, maxCreatedAlarms: 0,
maxEdgeEvents: 0,
maxDebugModeDurationMinutes: 15, maxDebugModeDurationMinutes: 15,
tenantServerRestLimitsConfiguration: '', tenantServerRestLimitsConfiguration: '',
customerServerRestLimitsConfiguration: '', customerServerRestLimitsConfiguration: '',

6
ui-ngx/src/assets/locale/locale.constant-en_US.json

@ -6263,6 +6263,7 @@
"advanced-settings": "Advanced settings", "advanced-settings": "Advanced settings",
"entities": "Entities", "entities": "Entities",
"rule-engine": "Rule Engine", "rule-engine": "Rule Engine",
"api-usage": "API Usage",
"time-to-live": "Time-to-live", "time-to-live": "Time-to-live",
"calculated-fields": "Calculated fields", "calculated-fields": "Calculated fields",
"alarms-and-notifications": "Alarms and notifications", "alarms-and-notifications": "Alarms and notifications",
@ -6398,7 +6399,10 @@
"max-sms-range": "SMS sent maximum number can't be negative", "max-sms-range": "SMS sent maximum number can't be negative",
"max-created-alarms": "Alarms created maximum number", "max-created-alarms": "Alarms created maximum number",
"max-created-alarms-required": "Alarms created maximum number is required.", "max-created-alarms-required": "Alarms created maximum number is required.",
"max-created-alarms-range": "Alarms created maximum number be negative", "max-created-alarms-range": "Alarms created maximum number can't be negative",
"max-edge-events": "Edge events maximum number",
"max-edge-events-required": "Edge events maximum number is required.",
"max-edge-events-range": "Edge events maximum number can't be negative",
"no-queue": "No Queue configured", "no-queue": "No Queue configured",
"add-queue": "Add Queue", "add-queue": "Add Queue",
"queues-with-count": "Queues ({{count}})", "queues-with-count": "Queues ({{count}})",

Loading…
Cancel
Save