diff --git a/application/pom.xml b/application/pom.xml index ec35ca3d77..9c6d05d31f 100644 --- a/application/pom.xml +++ b/application/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT thingsboard application diff --git a/application/src/main/java/org/thingsboard/server/controller/AdminController.java b/application/src/main/java/org/thingsboard/server/controller/AdminController.java index 0b04a6939f..8e8f4c4de3 100644 --- a/application/src/main/java/org/thingsboard/server/controller/AdminController.java +++ b/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 prevUrlOpt = CookieUtils.getCookie(request, PREV_URI_COOKIE_NAME); + String redirectUrl = getMailOAuth2RedirectUrl(request); Optional 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; } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java index 1e9a0ade87..aeab041c7f 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java +++ b/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 { diff --git a/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/AbstractOAuth2ClientMapper.java b/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/AbstractOAuth2ClientMapper.java index 9ac31a8813..4559c53862 100644 --- a/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/AbstractOAuth2ClientMapper.java +++ b/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) { diff --git a/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/CallbackUrlSchemeValidator.java b/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/CallbackUrlSchemeValidator.java new file mode 100644 index 0000000000..bdb4a6d4ab --- /dev/null +++ b/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 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]", "?"); + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/HttpCookieOAuth2AuthorizationRequestRepository.java b/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/HttpCookieOAuth2AuthorizationRequestRepository.java index b908a6c650..2b40e548e6 100644 --- a/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/HttpCookieOAuth2AuthorizationRequestRepository.java +++ b/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); } diff --git a/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationFailureHandler.java b/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationFailureHandler.java index 3b9d325c38..421be12c00 100644 --- a/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationFailureHandler.java +++ b/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="; diff --git a/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationSuccessHandler.java b/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationSuccessHandler.java index c22ceb944a..c99fb9378e 100644 --- a/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationSuccessHandler.java +++ b/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 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 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 += "/?"; } diff --git a/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/PrevUriValidator.java b/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/PrevUriValidator.java new file mode 100644 index 0000000000..14855ecda3 --- /dev/null +++ b/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]", "?"); + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/security/model/token/OAuth2AppTokenFactory.java b/application/src/main/java/org/thingsboard/server/service/security/model/token/OAuth2AppTokenFactory.java index a353aac86b..9e6b525364 100644 --- a/application/src/main/java/org/thingsboard/server/service/security/model/token/OAuth2AppTokenFactory.java +++ b/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; } diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/AbstractSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/AbstractSubscriptionService.java index 5df3e638f5..7b5a2a0fad 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/AbstractSubscriptionService.java +++ b/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 toSubscriptionManagerService, Supplier 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()); diff --git a/application/src/test/java/org/thingsboard/server/controller/AdminControllerMailOAuth2Test.java b/application/src/test/java/org/thingsboard/server/controller/AdminControllerMailOAuth2Test.java new file mode 100644 index 0000000000..b74acea73a --- /dev/null +++ b/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)}); + } + +} diff --git a/application/src/test/java/org/thingsboard/server/controller/AdminControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/AdminControllerTest.java index 21a8b3b2e9..bbdd324d82 100644 --- a/application/src/test/java/org/thingsboard/server/controller/AdminControllerTest.java +++ b/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()); diff --git a/application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/CallbackUrlSchemeValidatorTest.java b/application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/CallbackUrlSchemeValidatorTest.java new file mode 100644 index 0000000000..aa2eeb7982 --- /dev/null +++ b/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(); + } + +} diff --git a/application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/HttpCookieOAuth2AuthorizationRequestRepositoryTest.java b/application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/HttpCookieOAuth2AuthorizationRequestRepositoryTest.java new file mode 100644 index 0000000000..5cada8477f --- /dev/null +++ b/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 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); + } + +} diff --git a/application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/OAuth2ClientMapperTest.java b/application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/OAuth2ClientMapperTest.java new file mode 100644 index 0000000000..775ed5945e --- /dev/null +++ b/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(); + } + +} diff --git a/application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationFailureHandlerTest.java b/application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationFailureHandlerTest.java new file mode 100644 index 0000000000..15110aa46d --- /dev/null +++ b/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 redirectCaptor = ArgumentCaptor.forClass(String.class); + verify(response).sendRedirect(redirectCaptor.capture()); + return redirectCaptor.getValue(); + } + +} diff --git a/application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationSuccessHandlerTest.java b/application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationSuccessHandlerTest.java index 7e4c2645c7..eae9441711 100644 --- a/application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationSuccessHandlerTest.java +++ b/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 cookieCaptor = ArgumentCaptor.forClass(Cookie.class); + verify(response).addCookie(cookieCaptor.capture()); + assertThat(cookieCaptor.getValue().getName()).isEqualTo(PREV_URI_COOKIE_NAME); + assertThat(cookieCaptor.getValue().getMaxAge()).isZero(); } -} \ No newline at end of file + + @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 redirectCaptor = ArgumentCaptor.forClass(String.class); + verify(response).sendRedirect(redirectCaptor.capture()); + return redirectCaptor.getValue(); + } + +} diff --git a/application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/PrevUriValidatorTest.java b/application/src/test/java/org/thingsboard/server/service/security/auth/oauth2/PrevUriValidatorTest.java new file mode 100644 index 0000000000..487abffc24 --- /dev/null +++ b/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(); + } + +} diff --git a/application/src/test/java/org/thingsboard/server/service/security/model/token/OAuth2AppTokenFactoryTest.java b/application/src/test/java/org/thingsboard/server/service/security/model/token/OAuth2AppTokenFactoryTest.java new file mode 100644 index 0000000000..0130ddfff7 --- /dev/null +++ b/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)); + } + +} diff --git a/application/src/test/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionServiceTest.java b/application/src/test/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionServiceTest.java index dcfe9bd2bd..958a41449f 100644 --- a/application/src/test/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionServiceTest.java +++ b/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 listOfNNumbers(int N) { return LongStream.range(0, N).boxed().toList(); diff --git a/common/actor/pom.xml b/common/actor/pom.xml index 91e2e88f51..05dbdae924 100644 --- a/common/actor/pom.xml +++ b/common/actor/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT common org.thingsboard.common diff --git a/common/cache/pom.xml b/common/cache/pom.xml index 428acd9214..3413e7c22f 100644 --- a/common/cache/pom.xml +++ b/common/cache/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT common org.thingsboard.common diff --git a/common/cluster-api/pom.xml b/common/cluster-api/pom.xml index 230ba4fa1e..29736b874e 100644 --- a/common/cluster-api/pom.xml +++ b/common/cluster-api/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT common org.thingsboard.common diff --git a/common/coap-server/pom.xml b/common/coap-server/pom.xml index 6a8741fc89..f889c6bcaa 100644 --- a/common/coap-server/pom.xml +++ b/common/coap-server/pom.xml @@ -22,7 +22,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT common org.thingsboard.common diff --git a/common/dao-api/pom.xml b/common/dao-api/pom.xml index e3bfe0faa9..987d753d2c 100644 --- a/common/dao-api/pom.xml +++ b/common/dao-api/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT common org.thingsboard.common diff --git a/common/data/pom.xml b/common/data/pom.xml index 87ff452c82..527bc39881 100644 --- a/common/data/pom.xml +++ b/common/data/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT common org.thingsboard.common diff --git a/common/discovery-api/pom.xml b/common/discovery-api/pom.xml index 57664ffc3e..df8f86335f 100644 --- a/common/discovery-api/pom.xml +++ b/common/discovery-api/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT common org.thingsboard.common diff --git a/common/edge-api/pom.xml b/common/edge-api/pom.xml index 396e06c68b..2cbf7a2672 100644 --- a/common/edge-api/pom.xml +++ b/common/edge-api/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT common org.thingsboard.common diff --git a/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java b/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java index 7e58a2cb9f..2237e06f43 100644 --- a/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java +++ b/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java @@ -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 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 onEdgeUpdate, Consumer onDownlink, Consumer 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 channelType() { + return Epoll.isAvailable() ? EpollSocketChannel.class : NioSocketChannel.class; + } + + @PreDestroy + public void destroy() { + workerGroup.shutdownGracefully(); + } + private StreamObserver initOutputStream(String edgeKey, Consumer onUplinkResponse, Consumer 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) diff --git a/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java b/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java index 267fd2b3f2..d9a05ab2ce 100644 --- a/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java +++ b/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); diff --git a/common/edge-api/src/main/proto/edge.proto b/common/edge-api/src/main/proto/edge.proto index 376a4d8c5b..8650087f74 100644 --- a/common/edge-api/src/main/proto/edge.proto +++ b/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; } diff --git a/common/edge-api/src/test/java/org/thingsboard/edge/rpc/EdgeGrpcClientLeakTest.java b/common/edge-api/src/test/java/org/thingsboard/edge/rpc/EdgeGrpcClientLeakTest.java new file mode 100644 index 0000000000..9fc66ed33e --- /dev/null +++ b/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 handleMsgs(StreamObserver 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 inputStream = (StreamObserver) 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); + } + } + +} diff --git a/common/edqs/pom.xml b/common/edqs/pom.xml index f787a3be4e..f17be58106 100644 --- a/common/edqs/pom.xml +++ b/common/edqs/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT common org.thingsboard.common diff --git a/common/message/pom.xml b/common/message/pom.xml index e053213585..d92d489c84 100644 --- a/common/message/pom.xml +++ b/common/message/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT common org.thingsboard.common diff --git a/common/pom.xml b/common/pom.xml index ac34b5c5e8..665367ddc3 100644 --- a/common/pom.xml +++ b/common/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT thingsboard common diff --git a/common/proto/pom.xml b/common/proto/pom.xml index 2bdaa86165..eb10742275 100644 --- a/common/proto/pom.xml +++ b/common/proto/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT common org.thingsboard.common diff --git a/common/queue/pom.xml b/common/queue/pom.xml index 8e9207ab77..13a5ac837c 100644 --- a/common/queue/pom.xml +++ b/common/queue/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT common org.thingsboard.common diff --git a/common/script/pom.xml b/common/script/pom.xml index 0c276a4395..b2f1011c79 100644 --- a/common/script/pom.xml +++ b/common/script/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT common org.thingsboard.common diff --git a/common/script/remote-js-client/pom.xml b/common/script/remote-js-client/pom.xml index fc8229970e..7f285986ca 100644 --- a/common/script/remote-js-client/pom.xml +++ b/common/script/remote-js-client/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard.common - 4.2.2.4 + 4.2.2.5-SNAPSHOT script org.thingsboard.common.script diff --git a/common/script/script-api/pom.xml b/common/script/script-api/pom.xml index 103b994f01..826392ce6a 100644 --- a/common/script/script-api/pom.xml +++ b/common/script/script-api/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard.common - 4.2.2.4 + 4.2.2.5-SNAPSHOT script org.thingsboard.common.script diff --git a/common/stats/pom.xml b/common/stats/pom.xml index 67a59aa72c..e9340252e5 100644 --- a/common/stats/pom.xml +++ b/common/stats/pom.xml @@ -22,7 +22,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT common org.thingsboard.common diff --git a/common/transport/coap/pom.xml b/common/transport/coap/pom.xml index 83e08e9bb9..35a7b96c41 100644 --- a/common/transport/coap/pom.xml +++ b/common/transport/coap/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard.common - 4.2.2.4 + 4.2.2.5-SNAPSHOT transport org.thingsboard.common.transport diff --git a/common/transport/http/pom.xml b/common/transport/http/pom.xml index 05ccd94bcd..5cbccd7ede 100644 --- a/common/transport/http/pom.xml +++ b/common/transport/http/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard.common - 4.2.2.4 + 4.2.2.5-SNAPSHOT transport org.thingsboard.common.transport diff --git a/common/transport/lwm2m/pom.xml b/common/transport/lwm2m/pom.xml index 215dd9ba6d..9a7e8ddad5 100644 --- a/common/transport/lwm2m/pom.xml +++ b/common/transport/lwm2m/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard.common - 4.2.2.4 + 4.2.2.5-SNAPSHOT transport org.thingsboard.common.transport diff --git a/common/transport/mqtt/pom.xml b/common/transport/mqtt/pom.xml index 5539d95b8c..6664d14d1d 100644 --- a/common/transport/mqtt/pom.xml +++ b/common/transport/mqtt/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard.common - 4.2.2.4 + 4.2.2.5-SNAPSHOT transport org.thingsboard.common.transport diff --git a/common/transport/pom.xml b/common/transport/pom.xml index 7080890835..ec420af200 100644 --- a/common/transport/pom.xml +++ b/common/transport/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT common org.thingsboard.common diff --git a/common/transport/snmp/pom.xml b/common/transport/snmp/pom.xml index ce654035cd..ccabdf75af 100644 --- a/common/transport/snmp/pom.xml +++ b/common/transport/snmp/pom.xml @@ -21,7 +21,7 @@ org.thingsboard.common - 4.2.2.4 + 4.2.2.5-SNAPSHOT transport diff --git a/common/transport/transport-api/pom.xml b/common/transport/transport-api/pom.xml index 91ea55dd3c..723b95a74f 100644 --- a/common/transport/transport-api/pom.xml +++ b/common/transport/transport-api/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard.common - 4.2.2.4 + 4.2.2.5-SNAPSHOT transport org.thingsboard.common.transport diff --git a/common/util/pom.xml b/common/util/pom.xml index 90e579987d..c6584ec7b2 100644 --- a/common/util/pom.xml +++ b/common/util/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT common org.thingsboard.common diff --git a/common/version-control/pom.xml b/common/version-control/pom.xml index 353bed71e1..4c67a2c476 100644 --- a/common/version-control/pom.xml +++ b/common/version-control/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT common org.thingsboard.common diff --git a/dao/pom.xml b/dao/pom.xml index f75168e28e..83bc3ccf85 100644 --- a/dao/pom.xml +++ b/dao/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT thingsboard dao diff --git a/edqs/pom.xml b/edqs/pom.xml index f222c93d24..92f2541039 100644 --- a/edqs/pom.xml +++ b/edqs/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT thingsboard edqs diff --git a/monitoring/pom.xml b/monitoring/pom.xml index 2b083bf040..efd6dab9a3 100644 --- a/monitoring/pom.xml +++ b/monitoring/pom.xml @@ -21,7 +21,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT thingsboard diff --git a/msa/black-box-tests/pom.xml b/msa/black-box-tests/pom.xml index d49a78a8fc..bab206418b 100644 --- a/msa/black-box-tests/pom.xml +++ b/msa/black-box-tests/pom.xml @@ -21,7 +21,7 @@ org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT msa org.thingsboard.msa diff --git a/msa/edqs/pom.xml b/msa/edqs/pom.xml index d77d70dce4..247e1116b1 100644 --- a/msa/edqs/pom.xml +++ b/msa/edqs/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT msa org.thingsboard.msa diff --git a/msa/js-executor/package.json b/msa/js-executor/package.json index de6339b518..e88f5da2eb 100644 --- a/msa/js-executor/package.json +++ b/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", diff --git a/msa/js-executor/pom.xml b/msa/js-executor/pom.xml index 58ca037435..ca09388c19 100644 --- a/msa/js-executor/pom.xml +++ b/msa/js-executor/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT msa org.thingsboard.msa diff --git a/msa/monitoring/pom.xml b/msa/monitoring/pom.xml index 1bfa284dc8..451150aafc 100644 --- a/msa/monitoring/pom.xml +++ b/msa/monitoring/pom.xml @@ -22,7 +22,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT msa diff --git a/msa/pom.xml b/msa/pom.xml index caee13ef17..edaecf7223 100644 --- a/msa/pom.xml +++ b/msa/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT thingsboard msa diff --git a/msa/tb-node/pom.xml b/msa/tb-node/pom.xml index f945f16c69..70c12c90c1 100644 --- a/msa/tb-node/pom.xml +++ b/msa/tb-node/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT msa org.thingsboard.msa diff --git a/msa/tb/pom.xml b/msa/tb/pom.xml index 2e14052543..7cdda058e3 100644 --- a/msa/tb/pom.xml +++ b/msa/tb/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT msa org.thingsboard.msa diff --git a/msa/transport/coap/pom.xml b/msa/transport/coap/pom.xml index ec86562655..02b8a665b2 100644 --- a/msa/transport/coap/pom.xml +++ b/msa/transport/coap/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard.msa - 4.2.2.4 + 4.2.2.5-SNAPSHOT transport org.thingsboard.msa.transport diff --git a/msa/transport/http/pom.xml b/msa/transport/http/pom.xml index 0f758b4b4a..322e40f366 100644 --- a/msa/transport/http/pom.xml +++ b/msa/transport/http/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard.msa - 4.2.2.4 + 4.2.2.5-SNAPSHOT transport org.thingsboard.msa.transport diff --git a/msa/transport/lwm2m/pom.xml b/msa/transport/lwm2m/pom.xml index e3202ee67c..018afb6f23 100644 --- a/msa/transport/lwm2m/pom.xml +++ b/msa/transport/lwm2m/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard.msa - 4.2.2.4 + 4.2.2.5-SNAPSHOT transport org.thingsboard.msa.transport diff --git a/msa/transport/mqtt/pom.xml b/msa/transport/mqtt/pom.xml index 2cf1bc598f..0eac380740 100644 --- a/msa/transport/mqtt/pom.xml +++ b/msa/transport/mqtt/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard.msa - 4.2.2.4 + 4.2.2.5-SNAPSHOT transport org.thingsboard.msa.transport diff --git a/msa/transport/pom.xml b/msa/transport/pom.xml index 291ec91d7f..726d5fdb13 100644 --- a/msa/transport/pom.xml +++ b/msa/transport/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT msa org.thingsboard.msa diff --git a/msa/transport/snmp/pom.xml b/msa/transport/snmp/pom.xml index 05d04cad63..55e1f03df4 100644 --- a/msa/transport/snmp/pom.xml +++ b/msa/transport/snmp/pom.xml @@ -21,7 +21,7 @@ org.thingsboard.msa transport - 4.2.2.4 + 4.2.2.5-SNAPSHOT org.thingsboard.msa.transport diff --git a/msa/vc-executor-docker/pom.xml b/msa/vc-executor-docker/pom.xml index b4bda12d44..0c4c528f1a 100644 --- a/msa/vc-executor-docker/pom.xml +++ b/msa/vc-executor-docker/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT msa org.thingsboard.msa diff --git a/msa/vc-executor/pom.xml b/msa/vc-executor/pom.xml index ac96647732..c09b0dbb7f 100644 --- a/msa/vc-executor/pom.xml +++ b/msa/vc-executor/pom.xml @@ -21,7 +21,7 @@ org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT msa org.thingsboard.msa diff --git a/msa/web-ui/package.json b/msa/web-ui/package.json index e5acd2a6fa..4b1e489783 100644 --- a/msa/web-ui/package.json +++ b/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", diff --git a/msa/web-ui/pom.xml b/msa/web-ui/pom.xml index 2c9044bde6..79996bfdd4 100644 --- a/msa/web-ui/pom.xml +++ b/msa/web-ui/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT msa org.thingsboard.msa diff --git a/netty-mqtt/pom.xml b/netty-mqtt/pom.xml index e9220a3cda..cdcd9ef078 100644 --- a/netty-mqtt/pom.xml +++ b/netty-mqtt/pom.xml @@ -19,11 +19,11 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT thingsboard netty-mqtt - 4.2.2.4 + 4.2.2.5-SNAPSHOT jar Netty MQTT Client diff --git a/pom.xml b/pom.xml index 7858a15ea8..ecc5870817 100755 --- a/pom.xml +++ b/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT pom Thingsboard diff --git a/rest-client/pom.xml b/rest-client/pom.xml index 350ec1c208..3e07db05d5 100644 --- a/rest-client/pom.xml +++ b/rest-client/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT thingsboard rest-client diff --git a/rule-engine/pom.xml b/rule-engine/pom.xml index 9190a845ff..d2a7afffd9 100644 --- a/rule-engine/pom.xml +++ b/rule-engine/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT thingsboard rule-engine diff --git a/rule-engine/rule-engine-api/pom.xml b/rule-engine/rule-engine-api/pom.xml index f0e12c5e33..3112a8ba8c 100644 --- a/rule-engine/rule-engine-api/pom.xml +++ b/rule-engine/rule-engine-api/pom.xml @@ -22,7 +22,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT rule-engine org.thingsboard.rule-engine diff --git a/rule-engine/rule-engine-components/pom.xml b/rule-engine/rule-engine-components/pom.xml index 111e0ad89e..9bce009263 100644 --- a/rule-engine/rule-engine-components/pom.xml +++ b/rule-engine/rule-engine-components/pom.xml @@ -22,7 +22,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT rule-engine org.thingsboard.rule-engine diff --git a/tools/pom.xml b/tools/pom.xml index 925e28ab63..35de3d7584 100644 --- a/tools/pom.xml +++ b/tools/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT thingsboard tools diff --git a/tools/src/main/java/org/thingsboard/client/tools/migrator/DictionaryParser.java b/tools/src/main/java/org/thingsboard/client/tools/migrator/DictionaryParser.java deleted file mode 100644 index d7f38b388e..0000000000 --- a/tools/src/main/java/org/thingsboard/client/tools/migrator/DictionaryParser.java +++ /dev/null @@ -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 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]); - } - } -} diff --git a/tools/src/main/java/org/thingsboard/client/tools/migrator/MigratorTool.java b/tools/src/main/java/org/thingsboard/client/tools/migrator/MigratorTool.java deleted file mode 100644 index 8d73c9cdbc..0000000000 --- a/tools/src/main/java/org/thingsboard/client/tools/migrator/MigratorTool.java +++ /dev/null @@ -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; - } - -} diff --git a/tools/src/main/java/org/thingsboard/client/tools/migrator/PgCaMigrator.java b/tools/src/main/java/org/thingsboard/client/tools/migrator/PgCaMigrator.java deleted file mode 100644 index 2ef202dc9f..0000000000 --- a/tools/src/main/java/org/thingsboard/client/tools/migrator/PgCaMigrator.java +++ /dev/null @@ -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 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 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 result, List 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 result, List raw) { - long ts = Long.parseLong(raw.get(2)); - long partition = toPartitionTs(ts); - result.add(partition); - result.add(ts); - } - - private void addTimeseries(List result, List raw) { - result.add(Long.parseLong(raw.get(2))); - } - - private void addValues(List result, List 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 toValuesTs(List raw) { - - logLinesMigrated(linesTsMigrated++); - - List result = new ArrayList<>(); - - addTypeIdKey(result, raw); - addPartitions(result, raw); - addValues(result, raw); - - processPartitions(result); - - return result; - } - - private List toValuesLatest(List raw) { - logLinesMigrated(linesLatestMigrated++); - List 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 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> function) { - String currentLine; - long linesProcessed = 0; - while(iterator.hasNext()) { - logLinesProcessed(linesProcessed++); - currentLine = iterator.nextLine(); - if(isBlockFinished(currentLine)) { - return; - } - - try { - List raw = Arrays.stream(currentLine.trim().split("\t")) - .map(String::trim) - .collect(Collectors.toList()); - List 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 castToNumericIfPossible(List values) { - try { - if (values.get(6) != null && NumberUtils.isCreatable(values.get(6).toString())) { - Double casted = NumberUtils.createDouble(values.get(6).toString()); - List 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 ("); - } - -} diff --git a/tools/src/main/java/org/thingsboard/client/tools/migrator/README.md b/tools/src/main/java/org/thingsboard/client/tools/migrator/README.md deleted file mode 100644 index 6995632514..0000000000 --- a/tools/src/main/java/org/thingsboard/client/tools/migrator/README.md +++ /dev/null @@ -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 \ No newline at end of file diff --git a/tools/src/main/java/org/thingsboard/client/tools/migrator/RelatedEntitiesParser.java b/tools/src/main/java/org/thingsboard/client/tools/migrator/RelatedEntitiesParser.java deleted file mode 100644 index 59d28f09db..0000000000 --- a/tools/src/main/java/org/thingsboard/client/tools/migrator/RelatedEntitiesParser.java +++ /dev/null @@ -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 allEntityIdsAndTypes = new HashMap<>(); - - private final Map 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 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()); - } - } -} diff --git a/tools/src/main/java/org/thingsboard/client/tools/migrator/WriterBuilder.java b/tools/src/main/java/org/thingsboard/client/tools/migrator/WriterBuilder.java deleted file mode 100644 index e2d1571ade..0000000000 --- a/tools/src/main/java/org/thingsboard/client/tools/migrator/WriterBuilder.java +++ /dev/null @@ -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(); - } -} diff --git a/transport/coap/pom.xml b/transport/coap/pom.xml index 42c875888a..65b9af146d 100644 --- a/transport/coap/pom.xml +++ b/transport/coap/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT transport org.thingsboard.transport diff --git a/transport/http/pom.xml b/transport/http/pom.xml index f5055a2cae..2bbe2277cf 100644 --- a/transport/http/pom.xml +++ b/transport/http/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT transport org.thingsboard.transport diff --git a/transport/lwm2m/pom.xml b/transport/lwm2m/pom.xml index db657dd38d..8229e4d7d5 100644 --- a/transport/lwm2m/pom.xml +++ b/transport/lwm2m/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT transport org.thingsboard.transport diff --git a/transport/mqtt/pom.xml b/transport/mqtt/pom.xml index 1567b444f8..3f84fd44d0 100644 --- a/transport/mqtt/pom.xml +++ b/transport/mqtt/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT transport org.thingsboard.transport diff --git a/transport/pom.xml b/transport/pom.xml index f8ba556c92..b3279b32ed 100644 --- a/transport/pom.xml +++ b/transport/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT thingsboard transport diff --git a/transport/snmp/pom.xml b/transport/snmp/pom.xml index 90d178b7f4..d6f763a8a5 100644 --- a/transport/snmp/pom.xml +++ b/transport/snmp/pom.xml @@ -21,7 +21,7 @@ org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT transport diff --git a/ui-ngx/package.json b/ui-ngx/package.json index 023ede7bd1..5358528a6b 100644 --- a/ui-ngx/package.json +++ b/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", diff --git a/ui-ngx/pom.xml b/ui-ngx/pom.xml index 99712c03fa..9d4ec3f656 100644 --- a/ui-ngx/pom.xml +++ b/ui-ngx/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 4.2.2.4 + 4.2.2.5-SNAPSHOT thingsboard org.thingsboard