diff --git a/application/src/main/data/upgrade/basic/schema_update.sql b/application/src/main/data/upgrade/basic/schema_update.sql index 2f9c0cd302..59ff33a9fe 100644 --- a/application/src/main/data/upgrade/basic/schema_update.sql +++ b/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); -- 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 diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/BaseApiUsageState.java b/application/src/main/java/org/thingsboard/server/service/apiusage/BaseApiUsageState.java index 85f828973a..3320481e03 100644 --- a/application/src/main/java/org/thingsboard/server/service/apiusage/BaseApiUsageState.java +++ b/application/src/main/java/org/thingsboard/server/service/apiusage/BaseApiUsageState.java @@ -33,6 +33,7 @@ import java.util.Map; import java.util.concurrent.ConcurrentHashMap; public abstract class BaseApiUsageState { + private final Map currentCycleValues = new ConcurrentHashMap<>(); private final Map currentHourValues = new ConcurrentHashMap<>(); @@ -137,55 +138,31 @@ public abstract class BaseApiUsageState { } public ApiUsageStateValue getFeatureValue(ApiFeature feature) { - switch (feature) { - case TRANSPORT: - return apiUsageState.getTransportState(); - case RE: - return apiUsageState.getReExecState(); - case DB: - return apiUsageState.getDbStorageState(); - case JS: - return apiUsageState.getJsExecState(); - case TBEL: - return apiUsageState.getTbelExecState(); - case EMAIL: - return apiUsageState.getEmailExecState(); - case SMS: - return apiUsageState.getSmsExecState(); - case ALARM: - return apiUsageState.getAlarmExecState(); - default: - return ApiUsageStateValue.ENABLED; - } + return switch (feature) { + case TRANSPORT -> apiUsageState.getTransportState(); + case RE -> apiUsageState.getReExecState(); + case DB -> apiUsageState.getDbStorageState(); + case JS -> apiUsageState.getJsExecState(); + case TBEL -> apiUsageState.getTbelExecState(); + case EMAIL -> apiUsageState.getEmailExecState(); + case SMS -> apiUsageState.getSmsExecState(); + case ALARM -> apiUsageState.getAlarmExecState(); + case EDGE -> apiUsageState.getEdgeState(); + }; } public boolean setFeatureValue(ApiFeature feature, ApiUsageStateValue value) { ApiUsageStateValue currentValue = getFeatureValue(feature); switch (feature) { - case TRANSPORT: - apiUsageState.setTransportState(value); - break; - case RE: - apiUsageState.setReExecState(value); - break; - case DB: - apiUsageState.setDbStorageState(value); - break; - 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; + case TRANSPORT -> apiUsageState.setTransportState(value); + case RE -> apiUsageState.setReExecState(value); + case DB -> apiUsageState.setDbStorageState(value); + case JS -> apiUsageState.setJsExecState(value); + case TBEL -> apiUsageState.setTbelExecState(value); + case EMAIL -> apiUsageState.setEmailExecState(value); + case SMS -> apiUsageState.setSmsExecState(value); + case ALARM -> apiUsageState.setAlarmExecState(value); + case EDGE -> apiUsageState.setEdgeState(value); } return !currentValue.equals(value); } diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java index 0f16610c6f..be88335a92 100644 --- a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java @@ -59,7 +59,6 @@ import org.thingsboard.server.dao.tenant.TbTenantProfileCache; import org.thingsboard.server.dao.tenant.TenantService; import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.dao.usagerecord.ApiUsageStateService; -import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg; import org.thingsboard.server.gen.transport.TransportProtos.UsageStatsKVProto; import org.thingsboard.server.queue.common.TbProtoQueueMsg; @@ -149,27 +148,7 @@ public class DefaultTbApiUsageStateService extends AbstractPartitionBasedService public void process(TbProtoQueueMsg msgPack, TbCallback callback) { ToUsageStatsServiceMsg serviceMsg = msgPack.getValue(); String serviceId = serviceMsg.getServiceId(); - - List 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 -> { + serviceMsg.getMsgsList().forEach(msg -> { TenantId tenantId = TenantId.fromUUID(new UUID(msg.getTenantIdMSB(), msg.getTenantIdLSB())); EntityId ownerId; if (msg.getCustomerIdMSB() != 0 && msg.getCustomerIdLSB() != 0) { @@ -184,7 +163,9 @@ public class DefaultTbApiUsageStateService extends AbstractPartitionBasedService } private void processEntityUsageStats(TenantId tenantId, EntityId ownerId, List values, String serviceId) { - if (deletedEntities.contains(ownerId)) return; + if (deletedEntities.contains(ownerId)) { + return; + } BaseApiUsageState usageState; List updatedEntries; @@ -205,14 +186,7 @@ public class DefaultTbApiUsageStateService extends AbstractPartitionBasedService updatedEntries = new ArrayList<>(ApiUsageRecordKey.values().length); Set apiFeatures = new HashSet<>(); for (UsageStatsKVProto statsItem : values) { - ApiUsageRecordKey recordKey; - - //For backward compatibility, remove after release - if (StringUtils.isNotEmpty(statsItem.getKey())) { - recordKey = ApiUsageRecordKey.valueOf(statsItem.getKey()); - } else { - recordKey = ProtoUtils.fromProto(statsItem.getRecordKey()); - } + ApiUsageRecordKey recordKey = ProtoUtils.fromProto(statsItem.getRecordKey()); StatsCalculationResult calculationResult = usageState.calculate(recordKey, statsItem.getValue(), serviceId); if (calculationResult.isValueChanged()) { @@ -598,4 +572,5 @@ public class DefaultTbApiUsageStateService extends AbstractPartitionBasedService private void destroy() { super.stop(); } + } diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/TbApiUsageStateService.java b/application/src/main/java/org/thingsboard/server/service/apiusage/TbApiUsageStateService.java index 5a81bbb78d..d647c6553e 100644 --- a/application/src/main/java/org/thingsboard/server/service/apiusage/TbApiUsageStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/apiusage/TbApiUsageStateService.java @@ -39,4 +39,5 @@ public interface TbApiUsageStateService extends TbApiUsageStateClient, RuleEngin void onCustomerDelete(CustomerId customerId); void onApiUsageStateUpdate(TenantId tenantId); + } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java index e9ba942afb..ae58ea1fa4 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java @@ -15,8 +15,8 @@ */ 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.SettableFuture; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; 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.EdgeStatsKey; 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.discovery.TopicService; import org.thingsboard.server.queue.provider.TbQueueProducerProvider; @@ -51,10 +53,23 @@ public class KafkaEdgeEventService extends BaseEdgeEventService { TopicPartitionInfo tpi = topicService.getEdgeEventNotificationsTopic(edgeEvent.getTenantId(), edgeEvent.getEdgeId()); ToEdgeEventNotificationMsg msg = ToEdgeEventNotificationMsg.newBuilder().setEdgeEventMsg(ProtoUtils.toProto(edgeEvent)).build(); - producerProvider.getTbEdgeEventsMsgProducer().send(tpi, new TbProtoQueueMsg<>(UUID.randomUUID(), msg), null); + SettableFuture 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.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_ADDED, edgeEvent.getTenantId(), edgeEvent.getEdgeId(), 1)); - return Futures.immediateFuture(null); + return result; } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/service/EdgeGrpcService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/service/EdgeGrpcService.java index b50d643634..5f6d1d93bc 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/service/EdgeGrpcService.java +++ b/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; import com.fasterxml.jackson.databind.node.ObjectNode; +import com.google.common.util.concurrent.Striped; +import io.grpc.Status; import io.grpc.stub.StreamObserver; import jakarta.annotation.PostConstruct; 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.context.ApplicationContext; import org.springframework.context.annotation.Lazy; +import org.springframework.context.event.EventListener; import org.springframework.stereotype.Service; import org.thingsboard.common.util.JacksonUtil; 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.TimeseriesSaveRequest; import org.thingsboard.server.cache.TbTransactionalCache; 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.DataConstants; +import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.EdgeEvent; 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.msg.TbMsgType; 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.TbMsgDataType; 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.FromEdgeSyncResponse; 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.RequestMsg; 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.telemetry.TelemetrySubscriptionService; +import java.util.List; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ExecutorService; import java.util.concurrent.ScheduledExecutorService; 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.IntConsumer; import static org.thingsboard.server.service.state.DefaultDeviceStateService.ACTIVITY_STATE; 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 TelemetrySubscriptionService tsSubService; private final TbTransactionalCache edgeIdServiceIdCache; + private final TbApiUsageStateClient apiUsageStateClient; + private final Striped tenantLocks = Striped.lock(64); private final ConcurrentMap> localSyncEdgeRequests = new ConcurrentHashMap<>(); private ScheduledExecutorService executorService; private ScheduledExecutorService sendDownlinkExecutorService; @@ -117,6 +132,46 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i 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 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 public StreamObserver handleMsgs(StreamObserver outputStream) { EdgeGrpcSessionManager sessionManager = applicationContext.getBean(EdgeGrpcSessionManager.class); @@ -195,17 +250,27 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i EdgeSessionState state = edgeSession.getState(); Edge edge = state.getEdge(); TenantId tenantId = state.getTenantId(); - log.info("[{}][{}] edge [{}] connected successfully.", tenantId, state.getSessionId(), edgeId); - if (sessions.hasByEdgeId(edgeId)) { - EdgeGrpcSessionManager existingSession = sessions.getByEdgeId(edgeId); - if (existingSession != null) { - UUID sessionId = existingSession.getState().getSessionId(); - log.info("[{}][{}] Replacing existing session [{}] for edge [{}]", tenantId, state.getSessionId(), sessionId, edgeId); - existingSession.destroyAndMarkAsZombieIfFailed(); - sessions.removeBySessionId(sessionId); + EdgeGrpcSessionManager replaced; + Lock lock = tenantLocks.get(tenantId); + lock.lock(); + try { + if (!isEdgeEnabled(tenantId)) { + throw new EdgeFeatureDisabledException("Edge feature disabled due to API limits"); + } + 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); long lastConnectTs = System.currentTimeMillis(); save(tenantId, edgeId, LAST_CONNECT_TIME, lastConnectTs); @@ -389,4 +454,10 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i e.shutdown(); } } + + private boolean isEdgeEnabled(TenantId tenantId) { + ApiUsageState state = apiUsageStateClient.getApiUsageState(tenantId); + return state == null || state.isEdgeEnabled(); + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeGrpcSession.java index 73a3b53e4c..21b452e506 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeGrpcSession.java +++ b/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.ListenableFuture; import com.google.common.util.concurrent.SettableFuture; +import io.grpc.Status; import io.grpc.stub.StreamObserver; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.checkerframework.checker.nullness.qual.Nullable; import org.springframework.data.util.Pair; +import org.thingsboard.edge.exception.EdgeFeatureDisabledException; import org.thingsboard.rule.engine.api.AttributesSaveRequest; import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.DataConstants; @@ -74,6 +76,7 @@ import java.util.UUID; import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; import java.util.function.BiConsumer; @@ -99,6 +102,7 @@ public class EdgeGrpcSession implements EdgeSession { private final EdgeSessionState state = new EdgeSessionState(); private final Lock downlinkMsgLock = new ReentrantLock(); private final ConcurrentLinkedQueue highPriorityQueue = new ConcurrentLinkedQueue<>(); + private final AtomicBoolean closed = new AtomicBoolean(false); private int clientMaxInboundMessageSize; @@ -273,8 +277,26 @@ public class EdgeGrpcSession implements EdgeSession { 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 public void close() { + if (!closed.compareAndSet(false, true)) { + log.debug("[{}][{}] close skipped — session already closed", getTenantId(), getSessionId()); + return; + } log.debug("[{}][{}] Closing session", getTenantId(), getSessionId()); state.setConnected(false); try { @@ -421,7 +443,7 @@ public class EdgeGrpcSession implements EdgeSession { if (state.isConnected() && pageData.hasNext()) { fetchAndSendEdgeEvents(fetcher, pageLink.nextPageLink(), result); } else { - EdgeEvent latestEdgeEvent = pageData.getData().get(pageData.getData().size() - 1); + EdgeEvent latestEdgeEvent = pageData.getData().getLast(); UUID idOffset = latestEdgeEvent.getUuidId(); if (idOffset != null) { Long newStartTs = Uuids.unixTimestamp(idOffset); @@ -457,7 +479,7 @@ public class EdgeGrpcSession implements EdgeSession { stopCurrentSendDownlinkMsgsTask(true); return; } - if (!state.getPendingMsgsMap().values().isEmpty()) { + if (!state.getPendingMsgsMap().isEmpty()) { Edge edge = state.getEdge(); List copy = new ArrayList<>(state.getPendingMsgsMap().values()); if (attempt > 1) { @@ -551,45 +573,54 @@ public class EdgeGrpcSession implements EdgeSession { private ConnectResponseMsg processConnect(ConnectRequestMsg request) { log.trace("[{}] processConnect [{}]", getSessionId(), request); Optional optional = ctx.getEdgeService().findEdgeByRoutingKey(TenantId.SYS_TENANT_ID, request.getEdgeRoutingKey()); - if (optional.isPresent()) { - Edge edge = optional.get(); - TenantId tenantId = edge.getTenantId(); - state.setEdge(edge); - try { - if (edge.getSecret().equals(request.getEdgeSecret())) { - 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(); - } - 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(); + if (optional.isEmpty()) { + return buildErrorResponse(ConnectResponseCode.BAD_CREDENTIALS, "Failed to find the edge! Routing key: " + request.getEdgeRoutingKey()); + } + Edge edge = optional.get(); + TenantId tenantId = edge.getTenantId(); + try { + if (!edge.getSecret().equals(request.getEdgeSecret())) { + String failureMsg = "Failed to validate the edge! Provided request secret: " + request.getEdgeSecret(); + processConnectionFailure(edge, failureMsg, "Failed to validate the edge!"); + return buildErrorResponse(ConnectResponseCode.BAD_CREDENTIALS, failureMsg); } + 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() - .setResponseCode(ConnectResponseCode.BAD_CREDENTIALS) - .setErrorMsg("Failed to find the edge! Routing key: " + request.getEdgeRoutingKey()) - .setConfiguration(EdgeConfiguration.getDefaultInstance()).build(); + .setResponseCode(code) + .setErrorMsg(errorMsg) + .setConfiguration(EdgeConfiguration.getDefaultInstance()) + .build(); } private void processSaveEdgeVersionAsAttribute(String edgeVersion) { @@ -618,4 +649,5 @@ public class EdgeGrpcSession implements EdgeSession { private UUID getSessionId() { return state.getSessionId(); } + } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeGrpcSessionDelegate.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeGrpcSessionDelegate.java index 9c8c5eddb9..f1a17f9cad 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeGrpcSessionDelegate.java +++ b/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; +import io.grpc.Status; import org.thingsboard.server.common.data.edge.EdgeEvent; 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) { getSession().startSyncProcess(fullSync); } + + @Override + public void closeWithError(Status status, String errorMsg) { + getSession().closeWithError(status, errorMsg); + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeSession.java index 0b7158211d..91d59b6a3c 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeSession.java +++ b/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; import com.google.common.util.concurrent.ListenableFuture; +import io.grpc.Status; import io.grpc.stub.StreamObserver; import org.springframework.data.util.Pair; import org.thingsboard.server.common.data.edge.EdgeEvent; @@ -31,13 +32,23 @@ import java.util.List; public interface EdgeSession extends Closeable { StreamObserver initInputStream(); + EdgeSessionState getState(); + void startSyncProcess(boolean fullSync); + void sendDownlinkMsg(ResponseMsg responseMsg); + void addHighPriorityEvent(EdgeEvent edgeEvent); + void processHighPriorityEvents(); + boolean hasHighPriorityEvents(); + ListenableFuture> fetchAndSendEdgeEvents(EdgeEventFetcher fetcher); + ListenableFuture sendDownlinkMsgsPack(List downlinkMsgsPack); + void closeWithError(Status status, String errorMsg); + } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeSessionsHolder.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeSessionsHolder.java index 2d7c3dd8e3..7ec1eb9daa 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/EdgeSessionsHolder.java +++ b/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.stereotype.Component; 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.service.edge.rpc.EdgeSessionState; import org.thingsboard.server.service.edge.rpc.session.manager.EdgeGrpcSessionManager; import java.util.HashSet; +import java.util.List; import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; @@ -71,6 +73,10 @@ public class EdgeSessionsHolder { return sessionsById.remove(sessionId); } + public List getByTenantId(TenantId tenantId) { + return sessions.values().stream().filter(s -> tenantId.equals(s.getState().getTenantId())).toList(); + } + public void remove(EdgeGrpcSessionManager session) { if (session == null) { log.warn("Can't remove session from holder because it's null"); @@ -80,4 +86,5 @@ public class EdgeSessionsHolder { removeByEdgeId(sessionState.getEdgeId()); removeBySessionId(sessionState.getSessionId()); } + } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/AbstractEdgeGrpcSessionManager.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/AbstractEdgeGrpcSessionManager.java index 0df1af1ada..81a0c0921f 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/AbstractEdgeGrpcSessionManager.java +++ b/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(); state.setEdge(edge); if (stateCustomerId != null && !stateCustomerId.equals(edge.getCustomerId())) { - // do not send edge configuration message on customer update - // message send by separate flow from assign_to or unassing_from customer + // do not send an edge configuration message on a customer update + // message send by separate flow from assign_to or unassign_from customer return; } EdgeUpdateMsg edgeConfig = EdgeUpdateMsg.newBuilder() @@ -146,4 +146,5 @@ public abstract class AbstractEdgeGrpcSessionManager extends EdgeGrpcSessionDele } } } + } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/EdgeGrpcSessionManager.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/EdgeGrpcSessionManager.java index 9bf51107e4..b57a3a89df 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/EdgeGrpcSessionManager.java +++ b/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; +import io.grpc.Status; import io.grpc.stub.StreamObserver; import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.EdgeEvent; @@ -41,6 +42,7 @@ public interface EdgeGrpcSessionManager { void onConfigurationUpdate(Edge edge); void onEdgeDisconnect(); void onEdgeRemoval(); + void closeWithError(Status status, String errorMsg); void destroyAndMarkAsZombieIfFailed(); boolean destroy(); diff --git a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java index 48bd497c6e..4340b751a3 100644 --- a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java +++ b/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.VALID; -/** - * Created by ashvayka on 05.10.18. - */ @Slf4j @Service @TbCoreComponent diff --git a/application/src/test/java/org/thingsboard/server/edge/EdgeApiUsageDisabledTest.java b/application/src/test/java/org/thingsboard/server/edge/EdgeApiUsageDisabledTest.java new file mode 100644 index 0000000000..ee36b63cdc --- /dev/null +++ b/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()); + } + +} diff --git a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java b/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java index 8208dc4fc9..59e4f503ab 100644 --- a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java +++ b/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.gen.edge.v1.AdminSettingsUpdateMsg; 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.AlarmUpdateMsg; +import org.thingsboard.server.gen.edge.v1.ApiKeyUpdateMsg; import org.thingsboard.server.gen.edge.v1.AssetProfileUpdateMsg; import org.thingsboard.server.gen.edge.v1.AssetUpdateMsg; import org.thingsboard.server.gen.edge.v1.CalculatedFieldUpdateMsg; @@ -87,6 +87,7 @@ import java.util.stream.Collectors; public class EdgeImitator { private static final int MAX_DOWNLINK_FAILS = 2; + private final String routingKey; private final String routingSecret; @@ -106,6 +107,8 @@ public class EdgeImitator { @Getter private EdgeConfiguration configuration; + @Getter + private volatile Exception closeException; private final ConcurrentLinkedDeque downlinkMsgs; //Returns collection copy as Unmodifiable list @@ -195,6 +198,9 @@ public class EdgeImitator { private void onClose(Exception e) { log.info("onClose: {}", e.getMessage()); + if (this.closeException == null) { + this.closeException = e; + } } private ListenableFuture> processDownlinkMsg(DownlinkMsg downlinkMsg) { diff --git a/application/src/test/java/org/thingsboard/server/service/apiusage/ApiUsageTest.java b/application/src/test/java/org/thingsboard/server/service/apiusage/ApiUsageTest.java index 4cd41c314e..b250cd04a7 100644 --- a/application/src/test/java/org/thingsboard/server/service/apiusage/ApiUsageTest.java +++ b/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 { - private Tenant savedTenant; - private User tenantAdmin; - private static final int MAX_DP_ENABLE_VALUE = 12; 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; + @Autowired private ApiUsageStateService apiUsageStateService; @Autowired @@ -76,16 +75,16 @@ public class ApiUsageTest extends AbstractControllerTest { Tenant tenant = new Tenant(); tenant.setTitle("My tenant"); tenant.setTenantProfileId(savedTenantProfile.getId()); - savedTenant = saveTenant(tenant); + Tenant savedTenant = saveTenant(tenant); tenantId = savedTenant.getId(); assertNotNull(savedTenant); - tenantAdmin = new User(); + User tenantAdmin = new User(); tenantAdmin.setAuthority(Authority.TENANT_ADMIN); tenantAdmin.setTenantId(savedTenant.getId()); tenantAdmin.setEmail("tenant2@thingsboard.org"); - tenantAdmin = createUserAndLogin(tenantAdmin, "testPassword1"); + createUserAndLogin(tenantAdmin, "testPassword1"); } @Test @@ -137,6 +136,25 @@ public class ApiUsageTest extends AbstractControllerTest { 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() { return apiUsageStateService.findTenantApiUsageState(tenantId); } @@ -151,6 +169,7 @@ public class ApiUsageTest extends AbstractControllerTest { .maxDPStorageDays(MAX_DP_ENABLE_VALUE) .maxSms(MAX_SMS_ENABLE_VALUE) .smsEnabled(true) + .maxEdgeEvents(MAX_EDGE_ENABLE_VALUE) .warnThreshold(WARN_THRESHOLD_VALUE) .build(); diff --git a/application/src/test/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateServiceTest.java b/application/src/test/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateServiceTest.java index a358c0f12c..e3cb1d77f0 100644 --- a/application/src/test/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateServiceTest.java +++ b/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.TenantProfileData; 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.dao.service.DaoSqlTest; import org.thingsboard.server.dao.usagerecord.ApiUsageStateService; @@ -64,7 +65,6 @@ public class DefaultTbApiUsageStateServiceTest extends AbstractControllerTest { private ApiUsageStateService apiUsageStateService; private TenantId tenantId; - private Tenant savedTenant; private TenantProfile savedTenantProfile; private static final int MAX_ENABLE_VALUE = 5000; @@ -83,48 +83,20 @@ public class DefaultTbApiUsageStateServiceTest extends AbstractControllerTest { Tenant tenant = new Tenant(); tenant.setTitle("My tenant"); tenant.setTenantProfileId(savedTenantProfile.getId()); - savedTenant = saveTenant(tenant); + Tenant savedTenant = saveTenant(tenant); tenantId = savedTenant.getId(); Assert.assertNotNull(savedTenant); } @Test public void testProcess_transitionFromWarningToDisabled() { - TransportProtos.ToUsageStatsServiceMsg.Builder warningMsgBuilder = TransportProtos.ToUsageStatsServiceMsg.newBuilder() - .setTenantIdMSB(tenantId.getId().getMostSignificantBits()) - .setTenantIdLSB(tenantId.getId().getLeastSignificantBits()) - .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 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 disableMsg = new TbProtoQueueMsg<>(UUID.randomUUID(), disableStatsMsg); + sendUsageStats(tenantId, ApiUsageRecordKey.STORAGE_DP_COUNT, VALUE_WARNING); + await().atMost(5, TimeUnit.SECONDS).until(() -> + apiUsageStateService.findTenantApiUsageState(tenantId).getDbStorageState() == ApiUsageStateValue.WARNING); - service.process(disableMsg, TbCallback.EMPTY); - assertEquals(ApiUsageStateValue.DISABLED, apiUsageStateService.findTenantApiUsageState(tenantId).getDbStorageState()); + sendUsageStats(tenantId, ApiUsageRecordKey.STORAGE_DP_COUNT, VALUE_DISABLE); + await().atMost(5, TimeUnit.SECONDS).until(() -> + apiUsageStateService.findTenantApiUsageState(tenantId).getDbStorageState() == ApiUsageStateValue.DISABLED); } @Test @@ -138,6 +110,7 @@ public class DefaultTbApiUsageStateServiceTest extends AbstractControllerTest { apiUsageState.setTransportState(ApiUsageStateValue.ENABLED); apiUsageState.setEmailExecState(ApiUsageStateValue.ENABLED); apiUsageState.setJsExecState(ApiUsageStateValue.ENABLED); + apiUsageState.setEdgeState(ApiUsageStateValue.ENABLED); apiUsageState.setTenantId(tenantId); apiUsageState.setEntityId(tenantId); @@ -194,7 +167,7 @@ public class DefaultTbApiUsageStateServiceTest extends AbstractControllerTest { await().atMost(5, TimeUnit.SECONDS).until(() -> { Optional 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 @@ -216,10 +189,10 @@ public class DefaultTbApiUsageStateServiceTest extends AbstractControllerTest { await().atMost(5, TimeUnit.SECONDS).until(() -> { Optional 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() .smsEnabled(false) .build(); @@ -237,10 +210,98 @@ public class DefaultTbApiUsageStateServiceTest extends AbstractControllerTest { await().atMost(5, TimeUnit.SECONDS).until(() -> { Optional 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() { TenantProfile tenantProfile = new TenantProfile(); tenantProfile.setName("Tenant Profile"); @@ -257,4 +318,4 @@ public class DefaultTbApiUsageStateServiceTest extends AbstractControllerTest { return tenantProfile; } -} \ No newline at end of file +} diff --git a/application/src/test/java/org/thingsboard/server/transport/coap/CoapTransportFeatureDisabledTest.java b/application/src/test/java/org/thingsboard/server/transport/coap/CoapTransportFeatureDisabledTest.java new file mode 100644 index 0000000000..42064f6b21 --- /dev/null +++ b/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(); + } + +} diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/MqttTransportFeatureDisabledTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/MqttTransportFeatureDisabledTest.java new file mode 100644 index 0000000000..41f7cb611a --- /dev/null +++ b/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) {} + } + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/ApiFeature.java b/common/data/src/main/java/org/thingsboard/server/common/data/ApiFeature.java index 335b9efd16..2723f217c6 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/ApiFeature.java +++ b/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; public enum ApiFeature { + TRANSPORT("transportApiState", "Device API"), DB("dbApiState", "Telemetry persistence"), RE("ruleEngineApiState", "Rule Engine execution"), @@ -25,7 +26,8 @@ public enum ApiFeature { TBEL("tbelExecutionApiState", "Tbel functions execution"), EMAIL("emailApiState", "Email messages"), SMS("smsApiState", "SMS messages"), - ALARM("alarmApiState", "Alarms"); + ALARM("alarmApiState", "Alarms"), + EDGE("edgeApiState", "Edge"); @Getter private final String apiStateKey; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageRecordKey.java b/common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageRecordKey.java index ceb657374b..d954778b33 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageRecordKey.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageRecordKey.java @@ -28,6 +28,7 @@ public enum ApiUsageRecordKey { EMAIL_EXEC_COUNT(ApiFeature.EMAIL, "emailCount", "emailLimit", "email message", true, true), SMS_EXEC_COUNT(ApiFeature.SMS, "smsCount", "smsLimit", "SMS message", true, true), CREATED_ALARMS_COUNT(ApiFeature.ALARM, "createdAlarmsCount", "createdAlarmsLimit", "alarm"), + EDGE_EVENT_COUNT(ApiFeature.EDGE, "edgeEventCount", "edgeEventLimit", "edge event"), ACTIVE_DEVICES("activeDevicesCount"), INACTIVE_DEVICES("inactiveDevicesCount"); @@ -39,6 +40,7 @@ public enum ApiUsageRecordKey { private static final ApiUsageRecordKey[] EMAIL_RECORD_KEYS = {EMAIL_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[] EDGE_RECORD_KEYS = {EDGE_EVENT_COUNT}; @Getter private final ApiFeature apiFeature; @@ -71,26 +73,17 @@ public enum ApiUsageRecordKey { } public static ApiUsageRecordKey[] getKeys(ApiFeature feature) { - switch (feature) { - case TRANSPORT: - return TRANSPORT_RECORD_KEYS; - case DB: - return DB_RECORD_KEYS; - case RE: - return RE_RECORD_KEYS; - case JS: - return JS_RECORD_KEYS; - case TBEL: - 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[]{}; - } + return switch (feature) { + case TRANSPORT -> TRANSPORT_RECORD_KEYS; + case DB -> DB_RECORD_KEYS; + case RE -> RE_RECORD_KEYS; + case JS -> JS_RECORD_KEYS; + case TBEL -> TBEL_RECORD_KEYS; + case EMAIL -> EMAIL_RECORD_KEYS; + case SMS -> SMS_RECORD_KEYS; + case ALARM -> ALARM_RECORD_KEYS; + case EDGE -> EDGE_RECORD_KEYS; + }; } } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageState.java b/common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageState.java index 8c1e1ef85a..8b8b3a9b1d 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageState.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageState.java @@ -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.TenantId; +import java.io.Serial; + @ToString @EqualsAndHashCode(callSuper = true) @Getter @Setter public class ApiUsageState extends BaseData implements HasTenantId, HasVersion { + @Serial private static final long serialVersionUID = 8250339805336035966L; private TenantId tenantId; @@ -41,6 +44,7 @@ public class ApiUsageState extends BaseData implements HasTenan private ApiUsageStateValue emailExecState; private ApiUsageStateValue smsExecState; private ApiUsageStateValue alarmExecState; + private ApiUsageStateValue edgeState; private Long version; public ApiUsageState() { @@ -63,39 +67,44 @@ public class ApiUsageState extends BaseData implements HasTenan this.emailExecState = ur.getEmailExecState(); this.smsExecState = ur.getSmsExecState(); this.alarmExecState = ur.getAlarmExecState(); + this.edgeState = ur.getEdgeState(); this.version = ur.getVersion(); } public boolean isTransportEnabled() { - return !ApiUsageStateValue.DISABLED.equals(transportState); + return transportState != ApiUsageStateValue.DISABLED; } public boolean isReExecEnabled() { - return !ApiUsageStateValue.DISABLED.equals(reExecState); + return reExecState != ApiUsageStateValue.DISABLED; } public boolean isDbStorageEnabled() { - return !ApiUsageStateValue.DISABLED.equals(dbStorageState); + return dbStorageState != ApiUsageStateValue.DISABLED; } public boolean isJsExecEnabled() { - return !ApiUsageStateValue.DISABLED.equals(jsExecState); + return jsExecState != ApiUsageStateValue.DISABLED; } public boolean isTbelExecEnabled() { - return !ApiUsageStateValue.DISABLED.equals(tbelExecState); + return tbelExecState != ApiUsageStateValue.DISABLED; } public boolean isEmailSendEnabled() { - return !ApiUsageStateValue.DISABLED.equals(emailExecState); + return emailExecState != ApiUsageStateValue.DISABLED; } public boolean isSmsSendEnabled() { - return !ApiUsageStateValue.DISABLED.equals(smsExecState); + return smsExecState != ApiUsageStateValue.DISABLED; } public boolean isAlarmCreationEnabled() { return alarmExecState != ApiUsageStateValue.DISABLED; } + public boolean isEdgeEnabled() { + return edgeState != ApiUsageStateValue.DISABLED; + } + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageStateValue.java b/common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageStateValue.java index aa639f9c14..78b9d5db9e 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageStateValue.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageStateValue.java @@ -19,8 +19,8 @@ public enum ApiUsageStateValue { ENABLED, WARNING, DISABLED; - public static ApiUsageStateValue toMoreRestricted(ApiUsageStateValue a, ApiUsageStateValue b) { return a.ordinal() > b.ordinal() ? a : b; } + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/edqs/fields/ApiUsageStateFields.java b/common/data/src/main/java/org/thingsboard/server/common/data/edqs/fields/ApiUsageStateFields.java index b3d8492ef3..e8c5ae05f7 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/edqs/fields/ApiUsageStateFields.java +++ b/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 smsExecState; private ApiUsageStateValue alarmExecState; + private ApiUsageStateValue edgeState; public ApiUsageStateFields(UUID id, long createdTime, UUID tenantId, UUID entityId, String entityType, ApiUsageStateValue transportState, ApiUsageStateValue dbStorageState, ApiUsageStateValue reExecState, ApiUsageStateValue jsExecState, ApiUsageStateValue tbelExecState, ApiUsageStateValue emailExecState, ApiUsageStateValue smsExecState, ApiUsageStateValue alarmExecState, - Long version) { + ApiUsageStateValue edgeState, Long version) { super(id, createdTime, tenantId, null, null, version); this.entityId = (entityType != null && entityId != null) ? EntityIdFactory.getByTypeAndUuid(entityType, entityId) : null; this.transportState = transportState; @@ -55,5 +56,7 @@ public class ApiUsageStateFields extends AbstractEntityFields { this.emailExecState = emailExecState; this.smsExecState = smsExecState; this.alarmExecState = alarmExecState; + this.edgeState = edgeState; } + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/edqs/fields/FieldsUtil.java b/common/data/src/main/java/org/thingsboard/server/common/data/edqs/fields/FieldsUtil.java index ab5a6856a1..dd5ba9107c 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/edqs/fields/FieldsUtil.java +++ b/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 static EntityFields toFields(Object entity) { - if (entity instanceof Customer customer) { - return toFields(customer); - } else if (entity instanceof Tenant tenant) { - return toFields(tenant); - } else if (entity instanceof TenantProfile tenantProfile) { - return toFields(tenantProfile); - } else if (entity instanceof Device device) { - return toFields(device); - } else if (entity instanceof Asset asset) { - return toFields(asset); - } else if (entity instanceof Edge edge) { - return toFields(edge); - } else if (entity instanceof EntityView entityView) { - return toFields(entityView); - } else if (entity instanceof User user) { - return toFields(user); - } else if (entity instanceof Dashboard dashboard) { - return toFields(dashboard); - } else if (entity instanceof RuleChain ruleChain) { - 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()); - } + return switch (entity) { + case Customer customer -> toFields(customer); + case Tenant tenant -> toFields(tenant); + case TenantProfile tenantProfile -> toFields(tenantProfile); + case Device device -> toFields(device); + case Asset asset -> toFields(asset); + case Edge edge -> toFields(edge); + case EntityView entityView -> toFields(entityView); + case User user -> toFields(user); + case Dashboard dashboard -> toFields(dashboard); + case RuleChain ruleChain -> toFields(ruleChain); + case RuleNode ruleNode -> toFields(ruleNode); + case WidgetType widgetType -> toFields(widgetType); + case WidgetsBundle widgetsBundle -> toFields(widgetsBundle); + case DeviceProfile deviceProfile -> toFields(deviceProfile); + case AssetProfile assetProfile -> toFields(assetProfile); + case QueueStats queueStats -> toFields(queueStats); + case ApiUsageState apiUsageState -> toFields(apiUsageState); + default -> throw new IllegalArgumentException("Unsupported entity type: " + entity.getClass().getName()); + }; } private static CustomerFields toFields(Customer entity) { @@ -284,6 +267,7 @@ public class FieldsUtil { .emailExecState(entity.getEmailExecState()) .smsExecState(entity.getSmsExecState()) .alarmExecState(entity.getAlarmExecState()) + .edgeState(entity.getEdgeState()) .version(entity.getVersion()) .build(); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java index 297f92cffb..f8ee6c7fdd 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java +++ b/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; @Schema(example = "1000") private long maxCreatedAlarms; + @Schema(example = "10000000") + private long maxEdgeEvents; @RateLimit(fieldName = "REST requests for tenant") private String tenantServerRestLimitsConfiguration; @@ -218,6 +220,7 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura case EMAIL_EXEC_COUNT -> maxEmails; case SMS_EXEC_COUNT -> maxSms; case CREATED_ALARMS_COUNT -> maxCreatedAlarms; + case EDGE_EVENT_COUNT -> maxEdgeEvents; default -> 0L; }; } @@ -226,7 +229,7 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura public boolean getProfileFeatureEnabled(ApiUsageRecordKey key) { switch (key) { case SMS_EXEC_COUNT: - return smsEnabled == null || Boolean.TRUE.equals(smsEnabled); + return smsEnabled == null || smsEnabled; default: return true; } diff --git a/common/edge-api/src/main/java/org/thingsboard/edge/exception/EdgeConnectionException.java b/common/edge-api/src/main/java/org/thingsboard/edge/exception/EdgeConnectionException.java index 447b6f1f8d..d522da38b6 100644 --- a/common/edge-api/src/main/java/org/thingsboard/edge/exception/EdgeConnectionException.java +++ b/common/edge-api/src/main/java/org/thingsboard/edge/exception/EdgeConnectionException.java @@ -15,15 +15,15 @@ */ package org.thingsboard.edge.exception; +import java.io.Serial; + public class EdgeConnectionException extends RuntimeException { + @Serial private static final long serialVersionUID = -4372754681230555723L; public EdgeConnectionException(String message) { super(message); } - public EdgeConnectionException(String message, Throwable cause) { - super(message, cause); - } -} \ No newline at end of file +} diff --git a/common/edge-api/src/main/java/org/thingsboard/edge/exception/EdgeFeatureDisabledException.java b/common/edge-api/src/main/java/org/thingsboard/edge/exception/EdgeFeatureDisabledException.java new file mode 100644 index 0000000000..410295cc1a --- /dev/null +++ b/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); + } + +} diff --git a/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java b/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java index 7e58a2cb9f..1fae79cb28 100644 --- a/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java +++ b/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.stereotype.Service; 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.StringUtils; import org.thingsboard.server.gen.edge.v1.ConnectRequestMsg; @@ -170,7 +171,11 @@ public class EdgeGrpcClient implements EdgeRpcClient { } catch (InterruptedException 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()) { log.debug("[{}] Edge update message received {}", edgeKey, responseMsg.getEdgeUpdateMsg()); diff --git a/common/edge-api/src/main/proto/edge.proto b/common/edge-api/src/main/proto/edge.proto index bcea3c30d9..cae7d75f57 100644 --- a/common/edge-api/src/main/proto/edge.proto +++ b/common/edge-api/src/main/proto/edge.proto @@ -97,6 +97,7 @@ enum ConnectResponseCode { ACCEPTED = 0; BAD_CREDENTIALS = 1; SERVER_UNAVAILABLE = 2; + FEATURE_DISABLED = 3; } message ConnectResponseMsg { diff --git a/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java b/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java index 784bd97c79..7e45eb2b9d 100644 --- a/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java +++ b/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 ACTIVE_DEVICES -> ApiUsageRecordKeyProto.ACTIVE_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 ACTIVE_DEVICES -> ApiUsageRecordKey.ACTIVE_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()) .setSmsExecState(apiUsageState.getSmsExecState().name()) .setAlarmExecState(apiUsageState.getAlarmExecState().name()) + .setEdgeState(apiUsageState.getEdgeState().name()) .setVersion(apiUsageState.getVersion()) .build(); } @@ -1215,6 +1218,12 @@ public class ProtoUtils { apiUsageState.setEmailExecState(ApiUsageStateValue.valueOf(proto.getEmailExecState())); apiUsageState.setSmsExecState(ApiUsageStateValue.valueOf(proto.getSmsExecState())); 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()); return apiUsageState; } @@ -1309,18 +1318,13 @@ public class ProtoUtils { public static TransportProtos.EntityUpdateMsg toEntityUpdateProto(T entity) { var builder = TransportProtos.EntityUpdateMsg.newBuilder(); - if (entity instanceof Device) { - builder.setDevice(toProto((Device) entity)); - } else if (entity instanceof DeviceProfile) { - builder.setDeviceProfile(toProto((DeviceProfile) entity)); - } else if (entity instanceof Tenant) { - builder.setTenant(toProto((Tenant) entity)); - } else if (entity instanceof TenantProfile) { - 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()); + switch (entity) { + case Device device -> builder.setDevice(toProto(device)); + case DeviceProfile deviceProfile -> builder.setDeviceProfile(toProto(deviceProfile)); + case Tenant tenant -> builder.setTenant(toProto(tenant)); + case TenantProfile tenantProfile -> builder.setTenantProfile(toProto(tenantProfile)); + case ApiUsageState apiUsageState -> builder.setApiUsageState(toProto(apiUsageState)); + default -> log.warn("[{}] entity does not support toProto serialization .", entity.getClass().getSimpleName()); } return builder.build(); } diff --git a/common/proto/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto index 68f7052c69..297f062ded 100644 --- a/common/proto/src/main/proto/queue.proto +++ b/common/proto/src/main/proto/queue.proto @@ -81,6 +81,7 @@ enum ApiUsageRecordKeyProto { CREATED_ALARMS_COUNT = 8; ACTIVE_DEVICES = 9; INACTIVE_DEVICES = 10; + EDGE_EVENT_COUNT = 11; } message EntityIdProto { @@ -387,6 +388,7 @@ message ApiUsageStateProto { string smsExecState = 15; string alarmExecState = 16; int64 version = 17; + string edgeState = 18; } message RepositorySettingsProto { @@ -1783,19 +1785,15 @@ message ToTransportMsg { } message UsageStatsKVProto { - string key = 1 [deprecated=true]; + reserved 1; + reserved "key"; int64 value = 2; ApiUsageRecordKeyProto recordKey = 3; } message ToUsageStatsServiceMsg { - int64 tenantIdMSB = 1 [deprecated=true]; - int64 tenantIdLSB = 2 [deprecated=true]; - 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]; + reserved 1 to 7; + reserved "tenantIdMSB", "tenantIdLSB", "entityIdMSB", "entityIdLSB", "values", "customerIdMSB", "customerIdLSB"; string serviceId = 8; repeated UsageStatsServiceMsg msgs = 9; } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageReportClient.java b/common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageReportClient.java index 50e131103b..717e51d8b5 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageReportClient.java +++ b/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.util.ProtoUtils; 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.UsageStatsKVProto; +import org.thingsboard.server.gen.transport.TransportProtos.UsageStatsServiceMsg; import org.thingsboard.server.queue.TbQueueProducer; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.discovery.PartitionService; @@ -178,7 +178,9 @@ public class DefaultTbApiUsageReportClient implements TbApiUsageReportClient { @Override public void report(TenantId tenantId, CustomerId customerId, ApiUsageRecordKey key, long value) { - if (!enabled) return; + if (!enabled) { + return; + } ReportLevel[] reportLevels = new ReportLevel[3]; reportLevels[0] = ReportLevel.of(tenantId); @@ -199,7 +201,9 @@ public class DefaultTbApiUsageReportClient implements TbApiUsageReportClient { private void report(ApiUsageRecordKey key, long value, ReportLevel... levels) { ConcurrentMap statsForKey = stats.get(key); for (ReportLevel level : levels) { - if (level == null) continue; + if (level == null) { + continue; + } AtomicLong n = statsForKey.computeIfAbsent(level, k -> new AtomicLong()); if (key.isCounter()) { @@ -212,6 +216,7 @@ public class DefaultTbApiUsageReportClient implements TbApiUsageReportClient { @Data private static class ReportLevel { + private final TenantId tenantId; private final CustomerId customerId; @@ -231,12 +236,14 @@ public class DefaultTbApiUsageReportClient implements TbApiUsageReportClient { @Data private static class ParentEntity { + private final TenantId tenantId; private final CustomerId customerId; public EntityId getId() { return customerId != null ? customerId : tenantId; } + } } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java index 361453abff..36c81edf92 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java @@ -807,6 +807,31 @@ public class DefaultTransportService extends TransportActivityManager implements 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 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 public void log(TransportProtos.SessionInfoProto sessionInfo, String msg) { if (!logEnabled || sessionInfo == null || StringUtils.isEmpty(msg)) { @@ -1008,7 +1033,9 @@ public class DefaultTransportService extends TransportActivityManager implements case APIUSAGESTATE: ApiUsageState apiUsageState = ProtoUtils.fromProto(msg.getApiUsageState()); 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; case DEVICE: onDeviceUpdate(ProtoUtils.fromProto(msg.getDevice())); diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/BaseEdgeEventService.java b/dao/src/main/java/org/thingsboard/server/dao/edge/BaseEdgeEventService.java index 22de47e301..ee8e3aa45c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/BaseEdgeEventService.java +++ b/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.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.edge.EdgeEvent; 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.TimePageLink; import org.thingsboard.server.common.msg.tools.TbRateLimitsException; +import org.thingsboard.server.common.stats.TbApiUsageReportClient; import org.thingsboard.server.dao.service.DataValidator; public abstract class BaseEdgeEventService implements EdgeEventService { @@ -35,6 +37,8 @@ public abstract class BaseEdgeEventService implements EdgeEventService { private RateLimitService rateLimitService; @Autowired private DataValidator edgeEventValidator; + @Autowired(required = false) + private TbApiUsageReportClient apiUsageReportClient; @Override public PageData 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); } + protected void reportEdgeEventUsage(EdgeEvent edgeEvent) { + if (apiUsageReportClient != null) { + apiUsageReportClient.report(edgeEvent.getTenantId(), null, ApiUsageRecordKey.EDGE_EVENT_COUNT, 1); + } + } + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeDao.java b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeDao.java index 5813a3d89a..595f890963 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeDao.java +++ b/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.UUID; -/** - * The Interface EdgeDao. - * - */ public interface EdgeDao extends Dao, TenantEntityDao { Edge save(TenantId tenantId, Edge edge); diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeEventDao.java b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeEventDao.java index 870c4e5054..a6a9dc4062 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeEventDao.java +++ b/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; -/** - * The Interface EdgeEventDao. - */ public interface EdgeEventDao extends Dao { - /** - * Save or update edge event object - * - * @param edgeEvent the event object - * @return saved edge event object future - */ ListenableFuture 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 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); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/PostgresEdgeEventService.java b/dao/src/main/java/org/thingsboard/server/dao/edge/PostgresEdgeEventService.java index 646d7374e0..4bd2051cfa 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/PostgresEdgeEventService.java +++ b/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) { statsCounterService.ifPresent(statsCounterService -> statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_ADDED, edgeEvent.getTenantId(), edgeEvent.getEdgeId(), 1)); + reportEdgeEventUsage(edgeEvent); eventPublisher.publishEvent(SaveEntityEvent.builder() .tenantId(edgeEvent.getTenantId()) .entityId(edgeEvent.getEdgeId()) diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java index 6f971e61f7..4c2623e57a 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java @@ -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_SMS_EXEC_COLUMN = "sms_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. diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/ApiUsageStateEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sql/ApiUsageStateEntity.java index 93c7ee4e8f..aac44c8c7d 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sql/ApiUsageStateEntity.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sql/ApiUsageStateEntity.java @@ -69,6 +69,9 @@ public class ApiUsageStateEntity extends BaseVersionedEntity impl @Enumerated(EnumType.STRING) @Column(name = ModelConstants.API_USAGE_STATE_ALARM_EXEC_COLUMN) private ApiUsageStateValue alarmExecState = ApiUsageStateValue.ENABLED; + @Enumerated(EnumType.STRING) + @Column(name = ModelConstants.API_USAGE_STATE_EDGE_COLUMN) + private ApiUsageStateValue edgeState = ApiUsageStateValue.ENABLED; public ApiUsageStateEntity() { } @@ -90,6 +93,7 @@ public class ApiUsageStateEntity extends BaseVersionedEntity impl this.emailExecState = ur.getEmailExecState(); this.smsExecState = ur.getSmsExecState(); this.alarmExecState = ur.getAlarmExecState(); + this.edgeState = ur.getEdgeState(); } @Override @@ -110,6 +114,7 @@ public class ApiUsageStateEntity extends BaseVersionedEntity impl ur.setEmailExecState(emailExecState); ur.setSmsExecState(smsExecState); ur.setAlarmExecState(alarmExecState); + ur.setEdgeState(edgeState); ur.setVersion(version); return ur; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/usagerecord/ApiUsageStateRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/usagerecord/ApiUsageStateRepository.java index 6f5976745d..1fc73e1080 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/usagerecord/ApiUsageStateRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/usagerecord/ApiUsageStateRepository.java @@ -51,7 +51,7 @@ public interface ApiUsageStateRepository extends JpaRepository :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 findNextBatch(@Param("id") UUID id, Limit limit); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/usagerecord/ApiUsageStateServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/usagerecord/ApiUsageStateServiceImpl.java index 7f0e507f5f..94cd69f468 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/usagerecord/ApiUsageStateServiceImpl.java +++ b/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.setEmailExecState(ApiUsageStateValue.ENABLED); apiUsageState.setAlarmExecState(ApiUsageStateValue.ENABLED); + apiUsageState.setEdgeState(ApiUsageStateValue.ENABLED); apiUsageStateValidator.validate(apiUsageState, ApiUsageState::getTenantId); 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()))); apiUsageStates.add(new BasicTsKvEntry(saved.getCreatedTime(), 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); if (configuration != null) { diff --git a/dao/src/main/resources/sql/schema-entities.sql b/dao/src/main/resources/sql/schema-entities.sql index 6a66b203a4..66b2134803 100644 --- a/dao/src/main/resources/sql/schema-entities.sql +++ b/dao/src/main/resources/sql/schema-entities.sql @@ -706,6 +706,7 @@ CREATE TABLE IF NOT EXISTS api_usage_state ( email_exec varchar(32), sms_exec varchar(32), alarm_exec varchar(32), + edge varchar(32), version BIGINT DEFAULT 1, CONSTRAINT api_usage_state_unq_key UNIQUE (tenant_id, entity_id) ); diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/ApiUsageStateServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/ApiUsageStateServiceTest.java index e95c353687..276bf53f56 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/ApiUsageStateServiceTest.java +++ b/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.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 public class ApiUsageStateServiceTest extends AbstractServiceTest { @@ -33,7 +37,22 @@ public class ApiUsageStateServiceTest extends AbstractServiceTest { @Test public void testFindTenantApiUsageState() { 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 @@ -42,7 +61,7 @@ public class ApiUsageStateServiceTest extends AbstractServiceTest { state.setTransportState(ApiUsageStateValue.DISABLED); ApiUsageState updated = apiUsageStateService.update(state); - Assert.assertEquals(ApiUsageStateValue.DISABLED, updated.getTransportState()); + assertEquals(ApiUsageStateValue.DISABLED, updated.getTransportState()); } @Test @@ -53,20 +72,76 @@ public class ApiUsageStateServiceTest extends AbstractServiceTest { 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 public void testFindApiUsageStateByEntityId() { ApiUsageState state = apiUsageStateService.findApiUsageStateByEntityId(tenantId); - Assert.assertNotNull(state); + assertNotNull(state); } @Test public void testDeleteByTenantId() { ApiUsageState state = apiUsageStateService.findTenantApiUsageState(tenantId); - Assert.assertNotNull(state); + assertNotNull(state); apiUsageStateService.deleteByTenantId(tenantId); state = apiUsageStateService.findTenantApiUsageState(tenantId); - Assert.assertNull(state); + assertNull(state); } } diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/DeviceServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/DeviceServiceTest.java index b9fd609521..e7ce898614 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/DeviceServiceTest.java +++ b/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.mockito.Mockito; 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.TransactionStatus; 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.DeviceProfileService; 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.ota.OtaPackageService; import org.thingsboard.server.dao.service.validator.DeviceCredentialsDataValidator; import org.thingsboard.server.dao.tenant.TbTenantProfileCache; import org.thingsboard.server.dao.tenant.TenantProfileService; +import org.thingsboard.server.exception.DataValidationException; import java.nio.ByteBuffer; import java.util.ArrayList; @@ -111,10 +111,10 @@ public class DeviceServiceTest extends AbstractServiceTest { private CalculatedFieldService calculatedFieldService; @Autowired private PlatformTransactionManager platformTransactionManager; - @SpyBean + @MockitoSpyBean private DeviceCredentialsDataValidator validator; - private IdComparator idComparator = new IdComparator<>(); + private final IdComparator idComparator = new IdComparator<>(); private TenantId anotherTenantId; private static ListeningExecutorService executor; diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/EdgeEventServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/EdgeEventServiceTest.java index db9c074a1b..6b5d38f73d 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/EdgeEventServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/EdgeEventServiceTest.java @@ -74,7 +74,7 @@ public class EdgeEventServiceTest extends AbstractServiceTest { PageData edgeEvents = edgeEventService.findEdgeEvents(tenantId, edgeId, 0L, null, new TimePageLink(1)); 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.getEdgeId(), edgeEvent.getEdgeId()); Assert.assertEquals(saved.getEntityId(), edgeEvent.getEntityId()); @@ -124,7 +124,7 @@ public class EdgeEventServiceTest extends AbstractServiceTest { Assert.assertNotNull(edgeEvents.getData()); 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()); edgeEventDao.cleanupEvents(1); @@ -136,4 +136,4 @@ public class EdgeEventServiceTest extends AbstractServiceTest { return edgeEventService.saveAsync(edgeEvent); } -} \ No newline at end of file +} diff --git a/edqs/src/test/java/org/thingsboard/server/edqs/repo/ApiUsageStateFilterTest.java b/edqs/src/test/java/org/thingsboard/server/edqs/repo/ApiUsageStateFilterTest.java index 494ca24a9e..6839e135bf 100644 --- a/edqs/src/test/java/org/thingsboard/server/edqs/repo/ApiUsageStateFilterTest.java +++ b/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 java.util.Arrays; +import java.util.List; import java.util.UUID; 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); 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()); } @@ -81,6 +82,7 @@ public class ApiUsageStateFilterTest extends AbstractEDQTest { apiUsageState.setSmsExecState(ApiUsageStateValue.ENABLED); apiUsageState.setEmailExecState(ApiUsageStateValue.ENABLED); apiUsageState.setAlarmExecState(ApiUsageStateValue.ENABLED); + apiUsageState.setEdgeState(ApiUsageStateValue.ENABLED); return apiUsageState; } @@ -99,7 +101,7 @@ public class ApiUsageStateFilterTest extends AbstractEDQTest { nameFilter.setPredicate(predicate); nameFilter.setValueType(EntityKeyValueType.STRING); - return new EntityDataQuery(filter, pageLink, entityFields, null, Arrays.asList(nameFilter)); + return new EntityDataQuery(filter, pageLink, entityFields, null, List.of(nameFilter)); } } diff --git a/ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.html b/ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.html index ddd159388a..e2158854b7 100644 --- a/ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.html +++ b/ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.html @@ -182,18 +182,18 @@ - tenant-profile.max-transport-messages + tenant-profile.max-rule-node-executions-per-message - @if (tenantProfileConfigurationForm.get('maxTransportMessages').hasError('required')) { + @if (tenantProfileConfigurationForm.get('maxRuleNodeExecutionsPerMessage').hasError('required')) { - {{ 'tenant-profile.max-transport-messages-required' | translate}} + {{ 'tenant-profile.max-rule-node-executions-per-message-required' | translate}} } - @if (tenantProfileConfigurationForm.get('maxTransportMessages').hasError('min')) { + @if (tenantProfileConfigurationForm.get('maxRuleNodeExecutionsPerMessage').hasError('min')) { - {{ 'tenant-profile.max-transport-messages-range' | translate}} + {{ 'tenant-profile.max-rule-node-executions-per-message-range' | translate}} } @@ -242,24 +242,58 @@ + + + + +
+ + {{ 'tenant-profile.api-usage' | translate }} tenant-profile.unlimited + +
+ + tenant-profile.max-transport-messages + + @if (tenantProfileConfigurationForm.get('maxTransportMessages').hasError('required')) { + + {{ 'tenant-profile.max-transport-messages-required' | translate}} + + } + @if (tenantProfileConfigurationForm.get('maxTransportMessages').hasError('min')) { + + {{ 'tenant-profile.max-transport-messages-range' | translate}} + + } + + + + tenant-profile.max-edge-events + + @if (tenantProfileConfigurationForm.get('maxEdgeEvents').hasError('required')) { + + {{ 'tenant-profile.max-edge-events-required' | translate}} + + } + @if (tenantProfileConfigurationForm.get('maxEdgeEvents').hasError('min')) { + + {{ 'tenant-profile.max-edge-events-range' | translate}} + + } + + +
+ + + + tenant-profile.advanced-settings + + +
- - tenant-profile.max-rule-node-executions-per-message - - @if (tenantProfileConfigurationForm.get('maxRuleNodeExecutionsPerMessage').hasError('required')) { - - {{ 'tenant-profile.max-rule-node-executions-per-message-required' | translate}} - - } - @if (tenantProfileConfigurationForm.get('maxRuleNodeExecutionsPerMessage').hasError('min')) { - - {{ 'tenant-profile.max-rule-node-executions-per-message-range' | translate}} - - } - - tenant-profile.max-transport-data-points +
+
{{ 'tenant-profile.calculated-fields' | translate }} tenant-profile.unlimited diff --git a/ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.ts b/ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.ts index 9fe6232f4c..bf0bb4369a 100644 --- a/ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.ts +++ b/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], smsEnabled: [false], maxCreatedAlarms: [0, [Validators.required, Validators.min(0)]], + maxEdgeEvents: [0, [Validators.required, Validators.min(0)]], maxDebugModeDurationMinutes: [0, [Validators.min(0)]], defaultStorageTtlDays: [0, [Validators.required, Validators.min(0)]], alarmsTtlDays: [0, [Validators.required, Validators.min(0)]], diff --git a/ui-ngx/src/app/shared/models/tenant.model.ts b/ui-ngx/src/app/shared/models/tenant.model.ts index 262ecfc62a..df18312504 100644 --- a/ui-ngx/src/app/shared/models/tenant.model.ts +++ b/ui-ngx/src/app/shared/models/tenant.model.ts @@ -71,6 +71,7 @@ export interface DefaultTenantProfileConfiguration { maxSms: number; smsEnabled: boolean; maxCreatedAlarms: number; + maxEdgeEvents: number; maxDebugModeDurationMinutes: number; @@ -154,6 +155,7 @@ export function createTenantProfileConfiguration(type: TenantProfileType): Tenan maxSms: 0, smsEnabled: true, maxCreatedAlarms: 0, + maxEdgeEvents: 0, maxDebugModeDurationMinutes: 15, tenantServerRestLimitsConfiguration: '', customerServerRestLimitsConfiguration: '', diff --git a/ui-ngx/src/assets/locale/locale.constant-en_US.json b/ui-ngx/src/assets/locale/locale.constant-en_US.json index b1a3996e3c..950b92e247 100644 --- a/ui-ngx/src/assets/locale/locale.constant-en_US.json +++ b/ui-ngx/src/assets/locale/locale.constant-en_US.json @@ -6263,6 +6263,7 @@ "advanced-settings": "Advanced settings", "entities": "Entities", "rule-engine": "Rule Engine", + "api-usage": "API Usage", "time-to-live": "Time-to-live", "calculated-fields": "Calculated fields", "alarms-and-notifications": "Alarms and notifications", @@ -6398,7 +6399,10 @@ "max-sms-range": "SMS sent maximum number can't be negative", "max-created-alarms": "Alarms created maximum number", "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", "add-queue": "Add Queue", "queues-with-count": "Queues ({{count}})",