Browse Source

Merge branch 'lts-4.2' into release-4.2

# Conflicts:
#	application/pom.xml
#	common/actor/pom.xml
#	common/cache/pom.xml
#	common/cluster-api/pom.xml
#	common/coap-server/pom.xml
#	common/dao-api/pom.xml
#	common/data/pom.xml
#	common/discovery-api/pom.xml
#	common/edge-api/pom.xml
#	common/edqs/pom.xml
#	common/message/pom.xml
#	common/pom.xml
#	common/proto/pom.xml
#	common/queue/pom.xml
#	common/script/pom.xml
#	common/script/remote-js-client/pom.xml
#	common/script/script-api/pom.xml
#	common/stats/pom.xml
#	common/transport/coap/pom.xml
#	common/transport/http/pom.xml
#	common/transport/lwm2m/pom.xml
#	common/transport/mqtt/pom.xml
#	common/transport/pom.xml
#	common/transport/snmp/pom.xml
#	common/transport/transport-api/pom.xml
#	common/util/pom.xml
#	common/version-control/pom.xml
#	dao/pom.xml
#	edqs/pom.xml
#	monitoring/pom.xml
#	msa/black-box-tests/pom.xml
#	msa/edqs/pom.xml
#	msa/js-executor/pom.xml
#	msa/monitoring/pom.xml
#	msa/pom.xml
#	msa/tb-node/pom.xml
#	msa/tb/pom.xml
#	msa/transport/coap/pom.xml
#	msa/transport/http/pom.xml
#	msa/transport/lwm2m/pom.xml
#	msa/transport/mqtt/pom.xml
#	msa/transport/pom.xml
#	msa/transport/snmp/pom.xml
#	msa/vc-executor-docker/pom.xml
#	msa/vc-executor/pom.xml
#	msa/web-ui/pom.xml
#	netty-mqtt/pom.xml
#	pom.xml
#	rest-client/pom.xml
#	rule-engine/pom.xml
#	rule-engine/rule-engine-api/pom.xml
#	rule-engine/rule-engine-components/pom.xml
#	tools/pom.xml
#	transport/coap/pom.xml
#	transport/http/pom.xml
#	transport/lwm2m/pom.xml
#	transport/mqtt/pom.xml
#	transport/pom.xml
#	transport/snmp/pom.xml
#	ui-ngx/pom.xml
release-4.2
Viacheslav Klimov 3 weeks ago
parent
commit
13bc560104
Failed to extract signature
  1. 2
      application/pom.xml
  2. 27
      application/src/main/java/org/thingsboard/server/controller/AdminController.java
  3. 7
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  4. 28
      application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/AbstractOAuth2ClientMapper.java
  5. 66
      application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/CallbackUrlSchemeValidator.java
  6. 5
      application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/HttpCookieOAuth2AuthorizationRequestRepository.java
  7. 5
      application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationFailureHandler.java
  8. 44
      application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationSuccessHandler.java
  9. 66
      application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/PrevUriValidator.java
  10. 6
      application/src/main/java/org/thingsboard/server/service/security/model/token/OAuth2AppTokenFactory.java
  11. 11
      application/src/main/java/org/thingsboard/server/service/telemetry/AbstractSubscriptionService.java
  12. 81
      application/src/test/java/org/thingsboard/server/controller/AdminControllerMailOAuth2Test.java
  13. 33
      application/src/test/java/org/thingsboard/server/controller/AdminControllerTest.java
  14. 81
      application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/CallbackUrlSchemeValidatorTest.java
  15. 69
      application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/HttpCookieOAuth2AuthorizationRequestRepositoryTest.java
  16. 115
      application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/OAuth2ClientMapperTest.java
  17. 103
      application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationFailureHandlerTest.java
  18. 186
      application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationSuccessHandlerTest.java
  19. 74
      application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/PrevUriValidatorTest.java
  20. 72
      application/src/test/java/org/thingsboard/server/service/security/model/token/OAuth2AppTokenFactoryTest.java
  21. 22
      application/src/test/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionServiceTest.java
  22. 2
      common/actor/pom.xml
  23. 2
      common/cache/pom.xml
  24. 2
      common/cluster-api/pom.xml
  25. 2
      common/coap-server/pom.xml
  26. 2
      common/dao-api/pom.xml
  27. 2
      common/data/pom.xml
  28. 2
      common/discovery-api/pom.xml
  29. 2
      common/edge-api/pom.xml
  30. 59
      common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java
  31. 2
      common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java
  32. 1
      common/edge-api/src/main/proto/edge.proto
  33. 168
      common/edge-api/src/test/java/org/thingsboard/edge/rpc/EdgeGrpcClientLeakTest.java
  34. 2
      common/edqs/pom.xml
  35. 2
      common/message/pom.xml
  36. 2
      common/pom.xml
  37. 2
      common/proto/pom.xml
  38. 2
      common/queue/pom.xml
  39. 2
      common/script/pom.xml
  40. 2
      common/script/remote-js-client/pom.xml
  41. 2
      common/script/script-api/pom.xml
  42. 2
      common/stats/pom.xml
  43. 2
      common/transport/coap/pom.xml
  44. 2
      common/transport/http/pom.xml
  45. 2
      common/transport/lwm2m/pom.xml
  46. 2
      common/transport/mqtt/pom.xml
  47. 2
      common/transport/pom.xml
  48. 2
      common/transport/snmp/pom.xml
  49. 2
      common/transport/transport-api/pom.xml
  50. 2
      common/util/pom.xml
  51. 2
      common/version-control/pom.xml
  52. 2
      dao/pom.xml
  53. 2
      edqs/pom.xml
  54. 2
      monitoring/pom.xml
  55. 2
      msa/black-box-tests/pom.xml
  56. 2
      msa/edqs/pom.xml
  57. 2
      msa/js-executor/package.json
  58. 2
      msa/js-executor/pom.xml
  59. 2
      msa/monitoring/pom.xml
  60. 2
      msa/pom.xml
  61. 2
      msa/tb-node/pom.xml
  62. 2
      msa/tb/pom.xml
  63. 2
      msa/transport/coap/pom.xml
  64. 2
      msa/transport/http/pom.xml
  65. 2
      msa/transport/lwm2m/pom.xml
  66. 2
      msa/transport/mqtt/pom.xml
  67. 2
      msa/transport/pom.xml
  68. 2
      msa/transport/snmp/pom.xml
  69. 2
      msa/vc-executor-docker/pom.xml
  70. 2
      msa/vc-executor/pom.xml
  71. 2
      msa/web-ui/package.json
  72. 2
      msa/web-ui/pom.xml
  73. 4
      netty-mqtt/pom.xml
  74. 2
      pom.xml
  75. 2
      rest-client/pom.xml
  76. 2
      rule-engine/pom.xml
  77. 2
      rule-engine/rule-engine-api/pom.xml
  78. 2
      rule-engine/rule-engine-components/pom.xml
  79. 2
      tools/pom.xml
  80. 74
      tools/src/main/java/org/thingsboard/client/tools/migrator/DictionaryParser.java
  81. 102
      tools/src/main/java/org/thingsboard/client/tools/migrator/MigratorTool.java
  82. 282
      tools/src/main/java/org/thingsboard/client/tools/migrator/PgCaMigrator.java
  83. 92
      tools/src/main/java/org/thingsboard/client/tools/migrator/README.md
  84. 88
      tools/src/main/java/org/thingsboard/client/tools/migrator/RelatedEntitiesParser.java
  85. 86
      tools/src/main/java/org/thingsboard/client/tools/migrator/WriterBuilder.java
  86. 2
      transport/coap/pom.xml
  87. 2
      transport/http/pom.xml
  88. 2
      transport/lwm2m/pom.xml
  89. 2
      transport/mqtt/pom.xml
  90. 2
      transport/pom.xml
  91. 2
      transport/snmp/pom.xml
  92. 2
      ui-ngx/package.json
  93. 2
      ui-ngx/pom.xml

2
application/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>thingsboard</artifactId>
</parent>
<artifactId>application</artifactId>

27
application/src/main/java/org/thingsboard/server/controller/AdminController.java

@ -76,6 +76,7 @@ import org.thingsboard.server.dao.settings.SecuritySettingsService;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.security.auth.jwt.settings.JwtSettingsService;
import org.thingsboard.server.service.security.auth.oauth2.CookieUtils;
import org.thingsboard.server.service.security.auth.oauth2.PrevUriValidator;
import org.thingsboard.server.service.security.model.SecurityUser;
import org.thingsboard.server.service.security.model.token.JwtTokenFactory;
import org.thingsboard.server.service.security.permission.Operation;
@ -92,6 +93,8 @@ import java.util.Optional;
import static org.thingsboard.server.controller.ControllerConstants.SYSTEM_AUTHORITY_PARAGRAPH;
import static org.thingsboard.server.controller.ControllerConstants.TENANT_AUTHORITY_PARAGRAPH;
import static org.thingsboard.server.service.security.auth.oauth2.HttpCookieOAuth2AuthorizationRequestRepository.PREV_URI_COOKIE_NAME;
import static org.thingsboard.server.service.security.auth.oauth2.HttpCookieOAuth2AuthorizationRequestRepository.PREV_URI_PARAMETER;
@RestController
@TbCoreComponent
@ -100,8 +103,7 @@ import static org.thingsboard.server.controller.ControllerConstants.TENANT_AUTHO
@RequiredArgsConstructor
public class AdminController extends BaseController {
private static final String PREV_URI_PATH_PARAMETER = "prevUri";
private static final String PREV_URI_COOKIE_NAME = "prev_uri";
private static final String DEFAULT_PREV_URI = "/settings/outgoing-mail";
private static final String STATE_COOKIE_NAME = "state";
private static final String MAIL_SETTINGS_KEY = "mail";
@ -419,8 +421,9 @@ public class AdminController extends BaseController {
@GetMapping(value = "/mail/oauth2/authorize", produces = "application/text")
public String getAuthorizationUrl(HttpServletRequest request, HttpServletResponse response) throws ThingsboardException {
String state = StringUtils.generateSafeToken();
if (request.getParameter(PREV_URI_PATH_PARAMETER) != null) {
CookieUtils.addCookie(response, PREV_URI_COOKIE_NAME, request.getParameter(PREV_URI_PATH_PARAMETER), 180);
String prevUriParam = request.getParameter(PREV_URI_PARAMETER);
if (PrevUriValidator.isValid(prevUriParam)) {
CookieUtils.addCookie(response, PREV_URI_COOKIE_NAME, prevUriParam, 180);
}
CookieUtils.addCookie(response, STATE_COOKIE_NAME, state, 180);
@ -445,12 +448,9 @@ public class AdminController extends BaseController {
public void codeProcessingUrl(
@RequestParam(value = "code") String code, @RequestParam(value = "state") String state,
HttpServletRequest request, HttpServletResponse response) throws ThingsboardException, IOException {
Optional<Cookie> prevUrlOpt = CookieUtils.getCookie(request, PREV_URI_COOKIE_NAME);
String redirectUrl = getMailOAuth2RedirectUrl(request);
Optional<Cookie> cookieState = CookieUtils.getCookie(request, STATE_COOKIE_NAME);
String baseUrl = this.systemSecurityService.getBaseUrl(TenantId.SYS_TENANT_ID, new CustomerId(EntityId.NULL_UUID), request);
String prevUri = baseUrl + (prevUrlOpt.isPresent() ? prevUrlOpt.get().getValue() : "/settings/outgoing-mail");
if (cookieState.isEmpty() || !cookieState.get().getValue().equals(state)) {
CookieUtils.deleteCookie(request, response, STATE_COOKIE_NAME);
throw new ThingsboardException("Refresh token was not generated, invalid state param", ThingsboardErrorCode.BAD_REQUEST_PARAMS);
@ -480,7 +480,16 @@ public class AdminController extends BaseController {
((ObjectNode) jsonValue).put("tokenGenerated", true);
adminSettingsService.saveAdminSettings(TenantId.SYS_TENANT_ID, adminSettings);
response.sendRedirect(prevUri);
response.sendRedirect(redirectUrl);
}
String getMailOAuth2RedirectUrl(HttpServletRequest request) {
String baseUrl = this.systemSecurityService.getBaseUrl(TenantId.SYS_TENANT_ID, new CustomerId(EntityId.NULL_UUID), request);
String prevUri = CookieUtils.getCookie(request, PREV_URI_COOKIE_NAME)
.map(Cookie::getValue)
.filter(PrevUriValidator::isValid)
.orElse(DEFAULT_PREV_URI);
return baseUrl + prevUri;
}
}

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

@ -269,7 +269,12 @@ public abstract class EdgeGrpcSession implements Closeable {
@Override
public void onFailure(Throwable t) {
log.error("[{}][{}] Exception during sync process", tenantId, edge.getId(), t);
log.error("[{}][{}] Exception during sync process, skipping fetcher {} and continuing",
tenantId, edge.getId(), next.getClass().getSimpleName(), t);
// Keep walking the cursor: returning here leaves syncInProgress set for the life of the
// session, so the edge never receives SyncCompletedMsg and both general downlink delivery
// and uplink processing stay gated until the session is re-established.
doSync(cursor);
}
}, ctx.getGrpcCallbackExecutorService());
} else {

28
application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/AbstractOAuth2ClientMapper.java

@ -95,7 +95,7 @@ public abstract class AbstractOAuth2ClientMapper {
UserPrincipal principal = new UserPrincipal(UserPrincipal.Type.USER_NAME, oauth2User.getEmail());
User user = userService.findUserByEmail(TenantId.SYS_TENANT_ID, oauth2User.getEmail());
User user = findUserForClient(oauth2User.getEmail(), oAuth2Client);
if (user == null && !config.isAllowUserCreation()) {
throw new UsernameNotFoundException("User not found: " + oauth2User.getEmail());
@ -104,7 +104,7 @@ public abstract class AbstractOAuth2ClientMapper {
if (user == null) {
userCreationLock.lock();
try {
user = userService.findUserByEmail(TenantId.SYS_TENANT_ID, oauth2User.getEmail());
user = findUserForClient(oauth2User.getEmail(), oAuth2Client);
if (user == null) {
user = new User();
if (oauth2User.getCustomerId() == null && StringUtils.isEmpty(oauth2User.getCustomerName())) {
@ -112,7 +112,12 @@ public abstract class AbstractOAuth2ClientMapper {
} else {
user.setAuthority(Authority.CUSTOMER_USER);
}
TenantId tenantId = oauth2User.getTenantId() != null ? oauth2User.getTenantId() : getTenantId(oauth2User.getTenantName());
TenantId tenantId;
if (oAuth2Client.getTenantId().isSysTenantId()) {
tenantId = oauth2User.getTenantId() != null ? oauth2User.getTenantId() : getTenantId(oauth2User.getTenantName());
} else {
tenantId = oAuth2Client.getTenantId();
}
user.setTenantId(tenantId);
CustomerId customerId = oauth2User.getCustomerId() != null ?
oauth2User.getCustomerId() : getCustomerId(user.getTenantId(), oauth2User.getCustomerName());
@ -164,6 +169,23 @@ public abstract class AbstractOAuth2ClientMapper {
}
}
/**
* The user is matched by email alone, and the email is whatever the provider chose to send. A client registered by a
* tenant must therefore only ever resolve users of that same tenant; a system client is platform-wide by design.
*/
private User findUserForClient(String email, OAuth2Client oAuth2Client) {
User user = userService.findUserByEmail(TenantId.SYS_TENANT_ID, email);
if (user == null || oAuth2Client.getTenantId().isSysTenantId()) {
return user;
}
if (user.getTenantId().equals(oAuth2Client.getTenantId())) {
return user;
}
log.warn("OAuth2 client [{}] of tenant [{}] cannot resolve user [{}] [{}] of tenant [{}]: outside of the client tenant",
oAuth2Client.getId(), oAuth2Client.getTenantId(), user.getId(), email, user.getTenantId());
throw new UsernameNotFoundException("User not found: " + email);
}
private TenantId getTenantId(String name) throws Exception {
Tenant tenant = tenantService.findTenantByName(name);
if (tenant != null) {

66
application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/CallbackUrlSchemeValidator.java

@ -0,0 +1,66 @@
/**
* 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.security.auth.oauth2;
import lombok.extern.slf4j.Slf4j;
import org.springframework.security.oauth2.core.endpoint.OAuth2AuthorizationRequest;
import org.thingsboard.server.common.data.StringUtils;
import java.util.Locale;
import java.util.Set;
import java.util.regex.Pattern;
@Slf4j
public class CallbackUrlSchemeValidator {
// RFC 3986 scheme grammar, plus '_': mobile apps derive the scheme from their package name, which may contain one
private static final Pattern SCHEME_PATTERN = Pattern.compile("[a-zA-Z][a-zA-Z0-9+.\\-_]*");
private static final Set<String> FORBIDDEN_SCHEMES = Set.of("http", "https", "javascript", "data", "file", "vbscript");
private static final int MAX_LOGGED_LENGTH = 128;
/**
* The redirect carrying the access token is built as callbackUrlScheme + ':', so only a mobile app scheme may
* pass: a web scheme would send the token to whatever host follows it.
*/
public static boolean isValid(String callbackUrlScheme) {
return !StringUtils.isEmpty(callbackUrlScheme)
&& SCHEME_PATTERN.matcher(callbackUrlScheme).matches()
&& !FORBIDDEN_SCHEMES.contains(callbackUrlScheme.toLowerCase(Locale.ROOT));
}
/**
* The attribute is restored from the oauth2_auth_request cookie, which the client can replace, so the scheme is
* checked again on read and not only when the authorization request is built.
*/
public static String getCallbackUrlScheme(OAuth2AuthorizationRequest authorizationRequest) {
String callbackUrlScheme = authorizationRequest != null ?
authorizationRequest.getAttribute(TbOAuth2ParameterNames.CALLBACK_URL_SCHEME) : null;
if (StringUtils.isEmpty(callbackUrlScheme)) {
return null;
}
if (!isValid(callbackUrlScheme)) {
log.warn("Ignoring invalid callback url scheme: [{}]", forLog(callbackUrlScheme));
return null;
}
return callbackUrlScheme;
}
// a rejected value is attacker-controlled: it must not be able to forge log lines
private static String forLog(String value) {
return value.substring(0, Math.min(value.length(), MAX_LOGGED_LENGTH)).replaceAll("[^\\x20-\\x7E]", "?");
}
}

5
application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/HttpCookieOAuth2AuthorizationRequestRepository.java

@ -43,8 +43,9 @@ public class HttpCookieOAuth2AuthorizationRequestRepository implements Authoriza
CookieUtils.deleteCookie(request, response, OAUTH2_AUTHORIZATION_REQUEST_COOKIE_NAME);
return;
}
if (request.getParameter(PREV_URI_PARAMETER) != null) {
CookieUtils.addCookie(response, PREV_URI_COOKIE_NAME, request.getParameter(PREV_URI_PARAMETER), cookieExpireSeconds);
String prevUri = request.getParameter(PREV_URI_PARAMETER);
if (PrevUriValidator.isValid(prevUri)) {
CookieUtils.addCookie(response, PREV_URI_COOKIE_NAME, prevUri, cookieExpireSeconds);
}
CookieUtils.addCookie(response, OAUTH2_AUTHORIZATION_REQUEST_COOKIE_NAME, CookieUtils.serialize(authorizationRequest), cookieExpireSeconds);
}

5
application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationFailureHandler.java

@ -54,11 +54,8 @@ public class Oauth2AuthenticationFailureHandler extends SimpleUrlAuthenticationF
throws IOException, ServletException {
String baseUrl;
String errorPrefix;
String callbackUrlScheme = null;
OAuth2AuthorizationRequest authorizationRequest = httpCookieOAuth2AuthorizationRequestRepository.loadAuthorizationRequest(request);
if (authorizationRequest != null) {
callbackUrlScheme = authorizationRequest.getAttribute(TbOAuth2ParameterNames.CALLBACK_URL_SCHEME);
}
String callbackUrlScheme = CallbackUrlSchemeValidator.getCallbackUrlScheme(authorizationRequest);
if (!StringUtils.isEmpty(callbackUrlScheme)) {
baseUrl = callbackUrlScheme + ":";
errorPrefix = "/?error=";

44
application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationSuccessHandler.java

@ -82,18 +82,9 @@ public class Oauth2AuthenticationSuccessHandler extends SimpleUrlAuthenticationS
HttpServletResponse response,
Authentication authentication) throws IOException {
OAuth2AuthorizationRequest authorizationRequest = httpCookieOAuth2AuthorizationRequestRepository.loadAuthorizationRequest(request);
String callbackUrlScheme = authorizationRequest.getAttribute(TbOAuth2ParameterNames.CALLBACK_URL_SCHEME);
String baseUrl;
if (!StringUtils.isEmpty(callbackUrlScheme)) {
baseUrl = callbackUrlScheme + ":";
} else {
baseUrl = this.systemSecurityService.getBaseUrl(TenantId.SYS_TENANT_ID, new CustomerId(EntityId.NULL_UUID), request);
Optional<Cookie> prevUrlOpt = CookieUtils.getCookie(request, PREV_URI_COOKIE_NAME);
if (prevUrlOpt.isPresent()) {
baseUrl += prevUrlOpt.get().getValue();
CookieUtils.deleteCookie(request, response, PREV_URI_COOKIE_NAME);
}
}
String callbackUrlScheme = CallbackUrlSchemeValidator.getCallbackUrlScheme(authorizationRequest);
String baseUrl = getBaseUrl(request, callbackUrlScheme);
String prevUri = getPrevUri(request, response, callbackUrlScheme);
try {
OAuth2AuthenticationToken token = (OAuth2AuthenticationToken) authentication;
@ -108,7 +99,7 @@ public class Oauth2AuthenticationSuccessHandler extends SimpleUrlAuthenticationS
clearAuthenticationAttributes(request, response);
JwtPair tokenPair = tokenFactory.createTokenPair(securityUser);
getRedirectStrategy().sendRedirect(request, response, getRedirectUrl(baseUrl, tokenPair));
getRedirectStrategy().sendRedirect(request, response, getRedirectUrl(baseUrl + prevUri, tokenPair));
systemSecurityService.logLoginAction(securityUser, new RestAuthenticationDetails(request), ActionType.LOGIN, oauth2Client.getName(), null);
} catch (Exception e) {
log.debug("Error occurred during processing authentication success result. " +
@ -125,6 +116,31 @@ public class Oauth2AuthenticationSuccessHandler extends SimpleUrlAuthenticationS
}
}
String getBaseUrl(HttpServletRequest request, String callbackUrlScheme) {
if (!StringUtils.isEmpty(callbackUrlScheme)) {
return callbackUrlScheme + ":";
}
return this.systemSecurityService.getBaseUrl(TenantId.SYS_TENANT_ID, new CustomerId(EntityId.NULL_UUID), request);
}
/**
* The in-app path the user was on before the login, or an empty string. A present cookie is dropped whether or
* not its value passes validation - it is only meant to survive a single login round trip. The path is kept out
* of the base URL so that the error redirect, which appends its own path, stays routable.
*/
String getPrevUri(HttpServletRequest request, HttpServletResponse response, String callbackUrlScheme) {
if (!StringUtils.isEmpty(callbackUrlScheme)) {
return "";
}
Optional<Cookie> prevUriOpt = CookieUtils.getCookie(request, PREV_URI_COOKIE_NAME);
if (prevUriOpt.isEmpty()) {
return "";
}
String prevUri = prevUriOpt.get().getValue();
CookieUtils.deleteCookie(request, response, PREV_URI_COOKIE_NAME);
return PrevUriValidator.isValid(prevUri) ? prevUri : "";
}
protected void clearAuthenticationAttributes(HttpServletRequest request, HttpServletResponse response) {
super.clearAuthenticationAttributes(request);
httpCookieOAuth2AuthorizationRequestRepository.removeAuthorizationRequestCookies(request, response);
@ -133,6 +149,8 @@ public class Oauth2AuthenticationSuccessHandler extends SimpleUrlAuthenticationS
String getRedirectUrl(String baseUrl, JwtPair tokenPair) {
if (baseUrl.indexOf("?") > 0) {
baseUrl += "&";
} else if (baseUrl.endsWith("/")) {
baseUrl += "?";
} else {
baseUrl += "/?";
}

66
application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/PrevUriValidator.java

@ -0,0 +1,66 @@
/**
* 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.security.auth.oauth2;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.server.common.data.StringUtils;
import java.util.Locale;
@Slf4j
public class PrevUriValidator {
private static final int MAX_LENGTH = 2048;
private static final int MAX_LOGGED_LENGTH = 128;
public static boolean isValid(String prevUri) {
if (StringUtils.isEmpty(prevUri)) {
return false;
}
if (!isInAppPath(prevUri)) {
log.debug("Ignoring prevUri that is not an in-app path: [{}]", forLog(prevUri));
return false;
}
return true;
}
/**
* prevUri is appended to the platform base URL, which ends right after the authority, so the single leading '/'
* is what keeps the redirect on this host - it closes the authority before any of the value is read. The rest
* keeps an accepted value usable: it has to survive the cookie round trip (RFC 6265 allows neither control
* characters nor '"', ',', ';', '\' or non-ASCII) and to pass StrictHttpFirewall, which rejects '//', '%2f'
* and '%5c' in the path; a fragment would swallow the access token.
*/
private static boolean isInAppPath(String prevUri) {
if (prevUri.length() > MAX_LENGTH || prevUri.charAt(0) != '/') {
return false;
}
for (int i = 0; i < prevUri.length(); i++) {
char c = prevUri.charAt(i);
if (c <= ' ' || c >= 127 || c == '"' || c == ',' || c == ';' || c == '\\' || c == '#') {
return false;
}
}
String path = StringUtils.substringBefore(prevUri, "?").toLowerCase(Locale.ROOT);
return !path.contains("//") && !path.contains("%2f") && !path.contains("%5c");
}
// a rejected value is attacker-controlled: it must not be able to forge log lines
private static String forLog(String prevUri) {
return prevUri.substring(0, Math.min(prevUri.length(), MAX_LOGGED_LENGTH)).replaceAll("[^\\x20-\\x7E]", "?");
}
}

6
application/src/main/java/org/thingsboard/server/service/security/model/token/OAuth2AppTokenFactory.java

@ -26,6 +26,7 @@ import io.jsonwebtoken.security.SignatureException;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.service.security.auth.oauth2.CallbackUrlSchemeValidator;
import java.util.Base64;
import java.util.Date;
@ -58,13 +59,16 @@ public class OAuth2AppTokenFactory {
if (timeDiff > MAX_EXPIRATION_TIME_DIFF_MS) {
throw new IllegalArgumentException("Application token expiration time can't be longer than 5 minutes");
}
if (!claims.getIssuer().equals(appPackage)) {
if (!appPackage.equals(claims.getIssuer())) {
throw new IllegalArgumentException("Application token issuer doesn't match application package");
}
String callbackUrlScheme = claims.get(CALLBACK_URL_SCHEME, String.class);
if (StringUtils.isEmpty(callbackUrlScheme)) {
throw new IllegalArgumentException("Application token doesn't have callbackUrlScheme");
}
if (!CallbackUrlSchemeValidator.isValid(callbackUrlScheme)) {
throw new IllegalArgumentException("Application token has invalid callbackUrlScheme");
}
return callbackUrlScheme;
}

11
application/src/main/java/org/thingsboard/server/service/telemetry/AbstractSubscriptionService.java

@ -25,6 +25,7 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.common.data.exception.TenantNotFoundException;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.msg.queue.ServiceType;
@ -86,7 +87,15 @@ public abstract class AbstractSubscriptionService extends TbApplicationEventList
protected void forwardToSubscriptionManagerService(TenantId tenantId, EntityId entityId,
Consumer<SubscriptionManagerService> toSubscriptionManagerService,
Supplier<TransportProtos.ToCoreMsg> toCore) {
TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, entityId);
TopicPartitionInfo tpi;
try {
tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, entityId);
} catch (TenantNotFoundException e) {
// The tenant was deleted (e.g. concurrently with an in-flight asynchronous save callback),
// so there is no partition to route to and no subscribers to notify. Nothing to forward.
log.debug("[{}][{}] Skipping subscription update: tenant no longer exists.", tenantId, entityId);
return;
}
if (currentPartitions.contains(tpi)) {
if (subscriptionManagerService.isPresent()) {
toSubscriptionManagerService.accept(subscriptionManagerService.get());

81
application/src/test/java/org/thingsboard/server/controller/AdminControllerMailOAuth2Test.java

@ -0,0 +1,81 @@
/**
* 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.controller;
import jakarta.servlet.http.Cookie;
import jakarta.servlet.http.HttpServletRequest;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.ValueSource;
import org.mockito.InjectMocks;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.service.security.system.SystemSecurityService;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
import static org.thingsboard.server.service.security.auth.oauth2.HttpCookieOAuth2AuthorizationRequestRepository.PREV_URI_COOKIE_NAME;
@ExtendWith(MockitoExtension.class)
public class AdminControllerMailOAuth2Test {
private static final String BASE_URL = "https://thingsboard.example.com";
private static final String DEFAULT_PREV_URI = "/settings/outgoing-mail";
@Mock
private SystemSecurityService systemSecurityService;
@InjectMocks
private AdminController adminController;
private HttpServletRequest request;
@BeforeEach
public void before() {
request = mock(HttpServletRequest.class);
when(systemSecurityService.getBaseUrl(any(TenantId.class), any(CustomerId.class), any(HttpServletRequest.class))).thenReturn(BASE_URL);
}
@Test
public void testInAppPathIsTakenFromPrevUriCookie() {
givenPrevUriCookie("/settings/notifications?tab=1");
assertThat(adminController.getMailOAuth2RedirectUrl(request)).isEqualTo(BASE_URL + "/settings/notifications?tab=1");
}
@ParameterizedTest
@ValueSource(strings = {"@evil.com/", "//evil.com", "https://evil.com", "/\\evil.com", "/settings#fragment"})
public void testForgedPrevUriCookieIsIgnored(String prevUri) {
givenPrevUriCookie(prevUri);
assertThat(adminController.getMailOAuth2RedirectUrl(request)).isEqualTo(BASE_URL + DEFAULT_PREV_URI);
}
@Test
public void testRedirectUrlWithoutPrevUriCookie() {
when(request.getCookies()).thenReturn(null);
assertThat(adminController.getMailOAuth2RedirectUrl(request)).isEqualTo(BASE_URL + DEFAULT_PREV_URI);
}
private void givenPrevUriCookie(String prevUri) {
when(request.getCookies()).thenReturn(new Cookie[]{new Cookie(PREV_URI_COOKIE_NAME, prevUri)});
}
}

33
application/src/test/java/org/thingsboard/server/controller/AdminControllerTest.java

@ -17,6 +17,7 @@ package org.thingsboard.server.controller;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import jakarta.servlet.http.Cookie;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.RandomStringUtils;
import org.junit.Test;
@ -38,6 +39,8 @@ import static org.mockito.ArgumentMatchers.anyString;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.content;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.jsonPath;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
import static org.thingsboard.server.service.security.auth.oauth2.HttpCookieOAuth2AuthorizationRequestRepository.PREV_URI_COOKIE_NAME;
import static org.thingsboard.server.service.security.auth.oauth2.HttpCookieOAuth2AuthorizationRequestRepository.PREV_URI_PARAMETER;
@Slf4j
@DaoSqlTest
@ -111,6 +114,36 @@ public class AdminControllerTest extends AbstractControllerTest {
.andExpect(statusReason(containsString("is prohibited")));
}
@Test
public void testMailOAuth2AuthorizationStoresOnlyInAppPrevUri() throws Exception {
loginSysAdmin();
AdminSettings mailSettings = doGet("/api/admin/settings/mail", AdminSettings.class);
JsonNode originalJsonValue = mailSettings.getJsonValue();
try {
ObjectNode jsonValue = JacksonUtil.fromString(originalJsonValue.toString(), ObjectNode.class);
jsonValue.put("clientId", "clientId");
jsonValue.put("authUri", "https://accounts.google.com/o/oauth2/v2/auth");
jsonValue.put("redirectUri", "https://thingsboard.io/api/admin/mail/oauth2/code");
jsonValue.set("scope", JacksonUtil.newArrayNode().add("https://mail.google.com/"));
mailSettings.setJsonValue(jsonValue);
doPost("/api/admin/settings", mailSettings, AdminSettings.class);
Cookie prevUriCookie = doGet("/api/admin/mail/oauth2/authorize?" + PREV_URI_PARAMETER + "=@evil.com/")
.andExpect(status().isOk()).andReturn().getResponse().getCookie(PREV_URI_COOKIE_NAME);
assertThat(prevUriCookie).isNull();
prevUriCookie = doGet("/api/admin/mail/oauth2/authorize?" + PREV_URI_PARAMETER + "=/settings/outgoing-mail")
.andExpect(status().isOk()).andReturn().getResponse().getCookie(PREV_URI_COOKIE_NAME);
assertThat(prevUriCookie).isNotNull();
assertThat(prevUriCookie.getValue()).isEqualTo("/settings/outgoing-mail");
} finally {
// the mail settings are shared by the whole test context
AdminSettings currentSettings = doGet("/api/admin/settings/mail", AdminSettings.class);
currentSettings.setJsonValue(originalJsonValue);
doPost("/api/admin/settings", currentSettings, AdminSettings.class);
}
}
@Test
public void testSendTestMail() throws Exception {
Mockito.doNothing().when(mailService).sendTestMail(any(), anyString());

81
application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/CallbackUrlSchemeValidatorTest.java

@ -0,0 +1,81 @@
/**
* 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.security.auth.oauth2;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.NullAndEmptySource;
import org.junit.jupiter.params.provider.ValueSource;
import org.springframework.security.oauth2.core.endpoint.OAuth2AuthorizationRequest;
import static org.assertj.core.api.Assertions.assertThat;
public class CallbackUrlSchemeValidatorTest {
@ParameterizedTest
@ValueSource(strings = {"tbmobile", "tb-mobile.app1", "TbMobile+1", "org.mycompany.myapp.auth", "com.my_company.app.auth"})
public void testMobileAppSchemeIsValid(String callbackUrlScheme) {
assertThat(CallbackUrlSchemeValidator.isValid(callbackUrlScheme)).isTrue();
}
@ParameterizedTest
@NullAndEmptySource
@ValueSource(strings = {
"https://evil.com",
"http://evil.com",
"https",
"HTTPS",
"javascript",
"data",
"file",
"vbscript",
"//evil.com",
"tbmobile/evil.com",
"tbmobile:evil.com",
"tbmobile evil",
"1tbmobile",
"tbmobile@evil.com"
})
public void testInvalidSchemeIsRejected(String callbackUrlScheme) {
assertThat(CallbackUrlSchemeValidator.isValid(callbackUrlScheme)).isFalse();
}
@Test
public void testValidSchemeIsTakenFromAuthorizationRequest() {
assertThat(CallbackUrlSchemeValidator.getCallbackUrlScheme(givenAuthorizationRequest("tbmobile"))).isEqualTo("tbmobile");
}
@Test
public void testForgedSchemeFromAuthorizationRequestIsIgnored() {
assertThat(CallbackUrlSchemeValidator.getCallbackUrlScheme(givenAuthorizationRequest("https://evil.com"))).isNull();
}
@Test
public void testAuthorizationRequestWithoutScheme() {
assertThat(CallbackUrlSchemeValidator.getCallbackUrlScheme(givenAuthorizationRequest(null))).isNull();
assertThat(CallbackUrlSchemeValidator.getCallbackUrlScheme(null)).isNull();
}
private OAuth2AuthorizationRequest givenAuthorizationRequest(String callbackUrlScheme) {
OAuth2AuthorizationRequest.Builder builder = OAuth2AuthorizationRequest.authorizationCode()
.authorizationUri("testUri").clientId("testId");
if (callbackUrlScheme != null) {
builder.attributes(attributes -> attributes.put(TbOAuth2ParameterNames.CALLBACK_URL_SCHEME, callbackUrlScheme));
}
return builder.build();
}
}

69
application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/HttpCookieOAuth2AuthorizationRequestRepositoryTest.java

@ -0,0 +1,69 @@
/**
* 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.security.auth.oauth2;
import jakarta.servlet.http.Cookie;
import jakarta.servlet.http.HttpServletRequest;
import jakarta.servlet.http.HttpServletResponse;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.NullAndEmptySource;
import org.junit.jupiter.params.provider.ValueSource;
import org.mockito.ArgumentCaptor;
import org.springframework.security.oauth2.core.endpoint.OAuth2AuthorizationRequest;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.Mockito.atLeastOnce;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import static org.thingsboard.server.service.security.auth.oauth2.HttpCookieOAuth2AuthorizationRequestRepository.PREV_URI_COOKIE_NAME;
import static org.thingsboard.server.service.security.auth.oauth2.HttpCookieOAuth2AuthorizationRequestRepository.PREV_URI_PARAMETER;
public class HttpCookieOAuth2AuthorizationRequestRepositoryTest {
private final HttpCookieOAuth2AuthorizationRequestRepository repository = new HttpCookieOAuth2AuthorizationRequestRepository();
@Test
public void testPrevUriSavedForInAppPath() {
assertThat(savePrevUri("/dashboards/3fa13530-6597-11ed-bd76-8bd591f0ec3e?state=someState"))
.isEqualTo("/dashboards/3fa13530-6597-11ed-bd76-8bd591f0ec3e?state=someState");
}
@ParameterizedTest
@NullAndEmptySource
@ValueSource(strings = {"@evil.com/"})
public void testPrevUriNotSavedForInvalidValue(String prevUri) {
assertThat(savePrevUri(prevUri)).isNull();
}
private String savePrevUri(String prevUri) {
HttpServletRequest request = mock(HttpServletRequest.class);
HttpServletResponse response = mock(HttpServletResponse.class);
when(request.getParameter(PREV_URI_PARAMETER)).thenReturn(prevUri);
repository.saveAuthorizationRequest(OAuth2AuthorizationRequest.authorizationCode()
.authorizationUri("testUri").clientId("testId").build(), request, response);
ArgumentCaptor<Cookie> cookieCaptor = ArgumentCaptor.forClass(Cookie.class);
verify(response, atLeastOnce()).addCookie(cookieCaptor.capture());
return cookieCaptor.getAllValues().stream()
.filter(cookie -> PREV_URI_COOKIE_NAME.equals(cookie.getName()))
.map(Cookie::getValue)
.findAny().orElse(null);
}
}

115
application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/OAuth2ClientMapperTest.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.service.security.auth.oauth2;
import org.junit.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.security.core.userdetails.UsernameNotFoundException;
import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.oauth2.OAuth2Client;
import org.thingsboard.server.common.data.security.Authority;
import org.thingsboard.server.controller.AbstractControllerTest;
import org.thingsboard.server.dao.oauth2.OAuth2User;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.dao.user.UserService;
import org.thingsboard.server.service.security.model.SecurityUser;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertThrows;
@DaoSqlTest
public class OAuth2ClientMapperTest extends AbstractControllerTest {
@Autowired
private BasicOAuth2ClientMapper basicOAuth2ClientMapper;
@Autowired
private UserService userService;
@Test
public void testShouldFindUserOfOwnTenant() throws Exception {
loginTenantAdmin();
OAuth2Client tenantClient = doPost("/api/oauth2/client", createOauth2Client(tenantId, "tenant client"), OAuth2Client.class);
OAuth2User oAuth2User = new OAuth2User();
oAuth2User.setEmail(TENANT_ADMIN_EMAIL);
SecurityUser securityUser = basicOAuth2ClientMapper.getOrCreateSecurityUserFromOAuth2User(oAuth2User, tenantClient);
assertThat(securityUser.getTenantId()).isEqualTo(tenantId);
assertThat(securityUser.getAuthority()).isEqualTo(Authority.TENANT_ADMIN);
}
@Test
public void testShouldNotFindUserOfAnotherTenant() throws Exception {
loginDifferentTenant();
loginTenantAdmin();
OAuth2Client tenantClient = doPost("/api/oauth2/client", createOauth2Client(tenantId, "tenant client"), OAuth2Client.class);
// the email attribute is controlled by the identity provider behind the client
OAuth2User oAuth2User = new OAuth2User();
oAuth2User.setEmail(DIFFERENT_TENANT_ADMIN_EMAIL);
UsernameNotFoundException exception = assertThrows(
UsernameNotFoundException.class,
() -> basicOAuth2ClientMapper.getOrCreateSecurityUserFromOAuth2User(oAuth2User, tenantClient));
assertThat(exception.getMessage()).isEqualTo("User not found: " + DIFFERENT_TENANT_ADMIN_EMAIL);
User differentTenantAdmin = userService.findUserByEmail(TenantId.SYS_TENANT_ID, DIFFERENT_TENANT_ADMIN_EMAIL);
assertThat(differentTenantAdmin.getTenantId()).isEqualTo(differentTenantId);
loginSysAdmin();
deleteDifferentTenant();
}
@Test
public void testShouldNotFindSysAdmin() throws Exception {
loginTenantAdmin();
OAuth2Client tenantClient = doPost("/api/oauth2/client", createOauth2Client(tenantId, "tenant client"), OAuth2Client.class);
OAuth2User oAuth2User = new OAuth2User();
oAuth2User.setEmail(SYS_ADMIN_EMAIL);
UsernameNotFoundException exception = assertThrows(
UsernameNotFoundException.class,
() -> basicOAuth2ClientMapper.getOrCreateSecurityUserFromOAuth2User(oAuth2User, tenantClient));
assertThat(exception.getMessage()).isEqualTo("User not found: " + SYS_ADMIN_EMAIL);
User sysAdmin = userService.findUserByEmail(TenantId.SYS_TENANT_ID, SYS_ADMIN_EMAIL);
assertThat(sysAdmin.getAuthority()).isEqualTo(Authority.SYS_ADMIN);
}
@Test
public void testShouldCreateUserInClientTenant() throws Exception {
loginDifferentTenant();
loginTenantAdmin();
OAuth2Client tenantClient = doPost("/api/oauth2/client", createOauth2Client(tenantId, "tenant client"), OAuth2Client.class);
// a custom mapper endpoint may return any tenant id; the client's own tenant must win
String email = "userA@corporation.gmail.com";
OAuth2User oAuth2User = new OAuth2User();
oAuth2User.setEmail(email);
oAuth2User.setTenantId(differentTenantId);
basicOAuth2ClientMapper.getOrCreateSecurityUserFromOAuth2User(oAuth2User, tenantClient);
User created = userService.findUserByEmail(TenantId.SYS_TENANT_ID, email);
assertThat(created.getTenantId()).isEqualTo(tenantId);
loginSysAdmin();
deleteDifferentTenant();
}
}

103
application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationFailureHandlerTest.java

@ -0,0 +1,103 @@
/**
* 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.security.auth.oauth2;
import jakarta.servlet.http.HttpServletRequest;
import jakarta.servlet.http.HttpServletResponse;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.ValueSource;
import org.mockito.ArgumentCaptor;
import org.springframework.security.authentication.AuthenticationServiceException;
import org.springframework.security.oauth2.core.endpoint.OAuth2AuthorizationRequest;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.service.security.system.SystemSecurityService;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
public class Oauth2AuthenticationFailureHandlerTest {
private static final String BASE_URL = "https://thingsboard.example.com";
private final HttpCookieOAuth2AuthorizationRequestRepository authorizationRequestRepository =
mock(HttpCookieOAuth2AuthorizationRequestRepository.class);
private final SystemSecurityService systemSecurityService = mock(SystemSecurityService.class);
private final Oauth2AuthenticationFailureHandler failureHandler =
new Oauth2AuthenticationFailureHandler(authorizationRequestRepository, systemSecurityService);
private HttpServletRequest request;
private HttpServletResponse response;
@BeforeEach
public void before() {
request = mock(HttpServletRequest.class);
response = mock(HttpServletResponse.class);
when(request.getContextPath()).thenReturn("");
when(response.encodeRedirectURL(anyString())).thenAnswer(invocation -> invocation.getArgument(0));
when(systemSecurityService.getBaseUrl(any(TenantId.class), any(CustomerId.class), any(HttpServletRequest.class))).thenReturn(BASE_URL);
}
@Test
public void testErrorIsSentToMobileAppScheme() throws Exception {
givenCallbackUrlScheme("tbmobile");
assertThat(sendFailure()).isEqualTo("tbmobile:/?error=someError");
}
/**
* The scheme is restored from the oauth2_auth_request cookie, so a forged one must not turn the error redirect
* into a link to another host.
*/
@ParameterizedTest
@ValueSource(strings = {"https://evil.com", "javascript"})
public void testForgedCallbackUrlSchemeFallsBackToLoginPage(String callbackUrlScheme) throws Exception {
givenCallbackUrlScheme(callbackUrlScheme);
assertThat(sendFailure()).isEqualTo(BASE_URL + "/login?loginError=someError");
}
@Test
public void testErrorIsSentToLoginPageWithoutCallbackUrlScheme() throws Exception {
givenCallbackUrlScheme(null);
assertThat(sendFailure()).isEqualTo(BASE_URL + "/login?loginError=someError");
}
@Test
public void testErrorIsSentToLoginPageWithoutAuthorizationRequest() throws Exception {
assertThat(sendFailure()).isEqualTo(BASE_URL + "/login?loginError=someError");
}
private void givenCallbackUrlScheme(String callbackUrlScheme) {
when(authorizationRequestRepository.loadAuthorizationRequest(request)).thenReturn(
OAuth2AuthorizationRequest.authorizationCode().authorizationUri("testUri").clientId("testId")
.attributes(attributes -> attributes.put(TbOAuth2ParameterNames.CALLBACK_URL_SCHEME, callbackUrlScheme))
.build());
}
private String sendFailure() throws Exception {
failureHandler.onAuthenticationFailure(request, response, new AuthenticationServiceException("someError"));
ArgumentCaptor<String> redirectCaptor = ArgumentCaptor.forClass(String.class);
verify(response).sendRedirect(redirectCaptor.capture());
return redirectCaptor.getValue();
}
}

186
application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationSuccessHandlerTest.java

@ -15,54 +15,178 @@
*/
package org.thingsboard.server.service.security.auth.oauth2;
import org.junit.Before;
import org.junit.Test;
import org.mockito.Mock;
import org.springframework.beans.factory.annotation.Autowired;
import org.thingsboard.server.common.data.id.UserId;
import jakarta.servlet.http.Cookie;
import jakarta.servlet.http.HttpServletRequest;
import jakarta.servlet.http.HttpServletResponse;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.CsvSource;
import org.junit.jupiter.params.provider.ValueSource;
import org.mockito.ArgumentCaptor;
import org.springframework.security.oauth2.client.OAuth2AuthorizedClient;
import org.springframework.security.oauth2.client.OAuth2AuthorizedClientService;
import org.springframework.security.oauth2.client.authentication.OAuth2AuthenticationToken;
import org.springframework.security.oauth2.core.OAuth2AccessToken;
import org.springframework.security.oauth2.core.user.OAuth2User;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.oauth2.MapperType;
import org.thingsboard.server.common.data.oauth2.OAuth2Client;
import org.thingsboard.server.common.data.oauth2.OAuth2MapperConfig;
import org.thingsboard.server.common.data.security.model.JwtPair;
import org.thingsboard.server.controller.AbstractControllerTest;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.dao.oauth2.OAuth2ClientService;
import org.thingsboard.server.service.security.model.SecurityUser;
import org.thingsboard.server.service.security.model.token.JwtTokenFactory;
import org.thingsboard.server.service.security.system.SystemSecurityService;
import java.time.Instant;
import java.util.UUID;
import static org.junit.Assert.assertEquals;
import static org.mockito.ArgumentMatchers.eq;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import static org.thingsboard.server.service.security.auth.oauth2.HttpCookieOAuth2AuthorizationRequestRepository.PREV_URI_COOKIE_NAME;
@DaoSqlTest
public class Oauth2AuthenticationSuccessHandlerTest extends AbstractControllerTest {
public class Oauth2AuthenticationSuccessHandlerTest {
@Autowired
private Oauth2AuthenticationSuccessHandler oauth2AuthenticationSuccessHandler;
private static final String BASE_URL = "https://thingsboard.example.com";
private static final String PREV_URI = "/dashboards/3fa13530-6597-11ed-bd76-8bd591f0ec3e";
private static final JwtPair TOKEN_PAIR = new JwtPair("testAccessToken", "testRefreshToken");
@Mock
private JwtTokenFactory jwtTokenFactory;
private final JwtTokenFactory tokenFactory = mock(JwtTokenFactory.class);
private final OAuth2ClientMapperProvider oauth2ClientMapperProvider = mock(OAuth2ClientMapperProvider.class);
private final OAuth2ClientService oAuth2ClientService = mock(OAuth2ClientService.class);
private final OAuth2AuthorizedClientService oAuth2AuthorizedClientService = mock(OAuth2AuthorizedClientService.class);
private final SystemSecurityService systemSecurityService = mock(SystemSecurityService.class);
private final Oauth2AuthenticationSuccessHandler successHandler = new Oauth2AuthenticationSuccessHandler(
tokenFactory, oauth2ClientMapperProvider, oAuth2ClientService, oAuth2AuthorizedClientService,
mock(HttpCookieOAuth2AuthorizationRequestRepository.class), systemSecurityService);
private SecurityUser securityUser;
private HttpServletRequest request;
private HttpServletResponse response;
@Before
@BeforeEach
public void before() {
UserId userId = new UserId(UUID.randomUUID());
securityUser = new SecurityUser(userId);
when(jwtTokenFactory.createTokenPair(eq(securityUser))).thenReturn(new JwtPair("testAccessToken", "testRefreshToken"));
request = mock(HttpServletRequest.class);
response = mock(HttpServletResponse.class);
when(systemSecurityService.getBaseUrl(any(TenantId.class), any(CustomerId.class), any(HttpServletRequest.class))).thenReturn(BASE_URL);
when(response.encodeRedirectURL(anyString())).thenAnswer(invocation -> invocation.getArgument(0));
}
@Test
public void testGetRedirectUrl() {
JwtPair jwtPair = jwtTokenFactory.createTokenPair(securityUser);
public void testInAppPathIsTakenFromPrevUriCookie() {
givenPrevUriCookie(PREV_URI + "?state=someState");
assertThat(successHandler.getBaseUrl(request, null)).isEqualTo(BASE_URL);
assertThat(successHandler.getPrevUri(request, response, null)).isEqualTo(PREV_URI + "?state=someState");
}
String urlWithoutParams = "http://localhost:8080/dashboardGroups/3fa13530-6597-11ed-bd76-8bd591f0ec3e";
String urlWithParams = "http://localhost:8080/dashboardGroups/3fa13530-6597-11ed-bd76-8bd591f0ec3e?state=someState&page=1";
@ParameterizedTest
@ValueSource(strings = {"@evil.com/", "//evil.com", "https://evil.com", "/\\evil.com", "/dashboards#fragment"})
public void testForgedPrevUriCookieIsIgnored(String prevUri) {
givenPrevUriCookie(prevUri);
assertThat(successHandler.getPrevUri(request, response, null)).isEmpty();
}
String redirectUrl = oauth2AuthenticationSuccessHandler.getRedirectUrl(urlWithoutParams, jwtPair);
String expectedUrl = urlWithoutParams + "/?accessToken=" + jwtPair.getToken() + "&refreshToken=" + jwtPair.getRefreshToken();
assertEquals(expectedUrl, redirectUrl);
@Test
public void testForgedPrevUriCookieIsDeleted() {
givenPrevUriCookie("@evil.com/");
redirectUrl = oauth2AuthenticationSuccessHandler.getRedirectUrl(urlWithParams, jwtPair);
expectedUrl = urlWithParams + "&accessToken=" + jwtPair.getToken() + "&refreshToken=" + jwtPair.getRefreshToken();
assertEquals(expectedUrl, redirectUrl);
successHandler.getPrevUri(request, response, null);
ArgumentCaptor<Cookie> cookieCaptor = ArgumentCaptor.forClass(Cookie.class);
verify(response).addCookie(cookieCaptor.capture());
assertThat(cookieCaptor.getValue().getName()).isEqualTo(PREV_URI_COOKIE_NAME);
assertThat(cookieCaptor.getValue().getMaxAge()).isZero();
}
}
@Test
public void testBaseUrlWithoutPrevUriCookie() {
when(request.getCookies()).thenReturn(null);
assertThat(successHandler.getBaseUrl(request, null)).isEqualTo(BASE_URL);
assertThat(successHandler.getPrevUri(request, response, null)).isEmpty();
}
@Test
public void testCallbackUrlSchemeIgnoresPrevUri() {
givenPrevUriCookie(PREV_URI);
assertThat(successHandler.getBaseUrl(request, "tbmobile")).isEqualTo("tbmobile:");
assertThat(successHandler.getPrevUri(request, response, "tbmobile")).isEmpty();
}
@Test
public void testSuccessRedirectCarriesTokensToPrevUri() throws Exception {
givenPrevUriCookie(PREV_URI);
givenSuccessfulLogin();
successHandler.onAuthenticationSuccess(request, response, givenAuthentication());
assertThat(captureRedirect()).isEqualTo(BASE_URL + PREV_URI +
"/?accessToken=testAccessToken&refreshToken=testRefreshToken");
}
/**
* The error redirect appends its own path, so it must be built from the base URL alone - with prevUri in it the
* result would be an unroutable https://host/dashboards/x/login?loginError=...
*/
@Test
public void testErrorRedirectDropsPrevUri() throws Exception {
givenPrevUriCookie(PREV_URI);
when(oAuth2ClientService.findOAuth2ClientById(any(), any())).thenThrow(new RuntimeException("someError"));
successHandler.onAuthenticationSuccess(request, response, givenAuthentication());
assertThat(captureRedirect()).isEqualTo(BASE_URL + "/login?loginError=someError");
}
@ParameterizedTest
@CsvSource({
"https://thingsboard.example.com/dashboards/3fa13530-6597-11ed-bd76-8bd591f0ec3e, https://thingsboard.example.com/dashboards/3fa13530-6597-11ed-bd76-8bd591f0ec3e/?",
"https://thingsboard.example.com/dashboards/3fa13530-6597-11ed-bd76-8bd591f0ec3e?state=someState&page=1, https://thingsboard.example.com/dashboards/3fa13530-6597-11ed-bd76-8bd591f0ec3e?state=someState&page=1&",
"https://thingsboard.example.com/, https://thingsboard.example.com/?"
})
public void testGetRedirectUrl(String baseUrl, String expectedPrefix) {
assertThat(successHandler.getRedirectUrl(baseUrl, TOKEN_PAIR))
.isEqualTo(expectedPrefix + "accessToken=testAccessToken&refreshToken=testRefreshToken");
}
private void givenPrevUriCookie(String prevUri) {
when(request.getCookies()).thenReturn(new Cookie[]{new Cookie(PREV_URI_COOKIE_NAME, prevUri)});
}
private OAuth2AuthenticationToken givenAuthentication() {
OAuth2User principal = mock(OAuth2User.class);
when(principal.getName()).thenReturn("testUser");
OAuth2AuthenticationToken token = mock(OAuth2AuthenticationToken.class);
when(token.getAuthorizedClientRegistrationId()).thenReturn(UUID.randomUUID().toString());
when(token.getPrincipal()).thenReturn(principal);
return token;
}
private void givenSuccessfulLogin() {
OAuth2Client oauth2Client = new OAuth2Client();
oauth2Client.setMapperConfig(OAuth2MapperConfig.builder().type(MapperType.BASIC).build());
when(oAuth2ClientService.findOAuth2ClientById(any(), any())).thenReturn(oauth2Client);
OAuth2AuthorizedClient authorizedClient = mock(OAuth2AuthorizedClient.class);
when(authorizedClient.getAccessToken()).thenReturn(new OAuth2AccessToken(OAuth2AccessToken.TokenType.BEARER,
"testProviderAccessToken", Instant.now(), Instant.now().plusSeconds(60)));
when(oAuth2AuthorizedClientService.loadAuthorizedClient(anyString(), anyString())).thenReturn(authorizedClient);
SecurityUser securityUser = mock(SecurityUser.class);
OAuth2ClientMapper mapper = mock(OAuth2ClientMapper.class);
when(mapper.getOrCreateUserByClientPrincipal(any(), any(), anyString(), any())).thenReturn(securityUser);
when(oauth2ClientMapperProvider.getOAuth2ClientMapperByType(any())).thenReturn(mapper);
when(tokenFactory.createTokenPair(securityUser)).thenReturn(TOKEN_PAIR);
}
private String captureRedirect() throws Exception {
ArgumentCaptor<String> redirectCaptor = ArgumentCaptor.forClass(String.class);
verify(response).sendRedirect(redirectCaptor.capture());
return redirectCaptor.getValue();
}
}

74
application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/PrevUriValidatorTest.java

@ -0,0 +1,74 @@
/**
* 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.security.auth.oauth2;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.NullAndEmptySource;
import org.junit.jupiter.params.provider.ValueSource;
import static org.assertj.core.api.Assertions.assertThat;
public class PrevUriValidatorTest {
@ParameterizedTest
@ValueSource(strings = {
"/",
"/login",
"/dashboards/3fa13530-6597-11ed-bd76-8bd591f0ec3e",
"/dashboards/3fa13530-6597-11ed-bd76-8bd591f0ec3e?state=someState&page=1",
"/settings/outgoing-mail",
"/some%20path?q=a+b",
"/dashboards/3fa13530-6597-11ed-bd76-8bd591f0ec3e?state=W3siaWQiOiJhL2IifV0%3D"
})
public void testValidPrevUri(String prevUri) {
assertThat(PrevUriValidator.isValid(prevUri)).isTrue();
}
@ParameterizedTest
@NullAndEmptySource
@ValueSource(strings = {
"@evil.com/",
"evil.com",
"https://evil.com",
"//evil.com",
"/\\evil.com",
"/\tevil.com",
"/ evil.com",
"/dashboards\\..\\evil.com",
"/login\nLocation: https://evil.com",
"/login\r\nSet-Cookie: a=b",
"/dashboards//evil.com",
"/dashboards%2Fevil.com",
"/dashboards%5cevil.com",
"/dashboards;jsessionid=1",
"/dashboards?title=a,b",
"/dashboards,list",
"/dashboards\"list",
"/dashboards/3fa13530-6597-11ed-bd76-8bd591f0ec3e#fragment",
"/панель"
})
public void testInvalidPrevUri(String prevUri) {
assertThat(PrevUriValidator.isValid(prevUri)).isFalse();
}
@Test
public void testPrevUriLengthLimit() {
assertThat(PrevUriValidator.isValid("/" + "a".repeat(2047))).isTrue();
assertThat(PrevUriValidator.isValid("/" + "a".repeat(2048))).isFalse();
}
}

72
application/src/test/java/org/thingsboard/server/service/security/model/token/OAuth2AppTokenFactoryTest.java

@ -0,0 +1,72 @@
/**
* 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.security.model.token;
import io.jsonwebtoken.JwtBuilder;
import io.jsonwebtoken.Jwts;
import io.jsonwebtoken.security.Keys;
import org.junit.jupiter.api.Test;
import javax.crypto.SecretKey;
import java.util.Base64;
import java.util.Date;
import java.util.concurrent.TimeUnit;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
public class OAuth2AppTokenFactoryTest {
private static final String APP_PACKAGE = "org.thingsboard.demo.app";
private static final byte[] KEY_BYTES = "yjNyylzT1TmiVE2jV3YTnUpZzwLLLdPDJKmhLNyXDPnLtVCLcJIjIGmDPKHNoDMK".getBytes();
private final OAuth2AppTokenFactory tokenFactory = new OAuth2AppTokenFactory();
@Test
public void testMobileAppSchemeIsAccepted() {
assertThat(validate(appToken("tb-mobile.app1", APP_PACKAGE))).isEqualTo("tb-mobile.app1");
}
@Test
public void testInvalidCallbackUrlSchemeIsRejected() {
assertThatThrownBy(() -> validate(appToken("https://evil.com", APP_PACKAGE)))
.isInstanceOf(IllegalArgumentException.class)
.hasMessageContaining("callbackUrlScheme");
}
@Test
public void testTokenWithoutIssuerIsRejected() {
assertThatThrownBy(() -> validate(appToken("tbmobile", null)))
.isInstanceOf(IllegalArgumentException.class)
.hasMessageContaining("issuer");
}
private String appToken(String callbackUrlScheme, String issuer) {
JwtBuilder builder = Jwts.builder()
.expiration(new Date(System.currentTimeMillis() + TimeUnit.MINUTES.toMillis(1)))
.claim("callbackUrlScheme", callbackUrlScheme);
if (issuer != null) {
builder.issuer(issuer);
}
SecretKey key = Keys.hmacShaKeyFor(KEY_BYTES);
return builder.signWith(key).compact();
}
private String validate(String appToken) {
return tokenFactory.validateTokenAndGetCallbackUrlScheme(APP_PACKAGE, appToken, Base64.getEncoder().encodeToString(KEY_BYTES));
}
}

22
application/src/test/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionServiceTest.java

@ -40,6 +40,7 @@ import org.thingsboard.server.common.data.ApiUsageStateValue;
import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.EntityView;
import org.thingsboard.server.common.data.exception.TenantNotFoundException;
import org.thingsboard.server.common.data.id.ApiUsageStateId;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.DeviceId;
@ -89,6 +90,7 @@ import java.util.stream.Stream;
import static com.google.common.util.concurrent.Futures.immediateFailedFuture;
import static com.google.common.util.concurrent.Futures.immediateFuture;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatNoException;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
@ -1155,6 +1157,26 @@ class DefaultTelemetrySubscriptionServiceTest {
then(deviceStateManager).shouldHaveNoInteractions();
}
/* --- Subscription forwarding --- */
@Test
void shouldSkipSubscriptionForwardWhenTenantWasDeleted() {
// GIVEN the tenant was deleted concurrently, so partition resolution fails
given(partitionService.resolve(ServiceType.TB_CORE, tenantId, entityId))
.willThrow(new TenantNotFoundException(tenantId));
// WHEN forwarding a subscription update (e.g. from an in-flight async save callback)
// THEN it must not propagate the exception
assertThatNoException().isThrownBy(() -> telemetryService.forwardToSubscriptionManagerService(
tenantId, entityId,
sm -> sm.onAttributesUpdate(tenantId, entityId, AttributeScope.SERVER_SCOPE.name(), List.of(), TbCallback.EMPTY),
() -> null));
// AND nothing is forwarded, since there is no partition to route to and no subscribers to notify
then(subscriptionManagerService).shouldHaveNoInteractions();
then(clusterService).shouldHaveNoInteractions();
}
// used to emulate versions returned by save APIs
private static List<Long> listOfNNumbers(int N) {
return LongStream.range(0, N).boxed().toList();

2
common/actor/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>common</artifactId>
</parent>
<groupId>org.thingsboard.common</groupId>

2
common/cache/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>common</artifactId>
</parent>
<groupId>org.thingsboard.common</groupId>

2
common/cluster-api/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>common</artifactId>
</parent>
<groupId>org.thingsboard.common</groupId>

2
common/coap-server/pom.xml

@ -22,7 +22,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>common</artifactId>
</parent>
<groupId>org.thingsboard.common</groupId>

2
common/dao-api/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>common</artifactId>
</parent>
<groupId>org.thingsboard.common</groupId>

2
common/data/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>common</artifactId>
</parent>
<groupId>org.thingsboard.common</groupId>

2
common/discovery-api/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>common</artifactId>
</parent>
<groupId>org.thingsboard.common</groupId>

2
common/edge-api/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>common</artifactId>
</parent>
<groupId>org.thingsboard.common</groupId>

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

@ -19,8 +19,17 @@ import io.grpc.HttpConnectProxiedSocketAddress;
import io.grpc.ManagedChannel;
import io.grpc.netty.shaded.io.grpc.netty.GrpcSslContexts;
import io.grpc.netty.shaded.io.grpc.netty.NettyChannelBuilder;
import io.grpc.netty.shaded.io.netty.channel.Channel;
import io.grpc.netty.shaded.io.netty.channel.EventLoopGroup;
import io.grpc.netty.shaded.io.netty.channel.epoll.Epoll;
import io.grpc.netty.shaded.io.netty.channel.epoll.EpollEventLoopGroup;
import io.grpc.netty.shaded.io.netty.channel.epoll.EpollSocketChannel;
import io.grpc.netty.shaded.io.netty.channel.nio.NioEventLoopGroup;
import io.grpc.netty.shaded.io.netty.channel.socket.nio.NioSocketChannel;
import io.grpc.netty.shaded.io.netty.handler.ssl.SslContextBuilder;
import io.grpc.netty.shaded.io.netty.util.concurrent.DefaultThreadFactory;
import io.grpc.stub.StreamObserver;
import jakarta.annotation.PreDestroy;
import lombok.Getter;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
@ -84,8 +93,14 @@ public class EdgeGrpcClient implements EdgeRpcClient {
private ManagedChannel channel;
private final EventLoopGroup workerGroup = createWorkerGroup();
private StreamObserver<RequestMsg> inputStream;
private volatile boolean connected;
private volatile boolean streamActive;
private static final ReentrantLock uplinkMsgLock = new ReentrantLock();
@Override
@ -95,7 +110,11 @@ public class EdgeGrpcClient implements EdgeRpcClient {
Consumer<EdgeConfiguration> onEdgeUpdate,
Consumer<DownlinkMsg> onDownlink,
Consumer<Exception> onError) {
connected = false;
streamActive = false;
NettyChannelBuilder builder = NettyChannelBuilder.forAddress(rpcHost, rpcPort)
.eventLoopGroup(workerGroup)
.channelType(channelType())
.maxInboundMessageSize(maxInboundMessageSize)
.keepAliveTime(keepAliveTimeSec, TimeUnit.SECONDS)
.keepAliveTimeout(keepAliveTimeoutSec, TimeUnit.SECONDS)
@ -131,6 +150,7 @@ public class EdgeGrpcClient implements EdgeRpcClient {
EdgeRpcServiceGrpc.EdgeRpcServiceStub stub = EdgeRpcServiceGrpc.newStub(channel);
log.info("[{}] Sending a connect request to the TB!", edgeKey);
this.inputStream = stub.withCompression("gzip").handleMsgs(initOutputStream(edgeKey, onUplinkResponse, onEdgeUpdate, onDownlink, onError));
streamActive = true;
this.inputStream.onNext(RequestMsg.newBuilder()
.setMsgType(RequestMsgType.CONNECT_RPC_MESSAGE)
.setConnectRequestMsg(ConnectRequestMsg.newBuilder()
@ -146,6 +166,20 @@ public class EdgeGrpcClient implements EdgeRpcClient {
return EdgeVersionComparator.getNewestEdgeVersion();
}
private static EventLoopGroup createWorkerGroup() {
DefaultThreadFactory threadFactory = new DefaultThreadFactory("edge-grpc-worker", true);
return Epoll.isAvailable() ? new EpollEventLoopGroup(1, threadFactory) : new NioEventLoopGroup(1, threadFactory);
}
private static Class<? extends Channel> channelType() {
return Epoll.isAvailable() ? EpollSocketChannel.class : NioSocketChannel.class;
}
@PreDestroy
public void destroy() {
workerGroup.shutdownGracefully();
}
private StreamObserver<ResponseMsg> initOutputStream(String edgeKey,
Consumer<UplinkResponseMsg> onUplinkResponse,
Consumer<EdgeConfiguration> onEdgeUpdate,
@ -162,8 +196,10 @@ public class EdgeGrpcClient implements EdgeRpcClient {
serverMaxInboundMessageSize = connectResponseMsg.getMaxInboundMessageSize();
}
log.info("[{}] Configuration received: {}", edgeKey, connectResponseMsg.getConfiguration());
connected = true;
onEdgeUpdate.accept(connectResponseMsg.getConfiguration());
} else {
connected = false;
log.error("[{}] Failed to establish the connection! Code: {}. Error message: {}.", edgeKey, connectResponseMsg.getResponseCode(), connectResponseMsg.getErrorMsg());
try {
EdgeGrpcClient.this.disconnect(true);
@ -186,6 +222,8 @@ public class EdgeGrpcClient implements EdgeRpcClient {
@Override
public void onError(Throwable t) {
connected = false;
streamActive = false;
log.warn("[{}] Stream was terminated due to error:", edgeKey, t);
try {
EdgeGrpcClient.this.disconnect(true);
@ -197,6 +235,8 @@ public class EdgeGrpcClient implements EdgeRpcClient {
@Override
public void onCompleted() {
connected = false;
streamActive = false;
log.info("[{}] Stream was closed and completed successfully!", edgeKey);
}
};
@ -204,6 +244,8 @@ public class EdgeGrpcClient implements EdgeRpcClient {
@Override
public void disconnect(boolean onError) throws InterruptedException {
connected = false;
streamActive = false;
if (!onError) {
try {
if (inputStream != null) {
@ -236,10 +278,19 @@ public class EdgeGrpcClient implements EdgeRpcClient {
}
}
@Override
public boolean isConnected() {
return connected;
}
@Override
public void sendUplinkMsg(UplinkMsg msg) {
uplinkMsgLock.lock();
try {
if (!streamActive) {
log.debug("Uplink msg is skipped, the cloud session is not established: {}", msg);
return;
}
this.inputStream.onNext(RequestMsg.newBuilder()
.setMsgType(RequestMsgType.UPLINK_RPC_MESSAGE)
.setUplinkMsg(msg)
@ -253,6 +304,10 @@ public class EdgeGrpcClient implements EdgeRpcClient {
public void sendSyncRequestMsg(boolean fullSyncRequired) {
uplinkMsgLock.lock();
try {
if (!streamActive) {
log.debug("Sync request msg is skipped, the cloud session is not established");
return;
}
SyncRequestMsg syncRequestMsg = SyncRequestMsg.newBuilder()
.setFullSync(fullSyncRequired)
.build();
@ -269,6 +324,10 @@ public class EdgeGrpcClient implements EdgeRpcClient {
public void sendDownlinkResponseMsg(DownlinkResponseMsg downlinkResponseMsg) {
uplinkMsgLock.lock();
try {
if (!streamActive) {
log.debug("Downlink response msg is skipped, the cloud session is not established: {}", downlinkResponseMsg);
return;
}
this.inputStream.onNext(RequestMsg.newBuilder()
.setMsgType(RequestMsgType.UPLINK_RPC_MESSAGE)
.setDownlinkResponseMsg(downlinkResponseMsg)

2
common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java

@ -34,6 +34,8 @@ public interface EdgeRpcClient {
void disconnect(boolean onError) throws InterruptedException;
boolean isConnected();
void sendSyncRequestMsg(boolean fullSyncRequired);
void sendUplinkMsg(UplinkMsg uplinkMsg);

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

@ -50,6 +50,7 @@ enum EdgeVersion {
V_4_2_2_2 = 4222;
V_4_2_2_3 = 4223;
V_4_2_2_4 = 4224;
V_4_2_2_5 = 4225;
V_LATEST = 99999;
}

168
common/edge-api/src/test/java/org/thingsboard/edge/rpc/EdgeGrpcClientLeakTest.java

@ -0,0 +1,168 @@
/**
* Copyright © 2016-2026 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.edge.rpc;
import io.grpc.Server;
import io.grpc.netty.shaded.io.grpc.netty.NettyServerBuilder;
import io.grpc.netty.shaded.io.netty.buffer.PooledByteBufAllocator;
import io.grpc.stub.StreamObserver;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
import org.springframework.test.util.ReflectionTestUtils;
import org.thingsboard.server.gen.edge.v1.ConnectResponseCode;
import org.thingsboard.server.gen.edge.v1.ConnectResponseMsg;
import org.thingsboard.server.gen.edge.v1.EdgeRpcServiceGrpc;
import org.thingsboard.server.gen.edge.v1.RequestMsg;
import org.thingsboard.server.gen.edge.v1.RequestMsgType;
import org.thingsboard.server.gen.edge.v1.ResponseMsg;
import org.thingsboard.server.gen.edge.v1.UplinkMsg;
import java.lang.reflect.Method;
import java.util.concurrent.TimeUnit;
import java.util.function.BooleanSupplier;
import static org.junit.jupiter.api.Assertions.fail;
class EdgeGrpcClientLeakTest {
// The window in which the default shared event loop group would be destroyed after the channel
// terminates (SharedResourceHolder delays destruction by 1 second). If EdgeGrpcClient ever goes
// back to the shared group, writes after this window hit a terminated executor and every buffer
// committed to grpc-netty's WriteQueue is pinned forever (4112 bytes per message, silent after
// the first RejectedExecutionException).
private static final long SHARED_GROUP_DEATH_WINDOW_MS = 3000;
private static final int MSG_COUNT = 50;
private static final long AWAIT_TIMEOUT_MS = 15_000;
private Server server;
private EdgeGrpcClient client;
@Test
void uplinksSentAfterTransportDeathDoNotPinPooledBuffers() throws Exception {
server = NettyServerBuilder.forPort(0)
.addService(new EdgeRpcServiceGrpc.EdgeRpcServiceImplBase() {
@Override
public StreamObserver<RequestMsg> handleMsgs(StreamObserver<ResponseMsg> outputStream) {
return new StreamObserver<>() {
@Override
public void onNext(RequestMsg requestMsg) {
if (requestMsg.hasConnectRequestMsg()) {
outputStream.onNext(ResponseMsg.newBuilder()
.setConnectResponseMsg(ConnectResponseMsg.newBuilder()
.setResponseCode(ConnectResponseCode.ACCEPTED)
.build())
.build());
}
}
@Override
public void onError(Throwable t) {
}
@Override
public void onCompleted() {
}
};
}
})
.build()
.start();
client = new EdgeGrpcClient();
ReflectionTestUtils.setField(client, "rpcHost", "localhost");
ReflectionTestUtils.setField(client, "rpcPort", server.getPort());
ReflectionTestUtils.setField(client, "timeoutSecs", 1);
ReflectionTestUtils.setField(client, "keepAliveTimeSec", 10);
ReflectionTestUtils.setField(client, "keepAliveTimeoutSec", 5);
ReflectionTestUtils.setField(client, "maxInboundMessageSize", 4194304);
client.connect("leakTest", "leakTest", msg -> {}, cfg -> {}, msg -> {}, e -> {});
await("client to connect", () -> client.isConnected());
server.shutdownNow();
server.awaitTermination(10, TimeUnit.SECONDS);
await("client to observe the transport death", () -> !client.isConnected());
Thread.sleep(SHARED_GROUP_DEATH_WINDOW_MS);
long baseline = pinnedBytes();
@SuppressWarnings("unchecked")
StreamObserver<RequestMsg> inputStream = (StreamObserver<RequestMsg>) ReflectionTestUtils.getField(client, "inputStream");
RequestMsg uplink = RequestMsg.newBuilder()
.setMsgType(RequestMsgType.UPLINK_RPC_MESSAGE)
.setUplinkMsg(UplinkMsg.newBuilder().setUplinkMsgId(1).build())
.build();
// Bypasses the connected gate on purpose: this models the check-then-act straggler (and the
// pre-gate retry loop) writing to a stream whose transport is already gone. Exceptions are
// swallowed the same way the production retry loop survives them.
for (int i = 0; i < MSG_COUNT; i++) {
try {
inputStream.onNext(uplink);
} catch (RuntimeException ignored) {
}
}
long deadline = System.currentTimeMillis() + AWAIT_TIMEOUT_MS;
while (pinnedBytes() > baseline) {
if (System.currentTimeMillis() > deadline) {
fail("Pinned pooled memory did not return to baseline: " + (pinnedBytes() - baseline)
+ " bytes retained after " + MSG_COUNT + " uplinks to a dead stream");
}
Thread.sleep(50);
}
}
@AfterEach
void tearDown() throws Exception {
if (client != null) {
client.disconnect(true);
client.destroy();
}
if (server != null) {
server.shutdownNow();
}
}
private void await(String what, BooleanSupplier condition) throws InterruptedException {
long deadline = System.currentTimeMillis() + AWAIT_TIMEOUT_MS;
while (!condition.getAsBoolean()) {
if (System.currentTimeMillis() > deadline) {
fail("Timed out waiting for " + what);
}
Thread.sleep(50);
}
}
// grpc-netty builds its own PooledByteBufAllocator instances instead of using
// PooledByteBufAllocator.DEFAULT, and the factory that owns them is package private - so they
// have to be pulled out reflectively. Both variants are checked so the assertion holds no matter
// which one the transport picks on this platform/version.
private static long pinnedBytes() {
try {
Class<?> utils = Class.forName("io.grpc.netty.shaded.io.grpc.netty.Utils");
Method getByteBufAllocator = utils.getDeclaredMethod("getByteBufAllocator", boolean.class);
getByteBufAllocator.setAccessible(true);
long total = 0;
for (boolean forceHeapBuffer : new boolean[]{false, true}) {
PooledByteBufAllocator pooled = (PooledByteBufAllocator) getByteBufAllocator.invoke(null, forceHeapBuffer);
total += pooled.pinnedDirectMemory() + pooled.pinnedHeapMemory();
}
return total;
} catch (Exception e) {
throw new IllegalStateException("Failed to read the gRPC allocator metrics", e);
}
}
}

2
common/edqs/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>common</artifactId>
</parent>
<groupId>org.thingsboard.common</groupId>

2
common/message/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>common</artifactId>
</parent>
<groupId>org.thingsboard.common</groupId>

2
common/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>thingsboard</artifactId>
</parent>
<artifactId>common</artifactId>

2
common/proto/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>common</artifactId>
</parent>
<groupId>org.thingsboard.common</groupId>

2
common/queue/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>common</artifactId>
</parent>
<groupId>org.thingsboard.common</groupId>

2
common/script/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>common</artifactId>
</parent>
<groupId>org.thingsboard.common</groupId>

2
common/script/remote-js-client/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard.common</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>script</artifactId>
</parent>
<groupId>org.thingsboard.common.script</groupId>

2
common/script/script-api/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard.common</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>script</artifactId>
</parent>
<groupId>org.thingsboard.common.script</groupId>

2
common/stats/pom.xml

@ -22,7 +22,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>common</artifactId>
</parent>
<groupId>org.thingsboard.common</groupId>

2
common/transport/coap/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard.common</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>transport</artifactId>
</parent>
<groupId>org.thingsboard.common.transport</groupId>

2
common/transport/http/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard.common</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>transport</artifactId>
</parent>
<groupId>org.thingsboard.common.transport</groupId>

2
common/transport/lwm2m/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard.common</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>transport</artifactId>
</parent>
<groupId>org.thingsboard.common.transport</groupId>

2
common/transport/mqtt/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard.common</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>transport</artifactId>
</parent>
<groupId>org.thingsboard.common.transport</groupId>

2
common/transport/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>common</artifactId>
</parent>
<groupId>org.thingsboard.common</groupId>

2
common/transport/snmp/pom.xml

@ -21,7 +21,7 @@
<parent>
<groupId>org.thingsboard.common</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>transport</artifactId>
</parent>

2
common/transport/transport-api/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard.common</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>transport</artifactId>
</parent>
<groupId>org.thingsboard.common.transport</groupId>

2
common/util/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>common</artifactId>
</parent>
<groupId>org.thingsboard.common</groupId>

2
common/version-control/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>common</artifactId>
</parent>
<groupId>org.thingsboard.common</groupId>

2
dao/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>thingsboard</artifactId>
</parent>
<artifactId>dao</artifactId>

2
edqs/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>thingsboard</artifactId>
</parent>
<artifactId>edqs</artifactId>

2
monitoring/pom.xml

@ -21,7 +21,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>thingsboard</artifactId>
</parent>

2
msa/black-box-tests/pom.xml

@ -21,7 +21,7 @@
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>msa</artifactId>
</parent>
<groupId>org.thingsboard.msa</groupId>

2
msa/edqs/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>msa</artifactId>
</parent>
<groupId>org.thingsboard.msa</groupId>

2
msa/js-executor/package.json

@ -1,7 +1,7 @@
{
"name": "thingsboard-js-executor",
"private": true,
"version": "4.2.2.4",
"version": "4.2.2.5",
"description": "ThingsBoard JavaScript Executor Microservice",
"main": "server.ts",
"bin": "server.js",

2
msa/js-executor/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>msa</artifactId>
</parent>
<groupId>org.thingsboard.msa</groupId>

2
msa/monitoring/pom.xml

@ -22,7 +22,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>msa</artifactId>
</parent>

2
msa/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>thingsboard</artifactId>
</parent>
<artifactId>msa</artifactId>

2
msa/tb-node/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>msa</artifactId>
</parent>
<groupId>org.thingsboard.msa</groupId>

2
msa/tb/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>msa</artifactId>
</parent>
<groupId>org.thingsboard.msa</groupId>

2
msa/transport/coap/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard.msa</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>transport</artifactId>
</parent>
<groupId>org.thingsboard.msa.transport</groupId>

2
msa/transport/http/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard.msa</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>transport</artifactId>
</parent>
<groupId>org.thingsboard.msa.transport</groupId>

2
msa/transport/lwm2m/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard.msa</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>transport</artifactId>
</parent>
<groupId>org.thingsboard.msa.transport</groupId>

2
msa/transport/mqtt/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard.msa</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>transport</artifactId>
</parent>
<groupId>org.thingsboard.msa.transport</groupId>

2
msa/transport/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>msa</artifactId>
</parent>
<groupId>org.thingsboard.msa</groupId>

2
msa/transport/snmp/pom.xml

@ -21,7 +21,7 @@
<parent>
<groupId>org.thingsboard.msa</groupId>
<artifactId>transport</artifactId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
</parent>
<groupId>org.thingsboard.msa.transport</groupId>

2
msa/vc-executor-docker/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>msa</artifactId>
</parent>
<groupId>org.thingsboard.msa</groupId>

2
msa/vc-executor/pom.xml

@ -21,7 +21,7 @@
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>msa</artifactId>
</parent>
<groupId>org.thingsboard.msa</groupId>

2
msa/web-ui/package.json

@ -1,7 +1,7 @@
{
"name": "thingsboard-web-ui",
"private": true,
"version": "4.2.2.4",
"version": "4.2.2.5",
"description": "ThingsBoard Web UI Microservice",
"main": "server.ts",
"bin": "server.js",

2
msa/web-ui/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>msa</artifactId>
</parent>
<groupId>org.thingsboard.msa</groupId>

4
netty-mqtt/pom.xml

@ -19,11 +19,11 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>thingsboard</artifactId>
</parent>
<artifactId>netty-mqtt</artifactId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<packaging>jar</packaging>
<name>Netty MQTT Client</name>

2
pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<groupId>org.thingsboard</groupId>
<artifactId>thingsboard</artifactId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<packaging>pom</packaging>
<name>Thingsboard</name>

2
rest-client/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>thingsboard</artifactId>
</parent>
<artifactId>rest-client</artifactId>

2
rule-engine/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>thingsboard</artifactId>
</parent>
<artifactId>rule-engine</artifactId>

2
rule-engine/rule-engine-api/pom.xml

@ -22,7 +22,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>rule-engine</artifactId>
</parent>
<groupId>org.thingsboard.rule-engine</groupId>

2
rule-engine/rule-engine-components/pom.xml

@ -22,7 +22,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>rule-engine</artifactId>
</parent>
<groupId>org.thingsboard.rule-engine</groupId>

2
tools/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>thingsboard</artifactId>
</parent>
<artifactId>tools</artifactId>

74
tools/src/main/java/org/thingsboard/client/tools/migrator/DictionaryParser.java

@ -1,74 +0,0 @@
/**
* 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.client.tools.migrator;
import org.apache.commons.io.FileUtils;
import org.apache.commons.io.LineIterator;
import org.thingsboard.server.common.data.StringUtils;
import java.io.File;
import java.io.IOException;
import java.util.HashMap;
import java.util.Map;
public class DictionaryParser {
private Map<String, String> dictionaryParsed = new HashMap<>();
public DictionaryParser(File sourceFile) throws IOException {
parseDictionaryDump(FileUtils.lineIterator(sourceFile));
}
public String getKeyByKeyId(String keyId) {
return dictionaryParsed.get(keyId);
}
private boolean isBlockFinished(String line) {
return StringUtils.isBlank(line) || line.equals("\\.");
}
private boolean isBlockStarted(String line) {
return line.startsWith("COPY public.key_dictionary (");
}
private void parseDictionaryDump(LineIterator iterator) throws IOException {
try {
String tempLine;
while (iterator.hasNext()) {
tempLine = iterator.nextLine();
if (isBlockStarted(tempLine)) {
processBlock(iterator);
}
}
} finally {
iterator.close();
}
}
private void processBlock(LineIterator lineIterator) {
String tempLine;
String[] lineSplited;
while(lineIterator.hasNext()) {
tempLine = lineIterator.nextLine();
if(isBlockFinished(tempLine)) {
return;
}
lineSplited = tempLine.split("\t");
dictionaryParsed.put(lineSplited[1], lineSplited[0]);
}
}
}

102
tools/src/main/java/org/thingsboard/client/tools/migrator/MigratorTool.java

@ -1,102 +0,0 @@
/**
* 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.client.tools.migrator;
import org.apache.commons.cli.BasicParser;
import org.apache.commons.cli.CommandLine;
import org.apache.commons.cli.CommandLineParser;
import org.apache.commons.cli.HelpFormatter;
import org.apache.commons.cli.Option;
import org.apache.commons.cli.Options;
import org.apache.commons.cli.ParseException;
import java.io.File;
public class MigratorTool {
public static void main(String[] args) {
CommandLine cmd = parseArgs(args);
try {
boolean castEnable = Boolean.parseBoolean(cmd.getOptionValue("castEnable"));
File allTelemetrySource = new File(cmd.getOptionValue("telemetryFrom"));
File tsSaveDir = null;
File partitionsSaveDir = null;
File latestSaveDir = null;
RelatedEntitiesParser allEntityIdsAndTypes =
new RelatedEntitiesParser(new File(cmd.getOptionValue("relatedEntities")));
DictionaryParser dictionaryParser = new DictionaryParser(allTelemetrySource);
if(cmd.getOptionValue("latestTelemetryOut") != null) {
latestSaveDir = new File(cmd.getOptionValue("latestTelemetryOut"));
}
if(cmd.getOptionValue("telemetryOut") != null) {
tsSaveDir = new File(cmd.getOptionValue("telemetryOut"));
partitionsSaveDir = new File(cmd.getOptionValue("partitionsOut"));
}
new PgCaMigrator(allTelemetrySource, tsSaveDir, partitionsSaveDir, latestSaveDir, allEntityIdsAndTypes, dictionaryParser, castEnable).migrate();
} catch (Throwable th) {
th.printStackTrace();
throw new IllegalStateException("failed", th);
}
}
private static CommandLine parseArgs(String[] args) {
Options options = new Options();
Option telemetryAllFrom = new Option("telemetryFrom", "telemetryFrom", true, "telemetry source file");
telemetryAllFrom.setRequired(true);
options.addOption(telemetryAllFrom);
Option latestTsOutOpt = new Option("latestOut", "latestTelemetryOut", true, "latest telemetry save dir");
latestTsOutOpt.setRequired(false);
options.addOption(latestTsOutOpt);
Option tsOutOpt = new Option("tsOut", "telemetryOut", true, "sstable save dir");
tsOutOpt.setRequired(false);
options.addOption(tsOutOpt);
Option partitionOutOpt = new Option("partitionsOut", "partitionsOut", true, "partitions save dir");
partitionOutOpt.setRequired(false);
options.addOption(partitionOutOpt);
Option castOpt = new Option("castEnable", "castEnable", true, "cast String to Double if possible");
castOpt.setRequired(true);
options.addOption(castOpt);
Option relatedOpt = new Option("relatedEntities", "relatedEntities", true, "related entities source file path");
relatedOpt.setRequired(true);
options.addOption(relatedOpt);
HelpFormatter formatter = new HelpFormatter();
CommandLineParser parser = new BasicParser();
try {
return parser.parse(options, args);
} catch (ParseException e) {
System.out.println(e.getMessage());
formatter.printHelp("utility-name", options);
System.exit(1);
}
return null;
}
}

282
tools/src/main/java/org/thingsboard/client/tools/migrator/PgCaMigrator.java

@ -1,282 +0,0 @@
/**
* 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.client.tools.migrator;
import com.google.common.collect.Lists;
import org.apache.cassandra.io.sstable.CQLSSTableWriter;
import org.apache.commons.io.FileUtils;
import org.apache.commons.io.LineIterator;
import org.apache.commons.lang3.math.NumberUtils;
import org.thingsboard.server.common.data.StringUtils;
import java.io.File;
import java.io.IOException;
import java.time.Instant;
import java.time.LocalDateTime;
import java.time.ZoneOffset;
import java.time.temporal.ChronoUnit;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Date;
import java.util.HashSet;
import java.util.List;
import java.util.Set;
import java.util.UUID;
import java.util.function.Function;
import java.util.stream.Collectors;
public class PgCaMigrator {
private final long LOG_BATCH = 1000000;
private final long rowPerFile = 1000000;
private long linesTsMigrated = 0;
private long linesLatestMigrated = 0;
private long castErrors = 0;
private long castedOk = 0;
private long currentWriterCount = 1;
private final File sourceFile;
private final boolean castStringIfPossible;
private final RelatedEntitiesParser entityIdsAndTypes;
private final DictionaryParser keyParser;
private CQLSSTableWriter currentTsWriter;
private CQLSSTableWriter currentPartitionsWriter;
private CQLSSTableWriter currentTsLatestWriter;
private final Set<String> partitions = new HashSet<>();
private File outTsDir;
private File outTsLatestDir;
public PgCaMigrator(File sourceFile,
File ourTsDir,
File outTsPartitionDir,
File outTsLatestDir,
RelatedEntitiesParser allEntityIdsAndTypes,
DictionaryParser dictionaryParser,
boolean castStringsIfPossible) {
this.sourceFile = sourceFile;
this.entityIdsAndTypes = allEntityIdsAndTypes;
this.keyParser = dictionaryParser;
this.castStringIfPossible = castStringsIfPossible;
if(outTsLatestDir != null) {
this.currentTsLatestWriter = WriterBuilder.getLatestWriter(outTsLatestDir);
this.outTsLatestDir = outTsLatestDir;
}
if(ourTsDir != null) {
this.currentTsWriter = WriterBuilder.getTsWriter(ourTsDir);
this.currentPartitionsWriter = WriterBuilder.getPartitionWriter(outTsPartitionDir);
this.outTsDir = ourTsDir;
}
}
public void migrate() throws IOException {
boolean isTsDone = false;
boolean isLatestDone = false;
String line;
LineIterator iterator = FileUtils.lineIterator(this.sourceFile);
try {
while(iterator.hasNext()) {
line = iterator.nextLine();
if(!isLatestDone && isBlockLatestStarted(line)) {
System.out.println("START TO MIGRATE LATEST");
long start = System.currentTimeMillis();
processBlock(iterator, currentTsLatestWriter, outTsLatestDir, this::toValuesLatest);
System.out.println("TOTAL LINES MIGRATED: " + linesLatestMigrated + ", FORMING OF SSL FOR LATEST TS FINISHED WITH TIME: " + (System.currentTimeMillis() - start) + " ms.");
isLatestDone = true;
}
if(!isTsDone && isBlockTsStarted(line)) {
System.out.println("START TO MIGRATE TS");
long start = System.currentTimeMillis();
processBlock(iterator, currentTsWriter, outTsDir, this::toValuesTs);
System.out.println("TOTAL LINES MIGRATED: " + linesTsMigrated + ", FORMING OF SSL FOR TS FINISHED WITH TIME: " + (System.currentTimeMillis() - start) + " ms.");
isTsDone = true;
}
}
System.out.println("Partitions collected " + partitions.size());
long startTs = System.currentTimeMillis();
for (String partition : partitions) {
String[] split = partition.split("\\|");
List<Object> values = Lists.newArrayList();
values.add(split[0]);
values.add(UUID.fromString(split[1]));
values.add(split[2]);
values.add(Long.parseLong(split[3]));
currentPartitionsWriter.addRow(values);
}
System.out.println(new Date() + " Migrated partitions " + partitions.size() + " in " + (System.currentTimeMillis() - startTs));
System.out.println();
System.out.println("Finished migrate Telemetry");
} finally {
iterator.close();
currentTsLatestWriter.close();
currentTsWriter.close();
currentPartitionsWriter.close();
}
}
private void logLinesProcessed(long lines) {
if (lines % LOG_BATCH == 0) {
System.out.println(new Date() + " lines processed = " + lines + " in, castOk " + castedOk + " castErr " + castErrors);
}
}
private void logLinesMigrated(long lines) {
if(lines % LOG_BATCH == 0) {
System.out.println(new Date() + " lines migrated = " + lines + " in, castOk " + castedOk + " castErr " + castErrors);
}
}
private void addTypeIdKey(List<Object> result, List<String> raw) {
result.add(entityIdsAndTypes.getEntityType(raw.get(0)));
result.add(UUID.fromString(raw.get(0)));
result.add(keyParser.getKeyByKeyId(raw.get(1)));
}
private void addPartitions(List<Object> result, List<String> raw) {
long ts = Long.parseLong(raw.get(2));
long partition = toPartitionTs(ts);
result.add(partition);
result.add(ts);
}
private void addTimeseries(List<Object> result, List<String> raw) {
result.add(Long.parseLong(raw.get(2)));
}
private void addValues(List<Object> result, List<String> raw) {
result.add(raw.get(3).equals("\\N") ? null : raw.get(3).equals("t") ? Boolean.TRUE : Boolean.FALSE);
result.add(raw.get(4).equals("\\N") ? null : raw.get(4));
result.add(raw.get(5).equals("\\N") ? null : Long.parseLong(raw.get(5)));
result.add(raw.get(6).equals("\\N") ? null : Double.parseDouble(raw.get(6)));
result.add(raw.get(7).equals("\\N") ? null : raw.get(7));
}
private List<Object> toValuesTs(List<String> raw) {
logLinesMigrated(linesTsMigrated++);
List<Object> result = new ArrayList<>();
addTypeIdKey(result, raw);
addPartitions(result, raw);
addValues(result, raw);
processPartitions(result);
return result;
}
private List<Object> toValuesLatest(List<String> raw) {
logLinesMigrated(linesLatestMigrated++);
List<Object> result = new ArrayList<>();
addTypeIdKey(result, raw);
addTimeseries(result, raw);
addValues(result, raw);
return result;
}
private long toPartitionTs(long ts) {
LocalDateTime time = LocalDateTime.ofInstant(Instant.ofEpochMilli(ts), ZoneOffset.UTC);
return time.truncatedTo(ChronoUnit.DAYS).withDayOfMonth(1).toInstant(ZoneOffset.UTC).toEpochMilli();
}
private void processPartitions(List<Object> values) {
String key = values.get(0) + "|" + values.get(1) + "|" + values.get(2) + "|" + values.get(3);
partitions.add(key);
}
private void processBlock(LineIterator iterator, CQLSSTableWriter writer, File outDir, Function<List<String>, List<Object>> function) {
String currentLine;
long linesProcessed = 0;
while(iterator.hasNext()) {
logLinesProcessed(linesProcessed++);
currentLine = iterator.nextLine();
if(isBlockFinished(currentLine)) {
return;
}
try {
List<String> raw = Arrays.stream(currentLine.trim().split("\t"))
.map(String::trim)
.collect(Collectors.toList());
List<Object> values = function.apply(raw);
if (this.currentWriterCount == 0) {
System.out.println(new Date() + " close writer " + new Date());
writer.close();
writer = WriterBuilder.getLatestWriter(outDir);
}
if (this.castStringIfPossible) {
writer.addRow(castToNumericIfPossible(values));
} else {
writer.addRow(values);
}
currentWriterCount++;
if (currentWriterCount >= rowPerFile) {
currentWriterCount = 0;
}
} catch (Exception ex) {
System.out.println(ex.getMessage() + " -> " + currentLine);
}
}
}
private List<Object> castToNumericIfPossible(List<Object> values) {
try {
if (values.get(6) != null && NumberUtils.isCreatable(values.get(6).toString())) {
Double casted = NumberUtils.createDouble(values.get(6).toString());
List<Object> numeric = Lists.newArrayList();
numeric.addAll(values);
numeric.set(6, null);
numeric.set(8, casted);
castedOk++;
return numeric;
}
} catch (Throwable th) {
castErrors++;
}
processPartitions(values);
return values;
}
private boolean isBlockFinished(String line) {
return StringUtils.isBlank(line) || line.equals("\\.");
}
private boolean isBlockTsStarted(String line) {
return line.startsWith("COPY public.ts_kv (");
}
private boolean isBlockLatestStarted(String line) {
return line.startsWith("COPY public.ts_kv_latest (");
}
}

92
tools/src/main/java/org/thingsboard/client/tools/migrator/README.md

@ -1,92 +0,0 @@
# Description:
This tool used for migrating ThingsBoard into hybrid mode from Postgres.
Performance of this tool depends on disk type and instance type (mostly on CPU resources).
But in general here are few benchmarks:
1. Creating Dump of the postgres ts_kv table -> 100GB = 90 minutes
2. If postgres table has size 100GB then dump file will be about 30GB
3. Generation SSTables from dump -> 100GB = 3 hours
4. 100GB Dump file will be converted into SSTable with size about 18GB
# Tool build Instruction:
Switch to `tools` module in Command Line and execute
mvn clean compile assembly:single
It will generate single jar file with all required dependencies inside `target dir` -> `tools-2.4.1-SNAPSHOT-jar-with-dependencies.jar`.
# Prepare requred files and run Tool:
#### Dump data from the source Postgres Database
*Do not use compression if possible because Tool can only work with uncompressed file
1. Dump related tables that need to correct save telemetry
`pg_dump -h localhost -U postgres -d thingsboard -T admin_settings -T attribute_kv -T audit_log -T component_discriptor -T device_credentials -T event -T oauth2_client_registration -T oauth2_client_registration_info -T oauth2_client_registration_template -T relation -T rule_node_state tb_schema_settings -T user_credentials > related_entities.dmp`
2. Dump `ts_kv` and child:
`pg_dump -h localhost -U postgres -d thingsboard --load-via-partition-root --data-only -t ts_kv* > ts_kv_all.dmp`
3. [Optional] Move table dumps to the instance where cassandra will be hosted
#### Prepare directory structure for SSTables
Tool use 3 different directories for saving SSTables - `ts_kv_cf`, `ts_kv_latest_cf`, `ts_kv_partitions_cf`
Create 3 empty directories. For example:
/home/user/migration/ts
/home/user/migration/ts_latest
/home/user/migration/ts_partition
#### Run tool
**If you want to migrate just `ts_kv` without `ts_kv_latest` or vice versa don't use arguments (paths) for output files*
**Note: if you run this tool on remote instance - don't forget to execute this command in `screen` to avoid unexpected termination*
```
java -jar ./tools-3.2.2-SNAPSHOT-jar-with-dependencies.jar
-telemetryFrom /home/user/dump/ts_kv_all.dmp
-relatedEntities /home/user/dump/related_entities.dmp
-latestOut /home/user/migration/ts_latest
-tsOut /home/user/migration/ts
-partitionsOut /home/user/migration/ts_partition
-castEnable false
```
*Use your paths for program arguments*
Tool execution time depends on DB size, CPU resources and Disk throughput
## Adding SSTables into Cassandra
* Note that this this part works only for single node Cassandra Cluster. If you have more nodes - it is better to use `sstableloader` tool.
1. [Optional] install Cassandra on the instance
2. [Optional] Using `cqlsh` create `thingsboard` keyspace and requred tables from this files `schema-keyspace.cql`, `schema-ts.cql` and `schema-ts-latest.cql` using `source` command
3. Stop Cassandra
4. Look at `/var/lib/cassandra/data/thingsboard` and check for names of data folders
5. Copy generated SSTable files into cassandra data dir using next command:
```
sudo find /home/user/migration/ts -name '*.*' -exec mv {} /var/lib/cassandra/data/thingsboard/ts_kv_cf-0e9aaf00ee5511e9a5fa7d6f489ffd13/ \;
sudo find /home/user/migration/ts_latest -name '*.*' -exec mv {} /var/lib/cassandra/data/thingsboard/ts_kv_latest_cf-161449d0ee5511e9a5fa7d6f489ffd13/ \;
sudo find /home/user/migration/ts_partition -name '*.*' -exec mv {} /var/lib/cassandra/data/thingsboard/ts_kv_partitions_cf-12e8fa80ee5511e9a5fa7d6f489ffd13/ \;
```
*Pay attention! Data folders have similar name `ts_kv_cf-0e9aaf00ee5511e9a5fa7d6f489ffd13`, but you have to use own*
6. Start Cassandra service and trigger compaction
Trigger compactions: `nodetool compact thingsboard`
Check compaction status: `nodetool compactionstats`
## Switch Thignsboard into Hybrid Mode
Modify Thingsboard properites file `thingsboard.yml`
- DATABASE_TS_TYPE = cassandra
- TS_KV_PARTITIONING = MONTHS
# Final steps
Start Thingsboard and verify migration

88
tools/src/main/java/org/thingsboard/client/tools/migrator/RelatedEntitiesParser.java

@ -1,88 +0,0 @@
/**
* 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.client.tools.migrator;
import org.apache.commons.io.FileUtils;
import org.apache.commons.io.LineIterator;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.StringUtils;
import java.io.File;
import java.io.IOException;
import java.util.HashMap;
import java.util.Map;
public class RelatedEntitiesParser {
private final Map<String, String> allEntityIdsAndTypes = new HashMap<>();
private final Map<String, EntityType> tableNameAndEntityType = Map.ofEntries(
Map.entry("COPY public.alarm ", EntityType.ALARM),
Map.entry("COPY public.asset ", EntityType.ASSET),
Map.entry("COPY public.customer ", EntityType.CUSTOMER),
Map.entry("COPY public.dashboard ", EntityType.DASHBOARD),
Map.entry("COPY public.device ", EntityType.DEVICE),
Map.entry("COPY public.rule_chain ", EntityType.RULE_CHAIN),
Map.entry("COPY public.rule_node ", EntityType.RULE_NODE),
Map.entry("COPY public.tenant ", EntityType.TENANT),
Map.entry("COPY public.tb_user ", EntityType.USER),
Map.entry("COPY public.entity_view ", EntityType.ENTITY_VIEW),
Map.entry("COPY public.widgets_bundle ", EntityType.WIDGETS_BUNDLE),
Map.entry("COPY public.widget_type ", EntityType.WIDGET_TYPE),
Map.entry("COPY public.tenant_profile ", EntityType.TENANT_PROFILE),
Map.entry("COPY public.device_profile ", EntityType.DEVICE_PROFILE),
Map.entry("COPY public.asset_profile ", EntityType.ASSET_PROFILE),
Map.entry("COPY public.api_usage_state ", EntityType.API_USAGE_STATE)
);
public RelatedEntitiesParser(File source) throws IOException {
processAllTables(FileUtils.lineIterator(source));
}
public String getEntityType(String uuid) {
return this.allEntityIdsAndTypes.get(uuid);
}
private boolean isBlockFinished(String line) {
return StringUtils.isBlank(line) || line.equals("\\.");
}
private void processAllTables(LineIterator lineIterator) throws IOException {
String currentLine;
try {
while (lineIterator.hasNext()) {
currentLine = lineIterator.nextLine();
for(Map.Entry<String, EntityType> entry : tableNameAndEntityType.entrySet()) {
if(currentLine.startsWith(entry.getKey())) {
processBlock(lineIterator, entry.getValue());
}
}
}
} finally {
lineIterator.close();
}
}
private void processBlock(LineIterator lineIterator, EntityType entityType) {
String currentLine;
while(lineIterator.hasNext()) {
currentLine = lineIterator.nextLine();
if(isBlockFinished(currentLine)) {
return;
}
allEntityIdsAndTypes.put(currentLine.split("\t")[0], entityType.name());
}
}
}

86
tools/src/main/java/org/thingsboard/client/tools/migrator/WriterBuilder.java

@ -1,86 +0,0 @@
/**
* 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.client.tools.migrator;
import org.apache.cassandra.io.sstable.CQLSSTableWriter;
import java.io.File;
public class WriterBuilder {
private static final String tsSchema = "CREATE TABLE thingsboard.ts_kv_cf (\n" +
" entity_type text, // (DEVICE, CUSTOMER, TENANT)\n" +
" entity_id timeuuid,\n" +
" key text,\n" +
" partition bigint,\n" +
" ts bigint,\n" +
" bool_v boolean,\n" +
" str_v text,\n" +
" long_v bigint,\n" +
" dbl_v double,\n" +
" json_v text,\n" +
" PRIMARY KEY (( entity_type, entity_id, key, partition ), ts)\n" +
");";
private static final String latestSchema = "CREATE TABLE IF NOT EXISTS thingsboard.ts_kv_latest_cf (\n" +
" entity_type text, // (DEVICE, CUSTOMER, TENANT)\n" +
" entity_id timeuuid,\n" +
" key text,\n" +
" ts bigint,\n" +
" bool_v boolean,\n" +
" str_v text,\n" +
" long_v bigint,\n" +
" dbl_v double,\n" +
" json_v text,\n" +
" PRIMARY KEY (( entity_type, entity_id ), key)\n" +
") WITH compaction = { 'class' : 'LeveledCompactionStrategy' };";
private static final String partitionSchema = "CREATE TABLE IF NOT EXISTS thingsboard.ts_kv_partitions_cf (\n" +
" entity_type text, // (DEVICE, CUSTOMER, TENANT)\n" +
" entity_id timeuuid,\n" +
" key text,\n" +
" partition bigint,\n" +
" PRIMARY KEY (( entity_type, entity_id, key ), partition)\n" +
") WITH CLUSTERING ORDER BY ( partition ASC )\n" +
" AND compaction = { 'class' : 'LeveledCompactionStrategy' };";
public static CQLSSTableWriter getTsWriter(File dir) {
return CQLSSTableWriter.builder()
.inDirectory(dir.getAbsolutePath())
.forTable(tsSchema)
.using("INSERT INTO thingsboard.ts_kv_cf (entity_type, entity_id, key, partition, ts, bool_v, str_v, long_v, dbl_v, json_v) " +
"VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)")
.build();
}
public static CQLSSTableWriter getLatestWriter(File dir) {
return CQLSSTableWriter.builder()
.inDirectory(dir.getAbsolutePath())
.forTable(latestSchema)
.using("INSERT INTO thingsboard.ts_kv_latest_cf (entity_type, entity_id, key, ts, bool_v, str_v, long_v, dbl_v, json_v) " +
"VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)")
.build();
}
public static CQLSSTableWriter getPartitionWriter(File dir) {
return CQLSSTableWriter.builder()
.inDirectory(dir.getAbsolutePath())
.forTable(partitionSchema)
.using("INSERT INTO thingsboard.ts_kv_partitions_cf (entity_type, entity_id, key, partition) " +
"VALUES (?, ?, ?, ?)")
.build();
}
}

2
transport/coap/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>transport</artifactId>
</parent>
<groupId>org.thingsboard.transport</groupId>

2
transport/http/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>transport</artifactId>
</parent>
<groupId>org.thingsboard.transport</groupId>

2
transport/lwm2m/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>transport</artifactId>
</parent>
<groupId>org.thingsboard.transport</groupId>

2
transport/mqtt/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>transport</artifactId>
</parent>
<groupId>org.thingsboard.transport</groupId>

2
transport/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>thingsboard</artifactId>
</parent>
<artifactId>transport</artifactId>

2
transport/snmp/pom.xml

@ -21,7 +21,7 @@
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>transport</artifactId>
</parent>

2
ui-ngx/package.json

@ -1,6 +1,6 @@
{
"name": "thingsboard",
"version": "4.2.2.4",
"version": "4.2.2.5",
"scripts": {
"ng": "ng",
"start": "node --max_old_space_size=8048 ./node_modules/@angular/cli/bin/ng serve --configuration development --host 0.0.0.0 --open",

2
ui-ngx/pom.xml

@ -20,7 +20,7 @@
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.thingsboard</groupId>
<version>4.2.2.4</version>
<version>4.2.2.5-SNAPSHOT</version>
<artifactId>thingsboard</artifactId>
</parent>
<groupId>org.thingsboard</groupId>

Loading…
Cancel
Save