From d2a9291e525217d9d6e0f225702ca9506c648d01 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Mon, 3 Jun 2024 15:16:32 +0200 Subject: [PATCH] added rate limits for the gateway device --- .../controller/TenantProfileController.java | 3 + .../mqtt/MqttGatewayRateLimitsTest.java | 90 ++++++++++++++++--- .../server/common/data/limit/LimitedApi.java | 1 + .../DefaultTenantProfileConfiguration.java | 3 + .../server/common/util/ProtoUtils.java | 6 ++ common/proto/src/main/proto/queue.proto | 2 + .../transport/auth/SessionInfoCreator.java | 1 + .../transport/auth/TransportDeviceInfo.java | 1 + .../DefaultTransportRateLimitService.java | 55 +++++++++--- .../transport/limits/TransportLimitsType.java | 2 +- .../limits/TransportRateLimitService.java | 2 +- .../service/DefaultTransportService.java | 37 +++++--- ...enant-profile-configuration.component.html | 12 ++- ...-tenant-profile-configuration.component.ts | 3 + .../tenant/rate-limits/rate-limits.models.ts | 9 ++ .../app/shared/models/limited-api.models.ts | 4 + .../assets/locale/locale.constant-en_US.json | 11 +++ 17 files changed, 201 insertions(+), 41 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/controller/TenantProfileController.java b/application/src/main/java/org/thingsboard/server/controller/TenantProfileController.java index 097aec6306..9c42683b00 100644 --- a/application/src/main/java/org/thingsboard/server/controller/TenantProfileController.java +++ b/application/src/main/java/org/thingsboard/server/controller/TenantProfileController.java @@ -141,6 +141,9 @@ public class TenantProfileController extends BaseController { " \"transportGatewayMsgRateLimit\": \"20:1,600:60\",\n" + " \"transportGatewayTelemetryMsgRateLimit\": \"20:1,600:60\",\n" + " \"transportGatewayTelemetryDataPointsRateLimit\": \"20:1,600:60\",\n" + + " \"transportGatewayDeviceMsgRateLimit\": \"20:1,600:60\",\n" + + " \"transportGatewayDeviceTelemetryMsgRateLimit\": \"20:1,600:60\",\n" + + " \"transportGatewayDeviceTelemetryDataPointsRateLimit\": \"20:1,600:60\",\n" + " \"maxTransportMessages\": 10000000,\n" + " \"maxTransportDataPoints\": 10000000,\n" + " \"maxREExecutions\": 4000000,\n" + diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/MqttGatewayRateLimitsTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/MqttGatewayRateLimitsTest.java index a1986ca2c2..b7569bc816 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/MqttGatewayRateLimitsTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/MqttGatewayRateLimitsTest.java @@ -16,6 +16,7 @@ package org.thingsboard.server.transport.mqtt; import com.fasterxml.jackson.databind.node.ObjectNode; +import org.awaitility.Awaitility; import org.junit.Assert; import org.junit.Before; import org.junit.Test; @@ -25,7 +26,7 @@ import org.springframework.test.context.TestPropertySource; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.TenantProfile; -import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.limit.LimitedApi; import org.thingsboard.server.common.data.notification.rule.trigger.RateLimitsTrigger; import org.thingsboard.server.common.data.security.DeviceCredentials; import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; @@ -34,12 +35,14 @@ import org.thingsboard.server.controller.AbstractControllerTest; import org.thingsboard.server.dao.service.DaoSqlTest; import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient; +import java.util.concurrent.TimeUnit; import java.util.function.Consumer; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; import static org.mockito.ArgumentMatchers.eq; import static org.thingsboard.server.common.data.limit.LimitedApi.TRANSPORT_MESSAGES_PER_GATEWAY; +import static org.thingsboard.server.common.data.limit.LimitedApi.TRANSPORT_MESSAGES_PER_GATEWAY_DEVICE; @DaoSqlTest @TestPropertySource(properties = { @@ -48,11 +51,14 @@ import static org.thingsboard.server.common.data.limit.LimitedApi.TRANSPORT_MESS }) public class MqttGatewayRateLimitsTest extends AbstractControllerTest { - private static final String TOPIC = "v1/gateway/telemetry"; + private static final String GATEWAY_TOPIC = "v1/gateway/telemetry"; + private static final String DEVICE_TOPIC = "v1/devices/me/telemetry"; private static final String DEVICE_A = "DeviceA"; private static final String DEVICE_B = "DeviceB"; - private DeviceId gatewayId; + private static final String DEVICE_PAYLOAD = "{\"temperature\": 42}"; + + private Device gateway; private String gatewayAccessToken; @SpyBean @@ -71,6 +77,10 @@ public class MqttGatewayRateLimitsTest extends AbstractControllerTest { profileConfiguration.setTransportGatewayTelemetryMsgRateLimit(null); profileConfiguration.setTransportGatewayTelemetryDataPointsRateLimit(null); + profileConfiguration.setTransportGatewayDeviceMsgRateLimit(null); + profileConfiguration.setTransportGatewayDeviceTelemetryMsgRateLimit(null); + profileConfiguration.setTransportGatewayDeviceTelemetryDataPointsRateLimit(null); + doPost("/api/tenantProfile", tenantProfile); loginTenantAdmin(); @@ -107,24 +117,24 @@ public class MqttGatewayRateLimitsTest extends AbstractControllerTest { MqttTestClient client = new MqttTestClient(); client.connectAndWait(gatewayAccessToken); - client.publishAndWait(TOPIC, getGatewayPayload(DEVICE_A)); + client.publishAndWait(GATEWAY_TOPIC, getGatewayPayload(DEVICE_A)); loginTenantAdmin(); Device deviceA = getDeviceByName(DEVICE_A); - var deviceATrigger = createRateLimitsTrigger(deviceA); + var deviceATrigger = createRateLimitsTrigger(deviceA, TRANSPORT_MESSAGES_PER_GATEWAY); Mockito.verify(notificationRuleProcessor, Mockito.never()).process(eq(deviceATrigger)); try { - client.publishAndWait(TOPIC, getGatewayPayload(DEVICE_B)); + client.publishAndWait(GATEWAY_TOPIC, getGatewayPayload(DEVICE_B)); } catch (Exception t) { } Device deviceB = getDeviceByName(DEVICE_B); - var deviceBTrigger = createRateLimitsTrigger(deviceB); + var deviceBTrigger = createRateLimitsTrigger(deviceB, TRANSPORT_MESSAGES_PER_GATEWAY); Mockito.verify(notificationRuleProcessor, Mockito.times(1)).process(deviceBTrigger); @@ -133,6 +143,62 @@ public class MqttGatewayRateLimitsTest extends AbstractControllerTest { } } + @Test + public void transportGatewayDeviceMsgRateLimitTest() throws Exception { + transportGatewayDeviceRateLimitTest(profileConfiguration -> profileConfiguration.setTransportGatewayDeviceMsgRateLimit("3:600")); + } + + @Test + public void transportGatewayDeviceTelemetryMsgRateLimitTest() throws Exception { + transportGatewayDeviceRateLimitTest(profileConfiguration -> profileConfiguration.setTransportGatewayDeviceTelemetryMsgRateLimit("1:600")); + } + + @Test + public void transportGatewayDeviceTelemetryDataPointsRateLimitTest() throws Exception { + transportGatewayDeviceRateLimitTest(profileConfiguration -> profileConfiguration.setTransportGatewayDeviceTelemetryDataPointsRateLimit("1:600")); + } + + private void transportGatewayDeviceRateLimitTest(Consumer profileConfiguration) throws Exception { + loginSysAdmin(); + + TenantProfile tenantProfile = doGet("/api/tenantProfile/" + tenantProfileId, TenantProfile.class); + Assert.assertNotNull(tenantProfile); + + profileConfiguration.accept((DefaultTenantProfileConfiguration) tenantProfile.getProfileData().getConfiguration()); + + doPost("/api/tenantProfile", tenantProfile); + + MqttTestClient client = new MqttTestClient(); + client.connectAndWait(gatewayAccessToken); + client.publishAndWait(DEVICE_TOPIC, DEVICE_PAYLOAD.getBytes()); + client.disconnect(); + + var gatewayTrigger = createRateLimitsTrigger(gateway, TRANSPORT_MESSAGES_PER_GATEWAY_DEVICE); + + Mockito.verify(notificationRuleProcessor, Mockito.never()).process(eq(gatewayTrigger)); + + loginTenantAdmin(); + + client = new MqttTestClient(); + + try { + client.connectAndWait(gatewayAccessToken); + client.publishAndWait(DEVICE_TOPIC, DEVICE_PAYLOAD.getBytes()); + if (client.isConnected()) { + client.disconnect(); + } + } catch (Exception t) { + } + + Awaitility.await() + .atMost(2, TimeUnit.SECONDS) + .untilAsserted(() -> Mockito.verify(notificationRuleProcessor, Mockito.times(1)).process(gatewayTrigger)); + + if (client.isConnected()) { + client.disconnect(); + } + } + private void createGateway() throws Exception { Device device = new Device(); device.setName("gateway"); @@ -141,7 +207,7 @@ public class MqttGatewayRateLimitsTest extends AbstractControllerTest { device.setAdditionalInfo(additionalInfo); device = doPost("/api/device", device, Device.class); assertNotNull(device); - gatewayId = device.getId(); + var gatewayId = device.getId(); assertNotNull(gatewayId); DeviceCredentials deviceCredentials = doGet("/api/device/" + gatewayId + "/credentials", DeviceCredentials.class); @@ -149,6 +215,8 @@ public class MqttGatewayRateLimitsTest extends AbstractControllerTest { assertEquals(gatewayId, deviceCredentials.getDeviceId()); gatewayAccessToken = deviceCredentials.getCredentialsId(); assertNotNull(gatewayAccessToken); + + this.gateway = device; } private Device getDeviceByName(String deviceName) throws Exception { @@ -158,13 +226,13 @@ public class MqttGatewayRateLimitsTest extends AbstractControllerTest { } private byte[] getGatewayPayload(String deviceName) { - return String.format("{\"%s\": [{\"values\": {\"temperature\": 42}}]}", deviceName).getBytes(); + return String.format("{\"%s\": [{\"values\": %s}]}", deviceName, DEVICE_PAYLOAD).getBytes(); } - private RateLimitsTrigger createRateLimitsTrigger(Device device) { + private RateLimitsTrigger createRateLimitsTrigger(Device device, LimitedApi limitedApi) { return RateLimitsTrigger.builder() .tenantId(tenantId) - .api(TRANSPORT_MESSAGES_PER_GATEWAY) + .api(limitedApi) .limitLevel(device.getId()) .limitLevelEntityName(device.getName()) .build(); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java index 11702100d5..7aa472bea1 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java @@ -41,6 +41,7 @@ public enum LimitedApi { TRANSPORT_MESSAGES_PER_TENANT("transport messages", true), TRANSPORT_MESSAGES_PER_DEVICE("transport messages per device", false), TRANSPORT_MESSAGES_PER_GATEWAY("transport messages per gateway", false), + TRANSPORT_MESSAGES_PER_GATEWAY_DEVICE("transport messages per gateway device", false), EMAILS("emails sending", true); private Function configExtractor; 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 d990c300da..40f4b7f42b 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 @@ -50,6 +50,9 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura private String transportGatewayMsgRateLimit; private String transportGatewayTelemetryMsgRateLimit; private String transportGatewayTelemetryDataPointsRateLimit; + private String transportGatewayDeviceMsgRateLimit; + private String transportGatewayDeviceTelemetryMsgRateLimit; + private String transportGatewayDeviceTelemetryDataPointsRateLimit; private String tenantEntityExportRateLimit; private String tenantEntityImportRateLimit; 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 16bd2ff896..e3d587e915 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 @@ -90,6 +90,8 @@ import java.util.UUID; import java.util.function.Function; import java.util.stream.Collectors; +import static org.thingsboard.server.common.data.DataConstants.GATEWAY_PARAMETER; + @Slf4j public class ProtoUtils { @@ -1013,6 +1015,10 @@ public class ProtoUtils { .setDeviceProfileIdLSB(device.getDeviceProfileId().getId().getLeastSignificantBits()) .setAdditionalInfo(JacksonUtil.toString(device.getAdditionalInfo())); + if (device.getAdditionalInfo().has(GATEWAY_PARAMETER)) { + builder.setIsGateway(device.getAdditionalInfo().get(GATEWAY_PARAMETER).booleanValue()); + } + PowerSavingConfiguration psmConfiguration = switch (device.getDeviceData().getTransportConfiguration().getType()) { case LWM2M -> (Lwm2mDeviceTransportConfiguration) device.getDeviceData().getTransportConfiguration(); case COAP -> (CoapDeviceTransportConfiguration) device.getDeviceData().getTransportConfiguration(); diff --git a/common/proto/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto index 038fdce130..581b993120 100644 --- a/common/proto/src/main/proto/queue.proto +++ b/common/proto/src/main/proto/queue.proto @@ -98,6 +98,7 @@ message SessionInfoProto { int64 customerIdLSB = 15; optional int64 gatewayIdMSB = 16; optional int64 gatewayIdLSB = 17; + bool isGateway = 18; } enum SessionEvent { @@ -184,6 +185,7 @@ message DeviceInfoProto { int64 edrxCycle = 13; int64 psmActivityTimer = 14; int64 pagingTransmissionWindow = 15; + bool isGateway = 16; } message DeviceProto { diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/auth/SessionInfoCreator.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/auth/SessionInfoCreator.java index 6ba41c245e..e0b8e261f8 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/auth/SessionInfoCreator.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/auth/SessionInfoCreator.java @@ -46,6 +46,7 @@ public class SessionInfoCreator { .setDeviceType(msg.getDeviceInfo().getDeviceType()) .setDeviceProfileIdMSB(msg.getDeviceInfo().getDeviceProfileId().getId().getMostSignificantBits()) .setDeviceProfileIdLSB(msg.getDeviceInfo().getDeviceProfileId().getId().getLeastSignificantBits()) + .setIsGateway(msg.getDeviceInfo().isGateway()) .build(); } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/auth/TransportDeviceInfo.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/auth/TransportDeviceInfo.java index ae58c3f1a9..c7239673f6 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/auth/TransportDeviceInfo.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/auth/TransportDeviceInfo.java @@ -38,4 +38,5 @@ public class TransportDeviceInfo implements Serializable { private Long edrxCycle; private Long psmActivityTimer; private Long pagingTransmissionWindow; + private boolean gateway; } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/DefaultTransportRateLimitService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/DefaultTransportRateLimitService.java index 975dd21fab..d070465c23 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/DefaultTransportRateLimitService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/DefaultTransportRateLimitService.java @@ -41,6 +41,7 @@ import java.util.function.BiConsumer; import java.util.function.Function; import static org.thingsboard.server.common.transport.limits.TransportLimitsType.DEVICE_LIMITS; +import static org.thingsboard.server.common.transport.limits.TransportLimitsType.GATEWAY_DEVICE_LIMITS; import static org.thingsboard.server.common.transport.limits.TransportLimitsType.GATEWAY_LIMITS; import static org.thingsboard.server.common.transport.limits.TransportLimitsType.TENANT_LIMITS; @@ -53,9 +54,11 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi private final ConcurrentMap tenantAllowed = new ConcurrentHashMap<>(); private final ConcurrentMap> tenantDevices = new ConcurrentHashMap<>(); private final ConcurrentMap> tenantGateways = new ConcurrentHashMap<>(); + private final ConcurrentMap> tenantGatewayDevices = new ConcurrentHashMap<>(); private final ConcurrentMap perTenantLimits = new ConcurrentHashMap<>(); private final ConcurrentMap perDeviceLimits = new ConcurrentHashMap<>(); private final ConcurrentMap perGatewayLimits = new ConcurrentHashMap<>(); + private final ConcurrentMap perGatewayDeviceLimits = new ConcurrentHashMap<>(); private final Map ipMap = new ConcurrentHashMap<>(); private final TransportTenantProfileCache tenantProfileCache; @@ -72,21 +75,23 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi } @Override - public TbPair checkLimits(TenantId tenantId, DeviceId gatewayId, DeviceId deviceId, int dataPoints) { + public TbPair checkLimits(TenantId tenantId, DeviceId gatewayId, DeviceId deviceId, int dataPoints, boolean isGateway) { if (!tenantAllowed.getOrDefault(tenantId, Boolean.TRUE)) { return TbPair.of(EntityType.API_USAGE_STATE, false); } if (!checkEntityRateLimit(dataPoints, getTenantRateLimits(tenantId))) { return TbPair.of(EntityType.TENANT, false); } - + if (isGateway && !checkEntityRateLimit(dataPoints, getGatewayDeviceRateLimits(tenantId, deviceId))) { + return TbPair.of(EntityType.DEVICE, true); + } if (gatewayId != null && !checkEntityRateLimit(dataPoints, getGatewayRateLimits(tenantId, gatewayId))) { return TbPair.of(EntityType.DEVICE, true); } - - if (deviceId != null && !checkEntityRateLimit(dataPoints, getDeviceRateLimits(tenantId, deviceId))) { + if (!isGateway && deviceId != null && !checkEntityRateLimit(dataPoints, getDeviceRateLimits(tenantId, deviceId))) { return TbPair.of(EntityType.DEVICE, false); } + return null; } @@ -104,8 +109,9 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi EntityTransportRateLimits tenantRateLimitPrototype = createRateLimits(update.getProfile(), TENANT_LIMITS); EntityTransportRateLimits deviceRateLimitPrototype = createRateLimits(update.getProfile(), DEVICE_LIMITS); EntityTransportRateLimits gatewayRateLimitPrototype = createRateLimits(update.getProfile(), GATEWAY_LIMITS); + EntityTransportRateLimits gatewayDeviceRateLimitPrototype = createRateLimits(update.getProfile(), GATEWAY_DEVICE_LIMITS); for (TenantId tenantId : update.getAffectedTenants()) { - update(tenantId, tenantRateLimitPrototype, deviceRateLimitPrototype, gatewayRateLimitPrototype); + update(tenantId, tenantRateLimitPrototype, deviceRateLimitPrototype, gatewayRateLimitPrototype, gatewayDeviceRateLimitPrototype); } } @@ -114,26 +120,34 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi EntityTransportRateLimits tenantRateLimitPrototype = createRateLimits(tenantProfileCache.get(tenantId), TENANT_LIMITS); EntityTransportRateLimits deviceRateLimitPrototype = createRateLimits(tenantProfileCache.get(tenantId), DEVICE_LIMITS); EntityTransportRateLimits gatewayRateLimitPrototype = createRateLimits(tenantProfileCache.get(tenantId), GATEWAY_LIMITS); - update(tenantId, tenantRateLimitPrototype, deviceRateLimitPrototype, gatewayRateLimitPrototype); + EntityTransportRateLimits gatewayDeviceRateLimitPrototype = createRateLimits(tenantProfileCache.get(tenantId), GATEWAY_DEVICE_LIMITS); + update(tenantId, tenantRateLimitPrototype, deviceRateLimitPrototype, gatewayRateLimitPrototype, gatewayDeviceRateLimitPrototype); } - private void update(TenantId tenantId, EntityTransportRateLimits tenantRateLimitPrototype, - EntityTransportRateLimits deviceRateLimitPrototype, EntityTransportRateLimits gatewayRateLimitPrototype) { + private void update(TenantId tenantId, EntityTransportRateLimits tenantRateLimitPrototype, EntityTransportRateLimits deviceRateLimitPrototype, + EntityTransportRateLimits gatewayRateLimitPrototype, EntityTransportRateLimits gatewayDeviceRateLimitPrototype) { mergeLimits(tenantId, tenantRateLimitPrototype, perTenantLimits::get, perTenantLimits::put); getTenantDevices(tenantId).forEach(deviceId -> mergeLimits(deviceId, deviceRateLimitPrototype, perDeviceLimits::get, perDeviceLimits::put)); - getTenantGateways(tenantId).forEach(deviceId -> mergeLimits(deviceId, gatewayRateLimitPrototype, perGatewayLimits::get, perGatewayLimits::put)); + getTenantGateways(tenantId).forEach(gatewayId -> mergeLimits(gatewayId, gatewayRateLimitPrototype, perGatewayLimits::get, perGatewayLimits::put)); + getTenantGatewayDevices(tenantId).forEach(gatewayId -> mergeLimits(gatewayId, gatewayDeviceRateLimitPrototype, perGatewayDeviceLimits::get, perGatewayDeviceLimits::put)); } @Override public void remove(TenantId tenantId) { perTenantLimits.remove(tenantId); tenantDevices.remove(tenantId); + tenantGateways.remove(tenantId); + tenantGatewayDevices.remove(tenantId); } @Override public void remove(DeviceId deviceId) { perDeviceLimits.remove(deviceId); + perGatewayLimits.remove(deviceId); + perGatewayDeviceLimits.remove(deviceId); tenantDevices.values().forEach(set -> set.remove(deviceId)); + tenantGateways.values().forEach(set -> set.remove(deviceId)); + tenantGatewayDevices.values().forEach(set -> set.remove(deviceId)); } @Override @@ -273,6 +287,11 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi telemetryMsgRateLimit = newLimit(profile.getTransportGatewayTelemetryMsgRateLimit()); telemetryDpRateLimit = newLimit(profile.getTransportGatewayTelemetryDataPointsRateLimit()); } + case GATEWAY_DEVICE_LIMITS -> { + regularMsgRateLimit = newLimit(profile.getTransportGatewayDeviceMsgRateLimit()); + telemetryMsgRateLimit = newLimit(profile.getTransportGatewayDeviceTelemetryMsgRateLimit()); + telemetryDpRateLimit = newLimit(profile.getTransportGatewayDeviceTelemetryDataPointsRateLimit()); + } default -> throw new IllegalStateException("Unknown limits type: " + limitsType); } @@ -296,10 +315,18 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi }); } - private EntityTransportRateLimits getGatewayRateLimits(TenantId tenantId, DeviceId deviceId) { - return perGatewayLimits.computeIfAbsent(deviceId, k -> { + private EntityTransportRateLimits getGatewayRateLimits(TenantId tenantId, DeviceId gatewayId) { + return perGatewayLimits.computeIfAbsent(gatewayId, k -> { EntityTransportRateLimits limits = createRateLimits(tenantProfileCache.get(tenantId), GATEWAY_LIMITS); - getTenantGateways(tenantId).add(deviceId); + getTenantGateways(tenantId).add(gatewayId); + return limits; + }); + } + + private EntityTransportRateLimits getGatewayDeviceRateLimits(TenantId tenantId, DeviceId gatewayId) { + return perGatewayDeviceLimits.computeIfAbsent(gatewayId, k -> { + EntityTransportRateLimits limits = createRateLimits(tenantProfileCache.get(tenantId), GATEWAY_DEVICE_LIMITS); + getTenantGatewayDevices(tenantId).add(gatewayId); return limits; }); } @@ -312,4 +339,8 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi return tenantGateways.computeIfAbsent(tenantId, id -> ConcurrentHashMap.newKeySet()); } + private Set getTenantGatewayDevices(TenantId tenantId) { + return tenantGatewayDevices.computeIfAbsent(tenantId, id -> ConcurrentHashMap.newKeySet()); + } + } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportLimitsType.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportLimitsType.java index ab65b07414..0c822f9589 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportLimitsType.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportLimitsType.java @@ -16,5 +16,5 @@ package org.thingsboard.server.common.transport.limits; public enum TransportLimitsType { - TENANT_LIMITS, DEVICE_LIMITS, GATEWAY_LIMITS + TENANT_LIMITS, DEVICE_LIMITS, GATEWAY_LIMITS, GATEWAY_DEVICE_LIMITS } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportRateLimitService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportRateLimitService.java index e33f9db240..aefcecb630 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportRateLimitService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportRateLimitService.java @@ -25,7 +25,7 @@ import java.net.InetSocketAddress; public interface TransportRateLimitService { - TbPair checkLimits(TenantId tenantId, DeviceId gatewayId, DeviceId deviceId, int dataPoints); + TbPair checkLimits(TenantId tenantId, DeviceId gatewayId, DeviceId deviceId, int dataPoints, boolean isGateway); void update(TenantProfileUpdateResult update); 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 e82e7812fb..cc91baca0c 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 @@ -429,8 +429,8 @@ public class DefaultTransportService extends TransportActivityManager implements @Override public void process(TenantId tenantId, TransportProtos.GetOrCreateDeviceFromGatewayRequestMsg requestMsg, TransportServiceCallback callback) { log.trace("Processing msg: {}", requestMsg); - DeviceId gatewayid = new DeviceId(new UUID(requestMsg.getGatewayIdMSB(), requestMsg.getGatewayIdLSB())); - if (!checkLimits(tenantId, gatewayid, null, requestMsg.getDeviceName(), requestMsg, callback, 0)) { + DeviceId gatewayId = new DeviceId(new UUID(requestMsg.getGatewayIdMSB(), requestMsg.getGatewayIdLSB())); + if (!checkLimits(tenantId, gatewayId, null, requestMsg.getDeviceName(), requestMsg, callback, 0, false)) { return; } @@ -476,6 +476,7 @@ public class DefaultTransportService extends TransportActivityManager implements tdi.setAdditionalInfo(di.getAdditionalInfo()); tdi.setDeviceName(di.getDeviceName()); tdi.setDeviceType(di.getDeviceType()); + tdi.setGateway(di.getIsGateway()); if (StringUtils.isNotEmpty(di.getPowerMode())) { tdi.setPowerMode(PowerMode.valueOf(di.getPowerMode())); tdi.setEdrxCycle(di.getEdrxCycle()); @@ -838,15 +839,15 @@ public class DefaultTransportService extends TransportActivityManager implements gatewayId = new DeviceId(new UUID(sessionInfo.getGatewayIdMSB(), sessionInfo.getGatewayIdLSB())); } - return checkLimits(tenantId, gatewayId, deviceId, sessionInfo.getDeviceName(), msg, callback, dataPoints); + return checkLimits(tenantId, gatewayId, deviceId, sessionInfo.getDeviceName(), msg, callback, dataPoints, sessionInfo.getIsGateway()); } - private boolean checkLimits(TenantId tenantId, DeviceId gatewayId, DeviceId deviceId, String deviceName, Object msg, TransportServiceCallback callback, int dataPoints) { + private boolean checkLimits(TenantId tenantId, DeviceId gatewayId, DeviceId deviceId, String deviceName, Object msg, TransportServiceCallback callback, int dataPoints, boolean isGateway) { if (log.isTraceEnabled()) { log.trace("[{}][{}] Processing msg: {}", tenantId, deviceName, msg); } - var rateLimitedPair = rateLimitService.checkLimits(tenantId, gatewayId, deviceId, dataPoints); + var rateLimitedPair = rateLimitService.checkLimits(tenantId, gatewayId, deviceId, dataPoints, isGateway); if (rateLimitedPair == null) { return true; } else { @@ -856,9 +857,15 @@ public class DefaultTransportService extends TransportActivityManager implements } if (rateLimitedEntityType == EntityType.DEVICE || rateLimitedEntityType == EntityType.TENANT) { - LimitedApi limitedApi = - rateLimitedEntityType == EntityType.TENANT ? LimitedApi.TRANSPORT_MESSAGES_PER_TENANT : - rateLimitedPair.getSecond() ? LimitedApi.TRANSPORT_MESSAGES_PER_GATEWAY : LimitedApi.TRANSPORT_MESSAGES_PER_DEVICE; + LimitedApi limitedApi; + + if (rateLimitedEntityType == EntityType.TENANT) { + limitedApi = LimitedApi.TRANSPORT_MESSAGES_PER_TENANT; + } else if (rateLimitedPair.getSecond()) { + limitedApi = isGateway ? LimitedApi.TRANSPORT_MESSAGES_PER_GATEWAY_DEVICE : LimitedApi.TRANSPORT_MESSAGES_PER_GATEWAY; + } else { + limitedApi = LimitedApi.TRANSPORT_MESSAGES_PER_DEVICE; + } EntityId limitLevel = rateLimitedEntityType == EntityType.DEVICE ? deviceId == null ? gatewayId : deviceId : tenantId; @@ -1023,16 +1030,20 @@ public class DefaultTransportService extends TransportActivityManager implements } else { newDeviceProfile = null; } + + JsonNode deviceAdditionalInfo = device.getAdditionalInfo(); + boolean isGateway = deviceAdditionalInfo.has(DataConstants.GATEWAY_PARAMETER) + && deviceAdditionalInfo.get(DataConstants.GATEWAY_PARAMETER).asBoolean(); + TransportProtos.SessionInfoProto newSessionInfo = TransportProtos.SessionInfoProto.newBuilder() .mergeFrom(md.getSessionInfo()) .setDeviceProfileIdMSB(deviceProfileIdMSB) .setDeviceProfileIdLSB(deviceProfileIdLSB) .setDeviceName(device.getName()) - .setDeviceType(device.getType()).build(); - JsonNode deviceAdditionalInfo = device.getAdditionalInfo(); - if (deviceAdditionalInfo.has(DataConstants.GATEWAY_PARAMETER) - && deviceAdditionalInfo.get(DataConstants.GATEWAY_PARAMETER).asBoolean() - && deviceAdditionalInfo.has(DataConstants.OVERWRITE_ACTIVITY_TIME_PARAMETER) + .setDeviceType(device.getType()) + .setIsGateway(isGateway).build(); + + if (isGateway && deviceAdditionalInfo.has(DataConstants.OVERWRITE_ACTIVITY_TIME_PARAMETER) && deviceAdditionalInfo.get(DataConstants.OVERWRITE_ACTIVITY_TIME_PARAMETER).isBoolean()) { md.setOverwriteActivityTime(deviceAdditionalInfo.get(DataConstants.OVERWRITE_ACTIVITY_TIME_PARAMETER).asBoolean()); } 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 91f96c7aba..31724191cf 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 @@ -520,13 +520,17 @@ -
+ +
-
+ +
@@ -547,7 +551,9 @@ -
+ +
( [RateLimitsType.GATEWAY_MESSAGES, 'tenant-profile.rate-limits.transport-gateway-msg'], [RateLimitsType.GATEWAY_TELEMETRY_MESSAGES, 'tenant-profile.rate-limits.transport-gateway-telemetry-msg'], [RateLimitsType.GATEWAY_TELEMETRY_DATA_POINTS, 'tenant-profile.rate-limits.transport-gateway-telemetry-data-points'], + [RateLimitsType.GATEWAY_DEVICE_MESSAGES, 'tenant-profile.rate-limits.transport-gateway-device-msg'], + [RateLimitsType.GATEWAY_DEVICE_TELEMETRY_MESSAGES, 'tenant-profile.rate-limits.transport-gateway-device-telemetry-msg'], + [RateLimitsType.GATEWAY_DEVICE_TELEMETRY_DATA_POINTS, 'tenant-profile.rate-limits.transport-gateway-device-telemetry-data-points'], [RateLimitsType.TENANT_SERVER_REST_LIMITS_CONFIGURATION, 'tenant-profile.rest-requests-for-tenant'], [RateLimitsType.CUSTOMER_SERVER_REST_LIMITS_CONFIGURATION, 'tenant-profile.customer-rest-limits'], [RateLimitsType.WS_UPDATE_PER_SESSION_RATE_LIMIT, 'tenant-profile.ws-limit-updates-per-session'], @@ -83,6 +89,9 @@ export const rateLimitsDialogTitleTranslationMap = new Map( [LimitedApi.CASSANDRA_QUERIES, 'api-limit.cassandra-queries'], [LimitedApi.TRANSPORT_MESSAGES_PER_TENANT, 'api-limit.transport-messages'], [LimitedApi.TRANSPORT_MESSAGES_PER_DEVICE, 'api-limit.transport-messages-per-device'], + [LimitedApi.TRANSPORT_MESSAGES_PER_GATEWAY, 'api-limit.transport-messages-per-gateway'], + [LimitedApi.TRANSPORT_MESSAGES_PER_GATEWAY_DEVICE, 'api-limit.transport-messages-per-gateway_device'], [LimitedApi.EDGE_EVENTS, 'api-limit.edge-events'], [LimitedApi.EDGE_EVENTS_PER_EDGE, 'api-limit.edge-events-per-edge'], [LimitedApi.EDGE_UPLINK_MESSAGES, 'api-limit.edge-uplink-messages'], 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 170d4836ee..ac235ef89f 100644 --- a/ui-ngx/src/assets/locale/locale.constant-en_US.json +++ b/ui-ngx/src/assets/locale/locale.constant-en_US.json @@ -894,6 +894,8 @@ "rest-api-requests-per-customer": "REST API requests per customer", "transport-messages": "Transport messages", "transport-messages-per-device": "Transport messages per device", + "transport-messages-per-gateway": "Transport messages per gateway", + "transport-messages-per-gateway-device": "Transport messages per gateway device", "ws-updates-per-session": "WS updates per session", "edge-events": "Edge events", "edge-events-per-edge": "Edge events per edge", @@ -4448,6 +4450,9 @@ "transport-gateway-msg-rate-limit": "Transport gateway messages", "transport-gateway-telemetry-msg-rate-limit": "Transport gateway telemetry messages", "transport-gateway-telemetry-data-points-rate-limit": "Transport gateway telemetry data points", + "transport-gateway-device-msg-rate-limit": "Transport gateway device messages", + "transport-gateway-device-telemetry-msg-rate-limit": "Transport gateway device telemetry messages", + "transport-gateway-device-telemetry-data-points-rate-limit": "Transport gateway device telemetry data points", "tenant-entity-export-rate-limit": "Entity version creation", "tenant-entity-import-rate-limit": "Entity version load", "tenant-notification-request-rate-limit": "Notification requests", @@ -4532,6 +4537,9 @@ "edit-transport-gateway-msg-title": "Edit transport gateway messages rate limits", "edit-transport-gateway-telemetry-msg-title": "Edit transport gateway telemetry messages rate limits", "edit-transport-gateway-telemetry-data-points-title": "Edit transport gateway telemetry data points rate limits", + "edit-transport-gateway-device-msg-title": "Edit transport gateway device messages rate limits", + "edit-transport-gateway-device-telemetry-msg-title": "Edit transport gateway device telemetry messages rate limits", + "edit-transport-gateway-device-telemetry-data-points-title": "Edit transport gateway device telemetry data points rate limits", "edit-tenant-rest-limits-title": "Edit REST requests for tenant rate limits", "edit-customer-rest-limits-title": "Edit REST requests for customer rate limits", "edit-ws-limit-updates-per-session-title": "Edit WS updates per session rate limits", @@ -4568,6 +4576,9 @@ "transport-gateway-msg": "Transport gateway messages", "transport-gateway-telemetry-msg": "Transport gateway telemetry messages", "transport-gateway-telemetry-data-points": "Transport gateway telemetry data points", + "transport-gateway-device-msg": "Transport gateway device messages", + "transport-gateway-device-telemetry-msg": "Transport gateway device telemetry messages", + "transport-gateway-device-telemetry-data-points": "Transport gateway device telemetry data points", "sec": "sec" } },