Browse Source

Merge remote-tracking branch 'origin/lts-4.2' into lts-4.3

# Conflicts:
#	application/src/main/java/org/thingsboard/server/service/entitiy/alarm/DefaultTbAlarmCommentService.java
#	application/src/test/java/org/thingsboard/server/controller/AlarmCommentControllerTest.java
pull/15765/head
Viacheslav Klimov 2 months ago
parent
commit
b234daf6b8
Failed to extract signature
  1. 22
      application/src/main/java/org/thingsboard/server/controller/AlarmCommentController.java
  2. 6
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
  3. 4
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java
  4. 2
      application/src/main/resources/thingsboard.yml
  5. 2
      application/src/test/java/org/thingsboard/server/client/AlarmCommentApiClientTest.java
  6. 23
      application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java
  7. 47
      application/src/test/java/org/thingsboard/server/controller/AlarmCommentControllerTest.java
  8. 4
      application/src/test/java/org/thingsboard/server/controller/HomePageApiTest.java
  9. 1
      application/src/test/java/org/thingsboard/server/controller/UserControllerTest.java
  10. 32
      application/src/test/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerServiceTest.java
  11. 78
      application/src/test/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerServiceTest.java
  12. 62
      common/coap-server/src/main/java/org/thingsboard/server/coapserver/DefaultCoapServerService.java
  13. 145
      common/coap-server/src/test/java/org/thingsboard/server/coapserver/DefaultCoapServerServiceTest.java
  14. 7
      common/data/src/main/java/org/thingsboard/server/common/data/StringUtils.java
  15. 2
      common/data/src/main/java/org/thingsboard/server/common/data/alarm/AlarmCommentSubType.java
  16. 12
      common/data/src/main/java/org/thingsboard/server/common/data/rpc/RpcError.java
  17. 7
      common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java
  18. 2
      common/proto/src/main/proto/queue.proto
  19. 11
      common/proto/src/test/java/org/thingsboard/server/common/util/ProtoUtilsTest.java
  20. 25
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/LwM2MTransportBootstrapService.java
  21. 115
      common/transport/lwm2m/src/test/java/org/thingsboard/server/transport/lwm2m/bootstrap/LwM2MTransportBootstrapServiceTest.java
  22. 115
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/DefaultTransportRateLimitService.java
  23. 35
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportLimitsType.java
  24. 4
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java
  25. 76
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportTenantProfileCache.java
  26. 202
      common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/limits/DefaultTransportRateLimitServiceTest.java
  27. 191
      common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/service/DefaultTransportTenantProfileCacheTest.java
  28. 18
      dao/src/main/java/org/thingsboard/server/dao/service/validator/DeviceCredentialsDataValidator.java
  29. 40
      dao/src/main/java/org/thingsboard/server/dao/util/DeviceConnectivityUtil.java
  30. 138
      dao/src/test/java/org/thingsboard/server/dao/service/validator/DeviceCredentialsDataValidatorTest.java
  31. 123
      dao/src/test/java/org/thingsboard/server/dao/util/DeviceConnectivityUtilTest.java
  32. 2
      transport/coap/src/main/resources/tb-coap-transport.yml
  33. 2
      transport/http/src/main/resources/tb-http-transport.yml
  34. 2
      transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml
  35. 2
      transport/mqtt/src/main/resources/tb-mqtt-transport.yml
  36. 2
      transport/snmp/src/main/resources/tb-snmp-transport.yml
  37. 2
      ui-ngx/src/assets/locale/locale.constant-en_US.json

22
application/src/main/java/org/thingsboard/server/controller/AlarmCommentController.java

@ -31,6 +31,7 @@ import org.thingsboard.server.common.data.alarm.Alarm;
import org.thingsboard.server.common.data.alarm.AlarmComment;
import org.thingsboard.server.common.data.alarm.AlarmCommentInfo;
import org.thingsboard.server.common.data.alarm.AlarmCommentType;
import org.thingsboard.server.common.data.exception.ThingsboardErrorCode;
import org.thingsboard.server.common.data.exception.ThingsboardException;
import org.thingsboard.server.common.data.id.AlarmCommentId;
import org.thingsboard.server.common.data.id.AlarmId;
@ -39,6 +40,7 @@ import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.config.annotations.ApiOperation;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.entitiy.alarm.TbAlarmCommentService;
import org.thingsboard.server.service.security.model.SecurityUser;
import org.thingsboard.server.service.security.permission.Operation;
import static org.thingsboard.server.controller.ControllerConstants.ALARM_COMMENT_ID_PARAM_DESCRIPTION;
@ -77,9 +79,13 @@ public class AlarmCommentController extends BaseController {
checkParameter(ALARM_ID, strAlarmId);
AlarmId alarmId = new AlarmId(toUUID(strAlarmId));
Alarm alarm = checkAlarmInfoId(alarmId, Operation.WRITE);
SecurityUser currentUser = getCurrentUser();
if (alarmComment.getId() != null) {
checkUserPermission(alarmComment, alarmId, "edit", currentUser);
}
alarmComment.setAlarmId(alarmId);
alarmComment.setType(AlarmCommentType.OTHER);
return tbAlarmCommentService.saveAlarmComment(alarm, alarmComment, getCurrentUser());
return tbAlarmCommentService.saveAlarmComment(alarm, alarmComment, currentUser);
}
@ApiOperation(value = "Delete Alarm comment (deleteAlarmComment)",
@ -93,7 +99,11 @@ public class AlarmCommentController extends BaseController {
AlarmCommentId alarmCommentId = new AlarmCommentId(toUUID(strCommentId));
AlarmComment alarmComment = checkAlarmCommentId(alarmCommentId, alarmId);
tbAlarmCommentService.deleteAlarmComment(alarm, alarmComment, getCurrentUser());
SecurityUser currentUser = getCurrentUser();
if (!currentUser.isTenantAdmin()) {
checkUserPermission(alarmComment, alarmId, "delete", currentUser);
}
tbAlarmCommentService.deleteAlarmComment(alarm, alarmComment, currentUser);
}
@ApiOperation(value = "Get Alarm comments (getAlarmComments)",
@ -120,4 +130,12 @@ public class AlarmCommentController extends BaseController {
return checkNotNull(alarmCommentService.findAlarmComments(alarm.getTenantId(), alarmId, pageLink));
}
private void checkUserPermission(AlarmComment alarmComment, AlarmId alarmId, String operation, SecurityUser currentUser) throws ThingsboardException {
AlarmComment existingAlarmComment = checkAlarmCommentId(alarmComment.getId(), alarmId);
if (existingAlarmComment.getUserId() != null && !existingAlarmComment.getUserId().equals(currentUser.getId())) {
throw new ThingsboardException("User is not allowed to " + operation + " other user's comment",
ThingsboardErrorCode.PERMISSION_DENIED);
}
}
}

6
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java

@ -464,10 +464,10 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
return firmwareStateService.process(msg.getValue());
}
private void forwardToCoreRpcService(FromDeviceRPCResponseProto proto, TbCallback callback) {
RpcError error = proto.getError() > 0 ? RpcError.values()[proto.getError()] : null;
void forwardToCoreRpcService(FromDeviceRPCResponseProto proto, TbCallback callback) {
RpcError error = RpcError.fromProtoErrorCode(proto.getError());
FromDeviceRpcResponse response = new FromDeviceRpcResponse(new UUID(proto.getRequestIdMSB(), proto.getRequestIdLSB())
, proto.getResponse(), error);
, proto.hasResponse() ? proto.getResponse() : null, error);
tbCoreDeviceRpcService.processRpcResponseFromRuleEngine(response);
callback.onSuccess();
}

4
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java

@ -179,9 +179,9 @@ public class DefaultTbRuleEngineConsumerService extends AbstractPartitionBasedCo
callback.onSuccess();
} else if (nfMsg.hasFromDeviceRpcResponse()) {
TransportProtos.FromDeviceRPCResponseProto proto = nfMsg.getFromDeviceRpcResponse();
RpcError error = proto.getError() > 0 ? RpcError.values()[proto.getError()] : null;
RpcError error = RpcError.fromProtoErrorCode(proto.getError());
FromDeviceRpcResponse response = new FromDeviceRpcResponse(new UUID(proto.getRequestIdMSB(), proto.getRequestIdLSB())
, proto.getResponse(), error);
, proto.hasResponse() ? proto.getResponse() : null, error);
tbDeviceRpcService.processRpcResponseFromDevice(response);
callback.onSuccess();
} else if (nfMsg.getQueueUpdateMsgsCount() > 0) {

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

@ -1158,6 +1158,8 @@ transport:
timeout: "${CLIENT_SIDE_RPC_TIMEOUT:60000}"
# Enable/disable http/mqtt/coap/lwm2m transport protocols (has higher priority than certain protocol's 'enabled' property)
api_enabled: "${TB_TRANSPORT_API_ENABLED:true}"
# Size of the thread pool that executes transport API callbacks (session registration, telemetry/attribute and RPC responses, entity update notifications, and the tenant profile fetch on a cache miss). Bounds how many such callbacks - including those that block on a backend round-trip - can run concurrently.
callback_thread_pool_size: "${TB_TRANSPORT_CALLBACK_THREAD_POOL_SIZE:20}"
log:
# Enable/Disable log of transport messages to telemetry. For example, logging of LwM2M registration update
enabled: "${TB_TRANSPORT_LOG_ENABLED:true}"

2
application/src/test/java/org/thingsboard/server/client/AlarmCommentApiClientTest.java

@ -100,7 +100,7 @@ public class AlarmCommentApiClientTest extends AbstractApiClientTest {
.filter(alarmCommentInfo -> alarmCommentInfo.getId().getId().equals(commentToDeleteId))
.findFirst()
.get();
assertEquals("User " + clientTenantAdmin.getEmail() + " deleted his comment", deletedComment.getComment().get("text").asText());
assertEquals("Comment was deleted by user " + clientTenantAdmin.getEmail(), deletedComment.getComment().get("text").asText());
}
}

23
application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java

@ -229,6 +229,7 @@ public abstract class AbstractWebTest extends AbstractInMemoryStorageTest {
protected static final String DIFFERENT_TENANT_ADMIN_PASSWORD = "difftenant";
protected static final String CUSTOMER_USER_EMAIL = "testcustomer@thingsboard.org";
protected static final String SECOND_CUSTOMER_USER_EMAIL = "testsecondcustomer@thingsboard.org";
private static final String CUSTOMER_USER_PASSWORD = "customer";
protected static final String DIFFERENT_CUSTOMER_USER_EMAIL = "testdifferentcustomer@thingsboard.org";
@ -268,6 +269,7 @@ public abstract class AbstractWebTest extends AbstractInMemoryStorageTest {
protected CustomerId differentTenantCustomerId;
protected UserId customerUserId;
protected UserId secondCustomerUserId;
protected UserId differentCustomerUserId;
protected UserId differentTenantCustomerUserId;
@ -393,9 +395,17 @@ public abstract class AbstractWebTest extends AbstractInMemoryStorageTest {
customerUser.setCustomerId(savedCustomer.getId());
customerUser.setEmail(CUSTOMER_USER_EMAIL);
customerUser = createUserAndLogin(customerUser, CUSTOMER_USER_PASSWORD);
customerUser = createUserAndActivate(customerUser, CUSTOMER_USER_PASSWORD);
customerUserId = customerUser.getId();
User secondCustomerUser = new User();
secondCustomerUser.setAuthority(Authority.CUSTOMER_USER);
secondCustomerUser.setTenantId(tenantId);
secondCustomerUser.setCustomerId(customerId);
secondCustomerUser.setEmail(SECOND_CUSTOMER_USER_EMAIL);
secondCustomerUser = createUserAndActivate(secondCustomerUser, CUSTOMER_USER_PASSWORD);
secondCustomerUserId = secondCustomerUser.getId();
resetTokens();
log.debug("Executed web test setup");
@ -494,6 +504,10 @@ public abstract class AbstractWebTest extends AbstractInMemoryStorageTest {
login(CUSTOMER_USER_EMAIL, CUSTOMER_USER_PASSWORD);
}
protected void loginSecondCustomerUser() throws Exception {
login(SECOND_CUSTOMER_USER_EMAIL, CUSTOMER_USER_PASSWORD);
}
protected void loginUser(String userName, String password) throws Exception {
login(userName, password);
}
@ -608,6 +622,13 @@ public abstract class AbstractWebTest extends AbstractInMemoryStorageTest {
return savedUser;
}
protected User createUserAndActivate(User user, String password) throws Exception {
User savedUser = doPost("/api/user", user, User.class);
JsonNode activateRequest = getActivateRequest(password);
doPost("/api/noauth/activate", activateRequest).andExpect(status().isOk());
return savedUser;
}
protected User createUser(User user, String password) throws Exception {
User savedUser = doPost("/api/user", user, User.class);
JsonNode activateRequest = getActivateRequest(password);

47
application/src/test/java/org/thingsboard/server/controller/AlarmCommentControllerTest.java

@ -161,6 +161,25 @@ public class AlarmCommentControllerTest extends AbstractControllerTest {
testLogEntityActionEntityEqClass(alarm, alarm.getId(), tenantId, customerId, tenantAdminUserId, TENANT_ADMIN_EMAIL, ActionType.UPDATED_COMMENT, 1, updatedAlarmComment);
}
@Test
public void testEditOthersAlarmCommentIsProhibited() throws Exception {
loginCustomerUser();
AlarmComment alarmComment = createAlarmComment(alarm.getId());
JsonNode newComment = JacksonUtil.newObjectNode().set("text", new TextNode("Second customer rewrite"));
alarmComment.setComment(newComment);
loginSecondCustomerUser();
doPost("/api/alarm/" + alarm.getId() + "/comment", alarmComment)
.andExpect(status().isForbidden())
.andExpect(statusReason(containsString("User is not allowed to edit other user's comment")));
loginTenantAdmin();
doPost("/api/alarm/" + alarm.getId() + "/comment", alarmComment)
.andExpect(status().isForbidden())
.andExpect(statusReason(containsString("User is not allowed to edit other user's comment")));
}
@Test
public void testUpdateAlarmViaDifferentTenant() throws Exception {
loginTenantAdmin();
@ -218,6 +237,32 @@ public class AlarmCommentControllerTest extends AbstractControllerTest {
testLogEntityActionEntityEqClass(alarm, alarm.getId(), tenantId, customerId, customerUserId, CUSTOMER_USER_EMAIL, ActionType.DELETED_COMMENT, 1, expectedAlarmComment);
}
@Test
public void testDeleteOthersAlarmCommentIsAllowedForAuthorOrTenantAdmin() throws Exception {
loginCustomerUser();
AlarmComment alarmComment = createAlarmComment(alarm.getId());
loginSecondCustomerUser();
Mockito.reset(tbClusterService, auditLogService);
doDelete("/api/alarm/" + alarm.getId() + "/comment/" + alarmComment.getId())
.andExpect(status().isForbidden())
.andExpect(statusReason(containsString("User is not allowed to delete other user's comment")));
loginTenantAdmin();
doDelete("/api/alarm/" + alarm.getId() + "/comment/" + alarmComment.getId())
.andExpect(status().isOk());
AlarmComment expectedAlarmComment = AlarmComment.builder()
.alarmId(alarm.getId())
.type(AlarmCommentType.SYSTEM)
.comment(JacksonUtil.newObjectNode()
.put("text", String.format(COMMENT_DELETED.getText(), TENANT_ADMIN_EMAIL))
.put("subtype", COMMENT_DELETED.name())
.put("userName", TENANT_ADMIN_EMAIL))
.build();
testLogEntityActionEntityEqClass(alarm, alarm.getId(), tenantId, customerId, tenantAdminUserId, TENANT_ADMIN_EMAIL, ActionType.DELETED_COMMENT, 1, expectedAlarmComment);
}
@Test
public void testDeleteAlarmViaTenant() throws Exception {
loginTenantAdmin();
@ -237,7 +282,7 @@ public class AlarmCommentControllerTest extends AbstractControllerTest {
assertThat(systemComment.getId()).isEqualTo(alarmComment.getId());
assertThat(systemComment.getType()).isEqualTo(AlarmCommentType.SYSTEM);
assertThat(systemComment.getComment().get("text").asText()).isEqualTo(String.format("User %s deleted his comment",
assertThat(systemComment.getComment().get("text").asText()).isEqualTo(String.format("Comment was deleted by user %s",
TENANT_ADMIN_EMAIL));
AlarmComment expectedAlarmComment = AlarmComment.builder()

4
application/src/test/java/org/thingsboard/server/controller/HomePageApiTest.java

@ -410,7 +410,7 @@ public class HomePageApiTest extends AbstractControllerTest {
Assert.assertEquals(1, usageInfo.getCustomers());
Assert.assertEquals(configuration.getMaxCustomers(), usageInfo.getMaxCustomers());
Assert.assertEquals(2, usageInfo.getUsers());
Assert.assertEquals(3, usageInfo.getUsers());
Assert.assertEquals(configuration.getMaxUsers(), usageInfo.getMaxUsers());
Assert.assertEquals(DEFAULT_DASHBOARDS_COUNT, usageInfo.getDashboards());
@ -476,7 +476,7 @@ public class HomePageApiTest extends AbstractControllerTest {
}
usageInfo = doGet("/api/usage", UsageInfo.class);
Assert.assertEquals(users.size() + 2, usageInfo.getUsers());
Assert.assertEquals(users.size() + 3, usageInfo.getUsers());
List<Dashboard> dashboards = new ArrayList<>();
for (int i = 0; i < 97; i++) {

1
application/src/test/java/org/thingsboard/server/controller/UserControllerTest.java

@ -717,6 +717,7 @@ public class UserControllerTest extends AbstractControllerTest {
String email = "testEmail1";
List<UserId> expectedCustomerUserIds = new ArrayList<>();
expectedCustomerUserIds.add(customerUserId);
expectedCustomerUserIds.add(secondCustomerUserId);
for (int i = 0; i < 45; i++) {
User customerUser = createCustomerUser(customerId);
customerUser.setEmail(email + StringUtils.randomAlphanumeric((int) (5 + Math.random() * 10)) + "@thingsboard.org");

32
application/src/test/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerServiceTest.java

@ -27,8 +27,11 @@ import org.mockito.junit.jupiter.MockitoExtension;
import org.springframework.test.util.ReflectionTestUtils;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.rpc.RpcError;
import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.service.rpc.TbCoreDeviceRpcService;
import org.thingsboard.server.service.ruleengine.RuleEngineCallService;
import org.thingsboard.server.service.state.DeviceStateService;
@ -51,6 +54,8 @@ public class DefaultTbCoreConsumerServiceTest {
private TbCoreConsumerStats statsMock;
@Mock
private RuleEngineCallService ruleEngineCallServiceMock;
@Mock
private TbCoreDeviceRpcService tbCoreDeviceRpcServiceMock;
@Mock
private TbCallback tbCallbackMock;
@ -638,4 +643,31 @@ public class DefaultTbCoreConsumerServiceTest {
then(ruleEngineCallServiceMock).should().onQueueMsg(restApiCallResponseMsgProto, tbCallbackMock);
}
@Test
public void givenNotFoundErrorAndNoResponse_whenForwardToCoreRpcService_thenNotFoundAndNullResponseAreRecovered() {
// GIVEN
ReflectionTestUtils.setField(defaultTbCoreConsumerServiceMock, "tbCoreDeviceRpcService", tbCoreDeviceRpcServiceMock);
var requestId = UUID.randomUUID();
// error = NOT_FOUND.ordinal() (0) and response left unset: the previously broken combination
// ('error > 0' dropped NOT_FOUND, proto3 default collapsed a null response to "").
var proto = TransportProtos.FromDeviceRPCResponseProto.newBuilder()
.setRequestIdMSB(requestId.getMostSignificantBits())
.setRequestIdLSB(requestId.getLeastSignificantBits())
.setError(RpcError.NOT_FOUND.ordinal())
.build();
doCallRealMethod().when(defaultTbCoreConsumerServiceMock).forwardToCoreRpcService(proto, tbCallbackMock);
// WHEN
defaultTbCoreConsumerServiceMock.forwardToCoreRpcService(proto, tbCallbackMock);
// THEN
var responseCaptor = ArgumentCaptor.forClass(FromDeviceRpcResponse.class);
then(tbCoreDeviceRpcServiceMock).should().processRpcResponseFromRuleEngine(responseCaptor.capture());
var response = responseCaptor.getValue();
assertThat(response.getId()).isEqualTo(requestId);
assertThat(response.getError()).contains(RpcError.NOT_FOUND);
assertThat(response.getResponse()).isEmpty();
then(tbCallbackMock).should().onSuccess();
}
}

78
application/src/test/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerServiceTest.java

@ -0,0 +1,78 @@
/**
* Copyright © 2016-2026 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.service.queue;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.springframework.test.util.ReflectionTestUtils;
import org.thingsboard.server.common.data.rpc.RpcError;
import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotificationMsg;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.service.rpc.TbRuleEngineDeviceRpcService;
import java.util.UUID;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.BDDMockito.then;
import static org.mockito.Mockito.doCallRealMethod;
@ExtendWith(MockitoExtension.class)
public class DefaultTbRuleEngineConsumerServiceTest {
@Mock
private TbRuleEngineDeviceRpcService tbDeviceRpcServiceMock;
@Mock
private TbCallback tbCallbackMock;
@Mock
private DefaultTbRuleEngineConsumerService defaultTbRuleEngineConsumerServiceMock;
@Test
public void givenNotFoundErrorAndNoResponse_whenHandleFromDeviceRpcResponse_thenNotFoundAndNullResponseAreRecovered() {
// GIVEN
ReflectionTestUtils.setField(defaultTbRuleEngineConsumerServiceMock, "tbDeviceRpcService", tbDeviceRpcServiceMock);
var requestId = UUID.randomUUID();
// error = NOT_FOUND.ordinal() (0) and response left unset: the previously broken combination
// ('error > 0' dropped NOT_FOUND, proto3 default collapsed a null response to "").
var proto = TransportProtos.FromDeviceRPCResponseProto.newBuilder()
.setRequestIdMSB(requestId.getMostSignificantBits())
.setRequestIdLSB(requestId.getLeastSignificantBits())
.setError(RpcError.NOT_FOUND.ordinal())
.build();
var nfMsg = ToRuleEngineNotificationMsg.newBuilder().setFromDeviceRpcResponse(proto).build();
var queueMsg = new TbProtoQueueMsg<>(requestId, nfMsg);
doCallRealMethod().when(defaultTbRuleEngineConsumerServiceMock).handleNotification(requestId, queueMsg, tbCallbackMock);
// WHEN
defaultTbRuleEngineConsumerServiceMock.handleNotification(requestId, queueMsg, tbCallbackMock);
// THEN
var responseCaptor = ArgumentCaptor.forClass(FromDeviceRpcResponse.class);
then(tbDeviceRpcServiceMock).should().processRpcResponseFromDevice(responseCaptor.capture());
var response = responseCaptor.getValue();
assertThat(response.getId()).isEqualTo(requestId);
assertThat(response.getError()).contains(RpcError.NOT_FOUND);
assertThat(response.getResponse()).isEmpty();
then(tbCallbackMock).should().onSuccess();
}
}

62
common/coap-server/src/main/java/org/thingsboard/server/coapserver/DefaultCoapServerService.java

@ -85,7 +85,9 @@ public class DefaultCoapServerService implements CoapServerService, SmartInitial
dtlsSessionsExecutor.shutdownNow();
}
log.info("Stopping CoAP server!");
server.destroy();
if (server != null) {
server.destroy();
}
log.info("CoAP server stopped!");
}
@ -105,27 +107,47 @@ public class DefaultCoapServerService implements CoapServerService, SmartInitial
private CoapServer createCoapServer() throws UnknownHostException {
Configuration networkConfig = createNetworkConfiguration();
server = new CoapServer(networkConfig);
CoapEndpoint.Builder noSecCoapEndpointBuilder = new CoapEndpoint.Builder();
InetAddress addr = InetAddress.getByName(coapServerContext.getHost());
InetSocketAddress sockAddr = new InetSocketAddress(addr, coapServerContext.getPort());
noSecCoapEndpointBuilder.setInetSocketAddress(sockAddr);
try {
server = new CoapServer(networkConfig);
CoapEndpoint.Builder noSecCoapEndpointBuilder = new CoapEndpoint.Builder();
InetAddress addr = InetAddress.getByName(coapServerContext.getHost());
InetSocketAddress sockAddr = new InetSocketAddress(addr, coapServerContext.getPort());
noSecCoapEndpointBuilder.setInetSocketAddress(sockAddr);
noSecCoapEndpointBuilder.setConfiguration(networkConfig);
CoapEndpoint noSecCoapEndpoint = noSecCoapEndpointBuilder.build();
server.addEndpoint(noSecCoapEndpoint);
if (isDtlsEnabled()) {
createDtlsEndpoint(networkConfig);
dtlsSessionsExecutor = ThingsBoardExecutors.newSingleThreadScheduledExecutor(getClass().getSimpleName());
dtlsSessionsExecutor.scheduleAtFixedRate(this::evictTimeoutSessions, new Random().nextInt((int) getDtlsSessionReportTimeout()), getDtlsSessionReportTimeout(), TimeUnit.MILLISECONDS);
}
Resource root = server.getRoot();
TbCoapServerMessageDeliverer messageDeliverer = new TbCoapServerMessageDeliverer(root);
server.setMessageDeliverer(messageDeliverer);
noSecCoapEndpointBuilder.setConfiguration(networkConfig);
CoapEndpoint noSecCoapEndpoint = noSecCoapEndpointBuilder.build();
server.addEndpoint(noSecCoapEndpoint);
if (isDtlsEnabled()) {
createDtlsEndpoint(networkConfig);
dtlsSessionsExecutor = ThingsBoardExecutors.newSingleThreadScheduledExecutor(getClass().getSimpleName());
dtlsSessionsExecutor.scheduleAtFixedRate(this::evictTimeoutSessions, new Random().nextInt((int) getDtlsSessionReportTimeout()), getDtlsSessionReportTimeout(), TimeUnit.MILLISECONDS);
server.start();
return server;
} catch (RuntimeException | UnknownHostException e) {
log.error("Failed to start CoAP server, releasing resources", e);
try {
if (dtlsSessionsExecutor != null) {
dtlsSessionsExecutor.shutdownNow();
}
if (server != null) {
server.destroy();
}
} catch (Exception suppressed) {
e.addSuppressed(suppressed);
} finally {
server = null;
dtlsSessionsExecutor = null;
dtlsConnector = null;
dtlsCoapEndpoint = null;
tbDtlsCertificateVerifier = null;
}
throw e;
}
Resource root = server.getRoot();
TbCoapServerMessageDeliverer messageDeliverer = new TbCoapServerMessageDeliverer(root);
server.setMessageDeliverer(messageDeliverer);
server.start();
return server;
}
private boolean isDtlsEnabled() {

145
common/coap-server/src/test/java/org/thingsboard/server/coapserver/DefaultCoapServerServiceTest.java

@ -0,0 +1,145 @@
/**
* Copyright © 2016-2026 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.coapserver;
import org.eclipse.californium.core.CoapServer;
import org.eclipse.californium.core.network.CoapEndpoint;
import org.eclipse.californium.core.server.resources.Resource;
import org.eclipse.californium.scandium.DTLSConnector;
import org.eclipse.californium.scandium.config.DtlsConnectorConfig;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.Mock;
import org.mockito.MockedConstruction;
import org.mockito.MockedStatic;
import org.mockito.junit.jupiter.MockitoExtension;
import org.springframework.test.util.ReflectionTestUtils;
import org.thingsboard.common.util.ThingsBoardExecutors;
import java.net.DatagramSocket;
import java.net.InetAddress;
import java.net.InetSocketAddress;
import java.util.concurrent.ScheduledExecutorService;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.mockConstruction;
import static org.mockito.Mockito.mockStatic;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
public class DefaultCoapServerServiceTest {
private static final String HOST = "127.0.0.1";
@Mock
private CoapServerContext mockCoapServerContext;
private DefaultCoapServerService service;
private DatagramSocket occupiedSocket;
private int occupiedPort;
@BeforeEach
public void setUp() throws Exception {
occupiedSocket = new DatagramSocket(new InetSocketAddress(InetAddress.getByName(HOST), 0));
occupiedPort = occupiedSocket.getLocalPort();
service = new DefaultCoapServerService();
ReflectionTestUtils.setField(service, "coapServerContext", mockCoapServerContext);
when(mockCoapServerContext.getHost()).thenReturn(HOST);
when(mockCoapServerContext.getPort()).thenReturn(occupiedPort);
when(mockCoapServerContext.getDtlsSettings()).thenReturn(null);
}
@AfterEach
public void tearDown() {
if (occupiedSocket != null && !occupiedSocket.isClosed()) {
occupiedSocket.close();
}
}
@Test
public void whenPlainBindFails_thenInitThrowsAndReleasesCoapServer() {
assertThatThrownBy(() -> service.init())
.isInstanceOf(IllegalStateException.class)
.hasMessageContaining("None of the server endpoints could be started");
assertThat(ReflectionTestUtils.getField(service, "server")).isNull();
assertThat(ReflectionTestUtils.getField(service, "dtlsSessionsExecutor")).isNull();
assertThat(ReflectionTestUtils.getField(service, "dtlsConnector")).isNull();
assertThat(ReflectionTestUtils.getField(service, "dtlsCoapEndpoint")).isNull();
assertThat(ReflectionTestUtils.getField(service, "tbDtlsCertificateVerifier")).isNull();
}
@Test
public void whenDtlsEnabledAndStartFails_thenInitShutsDownDtlsExecutorAndReleasesCoapServer() throws Exception {
// DTLS enabled: the DTLS endpoint is created and dtlsSessionsExecutor is scheduled before server.start().
// This exercises the catch's dtlsSessionsExecutor.shutdownNow() branch, which the plain-bind test does not.
TbCoapDtlsSettings mockDtlsSettings = mock(TbCoapDtlsSettings.class);
when(mockCoapServerContext.getDtlsSettings()).thenReturn(mockDtlsSettings);
DtlsConnectorConfig mockDtlsConfig = mock(DtlsConnectorConfig.class);
when(mockDtlsConfig.getAddress()).thenReturn(new InetSocketAddress(InetAddress.getByName(HOST), occupiedPort + 1));
TbCoapDtlsCertificateVerifier mockVerifier = mock(TbCoapDtlsCertificateVerifier.class);
when(mockVerifier.getDtlsSessionReportTimeout()).thenReturn(1800000L);
when(mockDtlsConfig.getAdvancedCertificateVerifier()).thenReturn(mockVerifier);
when(mockDtlsSettings.dtlsConnectorConfig(any())).thenReturn(mockDtlsConfig);
ScheduledExecutorService mockExecutor = mock(ScheduledExecutorService.class);
Resource mockRoot = mock(Resource.class);
try (MockedStatic<ThingsBoardExecutors> executorsStatic = mockStatic(ThingsBoardExecutors.class);
MockedConstruction<CoapServer> serverMock = mockConstruction(CoapServer.class, (server, ctx) -> {
when(server.getRoot()).thenReturn(mockRoot);
doThrow(new IllegalStateException("None of the server endpoints could be started")).when(server).start();
});
MockedConstruction<DTLSConnector> dtlsMock = mockConstruction(DTLSConnector.class);
MockedConstruction<CoapEndpoint.Builder> builderMock = mockConstruction(CoapEndpoint.Builder.class, (builder, ctx) -> {
when(builder.setInetSocketAddress(any())).thenReturn(builder);
when(builder.setConfiguration(any())).thenReturn(builder);
when(builder.setConnector(any(DTLSConnector.class))).thenReturn(builder);
when(builder.build()).thenReturn(mock(CoapEndpoint.class));
})) {
executorsStatic.when(() -> ThingsBoardExecutors.newSingleThreadScheduledExecutor(anyString())).thenReturn(mockExecutor);
assertThatThrownBy(() -> service.init())
.isInstanceOf(IllegalStateException.class)
.hasMessageContaining("None of the server endpoints could be started");
// DTLS branch was actually entered and the executor was created...
verify(mockDtlsSettings).dtlsConnectorConfig(any());
// ...and the cleanup branch shut it down and destroyed the server.
verify(mockExecutor).shutdownNow();
verify(serverMock.constructed().get(0)).destroy();
}
assertThat(ReflectionTestUtils.getField(service, "server")).isNull();
assertThat(ReflectionTestUtils.getField(service, "dtlsSessionsExecutor")).isNull();
assertThat(ReflectionTestUtils.getField(service, "dtlsConnector")).isNull();
assertThat(ReflectionTestUtils.getField(service, "dtlsCoapEndpoint")).isNull();
assertThat(ReflectionTestUtils.getField(service, "tbDtlsCertificateVerifier")).isNull();
}
}

7
common/data/src/main/java/org/thingsboard/server/common/data/StringUtils.java

@ -26,6 +26,7 @@ import java.util.Base64;
import java.util.List;
import java.util.Objects;
import java.util.function.Function;
import java.util.regex.Pattern;
import static org.apache.commons.lang3.StringUtils.repeat;
@ -39,6 +40,12 @@ public class StringUtils {
public static final int INDEX_NOT_FOUND = -1;
public static final Pattern CONTROL_CHARS = Pattern.compile("[\\x00-\\x1F\\x7F]");
public static boolean containsControlChars(String source) {
return source != null && CONTROL_CHARS.matcher(source).find();
}
public static boolean isEmpty(String source) {
return source == null || source.isEmpty();
}

2
common/data/src/main/java/org/thingsboard/server/common/data/alarm/AlarmCommentSubType.java

@ -24,7 +24,7 @@ public enum AlarmCommentSubType {
ASSIGNED_TO_USER("Alarm was assigned by user %s to user %s"),
UNASSIGNED_BY_USER("Alarm was unassigned by user %s"),
UNASSIGNED_FROM_DELETED_USER("Alarm was unassigned because user %s - was deleted"),
COMMENT_DELETED("User %s deleted his comment"),
COMMENT_DELETED("Comment was deleted by user %s"),
SEVERITY_CHANGED("Alarm severity was updated from %s to %s");
@Getter

12
common/data/src/main/java/org/thingsboard/server/common/data/rpc/RpcError.java

@ -20,4 +20,16 @@ package org.thingsboard.server.common.data.rpc;
*/
public enum RpcError {
NOT_FOUND, FORBIDDEN, NO_ACTIVE_CONNECTION, TIMEOUT, INTERNAL;
private static final RpcError[] VALUES = values();
/**
* Resolves an {@link RpcError} from the proto {@code error} ordinal.
* Returns {@code null} both for the "no error" sentinel (negative value) and for unknown ordinals
* that a newer node in a mixed-version cluster might emit, so callers never hit an
* {@link ArrayIndexOutOfBoundsException}.
*/
public static RpcError fromProtoErrorCode(int errorCode) {
return errorCode >= 0 && errorCode < VALUES.length ? VALUES[errorCode] : null;
}
}

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

@ -585,10 +585,11 @@ public class ProtoUtils {
}
private static ToDeviceActorNotificationMsg fromProto(TransportProtos.FromDeviceRpcResponseActorMsgProto proto) {
TransportProtos.FromDeviceRPCResponseProto rpcResponse = proto.getRpcResponse();
FromDeviceRpcResponse fromDeviceRpcResponse = new FromDeviceRpcResponse(
new UUID(proto.getRpcResponse().getRequestIdMSB(), proto.getRpcResponse().getRequestIdLSB()),
proto.getRpcResponse().getResponse(),
proto.getRpcResponse().getError() >= 0 ? RpcError.values()[proto.getRpcResponse().getError()] : null);
new UUID(rpcResponse.getRequestIdMSB(), rpcResponse.getRequestIdLSB()),
rpcResponse.hasResponse() ? rpcResponse.getResponse() : null,
RpcError.fromProtoErrorCode(rpcResponse.getError()));
return new FromDeviceRpcResponseActorMsg(
proto.getRequestId(),
TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())),

2
common/proto/src/main/proto/queue.proto

@ -1273,7 +1273,7 @@ message LocalSubscriptionServiceMsgProto {
message FromDeviceRPCResponseProto {
int64 requestIdMSB = 1;
int64 requestIdLSB = 2;
string response = 3;
optional string response = 3;
int32 error = 4;
}

11
common/proto/src/test/java/org/thingsboard/server/common/util/ProtoUtilsTest.java

@ -228,6 +228,17 @@ class ProtoUtilsTest {
assertThat(ProtoUtils.fromProto(serializedMsg)).as("deserialized").isEqualTo(msg);
}
@Test
void protoFromDeviceRpcResponseOnewaySerialization() {
// Oneway RPC success: response and error are both null. Relies on the proto
// 'optional string response' presence bit so the receiver round-trips null
// rather than seeing the proto3 default "".
FromDeviceRpcResponseActorMsg msg = new FromDeviceRpcResponseActorMsg(23, tenantId, deviceId, new FromDeviceRpcResponse(id, null, null));
TransportProtos.ToDeviceActorNotificationMsgProto serializedMsg = ProtoUtils.toProto(msg);
Assertions.assertNotNull(serializedMsg);
assertThat(ProtoUtils.fromProto(serializedMsg)).as("deserialized").isEqualTo(msg);
}
@Test
void protoRemoveRpcActorSerialization() {
RemoveRpcActorMsg msg = new RemoveRpcActorMsg(tenantId, deviceId, id);

25
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/LwM2MTransportBootstrapService.java

@ -82,13 +82,32 @@ public class LwM2MTransportBootstrapService implements SmartInitializingSingleto
@PostConstruct
public void init() {
log.info("Starting LwM2M transport bootstrap server...");
this.server = getLhBootstrapServer();
this.server.start();
log.info("Started LwM2M transport bootstrap server.");
LeshanBootstrapServer bootstrapServer = null;
try {
bootstrapServer = getLhBootstrapServer();
this.server = bootstrapServer;
bootstrapServer.start();
log.info("Started LwM2M transport bootstrap server.");
} catch (RuntimeException e) {
log.error("Failed to start LwM2M transport bootstrap server, releasing resources", e);
try {
if (bootstrapServer != null) {
bootstrapServer.destroy();
}
} catch (Exception suppressed) {
e.addSuppressed(suppressed);
} finally {
this.server = null;
}
throw e;
}
}
@PreDestroy
public void shutdown() {
if (server == null) {
return;
}
try {
log.info("Stopping LwM2M transport bootstrap server!");
server.destroy();

115
common/transport/lwm2m/src/test/java/org/thingsboard/server/transport/lwm2m/bootstrap/LwM2MTransportBootstrapServiceTest.java

@ -0,0 +1,115 @@
/**
* Copyright © 2016-2026 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.transport.lwm2m.bootstrap;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.mockito.junit.jupiter.MockitoSettings;
import org.mockito.quality.Strictness;
import org.springframework.test.util.ReflectionTestUtils;
import org.thingsboard.server.common.transport.TransportService;
import org.thingsboard.server.transport.lwm2m.bootstrap.secure.TbLwM2MDtlsBootstrapCertificateVerifier;
import org.thingsboard.server.transport.lwm2m.bootstrap.store.LwM2MBootstrapSecurityStore;
import org.thingsboard.server.transport.lwm2m.bootstrap.store.LwM2MInMemoryBootstrapConfigStore;
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportBootstrapConfig;
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig;
import java.net.DatagramSocket;
import java.net.InetAddress;
import java.net.InetSocketAddress;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
@MockitoSettings(strictness = Strictness.LENIENT)
public class LwM2MTransportBootstrapServiceTest {
private static final String HOST = "127.0.0.1";
@Mock
private LwM2MTransportServerConfig serverConfig;
@Mock
private LwM2MTransportBootstrapConfig bootstrapConfig;
@Mock
private LwM2MBootstrapSecurityStore lwM2MBootstrapSecurityStore;
@Mock
private LwM2MInMemoryBootstrapConfigStore lwM2MInMemoryBootstrapConfigStore;
@Mock
private TransportService transportService;
@Mock
private TbLwM2MDtlsBootstrapCertificateVerifier certificateVerifier;
private LwM2MTransportBootstrapService service;
private DatagramSocket occupiedPlain;
private DatagramSocket occupiedSecure;
@BeforeEach
public void setUp() throws Exception {
occupiedPlain = new DatagramSocket(new InetSocketAddress(InetAddress.getByName(HOST), 0));
occupiedSecure = new DatagramSocket(new InetSocketAddress(InetAddress.getByName(HOST), 0));
when(bootstrapConfig.getHost()).thenReturn(HOST);
when(bootstrapConfig.getPort()).thenReturn(occupiedPlain.getLocalPort());
when(bootstrapConfig.getSecureHost()).thenReturn(HOST);
when(bootstrapConfig.getSecurePort()).thenReturn(occupiedSecure.getLocalPort());
when(bootstrapConfig.getSslCredentials()).thenReturn(null);
when(serverConfig.isRecommendedCiphers()).thenReturn(false);
when(serverConfig.isRecommendedSupportedGroups()).thenReturn(false);
when(serverConfig.getDtlsRetransmissionTimeout()).thenReturn(9000);
when(serverConfig.getDtlsCidLength()).thenReturn(null);
service = new LwM2MTransportBootstrapService(
serverConfig,
bootstrapConfig,
lwM2MBootstrapSecurityStore,
lwM2MInMemoryBootstrapConfigStore,
transportService,
certificateVerifier
);
}
@AfterEach
public void tearDown() {
if (occupiedPlain != null && !occupiedPlain.isClosed()) {
occupiedPlain.close();
}
if (occupiedSecure != null && !occupiedSecure.isClosed()) {
occupiedSecure.close();
}
}
@Test
public void whenEndpointsFailToStart_thenInitThrowsAndReleasesBootstrapServer() {
assertThatThrownBy(() -> service.init())
.isInstanceOf(IllegalStateException.class)
.hasMessageContaining("None of the server endpoints could be started");
assertThat(ReflectionTestUtils.getField(service, "server")).isNull();
}
}

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

@ -107,11 +107,12 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi
@Override
public void update(TenantProfileUpdateResult update) {
log.info("Received tenant profile update: {}", update.getProfile());
EntityTransportRateLimits tenantRateLimitPrototype = createRateLimits(update.getProfile(), TENANT_LIMITS);
EntityTransportRateLimits deviceRateLimitPrototype = createRateLimits(update.getProfile(), DEVICE_LIMITS);
EntityTransportRateLimits gatewayRateLimitPrototype = createRateLimits(update.getProfile(), GATEWAY_LIMITS);
EntityTransportRateLimits gatewayDeviceRateLimitPrototype = createRateLimits(update.getProfile(), GATEWAY_DEVICE_LIMITS);
TenantProfile profile = update.getProfile();
log.info("Received tenant profile update: {}", profile);
EntityTransportRateLimits tenantRateLimitPrototype = createRateLimits(profile, TENANT_LIMITS);
EntityTransportRateLimits deviceRateLimitPrototype = createRateLimits(profile, DEVICE_LIMITS);
EntityTransportRateLimits gatewayRateLimitPrototype = createRateLimits(profile, GATEWAY_LIMITS);
EntityTransportRateLimits gatewayDeviceRateLimitPrototype = createRateLimits(profile, GATEWAY_DEVICE_LIMITS);
for (TenantId tenantId : update.getAffectedTenants()) {
update(tenantId, tenantRateLimitPrototype, deviceRateLimitPrototype, gatewayRateLimitPrototype, gatewayDeviceRateLimitPrototype);
}
@ -119,11 +120,13 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi
@Override
public void update(TenantId tenantId) {
EntityTransportRateLimits tenantRateLimitPrototype = createRateLimits(tenantProfileCache.get(tenantId), TENANT_LIMITS);
EntityTransportRateLimits deviceRateLimitPrototype = createRateLimits(tenantProfileCache.get(tenantId), DEVICE_LIMITS);
EntityTransportRateLimits gatewayRateLimitPrototype = createRateLimits(tenantProfileCache.get(tenantId), GATEWAY_LIMITS);
EntityTransportRateLimits gatewayDeviceRateLimitPrototype = createRateLimits(tenantProfileCache.get(tenantId), GATEWAY_DEVICE_LIMITS);
update(tenantId, tenantRateLimitPrototype, deviceRateLimitPrototype, gatewayRateLimitPrototype, gatewayDeviceRateLimitPrototype);
TenantProfile profile = tenantProfileCache.get(tenantId);
update(tenantId,
createRateLimits(profile, TENANT_LIMITS),
createRateLimits(profile, DEVICE_LIMITS),
createRateLimits(profile, GATEWAY_LIMITS),
createRateLimits(profile, GATEWAY_DEVICE_LIMITS)
);
}
private void update(TenantId tenantId, EntityTransportRateLimits tenantRateLimitPrototype, EntityTransportRateLimits deviceRateLimitPrototype,
@ -231,25 +234,26 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi
BiConsumer<T, EntityTransportRateLimits> putFunction) {
EntityTransportRateLimits oldRateLimits = getFunction.apply(entityId);
if (oldRateLimits == null) {
if (EntityType.TENANT.equals(entityId.getEntityType())) {
log.info("[{}] New rate limits: {}", entityId, newRateLimits);
} else {
log.debug("[{}] New rate limits: {}", entityId, newRateLimits);
}
logLimits(entityId, "New", newRateLimits);
putFunction.accept(entityId, newRateLimits);
} else {
EntityTransportRateLimits updated = merge(oldRateLimits, newRateLimits);
if (updated != null) {
if (EntityType.TENANT.equals(entityId.getEntityType())) {
log.info("[{}] Updated rate limits: {}", entityId, updated);
} else {
log.debug("[{}] Updated rate limits: {}", entityId, updated);
}
logLimits(entityId, "Updated", updated);
putFunction.accept(entityId, updated);
}
}
}
private void logLimits(EntityId entityId, String action, EntityTransportRateLimits limits) {
// Tenant-level changes are logged at INFO; the much noisier per-device/gateway ones at DEBUG.
if (EntityType.TENANT.equals(entityId.getEntityType())) {
log.info("[{}] {} rate limits: {}", entityId, action, limits);
} else {
log.debug("[{}] {} rate limits: {}", entityId, action, limits);
}
}
private EntityTransportRateLimits merge(EntityTransportRateLimits oldRateLimits, EntityTransportRateLimits newRateLimits) {
boolean regularUpdate = !oldRateLimits.getRegularMsgRateLimit().getConfiguration().equals(newRateLimits.getRegularMsgRateLimit().getConfiguration());
boolean telemetryMsgRateUpdate = !oldRateLimits.getTelemetryMsgRateLimit().getConfiguration().equals(newRateLimits.getTelemetryMsgRateLimit().getConfiguration());
@ -269,36 +273,12 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi
DefaultTenantProfileConfiguration profile = (DefaultTenantProfileConfiguration) profileData.getConfiguration();
if (profile == null) {
return new EntityTransportRateLimits(ALLOW, ALLOW, ALLOW);
} else {
TransportRateLimit regularMsgRateLimit;
TransportRateLimit telemetryMsgRateLimit;
TransportRateLimit telemetryDpRateLimit;
switch (limitsType) {
case TENANT_LIMITS -> {
regularMsgRateLimit = newLimit(profile.getTransportTenantMsgRateLimit());
telemetryMsgRateLimit = newLimit(profile.getTransportTenantTelemetryMsgRateLimit());
telemetryDpRateLimit = newLimit(profile.getTransportTenantTelemetryDataPointsRateLimit());
}
case DEVICE_LIMITS -> {
regularMsgRateLimit = newLimit(profile.getTransportDeviceMsgRateLimit());
telemetryMsgRateLimit = newLimit(profile.getTransportDeviceTelemetryMsgRateLimit());
telemetryDpRateLimit = newLimit(profile.getTransportDeviceTelemetryDataPointsRateLimit());
}
case GATEWAY_LIMITS -> {
regularMsgRateLimit = newLimit(profile.getTransportGatewayMsgRateLimit());
telemetryMsgRateLimit = newLimit(profile.getTransportGatewayTelemetryMsgRateLimit());
telemetryDpRateLimit = newLimit(profile.getTransportGatewayTelemetryDataPointsRateLimit());
}
case GATEWAY_DEVICE_LIMITS -> {
regularMsgRateLimit = newLimit(profile.getTransportGatewayDeviceMsgRateLimit());
telemetryMsgRateLimit = newLimit(profile.getTransportGatewayDeviceTelemetryMsgRateLimit());
telemetryDpRateLimit = newLimit(profile.getTransportGatewayDeviceTelemetryDataPointsRateLimit());
}
default -> throw new IllegalStateException("Unknown limits type: " + limitsType);
}
return new EntityTransportRateLimits(regularMsgRateLimit, telemetryMsgRateLimit, telemetryDpRateLimit);
}
return new EntityTransportRateLimits(
newLimit(limitsType.getRegularMsgRateLimit().apply(profile)),
newLimit(limitsType.getTelemetryMsgRateLimit().apply(profile)),
newLimit(limitsType.getTelemetryDataPointsRateLimit().apply(profile))
);
}
private static TransportRateLimit newLimit(String config) {
@ -306,31 +286,36 @@ public class DefaultTransportRateLimitService implements TransportRateLimitServi
}
private EntityTransportRateLimits getTenantRateLimits(TenantId tenantId) {
return perTenantLimits.computeIfAbsent(tenantId, k -> createRateLimits(tenantProfileCache.get(tenantId), TENANT_LIMITS));
return getRateLimits(perTenantLimits, tenantId, tenantId, TENANT_LIMITS, null);
}
private EntityTransportRateLimits getDeviceRateLimits(TenantId tenantId, DeviceId deviceId) {
return perDeviceLimits.computeIfAbsent(deviceId, k -> {
EntityTransportRateLimits limits = createRateLimits(tenantProfileCache.get(tenantId), DEVICE_LIMITS);
getTenantDevices(tenantId).add(deviceId);
return limits;
});
return getRateLimits(perDeviceLimits, tenantId, deviceId, DEVICE_LIMITS, () -> getTenantDevices(tenantId).add(deviceId));
}
private EntityTransportRateLimits getGatewayRateLimits(TenantId tenantId, DeviceId gatewayId) {
return perGatewayLimits.computeIfAbsent(gatewayId, k -> {
EntityTransportRateLimits limits = createRateLimits(tenantProfileCache.get(tenantId), GATEWAY_LIMITS);
getTenantGateways(tenantId).add(gatewayId);
return limits;
});
return getRateLimits(perGatewayLimits, tenantId, gatewayId, GATEWAY_LIMITS, () -> getTenantGateways(tenantId).add(gatewayId));
}
private EntityTransportRateLimits getGatewayDeviceRateLimits(TenantId tenantId, DeviceId gatewayId) {
return perGatewayDeviceLimits.computeIfAbsent(gatewayId, k -> {
EntityTransportRateLimits limits = createRateLimits(tenantProfileCache.get(tenantId), GATEWAY_DEVICE_LIMITS);
getTenantGatewayDevices(tenantId).add(gatewayId);
return limits;
});
return getRateLimits(perGatewayDeviceLimits, tenantId, gatewayId, GATEWAY_DEVICE_LIMITS, () -> getTenantGatewayDevices(tenantId).add(gatewayId));
}
private <T extends EntityId> EntityTransportRateLimits getRateLimits(ConcurrentMap<T, EntityTransportRateLimits> limitsMap, TenantId tenantId,
T entityId, TransportLimitsType limitsType, Runnable onMiss) {
EntityTransportRateLimits limits = limitsMap.get(entityId);
if (limits == null) {
// Resolve the tenant profile WITHOUT holding the ConcurrentHashMap bin lock: the fetch may
// block on a cross-service round-trip, so it must run before computeIfAbsent's mapping function.
TenantProfile tenantProfile = tenantProfileCache.get(tenantId);
limits = limitsMap.computeIfAbsent(entityId, k -> createRateLimits(tenantProfile, limitsType));
// Runs on every observed miss, including callers that lost the computeIfAbsent race and got an
// existing value back - NOT only on actual creation, so the callback must be idempotent.
if (onMiss != null) {
onMiss.run();
}
}
return limits;
}
private Set<DeviceId> getTenantDevices(TenantId tenantId) {

35
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/limits/TransportLimitsType.java

@ -15,6 +15,39 @@
*/
package org.thingsboard.server.common.transport.limits;
import lombok.Getter;
import lombok.RequiredArgsConstructor;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
import java.util.function.Function;
@Getter
@RequiredArgsConstructor
public enum TransportLimitsType {
TENANT_LIMITS, DEVICE_LIMITS, GATEWAY_LIMITS, GATEWAY_DEVICE_LIMITS
TENANT_LIMITS(
DefaultTenantProfileConfiguration::getTransportTenantMsgRateLimit,
DefaultTenantProfileConfiguration::getTransportTenantTelemetryMsgRateLimit,
DefaultTenantProfileConfiguration::getTransportTenantTelemetryDataPointsRateLimit
),
DEVICE_LIMITS(
DefaultTenantProfileConfiguration::getTransportDeviceMsgRateLimit,
DefaultTenantProfileConfiguration::getTransportDeviceTelemetryMsgRateLimit,
DefaultTenantProfileConfiguration::getTransportDeviceTelemetryDataPointsRateLimit
),
GATEWAY_LIMITS(
DefaultTenantProfileConfiguration::getTransportGatewayMsgRateLimit,
DefaultTenantProfileConfiguration::getTransportGatewayTelemetryMsgRateLimit,
DefaultTenantProfileConfiguration::getTransportGatewayTelemetryDataPointsRateLimit
),
GATEWAY_DEVICE_LIMITS(
DefaultTenantProfileConfiguration::getTransportGatewayDeviceMsgRateLimit,
DefaultTenantProfileConfiguration::getTransportGatewayDeviceTelemetryMsgRateLimit,
DefaultTenantProfileConfiguration::getTransportGatewayDeviceTelemetryDataPointsRateLimit
);
private final Function<DefaultTenantProfileConfiguration, String> regularMsgRateLimit;
private final Function<DefaultTenantProfileConfiguration, String> telemetryMsgRateLimit;
private final Function<DefaultTenantProfileConfiguration, String> telemetryDataPointsRateLimit;
}

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

@ -153,6 +153,8 @@ public class DefaultTransportService extends TransportActivityManager implements
private int notificationsPollDuration;
@Value("${transport.stats.enabled:false}")
private boolean statsEnabled;
@Value("${transport.callback_thread_pool_size:20}")
private int callbackThreadPoolSize;
@Autowired
@Lazy
@ -198,7 +200,7 @@ public class DefaultTransportService extends TransportActivityManager implements
this.ruleEngineProducerStats = statsFactory.createMessagesStats(StatsType.RULE_ENGINE.getName() + ".producer");
this.tbCoreProducerStats = statsFactory.createMessagesStats(StatsType.CORE.getName() + ".producer");
this.transportApiStats = statsFactory.createMessagesStats(StatsType.TRANSPORT.getName() + ".producer");
this.transportCallbackExecutor = ThingsBoardExecutors.newWorkStealingPool(20, getClass());
this.transportCallbackExecutor = ThingsBoardExecutors.newWorkStealingPool(callbackThreadPoolSize, getClass());
this.scheduler.scheduleAtFixedRate(this::invalidateRateLimits, new Random().nextInt((int) sessionReportTimeout), sessionReportTimeout, TimeUnit.MILLISECONDS);
transportApiRequestTemplate = queueProvider.createTransportApiRequestTemplate();
transportApiRequestTemplate.setMessagesStats(transportApiStats);

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

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.common.transport.service;
import com.google.common.util.concurrent.Striped;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Lazy;
@ -37,14 +38,20 @@ 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
@TbTransportComponent
@Slf4j
public class DefaultTransportTenantProfileCache implements TransportTenantProfileCache {
private final Lock tenantProfileFetchLock = new ReentrantLock();
// Number of stripes for the per-tenant fetch locks. Only contended during concurrent cold-cache
// misses (cached tenants never take the lock), and concurrent fetches are already bounded by the
// transport callback pool, so this comfortably over-provisions the realistic concurrency.
private static final int TENANT_PROFILE_FETCH_LOCK_STRIPES = 1024;
// Bounded set of per-tenant locks: de-duplicates concurrent misses for the same tenant while
// letting different tenants fetch concurrently (eager array - no weak-ref overhead at this size).
private final Striped<Lock> tenantProfileFetchLocks = Striped.lock(TENANT_PROFILE_FETCH_LOCK_STRIPES);
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<>();
@ -103,43 +110,52 @@ public class DefaultTransportTenantProfileCache implements TransportTenantProfil
}
private TenantProfile getTenantProfile(TenantId tenantId) {
TenantProfile profile = null;
TenantProfileId tenantProfileId = tenantIds.get(tenantId);
if (tenantProfileId != null) {
profile = profiles.get(tenantProfileId);
}
TenantProfile profile = lookupCached(tenantId);
if (profile == null) {
tenantProfileFetchLock.lock();
// Per-tenant lock: de-duplicates concurrent misses for the SAME tenant while allowing
// different tenants to resolve their profiles concurrently.
Lock lock = tenantProfileFetchLocks.get(tenantId);
lock.lock();
try {
tenantProfileId = tenantIds.get(tenantId);
if (tenantProfileId != null) {
profile = profiles.get(tenantProfileId);
}
profile = lookupCached(tenantId);
if (profile == null) {
TransportProtos.GetEntityProfileRequestMsg msg = TransportProtos.GetEntityProfileRequestMsg.newBuilder()
.setEntityType(EntityType.TENANT.name())
.setEntityIdMSB(tenantId.getId().getMostSignificantBits())
.setEntityIdLSB(tenantId.getId().getLeastSignificantBits())
.build();
TransportProtos.GetEntityProfileResponseMsg entityProfileMsg = transportService.getEntityProfile(msg);
profile = ProtoUtils.fromProto(entityProfileMsg.getTenantProfile());
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());
ApiUsageState apiUsageState = ProtoUtils.fromProto(entityProfileMsg.getApiState());
rateLimitService.update(tenantId, apiUsageState.isTransportEnabled());
profile = fetchAndCacheTenantProfile(tenantId);
}
} finally {
tenantProfileFetchLock.unlock();
lock.unlock();
}
}
return profile;
}
private TenantProfile lookupCached(TenantId tenantId) {
TenantProfileId tenantProfileId = tenantIds.get(tenantId);
if (tenantProfileId != null) {
return profiles.get(tenantProfileId);
}
return null;
}
private TenantProfile fetchAndCacheTenantProfile(TenantId tenantId) {
TransportProtos.GetEntityProfileRequestMsg msg = TransportProtos.GetEntityProfileRequestMsg.newBuilder()
.setEntityType(EntityType.TENANT.name())
.setEntityIdMSB(tenantId.getId().getMostSignificantBits())
.setEntityIdLSB(tenantId.getId().getLeastSignificantBits())
.build();
TransportProtos.GetEntityProfileResponseMsg entityProfileMsg = transportService.getEntityProfile(msg);
TenantProfile profile = ProtoUtils.fromProto(entityProfileMsg.getTenantProfile());
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());
ApiUsageState apiUsageState = ProtoUtils.fromProto(entityProfileMsg.getApiState());
rateLimitService.update(tenantId, apiUsageState.isTransportEnabled());
return profile;
}
}

202
common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/limits/DefaultTransportRateLimitServiceTest.java

@ -0,0 +1,202 @@
/**
* Copyright © 2016-2026 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.limits;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.EnumSource;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.TenantProfileId;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
import org.thingsboard.server.common.data.tenant.profile.TenantProfileData;
import org.thingsboard.server.common.transport.TransportTenantProfileCache;
import org.thingsboard.server.common.transport.profile.TenantProfileUpdateResult;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
class DefaultTransportRateLimitServiceTest {
private TransportTenantProfileCache tenantProfileCache;
private ExecutorService executor;
private final TenantId tenant = TenantId.fromUUID(UUID.randomUUID());
@BeforeEach
void setUp() {
tenantProfileCache = mock(TransportTenantProfileCache.class);
executor = Executors.newCachedThreadPool();
}
@AfterEach
void tearDown() {
executor.shutdownNow();
}
@Test
void checkLimitsDoesNotHoldMapBinLockAcrossProfileFetch() throws Exception {
// Two concurrent rate-limit checks for the SAME tenant must both be able to reach
// the (blocking) tenant-profile fetch concurrently. If the blocking fetch runs inside
// ConcurrentHashMap.computeIfAbsent, the second caller is stuck on the bin reservation
// node and never reaches the fetch -> the latch never reaches zero.
CountDownLatch bothCallersReachedFetch = new CountDownLatch(2);
CountDownLatch releaseFetch = new CountDownLatch(1);
when(tenantProfileCache.get(tenant)).thenAnswer(invocation -> {
bothCallersReachedFetch.countDown();
releaseFetch.await(5, TimeUnit.SECONDS);
return tenantProfile();
});
DefaultTransportRateLimitService service = new DefaultTransportRateLimitService(tenantProfileCache);
Runnable check = () -> service.checkLimits(tenant, null, null, 1, false);
executor.submit(check);
executor.submit(check);
boolean bothReached = bothCallersReachedFetch.await(3, TimeUnit.SECONDS);
releaseFetch.countDown();
assertThat(bothReached)
.as("both checkLimits calls should reach the profile fetch concurrently (no bin lock across I/O)")
.isTrue();
}
@ParameterizedTest
@EnumSource(TransportLimitsType.class)
void eachLimitsTypeReadsItsOwnProfileFields(TransportLimitsType type) {
// Distinct sentinel per profile field so a transposed method reference (e.g. GATEWAY_DEVICE_LIMITS
// wired to the plain gateway getters) resolves to the wrong value and fails the assertion.
DefaultTenantProfileConfiguration config = new DefaultTenantProfileConfiguration();
config.setTransportTenantMsgRateLimit("tenant-msg");
config.setTransportTenantTelemetryMsgRateLimit("tenant-tele-msg");
config.setTransportTenantTelemetryDataPointsRateLimit("tenant-tele-dp");
config.setTransportDeviceMsgRateLimit("device-msg");
config.setTransportDeviceTelemetryMsgRateLimit("device-tele-msg");
config.setTransportDeviceTelemetryDataPointsRateLimit("device-tele-dp");
config.setTransportGatewayMsgRateLimit("gateway-msg");
config.setTransportGatewayTelemetryMsgRateLimit("gateway-tele-msg");
config.setTransportGatewayTelemetryDataPointsRateLimit("gateway-tele-dp");
config.setTransportGatewayDeviceMsgRateLimit("gateway-device-msg");
config.setTransportGatewayDeviceTelemetryMsgRateLimit("gateway-device-tele-msg");
config.setTransportGatewayDeviceTelemetryDataPointsRateLimit("gateway-device-tele-dp");
String prefix = switch (type) {
case TENANT_LIMITS -> "tenant";
case DEVICE_LIMITS -> "device";
case GATEWAY_LIMITS -> "gateway";
case GATEWAY_DEVICE_LIMITS -> "gateway-device";
};
assertThat(type.getRegularMsgRateLimit().apply(config)).isEqualTo(prefix + "-msg");
assertThat(type.getTelemetryMsgRateLimit().apply(config)).isEqualTo(prefix + "-tele-msg");
assertThat(type.getTelemetryDataPointsRateLimit().apply(config)).isEqualTo(prefix + "-tele-dp");
}
@ParameterizedTest
@EnumSource(EntityLevel.class)
void profileUpdateReachesEntityTrackedDuringFirstCheck(EntityLevel level) {
DeviceId entity = new DeviceId(UUID.randomUUID());
when(tenantProfileCache.get(tenant)).thenReturn(profileWithRegularMsgLimit(level, "100:600"));
DefaultTransportRateLimitService service = new DefaultTransportRateLimitService(tenantProfileCache);
// First check resolves the (permissive) limit and must register the entity into the per-tenant
// tracking set via the onMiss callback - otherwise a later update(tenantId) can't reach it.
assertThat(level.check(service, tenant, entity))
.as("permissive limit should allow the first %s check", level).isNull();
// Tighten the limit to a single message and push a profile update for this tenant.
service.update(new TenantProfileUpdateResult(profileWithRegularMsgLimit(level, "1:600"), Set.of(tenant)));
// The freshly merged "1:600" bucket allows exactly one message...
assertThat(level.check(service, tenant, entity)).isNull();
// ...and blocks the next one. This only happens if update(tenantId) reached the tracked entity.
assertThat(level.check(service, tenant, entity))
.as("update(tenantId) must reach the tracked %s so the tightened limit applies", level).isNotNull();
}
private TenantProfile tenantProfile() {
return profileWith(new DefaultTenantProfileConfiguration());
}
private TenantProfile profileWithRegularMsgLimit(EntityLevel level, String regularMsgRateLimit) {
DefaultTenantProfileConfiguration config = new DefaultTenantProfileConfiguration();
level.setRegularMsgRateLimit(config, regularMsgRateLimit);
return profileWith(config);
}
private TenantProfile profileWith(DefaultTenantProfileConfiguration config) {
TenantProfile profile = new TenantProfile(new TenantProfileId(UUID.randomUUID()));
profile.setName("test-profile");
TenantProfileData profileData = new TenantProfileData();
profileData.setConfiguration(config);
profile.setProfileData(profileData);
return profile;
}
private enum EntityLevel {
DEVICE {
@Override
void setRegularMsgRateLimit(DefaultTenantProfileConfiguration config, String value) {
config.setTransportDeviceMsgRateLimit(value);
}
@Override
Object check(DefaultTransportRateLimitService service, TenantId tenantId, DeviceId entityId) {
return service.checkLimits(tenantId, null, entityId, 0, false);
}
},
GATEWAY {
@Override
void setRegularMsgRateLimit(DefaultTenantProfileConfiguration config, String value) {
config.setTransportGatewayMsgRateLimit(value);
}
@Override
Object check(DefaultTransportRateLimitService service, TenantId tenantId, DeviceId entityId) {
return service.checkLimits(tenantId, entityId, null, 0, false);
}
},
GATEWAY_DEVICE {
@Override
void setRegularMsgRateLimit(DefaultTenantProfileConfiguration config, String value) {
config.setTransportGatewayDeviceMsgRateLimit(value);
}
@Override
Object check(DefaultTransportRateLimitService service, TenantId tenantId, DeviceId entityId) {
return service.checkLimits(tenantId, null, entityId, 0, true);
}
};
abstract void setRegularMsgRateLimit(DefaultTenantProfileConfiguration config, String value);
abstract Object check(DefaultTransportRateLimitService service, TenantId tenantId, DeviceId entityId);
}
}

191
common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/service/DefaultTransportTenantProfileCacheTest.java

@ -0,0 +1,191 @@
/**
* Copyright © 2016-2026 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.service;
import com.google.common.util.concurrent.Striped;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.thingsboard.server.common.data.ApiUsageState;
import org.thingsboard.server.common.data.ApiUsageStateValue;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.id.ApiUsageStateId;
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.limits.TransportRateLimitService;
import org.thingsboard.server.common.util.ProtoUtils;
import org.thingsboard.server.gen.transport.TransportProtos.GetEntityProfileRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.GetEntityProfileResponseMsg;
import java.util.ArrayList;
import java.util.List;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.Lock;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyBoolean;
import static org.mockito.Mockito.doNothing;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
class DefaultTransportTenantProfileCacheTest {
private DefaultTransportTenantProfileCache cache;
private TransportService transportService;
private TransportRateLimitService rateLimitService;
private ExecutorService executor;
// Must match DefaultTransportTenantProfileCache.TENANT_PROFILE_FETCH_LOCK_STRIPES.
private static final int STRIPE_COUNT = 1024;
private final TenantId tenantA = TenantId.fromUUID(UUID.randomUUID());
// Deterministically pick a tenant that maps to a DIFFERENT stripe than tenantA, so the cross-tenant
// test below cannot flake on the ~1/1024 chance two random UUIDs hash to the same stripe.
private final TenantId tenantB = differentStripeFrom(tenantA);
private static TenantId differentStripeFrom(TenantId other) {
Striped<Lock> probe = Striped.lock(STRIPE_COUNT);
TenantId candidate = TenantId.fromUUID(UUID.randomUUID());
while (probe.get(candidate) == probe.get(other)) {
candidate = TenantId.fromUUID(UUID.randomUUID());
}
return candidate;
}
@BeforeEach
void setUp() {
cache = new DefaultTransportTenantProfileCache();
transportService = mock(TransportService.class);
rateLimitService = mock(TransportRateLimitService.class);
doNothing().when(rateLimitService).update(any(TenantId.class), anyBoolean());
cache.setTransportService(transportService);
cache.setRateLimitService(rateLimitService);
executor = Executors.newCachedThreadPool();
}
@AfterEach
void tearDown() {
executor.shutdownNow();
}
@Test
void fetchForOneTenantDoesNotBlockResolutionOfAnotherTenant() throws Exception {
CountDownLatch tenantAFetchStarted = new CountDownLatch(1);
CountDownLatch releaseTenantA = new CountDownLatch(1);
GetEntityProfileResponseMsg responseA = responseFor(tenantA);
GetEntityProfileResponseMsg responseB = responseFor(tenantB);
when(transportService.getEntityProfile(any())).thenAnswer(invocation -> {
GetEntityProfileRequestMsg msg = invocation.getArgument(0);
TenantId requested = TenantId.fromUUID(new UUID(msg.getEntityIdMSB(), msg.getEntityIdLSB()));
if (requested.equals(tenantA)) {
tenantAFetchStarted.countDown();
releaseTenantA.await(5, TimeUnit.SECONDS);
return responseA;
}
return responseB;
});
// T1 starts fetching tenantA's profile and blocks inside the cross-service round-trip.
Future<TenantProfile> tenantAResult = executor.submit(() -> cache.get(tenantA));
assertThat(tenantAFetchStarted.await(5, TimeUnit.SECONDS))
.as("tenantA fetch should have started").isTrue();
// T2 resolves a different tenant - it must NOT wait for tenantA's in-flight fetch.
// Fails today (single global lock); passes once locking is per-tenant.
TenantProfile tenantBProfile = CompletableFuture
.supplyAsync(() -> cache.get(tenantB), executor)
.get(2, TimeUnit.SECONDS);
assertThat(tenantBProfile).isNotNull();
releaseTenantA.countDown();
assertThat(tenantAResult.get(5, TimeUnit.SECONDS)).isNotNull();
}
@Test
void concurrentMissesForSameTenantDedupeToSingleFetch() throws Exception {
// The per-tenant lock exists precisely so that concurrent cold misses for the SAME tenant collapse
// into a single cross-service fetch (the rest are served from cache). Assert that contract directly.
int callers = 8;
CountDownLatch fetchStarted = new CountDownLatch(1);
CountDownLatch releaseFetch = new CountDownLatch(1);
when(transportService.getEntityProfile(any())).thenAnswer(invocation -> {
fetchStarted.countDown();
// Hold the (single) in-flight fetch open while the other callers pile up on the per-tenant lock.
releaseFetch.await(5, TimeUnit.SECONDS);
return responseFor(tenantA);
});
CountDownLatch allSubmitted = new CountDownLatch(callers);
List<Future<TenantProfile>> results = new ArrayList<>();
for (int i = 0; i < callers; i++) {
results.add(executor.submit(() -> {
allSubmitted.countDown();
return cache.get(tenantA);
}));
}
assertThat(allSubmitted.await(5, TimeUnit.SECONDS)).as("all callers should start").isTrue();
assertThat(fetchStarted.await(5, TimeUnit.SECONDS)).as("the first fetch should start").isTrue();
releaseFetch.countDown();
for (Future<TenantProfile> result : results) {
assertThat(result.get(5, TimeUnit.SECONDS)).isNotNull();
}
// All 8 callers resolved the same tenant, but only one of them hit the backend.
verify(transportService, times(1)).getEntityProfile(any());
}
private GetEntityProfileResponseMsg responseFor(TenantId tenantId) {
TenantProfile profile = new TenantProfile(new TenantProfileId(UUID.randomUUID()));
profile.setName("profile-" + tenantId.getId());
return GetEntityProfileResponseMsg.newBuilder()
.setEntityType(EntityType.TENANT.name())
.setTenantProfile(ProtoUtils.toProto(profile))
.setApiState(ProtoUtils.toProto(enabledApiUsageState(tenantId)))
.build();
}
private ApiUsageState enabledApiUsageState(TenantId tenantId) {
ApiUsageState state = new ApiUsageState(new ApiUsageStateId(UUID.randomUUID()));
state.setTenantId(tenantId);
state.setEntityId(tenantId);
state.setTransportState(ApiUsageStateValue.ENABLED);
state.setDbStorageState(ApiUsageStateValue.ENABLED);
state.setReExecState(ApiUsageStateValue.ENABLED);
state.setJsExecState(ApiUsageStateValue.ENABLED);
state.setTbelExecState(ApiUsageStateValue.ENABLED);
state.setEmailExecState(ApiUsageStateValue.ENABLED);
state.setSmsExecState(ApiUsageStateValue.ENABLED);
state.setAlarmExecState(ApiUsageStateValue.ENABLED);
state.setVersion(1L);
return state;
}
}

18
dao/src/main/java/org/thingsboard/server/dao/service/validator/DeviceCredentialsDataValidator.java

@ -18,10 +18,13 @@ package org.thingsboard.server.dao.service.validator;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Component;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.device.credentials.BasicMqttCredentials;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.security.DeviceCredentials;
import org.thingsboard.server.common.data.security.DeviceCredentialsType;
import org.thingsboard.server.dao.device.DeviceCredentialsDao;
import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.exception.DeviceCredentialsValidationException;
@ -69,9 +72,24 @@ public class DeviceCredentialsDataValidator extends DataValidator<DeviceCredenti
if (StringUtils.isEmpty(deviceCredentials.getCredentialsId())) {
throw new DeviceCredentialsValidationException("Device credentials id should be specified!");
}
rejectControlChars(deviceCredentials.getCredentialsId(), "credentialsId");
if (deviceCredentials.getCredentialsType() == DeviceCredentialsType.MQTT_BASIC) {
BasicMqttCredentials mqtt = JacksonUtil.fromString(deviceCredentials.getCredentialsValue(), BasicMqttCredentials.class);
if (mqtt != null) {
rejectControlChars(mqtt.getClientId(), "clientId");
rejectControlChars(mqtt.getUserName(), "userName");
rejectControlChars(mqtt.getPassword(), "password");
}
}
Device device = deviceService.findDeviceById(tenantId, deviceCredentials.getDeviceId());
if (device == null) {
throw new DeviceCredentialsValidationException("Can't assign device credentials to non-existent device!");
}
}
private static void rejectControlChars(String value, String fieldName) {
if (StringUtils.containsControlChars(value)) {
throw new DeviceCredentialsValidationException(fieldName + " must not contain control characters!");
}
}
}

40
dao/src/main/java/org/thingsboard/server/dao/util/DeviceConnectivityUtil.java

@ -51,9 +51,27 @@ public class DeviceConnectivityUtil {
public static final String COAP_IMAGE = "thingsboard/coap-clients ";
private final static Pattern VALID_URL_PATTERN = Pattern.compile("^(https?)://[-a-zA-Z0-9+&@#/%?=~_|!:,.;]*[-a-zA-Z0-9+&@#/%=~_|]");
private static String stripControlChars(String value) {
return value == null ? null : StringUtils.CONTROL_CHARS.matcher(value).replaceAll("_");
}
// Escapes a value that is interpolated inside a double-quoted shell argument (e.g. -u "...") in the
// publish commands shown to the operator, so it cannot break out of the quotes or trigger command/
// variable substitution. Control chars are stripped first to keep the command on a single line.
private static String escapeShellArg(String value) {
if (value == null) {
return null;
}
return stripControlChars(value)
.replace("\\", "\\\\")
.replace("\"", "\\\"")
.replace("$", "\\$")
.replace("`", "\\`");
}
public static String getHttpPublishCommand(String protocol, String host, String port, DeviceCredentials deviceCredentials) {
return String.format("curl -v -X POST %s://%s%s/api/v1/%s/telemetry --header Content-Type:application/json --data " + JSON_EXAMPLE_PAYLOAD,
protocol, host, port, deviceCredentials.getCredentialsId());
protocol, host, port, stripControlChars(deviceCredentials.getCredentialsId()));
}
public static String getMqttPublishCommand(String protocol, String host, String port, String deviceTelemetryTopic, DeviceCredentials deviceCredentials) {
@ -66,20 +84,20 @@ public class DeviceConnectivityUtil {
switch (deviceCredentials.getCredentialsType()) {
case ACCESS_TOKEN:
command.append(" -u \"").append(deviceCredentials.getCredentialsId()).append("\"");
command.append(" -u \"").append(escapeShellArg(deviceCredentials.getCredentialsId())).append("\"");
break;
case MQTT_BASIC:
BasicMqttCredentials credentials = JacksonUtil.fromString(deviceCredentials.getCredentialsValue(),
BasicMqttCredentials.class);
if (credentials != null) {
if (StringUtils.isNotEmpty(credentials.getClientId())) {
command.append(" -i \"").append(credentials.getClientId()).append("\"");
command.append(" -i \"").append(escapeShellArg(credentials.getClientId())).append("\"");
}
if (StringUtils.isNotEmpty(credentials.getUserName())) {
command.append(" -u \"").append(credentials.getUserName()).append("\"");
command.append(" -u \"").append(escapeShellArg(credentials.getUserName())).append("\"");
}
if (StringUtils.isNotEmpty(credentials.getPassword())) {
command.append(" -P \"").append(credentials.getPassword()).append("\"");
command.append(" -P \"").append(escapeShellArg(credentials.getPassword())).append("\"");
}
} else {
return null;
@ -117,12 +135,12 @@ public class DeviceConnectivityUtil {
dockerComposeBuilder.append("\n");
dockerComposeBuilder.append(" # Environment variables\n");
dockerComposeBuilder.append(" environment:\n");
dockerComposeBuilder.append(" - TB_GW_HOST=").append(isLocalhost(host) ? HOST_DOCKER_INTERNAL : host).append("\n");
dockerComposeBuilder.append(" - TB_GW_HOST=").append(stripControlChars(isLocalhost(host) ? HOST_DOCKER_INTERNAL : host)).append("\n");
dockerComposeBuilder.append(" - TB_GW_PORT=1883\n");
switch (deviceCredentials.getCredentialsType()) {
case ACCESS_TOKEN:
dockerComposeBuilder.append(" - TB_GW_SECURITY_TYPE=accessToken\n");
dockerComposeBuilder.append(" - TB_GW_ACCESS_TOKEN=").append(deviceCredentials.getCredentialsId()).append("\n");
dockerComposeBuilder.append(" - TB_GW_ACCESS_TOKEN=").append(stripControlChars(deviceCredentials.getCredentialsId())).append("\n");
break;
case MQTT_BASIC:
dockerComposeBuilder.append(" - TB_GW_SECURITY_TYPE=usernamePassword\n");
@ -130,13 +148,13 @@ public class DeviceConnectivityUtil {
BasicMqttCredentials.class);
if (credentials != null) {
if (StringUtils.isNotEmpty(credentials.getClientId())) {
dockerComposeBuilder.append(" - TB_GW_CLIENT_ID=").append(credentials.getClientId()).append("\n");
dockerComposeBuilder.append(" - TB_GW_CLIENT_ID=").append(stripControlChars(credentials.getClientId())).append("\n");
}
if (StringUtils.isNotEmpty(credentials.getUserName())) {
dockerComposeBuilder.append(" - TB_GW_USERNAME=").append(credentials.getUserName()).append("\n");
dockerComposeBuilder.append(" - TB_GW_USERNAME=").append(stripControlChars(credentials.getUserName())).append("\n");
}
if (StringUtils.isNotEmpty(credentials.getPassword())) {
dockerComposeBuilder.append(" - TB_GW_PASSWORD=").append(credentials.getPassword()).append("\n");
dockerComposeBuilder.append(" - TB_GW_PASSWORD=").append(stripControlChars(credentials.getPassword())).append("\n");
}
}
break;
@ -201,7 +219,7 @@ public class DeviceConnectivityUtil {
String client = COAPS.equals(protocol) ? "coap-client-openssl" : "coap-client";
String certificate = COAPS.equals(protocol) ? " -R " + CA_ROOT_CERT_PEM : "";
return String.format("%s -v 6 -m POST%s -t \"application/json\" -e %s %s://%s%s/api/v1/%s/telemetry",
client, certificate, JSON_EXAMPLE_PAYLOAD, protocol, host, port, deviceCredentials.getCredentialsId());
client, certificate, JSON_EXAMPLE_PAYLOAD, protocol, host, port, stripControlChars(deviceCredentials.getCredentialsId()));
default:
return null;
}

138
dao/src/test/java/org/thingsboard/server/dao/service/validator/DeviceCredentialsDataValidatorTest.java

@ -0,0 +1,138 @@
/**
* Copyright © 2016-2026 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.dao.service.validator;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.InjectMocks;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.device.credentials.BasicMqttCredentials;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.security.DeviceCredentials;
import org.thingsboard.server.common.data.security.DeviceCredentialsType;
import org.thingsboard.server.dao.device.DeviceCredentialsDao;
import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.exception.DeviceCredentialsValidationException;
import java.util.UUID;
import static org.assertj.core.api.Assertions.assertThatCode;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.BDDMockito.willReturn;
@ExtendWith(MockitoExtension.class)
class DeviceCredentialsDataValidatorTest {
@Mock
DeviceCredentialsDao deviceCredentialsDao;
@Mock
DeviceService deviceService;
@InjectMocks
DeviceCredentialsDataValidator validator;
final TenantId tenantId = TenantId.fromUUID(UUID.fromString("9ef79cdf-37a8-4119-b682-2e7ed4e018da"));
final DeviceId deviceId = new DeviceId(UUID.fromString("11111111-1111-1111-1111-111111111111"));
@Test
void rejectsNewlineInAccessToken() {
DeviceCredentials creds = accessToken("safe_token\nentrypoint: [\"/bin/sh\"]");
assertThatThrownBy(() -> validator.validateDataImpl(tenantId, creds))
.isInstanceOf(DeviceCredentialsValidationException.class)
.hasMessageContaining("credentialsId")
.hasMessageContaining("control characters");
}
@Test
void rejectsCarriageReturnInAccessToken() {
DeviceCredentials creds = accessToken("token\rprivileged: true");
assertThatThrownBy(() -> validator.validateDataImpl(tenantId, creds))
.isInstanceOf(DeviceCredentialsValidationException.class)
.hasMessageContaining("control characters");
}
@Test
void rejectsNewlineInMqttClientId() {
DeviceCredentials creds = mqttBasic("cid\nentrypoint: x", "user", "pwd");
assertThatThrownBy(() -> validator.validateDataImpl(tenantId, creds))
.isInstanceOf(DeviceCredentialsValidationException.class)
.hasMessageContaining("clientId");
}
@Test
void rejectsNewlineInMqttUserName() {
DeviceCredentials creds = mqttBasic("cid", "user\nprivileged: true", "pwd");
assertThatThrownBy(() -> validator.validateDataImpl(tenantId, creds))
.isInstanceOf(DeviceCredentialsValidationException.class)
.hasMessageContaining("userName");
}
@Test
void rejectsNewlineInMqttPassword() {
DeviceCredentials creds = mqttBasic("cid", "user", "pwd\nentrypoint: x");
assertThatThrownBy(() -> validator.validateDataImpl(tenantId, creds))
.isInstanceOf(DeviceCredentialsValidationException.class)
.hasMessageContaining("password");
}
@Test
void acceptsValidCredentials() {
willReturn(new Device()).given(deviceService).findDeviceById(tenantId, deviceId);
DeviceCredentials creds = accessToken("safe_token_123");
assertThatCode(() -> validator.validateDataImpl(tenantId, creds))
.doesNotThrowAnyException();
}
@Test
void acceptsValidMqttBasicCredentials() {
willReturn(new Device()).given(deviceService).findDeviceById(tenantId, deviceId);
DeviceCredentials creds = mqttBasic("client-1", "user-1", "pwd-1");
assertThatCode(() -> validator.validateDataImpl(tenantId, creds))
.doesNotThrowAnyException();
}
private DeviceCredentials accessToken(String token) {
DeviceCredentials c = new DeviceCredentials();
c.setDeviceId(deviceId);
c.setCredentialsType(DeviceCredentialsType.ACCESS_TOKEN);
c.setCredentialsId(token);
return c;
}
private DeviceCredentials mqttBasic(String clientId, String userName, String password) {
BasicMqttCredentials inner = new BasicMqttCredentials();
inner.setClientId(clientId);
inner.setUserName(userName);
inner.setPassword(password);
DeviceCredentials c = new DeviceCredentials();
c.setDeviceId(deviceId);
c.setCredentialsType(DeviceCredentialsType.MQTT_BASIC);
c.setCredentialsId("mqtt-credentials-id");
c.setCredentialsValue(JacksonUtil.toString(inner));
return c;
}
}

123
dao/src/test/java/org/thingsboard/server/dao/util/DeviceConnectivityUtilTest.java

@ -16,6 +16,13 @@
package org.thingsboard.server.dao.util;
import org.junit.jupiter.api.Test;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.device.credentials.BasicMqttCredentials;
import org.thingsboard.server.common.data.security.DeviceCredentials;
import org.thingsboard.server.common.data.security.DeviceCredentialsType;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import static org.assertj.core.api.Assertions.assertThat;
@ -29,4 +36,120 @@ class DeviceConnectivityUtilTest {
assertThat(DeviceConnectivityUtil.CA_ROOT_CERT_PEM).doesNotContainAnyWhitespaces();
}
@Test
void validAccessTokenIsRenderedAsIs() throws Exception {
String yaml = renderCompose(accessToken("safe_token_123"));
assertThat(yaml).contains("- TB_GW_ACCESS_TOKEN=safe_token_123\n");
assertNoInjectedSiblingKeys(yaml);
}
@Test
void newlineInAccessTokenIsSanitized() throws Exception {
String malicious = "safe_token\n entrypoint: [\"/bin/bash\",\"-c\",\"id\"]";
String yaml = renderCompose(accessToken(malicious));
assertNoInjectedSiblingKeys(yaml);
}
@Test
void carriageReturnInAccessTokenIsSanitized() throws Exception {
String yaml = renderCompose(accessToken("token\rprivileged: true"));
assertNoInjectedSiblingKeys(yaml);
}
@Test
void newlineInMqttClientIdIsSanitized() throws Exception {
String yaml = renderCompose(mqttBasic("cid\n entrypoint: [\"/bin/sh\"]", "user", "pwd"));
assertNoInjectedSiblingKeys(yaml);
}
@Test
void newlineInMqttUserNameIsSanitized() throws Exception {
String yaml = renderCompose(mqttBasic("cid", "user\n privileged: true", "pwd"));
assertNoInjectedSiblingKeys(yaml);
}
@Test
void newlineInMqttPasswordIsSanitized() throws Exception {
String yaml = renderCompose(mqttBasic("cid", "user", "pwd\n entrypoint: [\"/bin/sh\"]"));
assertNoInjectedSiblingKeys(yaml);
}
@Test
void mqttBasicQuoteInUserNameIsEscapedInPublishCommand() {
String command = DeviceConnectivityUtil.getMqttPublishCommand(
"mqtt", "localhost", "1883", "v1/devices/me/telemetry",
mqttBasic("cid", "u\";touch pwned;echo \"", "pwd"));
// the double quote must be backslash-escaped so it cannot terminate the -u "..." argument
assertThat(command).contains("-u \"u\\\";touch pwned;echo \\\"\"");
assertThat(command).doesNotContain("-u \"u\";");
}
@Test
void controlCharsInMqttClientIdAreStrippedInPublishCommand() {
String command = DeviceConnectivityUtil.getMqttPublishCommand(
"mqtt", "localhost", "1883", "v1/devices/me/telemetry",
mqttBasic("c\nid", "user", "pwd"));
assertThat(command).doesNotContain("\n");
assertThat(command).contains("-i \"c_id\"");
}
@Test
void controlCharsInAccessTokenAreStrippedInHttpAndCoapCommands() {
DeviceCredentials creds = accessToken("tok\nen");
assertThat(DeviceConnectivityUtil.getHttpPublishCommand("http", "localhost", ":8080", creds))
.doesNotContain("\n")
.contains("/api/v1/tok_en/telemetry");
assertThat(DeviceConnectivityUtil.getCoapPublishCommand("coap", "localhost", ":5683", creds))
.doesNotContain("\n")
.contains("/api/v1/tok_en/telemetry");
}
private static String renderCompose(DeviceCredentials credentials) throws Exception {
var resource = DeviceConnectivityUtil.getGatewayDockerComposeFile(
"host.docker.internal", "3.8-stable", credentials);
try (var in = resource.getInputStream()) {
return new String(in.readAllBytes(), StandardCharsets.UTF_8);
}
}
private static DeviceCredentials accessToken(String token) {
DeviceCredentials c = new DeviceCredentials();
c.setCredentialsType(DeviceCredentialsType.ACCESS_TOKEN);
c.setCredentialsId(token);
return c;
}
private static DeviceCredentials mqttBasic(String clientId, String userName, String password) {
BasicMqttCredentials inner = new BasicMqttCredentials();
inner.setClientId(clientId);
inner.setUserName(userName);
inner.setPassword(password);
DeviceCredentials c = new DeviceCredentials();
c.setCredentialsType(DeviceCredentialsType.MQTT_BASIC);
c.setCredentialsId("mqtt-credentials-id");
c.setCredentialsValue(JacksonUtil.toString(inner));
return c;
}
private static void assertNoInjectedSiblingKeys(String yaml) throws IOException {
for (String line : yaml.split("\n")) {
String trimmed = line.replaceFirst("^\\s+", "");
assertThat(trimmed)
.as("unexpected sibling key — possible YAML injection: %s", line)
.doesNotStartWith("entrypoint:")
.doesNotStartWith("privileged:")
.doesNotStartWith("command:");
}
}
}

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

@ -133,6 +133,8 @@ redis:
blockWhenExhausted: "${REDIS_POOL_CONFIG_BLOCK_WHEN_EXHAUSTED:true}"
transport:
# Size of the thread pool that executes transport API callbacks (session registration, telemetry/attribute and RPC responses, entity update notifications, and the tenant profile fetch on a cache miss). Bounds how many such callbacks - including those that block on a backend round-trip - can run concurrently.
callback_thread_pool_size: "${TB_TRANSPORT_CALLBACK_THREAD_POOL_SIZE:20}"
# Local CoAP transport parameters
coap:
# CoaP processing timeout in milliseconds

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

@ -167,6 +167,8 @@ redis:
# HTTP server parameters
transport:
# Size of the thread pool that executes transport API callbacks (session registration, telemetry/attribute and RPC responses, entity update notifications, and the tenant profile fetch on a cache miss). Bounds how many such callbacks - including those that block on a backend round-trip - can run concurrently.
callback_thread_pool_size: "${TB_TRANSPORT_CALLBACK_THREAD_POOL_SIZE:20}"
http:
# HTTP request processing timeout in milliseconds
request_timeout: "${HTTP_REQUEST_TIMEOUT:60000}"

2
transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml

@ -134,6 +134,8 @@ redis:
# LWM2M server parameters
transport:
# Size of the thread pool that executes transport API callbacks (session registration, telemetry/attribute and RPC responses, entity update notifications, and the tenant profile fetch on a cache miss). Bounds how many such callbacks - including those that block on a backend round-trip - can run concurrently.
callback_thread_pool_size: "${TB_TRANSPORT_CALLBACK_THREAD_POOL_SIZE:20}"
sessions:
# Session inactivity timeout is a global configuration parameter that defines how long the device transport session will be opened after the last message arrives from the device.
# The parameter value is in milliseconds.

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

@ -135,6 +135,8 @@ redis:
# MQTT server parameters
transport:
# Size of the thread pool that executes transport API callbacks (session registration, telemetry/attribute and RPC responses, entity update notifications, and the tenant profile fetch on a cache miss). Bounds how many such callbacks - including those that block on a backend round-trip - can run concurrently.
callback_thread_pool_size: "${TB_TRANSPORT_CALLBACK_THREAD_POOL_SIZE:20}"
mqtt:
# MQTT bind-address
bind_address: "${MQTT_BIND_ADDRESS:0.0.0.0}"

2
transport/snmp/src/main/resources/tb-snmp-transport.yml

@ -134,6 +134,8 @@ redis:
# Snmp server parameters
transport:
# Size of the thread pool that executes transport API callbacks (session registration, telemetry/attribute and RPC responses, entity update notifications, and the tenant profile fetch on a cache miss). Bounds how many such callbacks - including those that block on a backend round-trip - can run concurrently.
callback_thread_pool_size: "${TB_TRANSPORT_CALLBACK_THREAD_POOL_SIZE:20}"
snmp:
# Enable/disable SNMP transport protocol
enabled: "${SNMP_ENABLED:true}"

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

@ -667,7 +667,7 @@
"assigned-to-user": "Alarm was assigned by user {{userName}} to user {{assigneeName}}",
"unassigned-to-user": "Alarm was unassigned by user {{userName}}",
"unassigned-from-deleted-user": "Alarm was unassigned because user {{userName}} - was deleted",
"comment-deleted": "User {{userName}} deleted his comment",
"comment-deleted": "Comment was deleted by user {{userName}}",
"severity-changed": "Alarm severity was updated from {{oldSeverity}} to {{newSeverity}}"
}
},

Loading…
Cancel
Save