Browse Source

Transport Rate Limits are now configurable via Tenant Profile

pull/3592/head
Andrii Shvaika 6 years ago
parent
commit
6ea39e835e
  1. 3
      application/src/main/data/json/tenant/device_profile/rule_chain_template.json
  2. 4
      application/src/main/java/org/thingsboard/server/controller/TenantController.java
  3. 5
      application/src/main/java/org/thingsboard/server/controller/TenantProfileController.java
  4. 16
      application/src/main/java/org/thingsboard/server/service/profile/DefaultTbDeviceProfileCache.java
  5. 3
      application/src/main/java/org/thingsboard/server/service/profile/TbDeviceProfileCache.java
  6. 56
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java
  7. 13
      application/src/main/java/org/thingsboard/server/service/queue/TbClusterService.java
  8. 63
      application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java
  9. 6
      application/src/main/resources/thingsboard.yml
  10. 47
      common/queue/src/main/proto/queue.proto
  11. 2
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportDeviceProfileCache.java
  12. 9
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java
  13. 38
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportTenantProfileCache.java
  14. 47
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/DefaultTransportRateLimitFactory.java
  15. 115
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/DefaultTransportRateLimitService.java
  16. 30
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/DummyTransportRateLimit.java
  17. 34
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/SimpleTransportRateLimit.java
  18. 24
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportRateLimit.java
  19. 24
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportRateLimitFactory.java
  20. 36
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportRateLimitService.java
  21. 33
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportRateLimitType.java
  22. 30
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/profile/TenantProfileUpdateResult.java
  23. 6
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportDeviceProfileCache.java
  24. 160
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java
  25. 154
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportTenantProfileCache.java
  26. 24
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/TransportTenantRoutingInfoService.java
  27. 2
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/util/DataDecodingEncodingService.java
  28. 4
      transport/coap/src/main/resources/tb-coap-transport.yml
  29. 4
      transport/http/src/main/resources/tb-http-transport.yml
  30. 4
      transport/mqtt/src/main/resources/tb-mqtt-transport.yml

3
application/src/main/data/json/tenant/device_profile/rule_chain_template.json

@ -94,7 +94,8 @@
"name": "Device Profile Node", "name": "Device Profile Node",
"debugMode": false, "debugMode": false,
"configuration": { "configuration": {
"persistAlarmRulesState": false "persistAlarmRulesState": false,
"fetchAlarmRulesStateOnStart": false
} }
} }
], ],

4
application/src/main/java/org/thingsboard/server/controller/TenantController.java

@ -92,6 +92,7 @@ public class TenantController extends BaseController {
installScripts.createDefaultRuleChains(tenant.getId()); installScripts.createDefaultRuleChains(tenant.getId());
} }
tenantProfileCache.evict(tenant.getId()); tenantProfileCache.evict(tenant.getId());
tbClusterService.onTenantChange(tenant, null);
return tenant; return tenant;
} catch (Exception e) { } catch (Exception e) {
throw handleException(e); throw handleException(e);
@ -105,9 +106,10 @@ public class TenantController extends BaseController {
checkParameter("tenantId", strTenantId); checkParameter("tenantId", strTenantId);
try { try {
TenantId tenantId = new TenantId(toUUID(strTenantId)); TenantId tenantId = new TenantId(toUUID(strTenantId));
checkTenantId(tenantId, Operation.DELETE); Tenant tenant = checkTenantId(tenantId, Operation.DELETE);
tenantService.deleteTenant(tenantId); tenantService.deleteTenant(tenantId);
tenantProfileCache.evict(tenantId); tenantProfileCache.evict(tenantId);
tbClusterService.onTenantDelete(tenant, null);
tbClusterService.onEntityStateChange(tenantId, tenantId, ComponentLifecycleEvent.DELETED); tbClusterService.onEntityStateChange(tenantId, tenantId, ComponentLifecycleEvent.DELETED);
} catch (Exception e) { } catch (Exception e) {
throw handleException(e); throw handleException(e);

5
application/src/main/java/org/thingsboard/server/controller/TenantProfileController.java

@ -34,6 +34,7 @@ import org.thingsboard.server.common.data.id.TenantProfileId;
import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent;
import org.thingsboard.server.dao.exception.DataValidationException;
import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.security.permission.Operation; import org.thingsboard.server.service.security.permission.Operation;
import org.thingsboard.server.service.security.permission.Resource; import org.thingsboard.server.service.security.permission.Resource;
@ -96,6 +97,7 @@ public class TenantProfileController extends BaseController {
tenantProfile = checkNotNull(tenantProfileService.saveTenantProfile(getTenantId(), tenantProfile)); tenantProfile = checkNotNull(tenantProfileService.saveTenantProfile(getTenantId(), tenantProfile));
tenantProfileCache.put(tenantProfile); tenantProfileCache.put(tenantProfile);
tbClusterService.onTenantProfileChange(tenantProfile, null);
tbClusterService.onEntityStateChange(TenantId.SYS_TENANT_ID, tenantProfile.getId(), tbClusterService.onEntityStateChange(TenantId.SYS_TENANT_ID, tenantProfile.getId(),
newTenantProfile ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); newTenantProfile ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED);
return tenantProfile; return tenantProfile;
@ -111,8 +113,9 @@ public class TenantProfileController extends BaseController {
checkParameter("tenantProfileId", strTenantProfileId); checkParameter("tenantProfileId", strTenantProfileId);
try { try {
TenantProfileId tenantProfileId = new TenantProfileId(toUUID(strTenantProfileId)); TenantProfileId tenantProfileId = new TenantProfileId(toUUID(strTenantProfileId));
checkTenantProfileId(tenantProfileId, Operation.DELETE); TenantProfile profile = checkTenantProfileId(tenantProfileId, Operation.DELETE);
tenantProfileService.deleteTenantProfile(getTenantId(), tenantProfileId); tenantProfileService.deleteTenantProfile(getTenantId(), tenantProfileId);
tbClusterService.onTenantProfileDelete(profile, null);
} catch (Exception e) { } catch (Exception e) {
throw handleException(e); throw handleException(e);
} }

16
application/src/main/java/org/thingsboard/server/service/profile/DefaultTbDeviceProfileCache.java

@ -61,7 +61,7 @@ public class DefaultTbDeviceProfileCache implements TbDeviceProfileCache {
profile = deviceProfileService.findDeviceProfileById(tenantId, deviceProfileId); profile = deviceProfileService.findDeviceProfileById(tenantId, deviceProfileId);
if (profile != null) { if (profile != null) {
deviceProfilesMap.put(deviceProfileId, profile); deviceProfilesMap.put(deviceProfileId, profile);
log.info("[{}] Fetch device profile into cache: {}", profile.getId(), profile); log.debug("[{}] Fetch device profile into cache: {}", profile.getId(), profile);
} }
} }
} finally { } finally {
@ -91,7 +91,7 @@ public class DefaultTbDeviceProfileCache implements TbDeviceProfileCache {
public void put(DeviceProfile profile) { public void put(DeviceProfile profile) {
if (profile.getId() != null) { if (profile.getId() != null) {
deviceProfilesMap.put(profile.getId(), profile); deviceProfilesMap.put(profile.getId(), profile);
log.info("[{}] pushed device profile to cache: {}", profile.getId(), profile); log.debug("[{}] pushed device profile to cache: {}", profile.getId(), profile);
notifyListeners(profile); notifyListeners(profile);
} }
} }
@ -99,7 +99,7 @@ public class DefaultTbDeviceProfileCache implements TbDeviceProfileCache {
@Override @Override
public void evict(TenantId tenantId, DeviceProfileId profileId) { public void evict(TenantId tenantId, DeviceProfileId profileId) {
DeviceProfile oldProfile = deviceProfilesMap.remove(profileId); DeviceProfile oldProfile = deviceProfilesMap.remove(profileId);
log.info("[{}] evict device profile from cache: {}", profileId, oldProfile); log.debug("[{}] evict device profile from cache: {}", profileId, oldProfile);
DeviceProfile newProfile = get(tenantId, profileId); DeviceProfile newProfile = get(tenantId, profileId);
if (newProfile != null) { if (newProfile != null) {
notifyListeners(newProfile); notifyListeners(newProfile);
@ -116,6 +116,16 @@ public class DefaultTbDeviceProfileCache implements TbDeviceProfileCache {
listeners.computeIfAbsent(tenantId, id -> new ConcurrentHashMap<>()).put(listenerId, listener); listeners.computeIfAbsent(tenantId, id -> new ConcurrentHashMap<>()).put(listenerId, listener);
} }
@Override
public DeviceProfile find(DeviceProfileId deviceProfileId) {
return deviceProfileService.findDeviceProfileById(TenantId.SYS_TENANT_ID, deviceProfileId);
}
@Override
public DeviceProfile findOrCreateDeviceProfile(TenantId tenantId, String profileName) {
return deviceProfileService.findOrCreateDeviceProfile(tenantId, profileName);
}
@Override @Override
public void removeListener(TenantId tenantId, EntityId listenerId) { public void removeListener(TenantId tenantId, EntityId listenerId) {
ConcurrentMap<EntityId, Consumer<DeviceProfile>> tenantListeners = listeners.get(tenantId); ConcurrentMap<EntityId, Consumer<DeviceProfile>> tenantListeners = listeners.get(tenantId);

3
application/src/main/java/org/thingsboard/server/service/profile/TbDeviceProfileCache.java

@ -29,4 +29,7 @@ public interface TbDeviceProfileCache extends RuleEngineDeviceProfileCache {
void evict(DeviceId id); void evict(DeviceId id);
DeviceProfile find(DeviceProfileId deviceProfileId);
DeviceProfile findOrCreateDeviceProfile(TenantId tenantId, String deviceType);
} }

56
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java

@ -23,11 +23,15 @@ import org.springframework.stereotype.Service;
import org.thingsboard.rule.engine.api.msg.ToDeviceActorNotificationMsg; import org.thingsboard.rule.engine.api.msg.ToDeviceActorNotificationMsg;
import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.HasName;
import org.thingsboard.server.common.data.Tenant;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.DeviceProfileId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.TenantProfileId;
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg;
@ -189,21 +193,51 @@ public class DefaultTbClusterService implements TbClusterService {
@Override @Override
public void onDeviceProfileChange(DeviceProfile deviceProfile, TbQueueCallback callback) { public void onDeviceProfileChange(DeviceProfile deviceProfile, TbQueueCallback callback) {
log.trace("[{}][{}] Processing device profile [{}] change event", deviceProfile.getTenantId(), deviceProfile.getId(), deviceProfile.getName()); onEntityChange(deviceProfile.getTenantId(), deviceProfile.getId(), deviceProfile, callback);
TransportProtos.DeviceProfileUpdateMsg profileUpdateMsg = TransportProtos.DeviceProfileUpdateMsg.newBuilder() }
.setData(ByteString.copyFrom(encodingService.encode(deviceProfile))).build();
ToTransportMsg transportMsg = ToTransportMsg.newBuilder().setDeviceProfileUpdateMsg(profileUpdateMsg).build(); @Override
broadcast(transportMsg); public void onTenantProfileChange(TenantProfile tenantProfile, TbQueueCallback callback) {
onEntityChange(TenantId.SYS_TENANT_ID, tenantProfile.getId(), tenantProfile, callback);
}
@Override
public void onTenantChange(Tenant tenant, TbQueueCallback callback) {
onEntityChange(TenantId.SYS_TENANT_ID, tenant.getId(), tenant, callback);
}
@Override
public void onDeviceProfileDelete(DeviceProfile entity, TbQueueCallback callback) {
onEntityDelete(entity.getTenantId(), entity.getId(), entity.getName(), callback);
} }
@Override @Override
public void onDeviceProfileDelete(DeviceProfile deviceProfile, TbQueueCallback callback) { public void onTenantProfileDelete(TenantProfile entity, TbQueueCallback callback) {
log.trace("[{}][{}] Processing device profile [{}] delete event", deviceProfile.getTenantId(), deviceProfile.getId(), deviceProfile.getName()); onEntityDelete(TenantId.SYS_TENANT_ID, entity.getId(), entity.getName(), callback);
TransportProtos.DeviceProfileDeleteMsg profileDeleteMsg = TransportProtos.DeviceProfileDeleteMsg.newBuilder() }
.setProfileIdMSB(deviceProfile.getId().getId().getMostSignificantBits())
.setProfileIdLSB(deviceProfile.getId().getId().getLeastSignificantBits()) @Override
public void onTenantDelete(Tenant entity, TbQueueCallback callback) {
onEntityDelete(TenantId.SYS_TENANT_ID, entity.getId(), entity.getName(), callback);
}
public <T extends HasName> void onEntityChange(TenantId tenantId, EntityId entityid, T entity, TbQueueCallback callback) {
log.trace("[{}][{}][{}] Processing [{}] change event", tenantId, entityid.getEntityType(), entityid.getId(), entity.getName());
TransportProtos.EntityUpdateMsg entityUpdateMsg = TransportProtos.EntityUpdateMsg.newBuilder()
.setEntityType(entityid.getEntityType().name())
.setData(ByteString.copyFrom(encodingService.encode(entity))).build();
ToTransportMsg transportMsg = ToTransportMsg.newBuilder().setEntityUpdateMsg(entityUpdateMsg).build();
broadcast(transportMsg);
}
private void onEntityDelete(TenantId tenantId, EntityId entityId, String name, TbQueueCallback callback) {
log.trace("[{}][{}][{}] Processing [{}] delete event", tenantId, entityId.getEntityType(), entityId.getId(), name);
TransportProtos.EntityDeleteMsg entityDeleteMsg = TransportProtos.EntityDeleteMsg.newBuilder()
.setEntityType(entityId.getEntityType().name())
.setEntityIdMSB(entityId.getId().getMostSignificantBits())
.setEntityIdLSB(entityId.getId().getLeastSignificantBits())
.build(); .build();
ToTransportMsg transportMsg = ToTransportMsg.newBuilder().setDeviceProfileDeleteMsg(profileDeleteMsg).build(); ToTransportMsg transportMsg = ToTransportMsg.newBuilder().setEntityDeleteMsg(entityDeleteMsg).build();
broadcast(transportMsg); broadcast(transportMsg);
} }

13
application/src/main/java/org/thingsboard/server/service/queue/TbClusterService.java

@ -17,7 +17,8 @@ package org.thingsboard.server.service.queue;
import org.thingsboard.rule.engine.api.msg.ToDeviceActorNotificationMsg; import org.thingsboard.rule.engine.api.msg.ToDeviceActorNotificationMsg;
import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.Tenant;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent;
@ -53,5 +54,13 @@ public interface TbClusterService {
void onDeviceProfileChange(DeviceProfile deviceProfile, TbQueueCallback callback); void onDeviceProfileChange(DeviceProfile deviceProfile, TbQueueCallback callback);
void onDeviceProfileDelete(DeviceProfile deviceProfileId, TbQueueCallback callback); void onDeviceProfileDelete(DeviceProfile deviceProfile, TbQueueCallback callback);
void onTenantProfileChange(TenantProfile tenantProfile, TbQueueCallback callback);
void onTenantProfileDelete(TenantProfile tenantProfile, TbQueueCallback callback);
void onTenantChange(Tenant tenant, TbQueueCallback callback);
void onTenantDelete(Tenant tenant, TbQueueCallback callback);
} }

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

@ -28,6 +28,7 @@ import org.springframework.util.StringUtils;
import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.TenantProfile; import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.device.credentials.BasicMqttCredentials; import org.thingsboard.server.common.data.device.credentials.BasicMqttCredentials;
import org.thingsboard.server.common.data.device.credentials.ProvisionDeviceCredentialsData; import org.thingsboard.server.common.data.device.credentials.ProvisionDeviceCredentialsData;
@ -45,21 +46,19 @@ import org.thingsboard.server.common.msg.TbMsgDataType;
import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.common.transport.util.DataDecodingEncodingService; import org.thingsboard.server.common.transport.util.DataDecodingEncodingService;
import org.thingsboard.server.dao.device.DeviceCredentialsService; import org.thingsboard.server.dao.device.DeviceCredentialsService;
import org.thingsboard.server.dao.device.DeviceProfileService;
import org.thingsboard.server.dao.device.DeviceProvisionService; import org.thingsboard.server.dao.device.DeviceProvisionService;
import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.device.provision.ProvisionRequest; import org.thingsboard.server.dao.device.provision.ProvisionRequest;
import org.thingsboard.server.dao.device.provision.ProvisionResponse; import org.thingsboard.server.dao.device.provision.ProvisionResponse;
import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.dao.relation.RelationService;
import org.thingsboard.server.dao.tenant.TenantProfileService;
import org.thingsboard.server.dao.tenant.TenantService; import org.thingsboard.server.dao.tenant.TenantService;
import org.thingsboard.server.dao.util.mapping.JacksonUtil; import org.thingsboard.server.dao.util.mapping.JacksonUtil;
import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.DeviceInfoProto; import org.thingsboard.server.gen.transport.TransportProtos.DeviceInfoProto;
import org.thingsboard.server.gen.transport.TransportProtos.GetOrCreateDeviceFromGatewayRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.GetOrCreateDeviceFromGatewayRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.GetOrCreateDeviceFromGatewayResponseMsg; import org.thingsboard.server.gen.transport.TransportProtos.GetOrCreateDeviceFromGatewayResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.GetTenantRoutingInfoRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.GetEntityProfileRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.GetTenantRoutingInfoResponseMsg; import org.thingsboard.server.gen.transport.TransportProtos.GetEntityProfileResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.TransportApiRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.TransportApiRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.TransportApiResponseMsg; import org.thingsboard.server.gen.transport.TransportProtos.TransportApiResponseMsg;
@ -70,6 +69,7 @@ import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.dao.device.provision.ProvisionFailedException; import org.thingsboard.server.dao.device.provision.ProvisionFailedException;
import org.thingsboard.server.service.executors.DbCallbackExecutorService; import org.thingsboard.server.service.executors.DbCallbackExecutorService;
import org.thingsboard.server.service.profile.TbDeviceProfileCache;
import org.thingsboard.server.service.profile.TbTenantProfileCache; import org.thingsboard.server.service.profile.TbTenantProfileCache;
import org.thingsboard.server.service.queue.TbClusterService; import org.thingsboard.server.service.queue.TbClusterService;
import org.thingsboard.server.service.state.DeviceStateService; import org.thingsboard.server.service.state.DeviceStateService;
@ -90,9 +90,7 @@ public class DefaultTransportApiService implements TransportApiService {
private static final ObjectMapper mapper = new ObjectMapper(); private static final ObjectMapper mapper = new ObjectMapper();
//TODO: Constructor dependencies; private final TbDeviceProfileCache deviceProfileCache;
private final DeviceProfileService deviceProfileService;
private final TenantService tenantService;
private final TbTenantProfileCache tenantProfileCache; private final TbTenantProfileCache tenantProfileCache;
private final DeviceService deviceService; private final DeviceService deviceService;
private final RelationService relationService; private final RelationService relationService;
@ -103,17 +101,15 @@ public class DefaultTransportApiService implements TransportApiService {
private final DataDecodingEncodingService dataDecodingEncodingService; private final DataDecodingEncodingService dataDecodingEncodingService;
private final DeviceProvisionService deviceProvisionService; private final DeviceProvisionService deviceProvisionService;
private final ConcurrentMap<String, ReentrantLock> deviceCreationLocks = new ConcurrentHashMap<>(); private final ConcurrentMap<String, ReentrantLock> deviceCreationLocks = new ConcurrentHashMap<>();
public DefaultTransportApiService(DeviceProfileService deviceProfileService, TenantService tenantService, public DefaultTransportApiService(TbDeviceProfileCache deviceProfileCache,
TbTenantProfileCache tenantProfileCache, DeviceService deviceService, TbTenantProfileCache tenantProfileCache, DeviceService deviceService,
RelationService relationService, DeviceCredentialsService deviceCredentialsService, RelationService relationService, DeviceCredentialsService deviceCredentialsService,
DeviceStateService deviceStateService, DbCallbackExecutorService dbCallbackExecutorService, DeviceStateService deviceStateService, DbCallbackExecutorService dbCallbackExecutorService,
TbClusterService tbClusterService, DataDecodingEncodingService dataDecodingEncodingService, TbClusterService tbClusterService, DataDecodingEncodingService dataDecodingEncodingService,
DeviceProvisionService deviceProvisionService) { DeviceProvisionService deviceProvisionService) {
this.deviceProfileService = deviceProfileService; this.deviceProfileCache = deviceProfileCache;
this.tenantService = tenantService;
this.tenantProfileCache = tenantProfileCache; this.tenantProfileCache = tenantProfileCache;
this.deviceService = deviceService; this.deviceService = deviceService;
this.relationService = relationService; this.relationService = relationService;
@ -143,11 +139,8 @@ public class DefaultTransportApiService implements TransportApiService {
} else if (transportApiRequestMsg.hasGetOrCreateDeviceRequestMsg()) { } else if (transportApiRequestMsg.hasGetOrCreateDeviceRequestMsg()) {
return Futures.transform(handle(transportApiRequestMsg.getGetOrCreateDeviceRequestMsg()), return Futures.transform(handle(transportApiRequestMsg.getGetOrCreateDeviceRequestMsg()),
value -> new TbProtoQueueMsg<>(tbProtoQueueMsg.getKey(), value, tbProtoQueueMsg.getHeaders()), MoreExecutors.directExecutor()); value -> new TbProtoQueueMsg<>(tbProtoQueueMsg.getKey(), value, tbProtoQueueMsg.getHeaders()), MoreExecutors.directExecutor());
} else if (transportApiRequestMsg.hasGetTenantRoutingInfoRequestMsg()) { } else if (transportApiRequestMsg.hasEntityProfileRequestMsg()) {
return Futures.transform(handle(transportApiRequestMsg.getGetTenantRoutingInfoRequestMsg()), return Futures.transform(handle(transportApiRequestMsg.getEntityProfileRequestMsg()),
value -> new TbProtoQueueMsg<>(tbProtoQueueMsg.getKey(), value, tbProtoQueueMsg.getHeaders()), MoreExecutors.directExecutor());
} else if (transportApiRequestMsg.hasGetDeviceProfileRequestMsg()) {
return Futures.transform(handle(transportApiRequestMsg.getGetDeviceProfileRequestMsg()),
value -> new TbProtoQueueMsg<>(tbProtoQueueMsg.getKey(), value, tbProtoQueueMsg.getHeaders()), MoreExecutors.directExecutor()); value -> new TbProtoQueueMsg<>(tbProtoQueueMsg.getKey(), value, tbProtoQueueMsg.getHeaders()), MoreExecutors.directExecutor());
} else if (transportApiRequestMsg.hasProvisionDeviceRequestMsg()) { } else if (transportApiRequestMsg.hasProvisionDeviceRequestMsg()) {
return Futures.transform(handle(transportApiRequestMsg.getProvisionDeviceRequestMsg()), return Futures.transform(handle(transportApiRequestMsg.getProvisionDeviceRequestMsg()),
@ -238,7 +231,7 @@ public class DefaultTransportApiService implements TransportApiService {
device.setName(requestMsg.getDeviceName()); device.setName(requestMsg.getDeviceName());
device.setType(requestMsg.getDeviceType()); device.setType(requestMsg.getDeviceType());
device.setCustomerId(gateway.getCustomerId()); device.setCustomerId(gateway.getCustomerId());
DeviceProfile deviceProfile = deviceProfileService.findOrCreateDeviceProfile(gateway.getTenantId(), requestMsg.getDeviceType()); DeviceProfile deviceProfile = deviceProfileCache.findOrCreateDeviceProfile(gateway.getTenantId(), requestMsg.getDeviceType());
device.setDeviceProfileId(deviceProfile.getId()); device.setDeviceProfileId(deviceProfile.getId());
device = deviceService.saveDevice(device); device = deviceService.saveDevice(device);
relationService.saveRelationAsync(TenantId.SYS_TENANT_ID, new EntityRelation(gateway.getId(), device.getId(), "Created")); relationService.saveRelationAsync(TenantId.SYS_TENANT_ID, new EntityRelation(gateway.getId(), device.getId(), "Created"));
@ -258,7 +251,7 @@ public class DefaultTransportApiService implements TransportApiService {
} }
GetOrCreateDeviceFromGatewayResponseMsg.Builder builder = GetOrCreateDeviceFromGatewayResponseMsg.newBuilder() GetOrCreateDeviceFromGatewayResponseMsg.Builder builder = GetOrCreateDeviceFromGatewayResponseMsg.newBuilder()
.setDeviceInfo(getDeviceInfoProto(device)); .setDeviceInfo(getDeviceInfoProto(device));
DeviceProfile deviceProfile = deviceProfileService.findDeviceProfileById(device.getTenantId(), device.getDeviceProfileId()); DeviceProfile deviceProfile = deviceProfileCache.get(device.getTenantId(), device.getDeviceProfileId());
if (deviceProfile != null) { if (deviceProfile != null) {
builder.setProfileBody(ByteString.copyFrom(dataDecodingEncodingService.encode(deviceProfile))); builder.setProfileBody(ByteString.copyFrom(dataDecodingEncodingService.encode(deviceProfile)));
} else { } else {
@ -320,23 +313,21 @@ public class DefaultTransportApiService implements TransportApiService {
.build(); .build();
} }
private ListenableFuture<TransportApiResponseMsg> handle(GetTenantRoutingInfoRequestMsg requestMsg) { private ListenableFuture<TransportApiResponseMsg> handle(GetEntityProfileRequestMsg requestMsg) {
TenantId tenantId = new TenantId(new UUID(requestMsg.getTenantIdMSB(), requestMsg.getTenantIdLSB())); EntityType entityType = EntityType.valueOf(requestMsg.getEntityType());
UUID entityUuid = new UUID(requestMsg.getEntityIdMSB(), requestMsg.getEntityIdLSB());
ListenableFuture<TenantProfile> tenantProfileFuture = Futures.immediateFuture(tenantProfileCache.get(tenantId)); ByteString data;
return Futures.transform(tenantProfileFuture, tenantProfile -> TransportApiResponseMsg.newBuilder() if (entityType.equals(EntityType.DEVICE_PROFILE)) {
.setGetTenantRoutingInfoResponseMsg(GetTenantRoutingInfoResponseMsg.newBuilder().setIsolatedTbCore(tenantProfile.isIsolatedTbCore()) DeviceProfileId deviceProfileId = new DeviceProfileId(entityUuid);
.setIsolatedTbRuleEngine(tenantProfile.isIsolatedTbRuleEngine()).build()).build(), dbCallbackExecutorService); DeviceProfile deviceProfile = deviceProfileCache.find(deviceProfileId);
} data = ByteString.copyFrom(dataDecodingEncodingService.encode(deviceProfile));
} else if (entityType.equals(EntityType.TENANT)) {
private ListenableFuture<TransportApiResponseMsg> handle(TransportProtos.GetDeviceProfileRequestMsg requestMsg) { TenantProfile tenantProfile = tenantProfileCache.get(new TenantId(entityUuid));
DeviceProfileId profileId = new DeviceProfileId(new UUID(requestMsg.getProfileIdMSB(), requestMsg.getProfileIdLSB())); data = ByteString.copyFrom(dataDecodingEncodingService.encode(tenantProfile));
DeviceProfile deviceProfile = deviceProfileService.findDeviceProfileById(TenantId.SYS_TENANT_ID, profileId); } else {
return Futures.immediateFuture(TransportApiResponseMsg.newBuilder() throw new RuntimeException("Invalid entity profile request: " + entityType);
.setGetDeviceProfileResponseMsg( }
TransportProtos.GetDeviceProfileResponseMsg.newBuilder() return Futures.immediateFuture(TransportApiResponseMsg.newBuilder().setEntityProfileResponseMsg(GetEntityProfileResponseMsg.newBuilder().setData(data).build()).build());
.setData(ByteString.copyFrom(dataDecodingEncodingService.encode(deviceProfile)))
.build()).build());
} }
private ListenableFuture<TransportApiResponseMsg> getDeviceInfo(DeviceId deviceId, DeviceCredentials credentials) { private ListenableFuture<TransportApiResponseMsg> getDeviceInfo(DeviceId deviceId, DeviceCredentials credentials) {
@ -348,7 +339,7 @@ public class DefaultTransportApiService implements TransportApiService {
try { try {
ValidateDeviceCredentialsResponseMsg.Builder builder = ValidateDeviceCredentialsResponseMsg.newBuilder(); ValidateDeviceCredentialsResponseMsg.Builder builder = ValidateDeviceCredentialsResponseMsg.newBuilder();
builder.setDeviceInfo(getDeviceInfoProto(device)); builder.setDeviceInfo(getDeviceInfoProto(device));
DeviceProfile deviceProfile = deviceProfileService.findDeviceProfileById(device.getTenantId(), device.getDeviceProfileId()); DeviceProfile deviceProfile = deviceProfileCache.get(device.getTenantId(), device.getDeviceProfileId());
if (deviceProfile != null) { if (deviceProfile != null) {
builder.setProfileBody(ByteString.copyFrom(dataDecodingEncodingService.encode(deviceProfile))); builder.setProfileBody(ByteString.copyFrom(dataDecodingEncodingService.encode(deviceProfile)));
} else { } else {

6
application/src/main/resources/thingsboard.yml

@ -502,7 +502,7 @@ js:
remote: remote:
# Maximum allowed JavaScript execution errors before JavaScript will be blacklisted # Maximum allowed JavaScript execution errors before JavaScript will be blacklisted
max_errors: "${REMOTE_JS_SANDBOX_MAX_ERRORS:3}" max_errors: "${REMOTE_JS_SANDBOX_MAX_ERRORS:3}"
# Maximum time in seconds for black listed function to stay in the list. # Maximum time in seconds for black listed function to stay in 1:the list.
max_black_list_duration_sec: "${REMOTE_JS_SANDBOX_MAX_BLACKLIST_DURATION_SEC:60}" max_black_list_duration_sec: "${REMOTE_JS_SANDBOX_MAX_BLACKLIST_DURATION_SEC:60}"
stats: stats:
enabled: "${TB_JS_REMOTE_STATS_ENABLED:false}" enabled: "${TB_JS_REMOTE_STATS_ENABLED:false}"
@ -512,10 +512,6 @@ transport:
sessions: sessions:
inactivity_timeout: "${TB_TRANSPORT_SESSIONS_INACTIVITY_TIMEOUT:300000}" inactivity_timeout: "${TB_TRANSPORT_SESSIONS_INACTIVITY_TIMEOUT:300000}"
report_timeout: "${TB_TRANSPORT_SESSIONS_REPORT_TIMEOUT:30000}" report_timeout: "${TB_TRANSPORT_SESSIONS_REPORT_TIMEOUT:30000}"
rate_limits:
enabled: "${TB_TRANSPORT_RATE_LIMITS_ENABLED:false}"
tenant: "${TB_TRANSPORT_RATE_LIMITS_TENANT:1000:1,20000:60}"
device: "${TB_TRANSPORT_RATE_LIMITS_DEVICE:10:1,300:60}"
json: json:
# Cast String data types to Numeric if possible when processing Telemetry/Attributes JSON # Cast String data types to Numeric if possible when processing Telemetry/Attributes JSON
type_cast_enabled: "${JSON_TYPE_CAST_ENABLED:true}" type_cast_enabled: "${JSON_TYPE_CAST_ENABLED:true}"

47
common/queue/src/main/proto/queue.proto

@ -177,32 +177,26 @@ message GetOrCreateDeviceFromGatewayResponseMsg {
bytes profileBody = 2; bytes profileBody = 2;
} }
message GetTenantRoutingInfoRequestMsg { message GetEntityProfileRequestMsg {
int64 tenantIdMSB = 1; string entityType = 1;
int64 tenantIdLSB = 2; int64 entityIdMSB = 2;
} int64 entityIdLSB = 3;
message GetTenantRoutingInfoResponseMsg {
bool isolatedTbCore = 1;
bool isolatedTbRuleEngine = 2;
}
message GetDeviceProfileRequestMsg {
int64 profileIdMSB = 1;
int64 profileIdLSB = 2;
} }
message GetDeviceProfileResponseMsg { message GetEntityProfileResponseMsg {
bytes data = 1; string entityType = 1;
bytes data = 2;
} }
message DeviceProfileUpdateMsg { message EntityUpdateMsg {
bytes data = 1; string entityType = 1;
bytes data = 2;
} }
message DeviceProfileDeleteMsg { message EntityDeleteMsg {
int64 profileIdMSB = 1; string entityType = 1;
int64 profileIdLSB = 2; int64 entityIdMSB = 2;
int64 entityIdLSB = 3;
} }
message SessionCloseNotificationProto { message SessionCloseNotificationProto {
@ -482,8 +476,7 @@ message TransportApiRequestMsg {
ValidateDeviceTokenRequestMsg validateTokenRequestMsg = 1; ValidateDeviceTokenRequestMsg validateTokenRequestMsg = 1;
ValidateDeviceX509CertRequestMsg validateX509CertRequestMsg = 2; ValidateDeviceX509CertRequestMsg validateX509CertRequestMsg = 2;
GetOrCreateDeviceFromGatewayRequestMsg getOrCreateDeviceRequestMsg = 3; GetOrCreateDeviceFromGatewayRequestMsg getOrCreateDeviceRequestMsg = 3;
GetTenantRoutingInfoRequestMsg getTenantRoutingInfoRequestMsg = 4; GetEntityProfileRequestMsg entityProfileRequestMsg = 4;
GetDeviceProfileRequestMsg getDeviceProfileRequestMsg = 5;
ValidateBasicMqttCredRequestMsg validateBasicMqttCredRequestMsg = 6; ValidateBasicMqttCredRequestMsg validateBasicMqttCredRequestMsg = 6;
ProvisionDeviceRequestMsg provisionDeviceRequestMsg = 7; ProvisionDeviceRequestMsg provisionDeviceRequestMsg = 7;
} }
@ -492,9 +485,8 @@ message TransportApiRequestMsg {
message TransportApiResponseMsg { message TransportApiResponseMsg {
ValidateDeviceCredentialsResponseMsg validateCredResponseMsg = 1; ValidateDeviceCredentialsResponseMsg validateCredResponseMsg = 1;
GetOrCreateDeviceFromGatewayResponseMsg getOrCreateDeviceResponseMsg = 2; GetOrCreateDeviceFromGatewayResponseMsg getOrCreateDeviceResponseMsg = 2;
GetTenantRoutingInfoResponseMsg getTenantRoutingInfoResponseMsg = 4; GetEntityProfileResponseMsg entityProfileResponseMsg = 3;
GetDeviceProfileResponseMsg getDeviceProfileResponseMsg = 5; ProvisionDeviceResponseMsg provisionDeviceResponseMsg = 4;
ProvisionDeviceResponseMsg provisionDeviceResponseMsg = 6;
} }
/* Messages that are handled by ThingsBoard Core Service */ /* Messages that are handled by ThingsBoard Core Service */
@ -535,7 +527,8 @@ message ToTransportMsg {
AttributeUpdateNotificationMsg attributeUpdateNotification = 5; AttributeUpdateNotificationMsg attributeUpdateNotification = 5;
ToDeviceRpcRequestMsg toDeviceRequest = 6; ToDeviceRpcRequestMsg toDeviceRequest = 6;
ToServerRpcResponseMsg toServerResponse = 7; ToServerRpcResponseMsg toServerResponse = 7;
DeviceProfileUpdateMsg deviceProfileUpdateMsg = 8; /* For Tenant, TenantProfile and DeviceProfile */
DeviceProfileDeleteMsg deviceProfileDeleteMsg = 9; EntityUpdateMsg entityUpdateMsg = 8;
EntityDeleteMsg entityDeleteMsg = 9;
ProvisionDeviceResponseMsg provisionResponse = 10; ProvisionDeviceResponseMsg provisionResponse = 10;
} }

2
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportProfileCache.java → common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportDeviceProfileCache.java

@ -21,7 +21,7 @@ import org.thingsboard.server.common.data.id.DeviceProfileId;
import java.util.Optional; import java.util.Optional;
public interface TransportProfileCache { public interface TransportDeviceProfileCache {
DeviceProfile getOrCreate(DeviceProfileId id, ByteString profileBody); DeviceProfile getOrCreate(DeviceProfileId id, ByteString profileBody);

9
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java

@ -20,11 +20,12 @@ import org.thingsboard.server.common.data.DeviceTransportType;
import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.DeviceProfileId;
import org.thingsboard.server.common.transport.auth.GetOrCreateDeviceFromGatewayResponse; import org.thingsboard.server.common.transport.auth.GetOrCreateDeviceFromGatewayResponse;
import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse; import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ClaimDeviceMsg; import org.thingsboard.server.gen.transport.TransportProtos.ClaimDeviceMsg;
import org.thingsboard.server.gen.transport.TransportProtos.GetAttributeRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.GetAttributeRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.GetOrCreateDeviceFromGatewayRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.GetOrCreateDeviceFromGatewayRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.GetTenantRoutingInfoRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.GetEntityProfileRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.GetTenantRoutingInfoResponseMsg; import org.thingsboard.server.gen.transport.TransportProtos.GetEntityProfileResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.PostAttributeMsg; import org.thingsboard.server.gen.transport.TransportProtos.PostAttributeMsg;
import org.thingsboard.server.gen.transport.TransportProtos.PostTelemetryMsg; import org.thingsboard.server.gen.transport.TransportProtos.PostTelemetryMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceRequestMsg;
@ -47,7 +48,7 @@ import java.util.concurrent.ScheduledExecutorService;
*/ */
public interface TransportService { public interface TransportService {
GetTenantRoutingInfoResponseMsg getRoutingInfo(GetTenantRoutingInfoRequestMsg msg); GetEntityProfileResponseMsg getRoutingInfo(GetEntityProfileRequestMsg msg);
void process(DeviceTransportType transportType, ValidateDeviceTokenRequestMsg msg, void process(DeviceTransportType transportType, ValidateDeviceTokenRequestMsg msg,
TransportServiceCallback<ValidateDeviceCredentialsResponse> callback); TransportServiceCallback<ValidateDeviceCredentialsResponse> callback);
@ -64,8 +65,6 @@ public interface TransportService {
void process(ProvisionDeviceRequestMsg msg, void process(ProvisionDeviceRequestMsg msg,
TransportServiceCallback<ProvisionDeviceResponseMsg> callback); TransportServiceCallback<ProvisionDeviceResponseMsg> callback);
void getDeviceProfile(DeviceProfileId deviceProfileId, TransportServiceCallback<DeviceProfile> callback);
void onProfileUpdate(DeviceProfile deviceProfile); void onProfileUpdate(DeviceProfile deviceProfile);
boolean checkLimits(SessionInfoProto sessionInfo, Object msg, TransportServiceCallback<Void> callback); boolean checkLimits(SessionInfoProto sessionInfo, Object msg, TransportServiceCallback<Void> callback);

38
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportTenantProfileCache.java

@ -0,0 +1,38 @@
/**
* Copyright © 2016-2020 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport;
import com.google.protobuf.ByteString;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.TenantProfileId;
import org.thingsboard.server.common.transport.profile.TenantProfileUpdateResult;
import org.thingsboard.server.queue.discovery.TenantRoutingInfo;
import org.thingsboard.server.queue.discovery.TenantRoutingInfoService;
import java.util.Set;
public interface TransportTenantProfileCache {
TenantProfile get(TenantId tenantId);
TenantProfileUpdateResult put(ByteString profileBody);
boolean put(TenantId tenantId, TenantProfileId profileId);
Set<TenantId> remove(TenantProfileId profileId);
}

47
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/DefaultTransportRateLimitFactory.java

@ -0,0 +1,47 @@
/**
* Copyright © 2016-2020 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.limits;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import org.springframework.util.StringUtils;
import org.thingsboard.server.common.msg.tools.TbRateLimits;
@Slf4j
@Component
public class DefaultTransportRateLimitFactory implements TransportRateLimitFactory {
private static final DummyTransportRateLimit ALWAYS_TRUE = new DummyTransportRateLimit();
@Override
public TransportRateLimit create(TransportRateLimitType type, Object configuration) {
if (!StringUtils.isEmpty(configuration)) {
try {
return new SimpleTransportRateLimit(new TbRateLimits(configuration.toString()), configuration.toString());
} catch (Exception e) {
log.warn("[{}] Failed to init rate limit with configuration: {}", type, configuration, e);
return ALWAYS_TRUE;
}
} else {
return ALWAYS_TRUE;
}
}
@Override
public TransportRateLimit createDefault(TransportRateLimitType type) {
return ALWAYS_TRUE;
}
}

115
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/DefaultTransportRateLimitService.java

@ -0,0 +1,115 @@
/**
* Copyright © 2016-2020 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.limits;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.TenantProfileData;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.transport.TransportTenantProfileCache;
import org.thingsboard.server.common.transport.profile.TenantProfileUpdateResult;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
@Service
@Slf4j
public class DefaultTransportRateLimitService implements TransportRateLimitService {
private final ConcurrentMap<TenantId, TransportRateLimit[]> perTenantLimits = new ConcurrentHashMap<>();
private final ConcurrentMap<DeviceId, TransportRateLimit[]> perDeviceLimits = new ConcurrentHashMap<>();
private final TransportRateLimitFactory rateLimitFactory;
private final TransportTenantProfileCache tenantProfileCache;
public DefaultTransportRateLimitService(TransportRateLimitFactory rateLimitFactory, TransportTenantProfileCache tenantProfileCache) {
this.rateLimitFactory = rateLimitFactory;
this.tenantProfileCache = tenantProfileCache;
}
@Override
public TransportRateLimit getRateLimit(TenantId tenantId, TransportRateLimitType limitType) {
TransportRateLimit[] limits = perTenantLimits.get(tenantId);
if (limits == null) {
limits = fetchProfileAndInit(tenantId);
perTenantLimits.put(tenantId, limits);
}
return limits[limitType.ordinal()];
}
@Override
public TransportRateLimit getRateLimit(TenantId tenantId, DeviceId deviceId, TransportRateLimitType limitType) {
TransportRateLimit[] limits = perDeviceLimits.get(deviceId);
if (limits == null) {
limits = fetchProfileAndInit(tenantId);
perDeviceLimits.put(deviceId, limits);
}
return limits[limitType.ordinal()];
}
@Override
public void update(TenantProfileUpdateResult update) {
TransportRateLimit[] newLimits = createTransportRateLimits(update.getProfile());
for (TenantId tenantId : update.getAffectedTenants()) {
mergeLimits(tenantId, newLimits);
}
}
@Override
public void update(TenantId tenantId) {
mergeLimits(tenantId, fetchProfileAndInit(tenantId));
}
public void mergeLimits(TenantId tenantId, TransportRateLimit[] newRateLimits) {
TransportRateLimit[] oldRateLimits = perTenantLimits.get(tenantId);
if (oldRateLimits == null) {
perTenantLimits.put(tenantId, newRateLimits);
} else {
for (int i = 0; i < TransportRateLimitType.values().length; i++) {
TransportRateLimit newLimit = newRateLimits[i];
TransportRateLimit oldLimit = oldRateLimits[i];
if (newLimit != null && (oldLimit == null || !oldLimit.getConfiguration().equals(newLimit.getConfiguration()))) {
oldRateLimits[i] = newLimit;
}
}
}
}
@Override
public void remove(TenantId tenantId) {
perTenantLimits.remove(tenantId);
}
@Override
public void remove(DeviceId deviceId) {
perDeviceLimits.remove(deviceId);
}
private TransportRateLimit[] fetchProfileAndInit(TenantId tenantId) {
return perTenantLimits.computeIfAbsent(tenantId, tmp -> createTransportRateLimits(tenantProfileCache.get(tenantId)));
}
private TransportRateLimit[] createTransportRateLimits(TenantProfile tenantProfile) {
TenantProfileData profileData = tenantProfile.getProfileData();
TransportRateLimit[] rateLimits = new TransportRateLimit[TransportRateLimitType.values().length];
for (TransportRateLimitType type : TransportRateLimitType.values()) {
rateLimits[type.ordinal()] = rateLimitFactory.create(type, profileData.getProperties().get(type.getConfigurationKey()));
}
return rateLimits;
}
}

30
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/DummyTransportRateLimit.java

@ -0,0 +1,30 @@
/**
* Copyright © 2016-2020 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.limits;
public class DummyTransportRateLimit implements TransportRateLimit {
@Override
public String getConfiguration() {
return "";
}
@Override
public boolean tryConsume() {
return true;
}
}

34
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/SimpleTransportRateLimit.java

@ -0,0 +1,34 @@
/**
* Copyright © 2016-2020 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.limits;
import lombok.Getter;
import lombok.RequiredArgsConstructor;
import org.thingsboard.server.common.msg.tools.TbRateLimits;
@RequiredArgsConstructor
public class SimpleTransportRateLimit implements TransportRateLimit {
private final TbRateLimits rateLimit;
@Getter
private final String configuration;
@Override
public boolean tryConsume() {
return rateLimit.tryConsume();
}
}

24
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportRateLimit.java

@ -0,0 +1,24 @@
/**
* Copyright © 2016-2020 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.limits;
public interface TransportRateLimit {
String getConfiguration();
boolean tryConsume();
}

24
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportRateLimitFactory.java

@ -0,0 +1,24 @@
/**
* Copyright © 2016-2020 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.limits;
public interface TransportRateLimitFactory {
TransportRateLimit create(TransportRateLimitType type, Object config);
TransportRateLimit createDefault(TransportRateLimitType type);
}

36
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportRateLimitService.java

@ -0,0 +1,36 @@
/**
* Copyright © 2016-2020 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.limits;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.transport.profile.TenantProfileUpdateResult;
public interface TransportRateLimitService {
TransportRateLimit getRateLimit(TenantId tenantId, TransportRateLimitType limit);
TransportRateLimit getRateLimit(TenantId tenantId, DeviceId deviceId, TransportRateLimitType limit);
void update(TenantProfileUpdateResult update);
void update(TenantId tenantId);
void remove(TenantId tenantId);
void remove(DeviceId deviceId);
}

33
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportRateLimitType.java

@ -0,0 +1,33 @@
/**
* Copyright © 2016-2020 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.limits;
import lombok.Getter;
public enum TransportRateLimitType {
TENANT_MAX_MSGS("transport.tenant.max.msg"),
TENANT_MAX_DATA_POINTS("transport.tenant.max.dataPoints"),
DEVICE_MAX_MSGS("transport.device.max.msg"),
DEVICE_MAX_DATA_POINTS("transport.device.max.dataPoints");
@Getter
private final String configurationKey;
TransportRateLimitType(String configurationKey) {
this.configurationKey = configurationKey;
}
}

30
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/profile/TenantProfileUpdateResult.java

@ -0,0 +1,30 @@
/**
* Copyright © 2016-2020 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.profile;
import lombok.Data;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.id.TenantId;
import java.util.Set;
@Data
public class TenantProfileUpdateResult {
private final TenantProfile profile;
private final Set<TenantId> affectedTenants;
}

6
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportProfileCache.java → common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportDeviceProfileCache.java

@ -21,7 +21,7 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.DeviceProfileId;
import org.thingsboard.server.common.transport.TransportProfileCache; import org.thingsboard.server.common.transport.TransportDeviceProfileCache;
import org.thingsboard.server.common.transport.util.DataDecodingEncodingService; import org.thingsboard.server.common.transport.util.DataDecodingEncodingService;
import java.util.Optional; import java.util.Optional;
@ -31,13 +31,13 @@ import java.util.concurrent.ConcurrentMap;
@Slf4j @Slf4j
@Component @Component
@ConditionalOnExpression("('${service.type:null}'=='monolith' && '${transport.api_enabled:true}'=='true') || '${service.type:null}'=='tb-transport'") @ConditionalOnExpression("('${service.type:null}'=='monolith' && '${transport.api_enabled:true}'=='true') || '${service.type:null}'=='tb-transport'")
public class DefaultTransportProfileCache implements TransportProfileCache { public class DefaultTransportDeviceProfileCache implements TransportDeviceProfileCache {
private final ConcurrentMap<DeviceProfileId, DeviceProfile> deviceProfiles = new ConcurrentHashMap<>(); private final ConcurrentMap<DeviceProfileId, DeviceProfile> deviceProfiles = new ConcurrentHashMap<>();
private final DataDecodingEncodingService dataDecodingEncodingService; private final DataDecodingEncodingService dataDecodingEncodingService;
public DefaultTransportProfileCache(DataDecodingEncodingService dataDecodingEncodingService) { public DefaultTransportDeviceProfileCache(DataDecodingEncodingService dataDecodingEncodingService) {
this.dataDecodingEncodingService = dataDecodingEncodingService; this.dataDecodingEncodingService = dataDecodingEncodingService;
} }

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

@ -29,25 +29,35 @@ import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.DeviceTransportType; import org.thingsboard.server.common.data.DeviceTransportType;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.Tenant;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.DeviceProfileId;
import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.TenantProfileId;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.common.msg.queue.ServiceQueue; import org.thingsboard.server.common.msg.queue.ServiceQueue;
import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.common.msg.session.SessionMsgType; import org.thingsboard.server.common.msg.session.SessionMsgType;
import org.thingsboard.server.common.msg.tools.TbRateLimits;
import org.thingsboard.server.common.msg.tools.TbRateLimitsException; import org.thingsboard.server.common.msg.tools.TbRateLimitsException;
import org.thingsboard.server.common.stats.MessagesStats;
import org.thingsboard.server.common.stats.StatsFactory;
import org.thingsboard.server.common.stats.StatsType;
import org.thingsboard.server.common.transport.SessionMsgListener; import org.thingsboard.server.common.transport.SessionMsgListener;
import org.thingsboard.server.common.transport.TransportProfileCache; import org.thingsboard.server.common.transport.TransportDeviceProfileCache;
import org.thingsboard.server.common.transport.TransportService; import org.thingsboard.server.common.transport.TransportService;
import org.thingsboard.server.common.transport.TransportServiceCallback; import org.thingsboard.server.common.transport.TransportServiceCallback;
import org.thingsboard.server.common.transport.TransportTenantProfileCache;
import org.thingsboard.server.common.transport.auth.GetOrCreateDeviceFromGatewayResponse; import org.thingsboard.server.common.transport.auth.GetOrCreateDeviceFromGatewayResponse;
import org.thingsboard.server.common.transport.auth.TransportDeviceInfo; import org.thingsboard.server.common.transport.auth.TransportDeviceInfo;
import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse; import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse;
import org.thingsboard.server.common.transport.limits.TransportRateLimit;
import org.thingsboard.server.common.transport.limits.TransportRateLimitService;
import org.thingsboard.server.common.transport.limits.TransportRateLimitType;
import org.thingsboard.server.common.transport.profile.TenantProfileUpdateResult;
import org.thingsboard.server.common.transport.util.DataDecodingEncodingService;
import org.thingsboard.server.common.transport.util.JsonUtils; import org.thingsboard.server.common.transport.util.JsonUtils;
import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceRequestMsg;
@ -69,15 +79,13 @@ import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
import org.thingsboard.server.queue.provider.TbQueueProducerProvider; import org.thingsboard.server.queue.provider.TbQueueProducerProvider;
import org.thingsboard.server.queue.provider.TbTransportQueueFactory; import org.thingsboard.server.queue.provider.TbTransportQueueFactory;
import org.thingsboard.server.common.stats.MessagesStats;
import org.thingsboard.server.common.stats.StatsFactory;
import org.thingsboard.server.common.stats.StatsType;
import javax.annotation.PostConstruct; import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy; import javax.annotation.PreDestroy;
import java.util.Collections; import java.util.Collections;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.Optional;
import java.util.Random; import java.util.Random;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
@ -98,12 +106,6 @@ import java.util.concurrent.atomic.AtomicInteger;
@ConditionalOnExpression("('${service.type:null}'=='monolith' && '${transport.api_enabled:true}'=='true') || '${service.type:null}'=='tb-transport'") @ConditionalOnExpression("('${service.type:null}'=='monolith' && '${transport.api_enabled:true}'=='true') || '${service.type:null}'=='tb-transport'")
public class DefaultTransportService implements TransportService { public class DefaultTransportService implements TransportService {
@Value("${transport.rate_limits.enabled}")
private boolean rateLimitEnabled;
@Value("${transport.rate_limits.tenant}")
private String perTenantLimitsConf;
@Value("${transport.rate_limits.device}")
private String perDevicesLimitsConf;
@Value("${transport.sessions.inactivity_timeout}") @Value("${transport.sessions.inactivity_timeout}")
private long sessionInactivityTimeout; private long sessionInactivityTimeout;
@Value("${transport.sessions.report_timeout}") @Value("${transport.sessions.report_timeout}")
@ -119,7 +121,10 @@ public class DefaultTransportService implements TransportService {
private final PartitionService partitionService; private final PartitionService partitionService;
private final TbServiceInfoProvider serviceInfoProvider; private final TbServiceInfoProvider serviceInfoProvider;
private final StatsFactory statsFactory; private final StatsFactory statsFactory;
private final TransportProfileCache transportProfileCache; private final TransportDeviceProfileCache deviceProfileCache;
private final TransportTenantProfileCache tenantProfileCache;
private final TransportRateLimitService rateLimitService;
private final DataDecodingEncodingService dataDecodingEncodingService;
protected TbQueueRequestTemplate<TbProtoQueueMsg<TransportApiRequestMsg>, TbProtoQueueMsg<TransportApiResponseMsg>> transportApiRequestTemplate; protected TbQueueRequestTemplate<TbProtoQueueMsg<TransportApiRequestMsg>, TbProtoQueueMsg<TransportApiResponseMsg>> transportApiRequestTemplate;
protected TbQueueProducer<TbProtoQueueMsg<ToRuleEngineMsg>> ruleEngineMsgProducer; protected TbQueueProducer<TbProtoQueueMsg<ToRuleEngineMsg>> ruleEngineMsgProducer;
@ -132,14 +137,11 @@ public class DefaultTransportService implements TransportService {
protected ScheduledExecutorService schedulerExecutor; protected ScheduledExecutorService schedulerExecutor;
protected ExecutorService transportCallbackExecutor; protected ExecutorService transportCallbackExecutor;
private ExecutorService mainConsumerExecutor;
private final ConcurrentMap<UUID, SessionMetaData> sessions = new ConcurrentHashMap<>(); private final ConcurrentMap<UUID, SessionMetaData> sessions = new ConcurrentHashMap<>();
private final Map<String, RpcRequestMetadata> toServerRpcPendingMap = new ConcurrentHashMap<>(); private final Map<String, RpcRequestMetadata> toServerRpcPendingMap = new ConcurrentHashMap<>();
//TODO 3.2: @ybondarenko Implement cleanup of this maps.
private final ConcurrentMap<TenantId, TbRateLimits> perTenantLimits = new ConcurrentHashMap<>();
private final ConcurrentMap<DeviceId, TbRateLimits> perDeviceLimits = new ConcurrentHashMap<>();
private ExecutorService mainConsumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("transport-consumer"));
private volatile boolean stopped = false; private volatile boolean stopped = false;
public DefaultTransportService(TbServiceInfoProvider serviceInfoProvider, public DefaultTransportService(TbServiceInfoProvider serviceInfoProvider,
@ -147,22 +149,22 @@ public class DefaultTransportService implements TransportService {
TbQueueProducerProvider producerProvider, TbQueueProducerProvider producerProvider,
PartitionService partitionService, PartitionService partitionService,
StatsFactory statsFactory, StatsFactory statsFactory,
TransportProfileCache transportProfileCache) { TransportDeviceProfileCache deviceProfileCache,
TransportTenantProfileCache tenantProfileCache,
TransportRateLimitService rateLimitService, DataDecodingEncodingService dataDecodingEncodingService) {
this.serviceInfoProvider = serviceInfoProvider; this.serviceInfoProvider = serviceInfoProvider;
this.queueProvider = queueProvider; this.queueProvider = queueProvider;
this.producerProvider = producerProvider; this.producerProvider = producerProvider;
this.partitionService = partitionService; this.partitionService = partitionService;
this.statsFactory = statsFactory; this.statsFactory = statsFactory;
this.transportProfileCache = transportProfileCache; this.deviceProfileCache = deviceProfileCache;
this.tenantProfileCache = tenantProfileCache;
this.rateLimitService = rateLimitService;
this.dataDecodingEncodingService = dataDecodingEncodingService;
} }
@PostConstruct @PostConstruct
public void init() { public void init() {
if (rateLimitEnabled) {
//Just checking the configuration parameters
new TbRateLimits(perTenantLimitsConf);
new TbRateLimits(perDevicesLimitsConf);
}
this.ruleEngineProducerStats = statsFactory.createMessagesStats(StatsType.RULE_ENGINE.getName() + ".producer"); this.ruleEngineProducerStats = statsFactory.createMessagesStats(StatsType.RULE_ENGINE.getName() + ".producer");
this.tbCoreProducerStats = statsFactory.createMessagesStats(StatsType.CORE.getName() + ".producer"); this.tbCoreProducerStats = statsFactory.createMessagesStats(StatsType.CORE.getName() + ".producer");
this.transportApiStats = statsFactory.createMessagesStats(StatsType.TRANSPORT.getName() + ".producer"); this.transportApiStats = statsFactory.createMessagesStats(StatsType.TRANSPORT.getName() + ".producer");
@ -177,6 +179,7 @@ public class DefaultTransportService implements TransportService {
TopicPartitionInfo tpi = partitionService.getNotificationsTopic(ServiceType.TB_TRANSPORT, serviceInfoProvider.getServiceId()); TopicPartitionInfo tpi = partitionService.getNotificationsTopic(ServiceType.TB_TRANSPORT, serviceInfoProvider.getServiceId());
transportNotificationsConsumer.subscribe(Collections.singleton(tpi)); transportNotificationsConsumer.subscribe(Collections.singleton(tpi));
transportApiRequestTemplate.init(); transportApiRequestTemplate.init();
mainConsumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("transport-consumer"));
mainConsumerExecutor.execute(() -> { mainConsumerExecutor.execute(() -> {
while (!stopped) { while (!stopped) {
try { try {
@ -208,10 +211,6 @@ public class DefaultTransportService implements TransportService {
@PreDestroy @PreDestroy
public void destroy() { public void destroy() {
if (rateLimitEnabled) {
perTenantLimits.clear();
perDeviceLimits.clear();
}
stopped = true; stopped = true;
if (transportNotificationsConsumer != null) { if (transportNotificationsConsumer != null) {
@ -232,7 +231,7 @@ public class DefaultTransportService implements TransportService {
} }
@Override @Override
public ScheduledExecutorService getSchedulerExecutor(){ public ScheduledExecutorService getSchedulerExecutor() {
return this.schedulerExecutor; return this.schedulerExecutor;
} }
@ -242,12 +241,12 @@ public class DefaultTransportService implements TransportService {
} }
@Override @Override
public TransportProtos.GetTenantRoutingInfoResponseMsg getRoutingInfo(TransportProtos.GetTenantRoutingInfoRequestMsg msg) { public TransportProtos.GetEntityProfileResponseMsg getRoutingInfo(TransportProtos.GetEntityProfileRequestMsg msg) {
TbProtoQueueMsg<TransportProtos.TransportApiRequestMsg> protoMsg = TbProtoQueueMsg<TransportProtos.TransportApiRequestMsg> protoMsg =
new TbProtoQueueMsg<>(UUID.randomUUID(), TransportProtos.TransportApiRequestMsg.newBuilder().setGetTenantRoutingInfoRequestMsg(msg).build()); new TbProtoQueueMsg<>(UUID.randomUUID(), TransportProtos.TransportApiRequestMsg.newBuilder().setEntityProfileRequestMsg(msg).build());
try { try {
TbProtoQueueMsg<TransportApiResponseMsg> response = transportApiRequestTemplate.send(protoMsg).get(); TbProtoQueueMsg<TransportApiResponseMsg> response = transportApiRequestTemplate.send(protoMsg).get();
return response.getValue().getGetTenantRoutingInfoResponseMsg(); return response.getValue().getEntityProfileResponseMsg();
} catch (InterruptedException | ExecutionException e) { } catch (InterruptedException | ExecutionException e) {
throw new RuntimeException(e); throw new RuntimeException(e);
} }
@ -289,7 +288,7 @@ public class DefaultTransportService implements TransportService {
result.deviceInfo(tdi); result.deviceInfo(tdi);
ByteString profileBody = msg.getProfileBody(); ByteString profileBody = msg.getProfileBody();
if (profileBody != null && !profileBody.isEmpty()) { if (profileBody != null && !profileBody.isEmpty()) {
DeviceProfile profile = transportProfileCache.getOrCreate(tdi.getDeviceProfileId(), profileBody); DeviceProfile profile = deviceProfileCache.getOrCreate(tdi.getDeviceProfileId(), profileBody);
if (transportType != DeviceTransportType.DEFAULT if (transportType != DeviceTransportType.DEFAULT
&& profile != null && profile.getTransportType() != DeviceTransportType.DEFAULT && profile.getTransportType() != transportType) { && profile != null && profile.getTransportType() != DeviceTransportType.DEFAULT && profile.getTransportType() != transportType) {
log.debug("[{}] Device profile [{}] has different transport type: {}, expected: {}", tdi.getDeviceId(), tdi.getDeviceProfileId(), profile.getTransportType(), transportType); log.debug("[{}] Device profile [{}] has different transport type: {}, expected: {}", tdi.getDeviceId(), tdi.getDeviceProfileId(), profile.getTransportType(), transportType);
@ -315,7 +314,7 @@ public class DefaultTransportService implements TransportService {
result.deviceInfo(tdi); result.deviceInfo(tdi);
ByteString profileBody = msg.getProfileBody(); ByteString profileBody = msg.getProfileBody();
if (profileBody != null && !profileBody.isEmpty()) { if (profileBody != null && !profileBody.isEmpty()) {
result.deviceProfile(transportProfileCache.getOrCreate(tdi.getDeviceProfileId(), profileBody)); result.deviceProfile(deviceProfileCache.getOrCreate(tdi.getDeviceProfileId(), profileBody));
} }
} }
return result.build(); return result.build();
@ -339,8 +338,8 @@ public class DefaultTransportService implements TransportService {
log.trace("Processing msg: {}", requestMsg); log.trace("Processing msg: {}", requestMsg);
TbProtoQueueMsg<TransportApiRequestMsg> protoMsg = new TbProtoQueueMsg<>(UUID.randomUUID(), TransportApiRequestMsg.newBuilder().setProvisionDeviceRequestMsg(requestMsg).build()); TbProtoQueueMsg<TransportApiRequestMsg> protoMsg = new TbProtoQueueMsg<>(UUID.randomUUID(), TransportApiRequestMsg.newBuilder().setProvisionDeviceRequestMsg(requestMsg).build());
ListenableFuture<ProvisionDeviceResponseMsg> response = Futures.transform(transportApiRequestTemplate.send(protoMsg), tmp -> ListenableFuture<ProvisionDeviceResponseMsg> response = Futures.transform(transportApiRequestTemplate.send(protoMsg), tmp ->
tmp.getValue().getProvisionDeviceResponseMsg() tmp.getValue().getProvisionDeviceResponseMsg()
, MoreExecutors.directExecutor()); , MoreExecutors.directExecutor());
AsyncCallbackTemplate.withCallback(response, callback::onSuccess, callback::onError, transportCallbackExecutor); AsyncCallbackTemplate.withCallback(response, callback::onSuccess, callback::onError, transportCallbackExecutor);
} }
@ -580,12 +579,11 @@ public class DefaultTransportService implements TransportService {
if (log.isTraceEnabled()) { if (log.isTraceEnabled()) {
log.trace("[{}] Processing msg: {}", toSessionId(sessionInfo), msg); log.trace("[{}] Processing msg: {}", toSessionId(sessionInfo), msg);
} }
if (!rateLimitEnabled) {
return true;
}
TenantId tenantId = new TenantId(new UUID(sessionInfo.getTenantIdMSB(), sessionInfo.getTenantIdLSB())); TenantId tenantId = new TenantId(new UUID(sessionInfo.getTenantIdMSB(), sessionInfo.getTenantIdLSB()));
TbRateLimits rateLimits = perTenantLimits.computeIfAbsent(tenantId, id -> new TbRateLimits(perTenantLimitsConf));
if (!rateLimits.tryConsume()) { TransportRateLimit tenantRateLimit = rateLimitService.getRateLimit(tenantId, TransportRateLimitType.TENANT_MAX_MSGS);
if (!tenantRateLimit.tryConsume()) {
if (callback != null) { if (callback != null) {
callback.onError(new TbRateLimitsException(EntityType.TENANT)); callback.onError(new TbRateLimitsException(EntityType.TENANT));
} }
@ -595,8 +593,8 @@ public class DefaultTransportService implements TransportService {
return false; return false;
} }
DeviceId deviceId = new DeviceId(new UUID(sessionInfo.getDeviceIdMSB(), sessionInfo.getDeviceIdLSB())); DeviceId deviceId = new DeviceId(new UUID(sessionInfo.getDeviceIdMSB(), sessionInfo.getDeviceIdLSB()));
rateLimits = perDeviceLimits.computeIfAbsent(deviceId, id -> new TbRateLimits(perDevicesLimitsConf)); TransportRateLimit deviceRateLimit = rateLimitService.getRateLimit(tenantId, deviceId, TransportRateLimitType.DEVICE_MAX_MSGS);
if (!rateLimits.tryConsume()) { if (!deviceRateLimit.tryConsume()) {
if (callback != null) { if (callback != null) {
callback.onError(new TbRateLimitsException(EntityType.DEVICE)); callback.onError(new TbRateLimitsException(EntityType.DEVICE));
} }
@ -637,16 +635,40 @@ public class DefaultTransportService implements TransportService {
deregisterSession(md.getSessionInfo()); deregisterSession(md.getSessionInfo());
} }
} else { } else {
if (toSessionMsg.hasDeviceProfileUpdateMsg()) { if (toSessionMsg.hasEntityUpdateMsg()) {
DeviceProfile deviceProfile = transportProfileCache.put(toSessionMsg.getDeviceProfileUpdateMsg().getData()); TransportProtos.EntityUpdateMsg msg = toSessionMsg.getEntityUpdateMsg();
if (deviceProfile != null) { EntityType entityType = EntityType.valueOf(msg.getEntityType());
onProfileUpdate(deviceProfile); if (EntityType.DEVICE_PROFILE.equals(entityType)) {
DeviceProfile deviceProfile = deviceProfileCache.put(msg.getData());
if (deviceProfile != null) {
onProfileUpdate(deviceProfile);
}
} else if (EntityType.TENANT_PROFILE.equals(entityType)) {
TenantProfileUpdateResult update = tenantProfileCache.put(msg.getData());
rateLimitService.update(update);
} else if (EntityType.TENANT.equals(entityType)) {
Optional<Tenant> profileOpt = dataDecodingEncodingService.decode(msg.getData().toByteArray());
if (profileOpt.isPresent()) {
Tenant tenant = profileOpt.get();
boolean updated = tenantProfileCache.put(tenant.getId(), tenant.getTenantProfileId());
if (updated) {
rateLimitService.update(tenant.getId());
}
}
}
} else if (toSessionMsg.hasEntityDeleteMsg()) {
TransportProtos.EntityDeleteMsg msg = toSessionMsg.getEntityDeleteMsg();
EntityType entityType = EntityType.valueOf(msg.getEntityType());
UUID entityUuid = new UUID(msg.getEntityIdMSB(), msg.getEntityIdLSB());
if (EntityType.DEVICE_PROFILE.equals(entityType)) {
deviceProfileCache.evict(new DeviceProfileId(new UUID(msg.getEntityIdMSB(), msg.getEntityIdLSB())));
} else if (EntityType.TENANT_PROFILE.equals(entityType)) {
tenantProfileCache.remove(new TenantProfileId(entityUuid));
} else if (EntityType.TENANT.equals(entityType)) {
rateLimitService.remove(new TenantId(entityUuid));
} else if (EntityType.DEVICE.equals(entityType)) {
rateLimitService.remove(new DeviceId(entityUuid));
} }
} else if (toSessionMsg.hasDeviceProfileDeleteMsg()) {
transportProfileCache.evict(new DeviceProfileId(new UUID(
toSessionMsg.getDeviceProfileDeleteMsg().getProfileIdMSB(),
toSessionMsg.getDeviceProfileDeleteMsg().getProfileIdLSB()
)));
} else { } else {
//TODO: should we notify the device actor about missed session? //TODO: should we notify the device actor about missed session?
log.debug("[{}] Missing session.", sessionId); log.debug("[{}] Missing session.", sessionId);
@ -654,38 +676,6 @@ public class DefaultTransportService implements TransportService {
} }
} }
@Override
public void getDeviceProfile(DeviceProfileId deviceProfileId, TransportServiceCallback<DeviceProfile> callback) {
DeviceProfile deviceProfile = transportProfileCache.get(deviceProfileId);
if (deviceProfile != null) {
callback.onSuccess(deviceProfile);
} else {
log.trace("Processing device profile request: [{}]", deviceProfileId);
TransportProtos.GetDeviceProfileRequestMsg msg = TransportProtos.GetDeviceProfileRequestMsg.newBuilder()
.setProfileIdMSB(deviceProfileId.getId().getMostSignificantBits())
.setProfileIdLSB(deviceProfileId.getId().getLeastSignificantBits())
.build();
TbProtoQueueMsg<TransportApiRequestMsg> protoMsg = new TbProtoQueueMsg<>(UUID.randomUUID(),
TransportApiRequestMsg.newBuilder().setGetDeviceProfileRequestMsg(msg).build());
AsyncCallbackTemplate.withCallback(transportApiRequestTemplate.send(protoMsg),
response -> {
ByteString devProfileBody = response.getValue().getGetDeviceProfileResponseMsg().getData();
if (devProfileBody != null && !devProfileBody.isEmpty()) {
DeviceProfile profile = transportProfileCache.put(devProfileBody);
if (profile != null) {
callback.onSuccess(profile);
} else {
log.warn("Failed to decode device profile: {}", devProfileBody);
callback.onError(new IllegalArgumentException("Failed to decode device profile!"));
}
} else {
log.warn("Failed to find device profile: [{}]", deviceProfileId);
callback.onError(new IllegalArgumentException("Failed to find device profile!"));
}
}, callback::onError, transportCallbackExecutor);
}
}
@Override @Override
public void onProfileUpdate(DeviceProfile deviceProfile) { public void onProfileUpdate(DeviceProfile deviceProfile) {
long deviceProfileIdMSB = deviceProfile.getId().getId().getMostSignificantBits(); long deviceProfileIdMSB = deviceProfile.getId().getId().getMostSignificantBits();
@ -750,7 +740,7 @@ public class DefaultTransportService implements TransportService {
private RuleChainId resolveRuleChainId(TransportProtos.SessionInfoProto sessionInfo) { private RuleChainId resolveRuleChainId(TransportProtos.SessionInfoProto sessionInfo) {
DeviceProfileId deviceProfileId = new DeviceProfileId(new UUID(sessionInfo.getDeviceProfileIdMSB(), sessionInfo.getDeviceProfileIdLSB())); DeviceProfileId deviceProfileId = new DeviceProfileId(new UUID(sessionInfo.getDeviceProfileIdMSB(), sessionInfo.getDeviceProfileIdLSB()));
DeviceProfile deviceProfile = transportProfileCache.get(deviceProfileId); DeviceProfile deviceProfile = deviceProfileCache.get(deviceProfileId);
RuleChainId ruleChainId; RuleChainId ruleChainId;
if (deviceProfile == null) { if (deviceProfile == null) {
log.warn("[{}] Device profile is null!", deviceProfileId); log.warn("[{}] Device profile is null!", deviceProfileId);
@ -779,7 +769,7 @@ public class DefaultTransportService implements TransportService {
} }
} }
private class StatsCallback implements TbQueueCallback { private static class StatsCallback implements TbQueueCallback {
private final TbQueueCallback callback; private final TbQueueCallback callback;
private final MessagesStats stats; private final MessagesStats stats;

154
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportTenantProfileCache.java

@ -0,0 +1,154 @@
/**
* Copyright © 2016-2020 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.service;
import com.google.protobuf.ByteString;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.TenantProfileId;
import org.thingsboard.server.common.transport.TransportService;
import org.thingsboard.server.common.transport.TransportTenantProfileCache;
import org.thingsboard.server.common.transport.profile.TenantProfileUpdateResult;
import org.thingsboard.server.common.transport.util.DataDecodingEncodingService;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.queue.discovery.TenantRoutingInfo;
import org.thingsboard.server.queue.discovery.TenantRoutingInfoService;
import java.util.Collections;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
@Component
@ConditionalOnExpression("('${service.type:null}'=='monolith' && '${transport.api_enabled:true}'=='true') || '${service.type:null}'=='tb-transport'")
@Slf4j
public class DefaultTransportTenantProfileCache implements TransportTenantProfileCache {
private final Lock tenantProfileFetchLock = new ReentrantLock();
private final ConcurrentMap<TenantProfileId, TenantProfile> profiles = new ConcurrentHashMap<>();
private final ConcurrentMap<TenantId, TenantProfileId> tenantIds = new ConcurrentHashMap<>();
private final ConcurrentMap<TenantProfileId, Set<TenantId>> tenantProfileIds = new ConcurrentHashMap<>();
private final DataDecodingEncodingService dataDecodingEncodingService;
private TransportService transportService;
@Lazy
@Autowired
public void setTransportService(TransportService transportService) {
this.transportService = transportService;
}
public DefaultTransportTenantProfileCache(DataDecodingEncodingService dataDecodingEncodingService) {
this.dataDecodingEncodingService = dataDecodingEncodingService;
}
@Override
public TenantProfile get(TenantId tenantId) {
return getTenantProfile(tenantId);
}
@Override
public TenantProfileUpdateResult put(ByteString profileBody) {
Optional<TenantProfile> profileOpt = dataDecodingEncodingService.decode(profileBody.toByteArray());
if (profileOpt.isPresent()) {
TenantProfile newProfile = profileOpt.get();
log.trace("[{}] put: {}", newProfile.getId(), newProfile);
return new TenantProfileUpdateResult(newProfile, tenantProfileIds.get(newProfile.getId()));
} else {
log.warn("Failed to decode profile: {}", profileBody.toString());
return new TenantProfileUpdateResult(null, Collections.emptySet());
}
}
@Override
public boolean put(TenantId tenantId, TenantProfileId profileId) {
log.trace("[{}] put: {}", tenantId, profileId);
TenantProfileId oldProfileId = tenantIds.get(tenantId);
if (oldProfileId != null && !oldProfileId.equals(profileId)) {
tenantProfileIds.computeIfAbsent(oldProfileId, id -> ConcurrentHashMap.newKeySet()).remove(tenantId);
tenantIds.put(tenantId, profileId);
tenantProfileIds.computeIfAbsent(profileId, id -> ConcurrentHashMap.newKeySet()).add(tenantId);
return true;
} else {
return false;
}
}
@Override
public Set<TenantId> remove(TenantProfileId profileId) {
Set<TenantId> tenants = tenantProfileIds.remove(profileId);
if (tenants != null) {
tenants.forEach(tenantIds::remove);
}
profiles.remove(profileId);
return tenants;
}
private TenantProfile getTenantProfile(TenantId tenantId) {
TenantProfile profile = null;
TenantProfileId tenantProfileId = tenantIds.get(tenantId);
if (tenantProfileId != null) {
profile = profiles.get(tenantProfileId);
}
if (profile == null) {
tenantProfileFetchLock.lock();
try {
tenantProfileId = tenantIds.get(tenantId);
if (tenantProfileId != null) {
profile = profiles.get(tenantProfileId);
}
if (profile == null) {
TransportProtos.GetEntityProfileRequestMsg msg = TransportProtos.GetEntityProfileRequestMsg.newBuilder()
.setEntityType(EntityType.TENANT.name())
.setEntityIdMSB(tenantId.getId().getMostSignificantBits())
.setEntityIdLSB(tenantId.getId().getLeastSignificantBits())
.build();
TransportProtos.GetEntityProfileResponseMsg routingInfo = transportService.getRoutingInfo(msg);
Optional<TenantProfile> profileOpt = dataDecodingEncodingService.decode(routingInfo.getData().toByteArray());
if (profileOpt.isPresent()) {
profile = profileOpt.get();
TenantProfile existingProfile = profiles.get(profile.getId());
if (existingProfile != null) {
profile = existingProfile;
} else {
profiles.put(profile.getId(), profile);
}
tenantProfileIds.computeIfAbsent(profile.getId(), id -> ConcurrentHashMap.newKeySet()).add(tenantId);
tenantIds.put(tenantId, profile.getId());
} else {
log.warn("[{}] Can't decode tenant profile: {}", tenantId, routingInfo.getData());
throw new RuntimeException("Can't decode tenant profile!");
}
}
} finally {
tenantProfileFetchLock.unlock();
}
}
return profile;
}
}

24
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/TransportTenantRoutingInfoService.java

@ -16,14 +16,11 @@
package org.thingsboard.server.common.transport.service; package org.thingsboard.server.common.transport.service;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.transport.TransportService; import org.thingsboard.server.common.transport.TransportTenantProfileCache;
import org.thingsboard.server.gen.transport.TransportProtos.GetTenantRoutingInfoRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.GetTenantRoutingInfoResponseMsg;
import org.thingsboard.server.queue.discovery.TenantRoutingInfo; import org.thingsboard.server.queue.discovery.TenantRoutingInfo;
import org.thingsboard.server.queue.discovery.TenantRoutingInfoService; import org.thingsboard.server.queue.discovery.TenantRoutingInfoService;
@ -32,21 +29,16 @@ import org.thingsboard.server.queue.discovery.TenantRoutingInfoService;
@ConditionalOnExpression("'${service.type:null}'=='tb-transport'") @ConditionalOnExpression("'${service.type:null}'=='tb-transport'")
public class TransportTenantRoutingInfoService implements TenantRoutingInfoService { public class TransportTenantRoutingInfoService implements TenantRoutingInfoService {
private TransportService transportService; private TransportTenantProfileCache tenantProfileCache;
@Lazy public TransportTenantRoutingInfoService(TransportTenantProfileCache tenantProfileCache) {
@Autowired this.tenantProfileCache = tenantProfileCache;
public void setTransportService(TransportService transportService) {
this.transportService = transportService;
} }
@Override @Override
public TenantRoutingInfo getRoutingInfo(TenantId tenantId) { public TenantRoutingInfo getRoutingInfo(TenantId tenantId) {
GetTenantRoutingInfoRequestMsg msg = GetTenantRoutingInfoRequestMsg.newBuilder() TenantProfile profile = tenantProfileCache.get(tenantId);
.setTenantIdMSB(tenantId.getId().getMostSignificantBits()) return new TenantRoutingInfo(tenantId, profile.isIsolatedTbCore(), profile.isIsolatedTbRuleEngine());
.setTenantIdLSB(tenantId.getId().getLeastSignificantBits())
.build();
GetTenantRoutingInfoResponseMsg routingInfo = transportService.getRoutingInfo(msg);
return new TenantRoutingInfo(tenantId, routingInfo.getIsolatedTbCore(), routingInfo.getIsolatedTbRuleEngine());
} }
} }

2
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/util/DataDecodingEncodingService.java

@ -15,8 +15,6 @@
*/ */
package org.thingsboard.server.common.transport.util; package org.thingsboard.server.common.transport.util;
import org.thingsboard.server.common.msg.TbActorMsg;
import java.util.Optional; import java.util.Optional;
public interface DataDecodingEncodingService { public interface DataDecodingEncodingService {

4
transport/coap/src/main/resources/tb-coap-transport.yml

@ -49,10 +49,6 @@ transport:
sessions: sessions:
inactivity_timeout: "${TB_TRANSPORT_SESSIONS_INACTIVITY_TIMEOUT:300000}" inactivity_timeout: "${TB_TRANSPORT_SESSIONS_INACTIVITY_TIMEOUT:300000}"
report_timeout: "${TB_TRANSPORT_SESSIONS_REPORT_TIMEOUT:30000}" report_timeout: "${TB_TRANSPORT_SESSIONS_REPORT_TIMEOUT:30000}"
rate_limits:
enabled: "${TB_TRANSPORT_RATE_LIMITS_ENABLED:false}"
tenant: "${TB_TRANSPORT_RATE_LIMITS_TENANT:1000:1,20000:60}"
device: "${TB_TRANSPORT_RATE_LIMITS_DEVICE:10:1,300:60}"
json: json:
# Cast String data types to Numeric if possible when processing Telemetry/Attributes JSON # Cast String data types to Numeric if possible when processing Telemetry/Attributes JSON
type_cast_enabled: "${JSON_TYPE_CAST_ENABLED:true}" type_cast_enabled: "${JSON_TYPE_CAST_ENABLED:true}"

4
transport/http/src/main/resources/tb-http-transport.yml

@ -42,10 +42,6 @@ transport:
sessions: sessions:
inactivity_timeout: "${TB_TRANSPORT_SESSIONS_INACTIVITY_TIMEOUT:300000}" inactivity_timeout: "${TB_TRANSPORT_SESSIONS_INACTIVITY_TIMEOUT:300000}"
report_timeout: "${TB_TRANSPORT_SESSIONS_REPORT_TIMEOUT:30000}" report_timeout: "${TB_TRANSPORT_SESSIONS_REPORT_TIMEOUT:30000}"
rate_limits:
enabled: "${TB_TRANSPORT_RATE_LIMITS_ENABLED:false}"
tenant: "${TB_TRANSPORT_RATE_LIMITS_TENANT:1000:1,20000:60}"
device: "${TB_TRANSPORT_RATE_LIMITS_DEVICE:10:1,300:60}"
json: json:
# Cast String data types to Numeric if possible when processing Telemetry/Attributes JSON # Cast String data types to Numeric if possible when processing Telemetry/Attributes JSON
type_cast_enabled: "${JSON_TYPE_CAST_ENABLED:true}" type_cast_enabled: "${JSON_TYPE_CAST_ENABLED:true}"

4
transport/mqtt/src/main/resources/tb-mqtt-transport.yml

@ -71,10 +71,6 @@ transport:
sessions: sessions:
inactivity_timeout: "${TB_TRANSPORT_SESSIONS_INACTIVITY_TIMEOUT:300000}" inactivity_timeout: "${TB_TRANSPORT_SESSIONS_INACTIVITY_TIMEOUT:300000}"
report_timeout: "${TB_TRANSPORT_SESSIONS_REPORT_TIMEOUT:30000}" report_timeout: "${TB_TRANSPORT_SESSIONS_REPORT_TIMEOUT:30000}"
rate_limits:
enabled: "${TB_TRANSPORT_RATE_LIMITS_ENABLED:false}"
tenant: "${TB_TRANSPORT_RATE_LIMITS_TENANT:1000:1,20000:60}"
device: "${TB_TRANSPORT_RATE_LIMITS_DEVICE:10:1,300:60}"
json: json:
# Cast String data types to Numeric if possible when processing Telemetry/Attributes JSON # Cast String data types to Numeric if possible when processing Telemetry/Attributes JSON
type_cast_enabled: "${JSON_TYPE_CAST_ENABLED:true}" type_cast_enabled: "${JSON_TYPE_CAST_ENABLED:true}"

Loading…
Cancel
Save