Browse Source
# Conflicts: # application/pom.xml # common/actor/pom.xml # common/cache/pom.xml # common/cluster-api/pom.xml # common/coap-server/pom.xml # common/dao-api/pom.xml # common/data/pom.xml # common/discovery-api/pom.xml # common/edge-api/pom.xml # common/edqs/pom.xml # common/message/pom.xml # common/pom.xml # common/proto/pom.xml # common/queue/pom.xml # common/script/pom.xml # common/script/remote-js-client/pom.xml # common/script/script-api/pom.xml # common/stats/pom.xml # common/transport/coap/pom.xml # common/transport/http/pom.xml # common/transport/lwm2m/pom.xml # common/transport/mqtt/pom.xml # common/transport/pom.xml # common/transport/snmp/pom.xml # common/transport/transport-api/pom.xml # common/util/pom.xml # common/version-control/pom.xml # dao/pom.xml # edqs/pom.xml # monitoring/pom.xml # msa/black-box-tests/pom.xml # msa/edqs/pom.xml # msa/js-executor/pom.xml # msa/monitoring/pom.xml # msa/pom.xml # msa/tb-node/pom.xml # msa/tb/pom.xml # msa/transport/coap/pom.xml # msa/transport/http/pom.xml # msa/transport/lwm2m/pom.xml # msa/transport/mqtt/pom.xml # msa/transport/pom.xml # msa/transport/snmp/pom.xml # msa/vc-executor-docker/pom.xml # msa/vc-executor/pom.xml # msa/web-ui/pom.xml # netty-mqtt/pom.xml # pom.xml # rest-client/pom.xml # rule-engine/pom.xml # rule-engine/rule-engine-api/pom.xml # rule-engine/rule-engine-components/pom.xml # tools/pom.xml # transport/coap/pom.xml # transport/http/pom.xml # transport/lwm2m/pom.xml # transport/mqtt/pom.xml # transport/pom.xml # transport/snmp/pom.xml # ui-ngx/pom.xmlrelease-4.3
102 changed files with 2088 additions and 906 deletions
@ -0,0 +1,66 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2026 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.service.security.auth.oauth2; |
||||
|
|
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.springframework.security.oauth2.core.endpoint.OAuth2AuthorizationRequest; |
||||
|
import org.thingsboard.server.common.data.StringUtils; |
||||
|
|
||||
|
import java.util.Locale; |
||||
|
import java.util.Set; |
||||
|
import java.util.regex.Pattern; |
||||
|
|
||||
|
@Slf4j |
||||
|
public class CallbackUrlSchemeValidator { |
||||
|
|
||||
|
// RFC 3986 scheme grammar, plus '_': mobile apps derive the scheme from their package name, which may contain one
|
||||
|
private static final Pattern SCHEME_PATTERN = Pattern.compile("[a-zA-Z][a-zA-Z0-9+.\\-_]*"); |
||||
|
private static final Set<String> FORBIDDEN_SCHEMES = Set.of("http", "https", "javascript", "data", "file", "vbscript"); |
||||
|
private static final int MAX_LOGGED_LENGTH = 128; |
||||
|
|
||||
|
/** |
||||
|
* The redirect carrying the access token is built as callbackUrlScheme + ':', so only a mobile app scheme may |
||||
|
* pass: a web scheme would send the token to whatever host follows it. |
||||
|
*/ |
||||
|
public static boolean isValid(String callbackUrlScheme) { |
||||
|
return !StringUtils.isEmpty(callbackUrlScheme) |
||||
|
&& SCHEME_PATTERN.matcher(callbackUrlScheme).matches() |
||||
|
&& !FORBIDDEN_SCHEMES.contains(callbackUrlScheme.toLowerCase(Locale.ROOT)); |
||||
|
} |
||||
|
|
||||
|
/** |
||||
|
* The attribute is restored from the oauth2_auth_request cookie, which the client can replace, so the scheme is |
||||
|
* checked again on read and not only when the authorization request is built. |
||||
|
*/ |
||||
|
public static String getCallbackUrlScheme(OAuth2AuthorizationRequest authorizationRequest) { |
||||
|
String callbackUrlScheme = authorizationRequest != null ? |
||||
|
authorizationRequest.getAttribute(TbOAuth2ParameterNames.CALLBACK_URL_SCHEME) : null; |
||||
|
if (StringUtils.isEmpty(callbackUrlScheme)) { |
||||
|
return null; |
||||
|
} |
||||
|
if (!isValid(callbackUrlScheme)) { |
||||
|
log.warn("Ignoring invalid callback url scheme: [{}]", forLog(callbackUrlScheme)); |
||||
|
return null; |
||||
|
} |
||||
|
return callbackUrlScheme; |
||||
|
} |
||||
|
|
||||
|
// a rejected value is attacker-controlled: it must not be able to forge log lines
|
||||
|
private static String forLog(String value) { |
||||
|
return value.substring(0, Math.min(value.length(), MAX_LOGGED_LENGTH)).replaceAll("[^\\x20-\\x7E]", "?"); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -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]", "?"); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -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)}); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,133 @@ |
|||||
|
/** |
||||
|
* 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.edge; |
||||
|
|
||||
|
import org.junit.After; |
||||
|
import org.junit.Before; |
||||
|
import org.junit.Test; |
||||
|
import org.mockito.ArgumentMatcher; |
||||
|
import org.springframework.beans.factory.annotation.Autowired; |
||||
|
import org.springframework.test.context.bean.override.mockito.MockitoSpyBean; |
||||
|
import org.springframework.test.util.ReflectionTestUtils; |
||||
|
import org.thingsboard.server.common.data.notification.rule.trigger.EdgeConnectionTrigger; |
||||
|
import org.thingsboard.server.common.data.notification.rule.trigger.NotificationRuleTrigger; |
||||
|
import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor; |
||||
|
import org.thingsboard.server.controller.AbstractWebTest; |
||||
|
import org.thingsboard.server.dao.service.DaoSqlTest; |
||||
|
import org.thingsboard.server.edge.imitator.EdgeImitator; |
||||
|
import org.thingsboard.server.service.edge.rpc.EdgeGrpcService; |
||||
|
|
||||
|
import java.util.Map; |
||||
|
import java.util.concurrent.TimeUnit; |
||||
|
|
||||
|
import static org.awaitility.Awaitility.await; |
||||
|
import static org.mockito.ArgumentMatchers.argThat; |
||||
|
import static org.mockito.Mockito.atLeastOnce; |
||||
|
import static org.mockito.Mockito.clearInvocations; |
||||
|
import static org.mockito.Mockito.never; |
||||
|
import static org.mockito.Mockito.times; |
||||
|
import static org.mockito.Mockito.verify; |
||||
|
|
||||
|
@DaoSqlTest |
||||
|
public class EdgeConnectionNotificationTest extends AbstractEdgeTest { |
||||
|
|
||||
|
private static final long DELAY_MS = 1500L; |
||||
|
|
||||
|
@MockitoSpyBean |
||||
|
private NotificationRuleProcessor notificationRuleProcessor; |
||||
|
|
||||
|
@Autowired |
||||
|
private EdgeGrpcService edgeGrpcService; |
||||
|
|
||||
|
private long originalDisconnectNotificationDelayMs; |
||||
|
|
||||
|
// Capture the bean's configured delay before each test and restore it after, so a test method that forgets
|
||||
|
// to call setDisconnectNotificationDelayMs can't silently inherit the previous method's mutated value
|
||||
|
// (the shared context means the mutation would otherwise persist across methods).
|
||||
|
@Before |
||||
|
public void captureDisconnectNotificationDelay() { |
||||
|
originalDisconnectNotificationDelayMs = (long) ReflectionTestUtils.getField(edgeGrpcService, "disconnectNotificationDelayMs"); |
||||
|
} |
||||
|
|
||||
|
@After |
||||
|
public void restoreDisconnectNotificationDelay() { |
||||
|
setDisconnectNotificationDelayMs(originalDisconnectNotificationDelayMs); |
||||
|
} |
||||
|
|
||||
|
// The delay is overridden per test (rather than via a per-class @TestPropertySource) so all cases share a
|
||||
|
// single Spring application context instead of booting a separate heavy context per delay value.
|
||||
|
private void setDisconnectNotificationDelayMs(long delayMs) { |
||||
|
ReflectionTestUtils.setField(edgeGrpcService, "disconnectNotificationDelayMs", delayMs); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenEdgeStaysDisconnected_whenDelayElapses_thenDisconnectNotificationSent() throws Exception { |
||||
|
// After the configured delay, the "disconnected" notification is sent exactly once.
|
||||
|
assertDisconnectNotificationSentOnce(DELAY_MS); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenZeroDelay_whenEdgeDisconnects_thenDisconnectNotificationSentImmediately() throws Exception { |
||||
|
// With a zero delay there is no debounce window - the "disconnected" notification fires right away.
|
||||
|
assertDisconnectNotificationSentOnce(0); |
||||
|
} |
||||
|
|
||||
|
private void assertDisconnectNotificationSentOnce(long delayMs) throws Exception { |
||||
|
setDisconnectNotificationDelayMs(delayMs); |
||||
|
clearInvocations(notificationRuleProcessor); |
||||
|
|
||||
|
edgeImitator.disconnect(); |
||||
|
|
||||
|
await().atMost(AbstractWebTest.TIMEOUT, TimeUnit.SECONDS).untilAsserted(() -> |
||||
|
verify(notificationRuleProcessor, times(1)).process(argThat(edgeConnectionTrigger(false)))); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenEdgeReconnectsWithinDelay_whenEdgeFlaps_thenDisconnectNotificationSuppressed() throws Exception { |
||||
|
setDisconnectNotificationDelayMs(DELAY_MS); |
||||
|
clearInvocations(notificationRuleProcessor); |
||||
|
|
||||
|
// Edge drops...
|
||||
|
edgeImitator.disconnect(); |
||||
|
// Wait until the server has processed the disconnect and scheduled the pending notification
|
||||
|
// (the edge id appears in the pendingDisconnectNotifications map) before reconnecting.
|
||||
|
await().atMost(AbstractWebTest.TIMEOUT, TimeUnit.SECONDS).until(() -> { |
||||
|
Map<?, ?> pending = (Map<?, ?>) ReflectionTestUtils.getField(edgeGrpcService, "pendingDisconnectNotifications"); |
||||
|
return pending != null && pending.containsKey(edge.getId()); |
||||
|
}); |
||||
|
|
||||
|
// ...and reconnects within the delay window, which must cancel the pending "disconnected" notification.
|
||||
|
EdgeImitator reconnected = createEdgeImitator(); |
||||
|
reconnected.connect(); |
||||
|
edgeImitator = reconnected; // let teardown clean up the live session
|
||||
|
|
||||
|
// The "connected" notification still fires immediately on reconnect (we suppress the disconnect only).
|
||||
|
await().atMost(AbstractWebTest.TIMEOUT, TimeUnit.SECONDS).untilAsserted(() -> |
||||
|
verify(notificationRuleProcessor, atLeastOnce()).process(argThat(edgeConnectionTrigger(true)))); |
||||
|
|
||||
|
// The "disconnected" notification must never be sent throughout the full delay window.
|
||||
|
await().during(DELAY_MS + 500, TimeUnit.MILLISECONDS) |
||||
|
.atMost(DELAY_MS + 2000, TimeUnit.MILLISECONDS) |
||||
|
.untilAsserted(() -> verify(notificationRuleProcessor, never()).process(argThat(edgeConnectionTrigger(false)))); |
||||
|
} |
||||
|
|
||||
|
private ArgumentMatcher<NotificationRuleTrigger> edgeConnectionTrigger(boolean connected) { |
||||
|
return trigger -> trigger instanceof EdgeConnectionTrigger edgeTrigger |
||||
|
&& edge.getId().equals(edgeTrigger.getEdgeId()) |
||||
|
&& edgeTrigger.isConnected() == connected; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,186 @@ |
|||||
|
/** |
||||
|
* 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.cf.ctx.state.aggregation.single; |
||||
|
|
||||
|
import com.fasterxml.jackson.databind.JsonNode; |
||||
|
import com.fasterxml.jackson.databind.node.ArrayNode; |
||||
|
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.Arguments; |
||||
|
import org.junit.jupiter.params.provider.MethodSource; |
||||
|
import org.mockito.InjectMocks; |
||||
|
import org.mockito.Mock; |
||||
|
import org.mockito.junit.jupiter.MockitoExtension; |
||||
|
import org.thingsboard.server.actors.ActorSystemContext; |
||||
|
import org.thingsboard.server.common.data.TenantProfile; |
||||
|
import org.thingsboard.server.common.data.cf.CalculatedField; |
||||
|
import org.thingsboard.server.common.data.cf.CalculatedFieldType; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.Argument; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.ArgumentType; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.TimeSeriesOutput; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggFunction; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggKeyInput; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggMetric; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.aggregation.single.EntityAggregationCalculatedFieldConfiguration; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval.CustomInterval; |
||||
|
import org.thingsboard.server.common.data.id.DeviceId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.data.kv.BasicKvEntry; |
||||
|
import org.thingsboard.server.common.data.kv.DoubleDataEntry; |
||||
|
import org.thingsboard.server.common.data.kv.StringDataEntry; |
||||
|
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; |
||||
|
import org.thingsboard.server.dao.tenant.TbTenantProfileCache; |
||||
|
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; |
||||
|
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; |
||||
|
import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry; |
||||
|
|
||||
|
import java.util.HashMap; |
||||
|
import java.util.Map; |
||||
|
import java.util.Optional; |
||||
|
import java.util.UUID; |
||||
|
import java.util.stream.Stream; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.junit.jupiter.params.provider.Arguments.arguments; |
||||
|
import static org.mockito.ArgumentMatchers.any; |
||||
|
import static org.mockito.Mockito.when; |
||||
|
|
||||
|
@ExtendWith(MockitoExtension.class) |
||||
|
public class EntityAggregationCalculatedFieldStateTest { |
||||
|
|
||||
|
private static final long INTERVAL_START_TS = 1_000L; |
||||
|
private static final long INTERVAL_END_TS = 2_000L; |
||||
|
|
||||
|
private final TenantId TENANT_ID = TenantId.fromUUID(UUID.fromString("80ee80ef-019f-46b1-80ba-22f3ef1b094c")); |
||||
|
private final DeviceId DEVICE_ID = new DeviceId(UUID.fromString("fc83d188-9cf5-4919-a774-d5c56bba2d27")); |
||||
|
|
||||
|
private EntityAggregationCalculatedFieldState state; |
||||
|
private CalculatedFieldCtx ctx; |
||||
|
|
||||
|
@Mock |
||||
|
private TenantProfile tenantProfile; |
||||
|
@Mock |
||||
|
private TbTenantProfileCache tenantProfileCache; |
||||
|
@InjectMocks |
||||
|
private ActorSystemContext systemContext; |
||||
|
|
||||
|
@BeforeEach |
||||
|
void setUp() { |
||||
|
when(tenantProfileCache.get(any(TenantId.class))).thenReturn(tenantProfile); |
||||
|
when(tenantProfile.getProfileConfiguration()).thenReturn(Optional.of(new DefaultTenantProfileConfiguration())); |
||||
|
|
||||
|
ctx = new CalculatedFieldCtx(getCalculatedField(), systemContext); |
||||
|
ctx.init(); |
||||
|
state = new EntityAggregationCalculatedFieldState(DEVICE_ID); |
||||
|
state.setCtx(ctx, null); |
||||
|
state.init(false); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testType() { |
||||
|
assertThat(state.getType()).isEqualTo(CalculatedFieldType.ENTITY_AGGREGATION); |
||||
|
} |
||||
|
|
||||
|
// A numeric aggregation result (SUM/AVG/COUNT/..., numeric MIN/MAX) must be serialized as a numeric
|
||||
|
// JSON node; a genuine string result (lexical MIN/MAX over string telemetry, e.g. a zero-padded code)
|
||||
|
// must stay a string node. The node type is asserted explicitly, so asText() is only used to verify the
|
||||
|
// value once the type is already pinned - it is not relied on to distinguish the types (that blindness
|
||||
|
// is what hid the bug).
|
||||
|
@ParameterizedTest(name = "{0} (precision {2}) -> numeric={3}") |
||||
|
@MethodSource("toResultSerializationCases") |
||||
|
void toResultSerializesResultWithTypePreservingNode(String metricName, BasicKvEntry kvEntry, Integer precision, |
||||
|
boolean expectNumeric, String expectedText) { |
||||
|
JsonNode value = toResultValue(metricName, kvEntry, precision); |
||||
|
|
||||
|
assertThat(value.isNumber()).isEqualTo(expectNumeric); |
||||
|
assertThat(value.isTextual()).isEqualTo(!expectNumeric); |
||||
|
assertThat(value.asText()).isEqualTo(expectedText); |
||||
|
} |
||||
|
|
||||
|
private static Stream<Arguments> toResultSerializationCases() { |
||||
|
return Stream.of( |
||||
|
// SUM: Number result, precision 0 -> whole-number (long) node
|
||||
|
arguments("consumption", new DoubleDataEntry("consumption", 400.0), 0, true, "400"), |
||||
|
// AVG: Number result, precision 2 -> half-up rounded double node
|
||||
|
arguments("avgConsumption", new DoubleDataEntry("avgConsumption", 133.335), 2, true, "133.34"), |
||||
|
// MAX over string telemetry: a zero-padded code is a genuine String result and must stay a
|
||||
|
// string node - as a number it would lose its padding ("0009" -> 9).
|
||||
|
arguments("maxCode", new StringDataEntry("maxCode", "0009"), 0, false, "0009") |
||||
|
); |
||||
|
} |
||||
|
|
||||
|
private JsonNode toResultValue(String metricName, BasicKvEntry kvEntry, Integer precision) { |
||||
|
AggIntervalEntry interval = new AggIntervalEntry(INTERVAL_START_TS, INTERVAL_END_TS); |
||||
|
ArgumentEntry argumentEntry = new SingleValueArgumentEntry(INTERVAL_START_TS, kvEntry, SingleValueArgumentEntry.DEFAULT_VERSION); |
||||
|
Map<AggIntervalEntry, Map<String, ArgumentEntry>> results = new HashMap<>(); |
||||
|
results.put(interval, Map.of(metricName, argumentEntry)); |
||||
|
|
||||
|
ArrayNode result = state.toResult(results, precision); |
||||
|
|
||||
|
assertThat(result.size()).isEqualTo(1); |
||||
|
assertThat(result.get(0).get("ts").asLong()).isEqualTo(INTERVAL_START_TS); |
||||
|
return result.get(0).get("values").get(metricName); |
||||
|
} |
||||
|
|
||||
|
private CalculatedField getCalculatedField() { |
||||
|
CalculatedField calculatedField = new CalculatedField(); |
||||
|
calculatedField.setTenantId(TENANT_ID); |
||||
|
calculatedField.setEntityId(DEVICE_ID); |
||||
|
calculatedField.setType(CalculatedFieldType.ENTITY_AGGREGATION); |
||||
|
calculatedField.setName("Test Entity Aggregation CF"); |
||||
|
calculatedField.setConfigurationVersion(1); |
||||
|
calculatedField.setConfiguration(getConfiguration()); |
||||
|
calculatedField.setVersion(1L); |
||||
|
return calculatedField; |
||||
|
} |
||||
|
|
||||
|
private EntityAggregationCalculatedFieldConfiguration getConfiguration() { |
||||
|
EntityAggregationCalculatedFieldConfiguration configuration = new EntityAggregationCalculatedFieldConfiguration(); |
||||
|
|
||||
|
Argument energy = new Argument(); |
||||
|
energy.setRefEntityKey(new ReferencedEntityKey("energy", ArgumentType.TS_LATEST, null)); |
||||
|
configuration.setArguments(Map.of("en", energy)); |
||||
|
|
||||
|
Map<String, AggMetric> metrics = new HashMap<>(); |
||||
|
AggMetric consumption = new AggMetric(); |
||||
|
consumption.setFunction(AggFunction.SUM); |
||||
|
consumption.setInput(new AggKeyInput("en")); |
||||
|
metrics.put("consumption", consumption); |
||||
|
|
||||
|
AggMetric avgConsumption = new AggMetric(); |
||||
|
avgConsumption.setFunction(AggFunction.AVG); |
||||
|
avgConsumption.setInput(new AggKeyInput("en")); |
||||
|
metrics.put("avgConsumption", avgConsumption); |
||||
|
|
||||
|
AggMetric maxCode = new AggMetric(); |
||||
|
maxCode.setFunction(AggFunction.MAX); |
||||
|
maxCode.setInput(new AggKeyInput("en")); |
||||
|
metrics.put("maxCode", maxCode); |
||||
|
configuration.setMetrics(metrics); |
||||
|
|
||||
|
configuration.setInterval(new CustomInterval("UTC", 0L, 5L)); |
||||
|
|
||||
|
TimeSeriesOutput output = new TimeSeriesOutput(); |
||||
|
output.setDecimalsByDefault(0); |
||||
|
configuration.setOutput(output); |
||||
|
|
||||
|
return configuration; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,194 @@ |
|||||
|
/** |
||||
|
* 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.edge.rpc; |
||||
|
|
||||
|
import org.junit.jupiter.api.BeforeEach; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.junit.jupiter.api.extension.ExtendWith; |
||||
|
import org.mockito.ArgumentMatcher; |
||||
|
import org.mockito.InjectMocks; |
||||
|
import org.mockito.Mock; |
||||
|
import org.mockito.junit.jupiter.MockitoExtension; |
||||
|
import org.springframework.test.util.ReflectionTestUtils; |
||||
|
import org.thingsboard.server.cache.SimpleTbCacheValueWrapper; |
||||
|
import org.thingsboard.server.cache.TbTransactionalCache; |
||||
|
import org.thingsboard.server.common.data.edge.Edge; |
||||
|
import org.thingsboard.server.common.data.id.EdgeId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.data.notification.rule.trigger.EdgeConnectionTrigger; |
||||
|
import org.thingsboard.server.common.data.notification.rule.trigger.NotificationRuleTrigger; |
||||
|
import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor; |
||||
|
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; |
||||
|
import org.thingsboard.server.service.edge.EdgeContextComponent; |
||||
|
|
||||
|
import java.util.UUID; |
||||
|
import java.util.concurrent.ConcurrentMap; |
||||
|
|
||||
|
import static org.mockito.ArgumentMatchers.any; |
||||
|
import static org.mockito.ArgumentMatchers.argThat; |
||||
|
import static org.mockito.Mockito.mock; |
||||
|
import static org.mockito.Mockito.never; |
||||
|
import static org.mockito.Mockito.times; |
||||
|
import static org.mockito.Mockito.verify; |
||||
|
import static org.mockito.Mockito.when; |
||||
|
|
||||
|
@ExtendWith(MockitoExtension.class) |
||||
|
public class EdgeGrpcServiceTest { |
||||
|
|
||||
|
private static final String THIS_NODE = "tb-core-1"; |
||||
|
private static final String OTHER_NODE = "tb-core-2"; |
||||
|
|
||||
|
@Mock |
||||
|
private TbTransactionalCache<EdgeId, String> edgeIdServiceIdCache; |
||||
|
|
||||
|
@Mock |
||||
|
private TbServiceInfoProvider serviceInfoProvider; |
||||
|
|
||||
|
@Mock |
||||
|
private EdgeContextComponent ctx; |
||||
|
|
||||
|
@Mock |
||||
|
private NotificationRuleProcessor ruleProcessor; |
||||
|
|
||||
|
@InjectMocks |
||||
|
private EdgeGrpcService edgeGrpcService; |
||||
|
|
||||
|
private EdgeId edgeId; |
||||
|
private Edge edge; |
||||
|
|
||||
|
@BeforeEach |
||||
|
public void setUp() { |
||||
|
edgeId = new EdgeId(UUID.randomUUID()); |
||||
|
edge = new Edge(edgeId); |
||||
|
edge.setTenantId(TenantId.fromUUID(UUID.randomUUID())); |
||||
|
edge.setName("test-edge"); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenCacheOwnedByThisNode_whenEvict_thenEntryIsEvicted() { |
||||
|
when(serviceInfoProvider.getServiceId()).thenReturn(THIS_NODE); |
||||
|
when(edgeIdServiceIdCache.get(edgeId)).thenReturn(SimpleTbCacheValueWrapper.wrap(THIS_NODE)); |
||||
|
|
||||
|
evictServiceIdCacheIfOwnedByThisNode(); |
||||
|
|
||||
|
verify(edgeIdServiceIdCache, times(1)).evict(edgeId); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenCacheOwnedByAnotherNode_whenEvict_thenEntryIsKept() { |
||||
|
// The edge already reconnected to another node within the keep-alive window: must NOT wipe the live owner.
|
||||
|
when(serviceInfoProvider.getServiceId()).thenReturn(THIS_NODE); |
||||
|
when(edgeIdServiceIdCache.get(edgeId)).thenReturn(SimpleTbCacheValueWrapper.wrap(OTHER_NODE)); |
||||
|
|
||||
|
evictServiceIdCacheIfOwnedByThisNode(); |
||||
|
|
||||
|
verify(edgeIdServiceIdCache, never()).evict(edgeId); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenEmptyCache_whenEvict_thenNothingEvicted() { |
||||
|
when(edgeIdServiceIdCache.get(edgeId)).thenReturn(null); |
||||
|
|
||||
|
evictServiceIdCacheIfOwnedByThisNode(); |
||||
|
|
||||
|
verify(edgeIdServiceIdCache, never()).evict(edgeId); |
||||
|
} |
||||
|
|
||||
|
// --- fireDelayedDisconnectNotification re-verify guard ---
|
||||
|
|
||||
|
@Test |
||||
|
public void givenEdgeReconnectedToThisNode_whenDelayFires_thenNotificationSuppressed() { |
||||
|
// A live session exists again on this node - suppress the stale disconnect notification.
|
||||
|
sessions().put(edgeId, mock(EdgeGrpcSession.class)); |
||||
|
|
||||
|
fireDelayedDisconnectNotification(); |
||||
|
|
||||
|
verify(ruleProcessor, never()).process(any()); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenEdgeReconnectedToAnotherNode_whenDelayFires_thenNotificationSuppressed() { |
||||
|
// No local session, but the cluster cache still points at some node: the edge is connected elsewhere.
|
||||
|
when(edgeIdServiceIdCache.get(edgeId)).thenReturn(SimpleTbCacheValueWrapper.wrap(OTHER_NODE)); |
||||
|
|
||||
|
fireDelayedDisconnectNotification(); |
||||
|
|
||||
|
verify(ruleProcessor, never()).process(any()); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenEdgeStaysDisconnectedClusterWide_whenDelayFires_thenNotificationSent() { |
||||
|
// No local session and no cache entry on any node: the edge is genuinely down - fire the notification.
|
||||
|
when(edgeIdServiceIdCache.get(edgeId)).thenReturn(null); |
||||
|
when(ctx.getRuleProcessor()).thenReturn(ruleProcessor); |
||||
|
|
||||
|
fireDelayedDisconnectNotification(); |
||||
|
|
||||
|
verify(ruleProcessor, times(1)).process(argThat(disconnectTrigger())); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenPendingDisconnect_whenDestroy_thenNotificationFlushed() { |
||||
|
// The edge is genuinely down (no session, no cache). A graceful shutdown must flush the pending
|
||||
|
// notification rather than drop it, otherwise a restart within the delay window swallows the alert.
|
||||
|
when(edgeIdServiceIdCache.get(edgeId)).thenReturn(null); |
||||
|
when(ctx.getRuleProcessor()).thenReturn(ruleProcessor); |
||||
|
pendingDisconnects().put(edgeId, new EdgeGrpcService.PendingDisconnect(edge)); |
||||
|
|
||||
|
destroy(); |
||||
|
|
||||
|
verify(ruleProcessor, times(1)).process(argThat(disconnectTrigger())); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenPendingDisconnectButReconnectedElsewhere_whenDestroy_thenNotificationSuppressed() { |
||||
|
// The flush still honors the re-verify guard: an edge that reconnected to another node must not alert.
|
||||
|
when(edgeIdServiceIdCache.get(edgeId)).thenReturn(SimpleTbCacheValueWrapper.wrap(OTHER_NODE)); |
||||
|
pendingDisconnects().put(edgeId, new EdgeGrpcService.PendingDisconnect(edge)); |
||||
|
|
||||
|
destroy(); |
||||
|
|
||||
|
verify(ruleProcessor, never()).process(any()); |
||||
|
} |
||||
|
|
||||
|
private void destroy() { |
||||
|
ReflectionTestUtils.invokeMethod(edgeGrpcService, "destroy"); |
||||
|
} |
||||
|
|
||||
|
private void evictServiceIdCacheIfOwnedByThisNode() { |
||||
|
ReflectionTestUtils.invokeMethod(edgeGrpcService, "evictServiceIdCacheIfOwnedByThisNode", edgeId); |
||||
|
} |
||||
|
|
||||
|
private void fireDelayedDisconnectNotification() { |
||||
|
EdgeGrpcService.PendingDisconnect pending = new EdgeGrpcService.PendingDisconnect(edge); |
||||
|
ReflectionTestUtils.invokeMethod(edgeGrpcService, "fireDelayedDisconnectNotification", pending); |
||||
|
} |
||||
|
|
||||
|
@SuppressWarnings("unchecked") |
||||
|
private ConcurrentMap<EdgeId, EdgeGrpcSession> sessions() { |
||||
|
return (ConcurrentMap<EdgeId, EdgeGrpcSession>) ReflectionTestUtils.getField(edgeGrpcService, "sessions"); |
||||
|
} |
||||
|
|
||||
|
@SuppressWarnings("unchecked") |
||||
|
private ConcurrentMap<EdgeId, EdgeGrpcService.PendingDisconnect> pendingDisconnects() { |
||||
|
return (ConcurrentMap<EdgeId, EdgeGrpcService.PendingDisconnect>) ReflectionTestUtils.getField(edgeGrpcService, "pendingDisconnectNotifications"); |
||||
|
} |
||||
|
|
||||
|
private static ArgumentMatcher<NotificationRuleTrigger> disconnectTrigger() { |
||||
|
return trigger -> trigger instanceof EdgeConnectionTrigger edgeTrigger && !edgeTrigger.isConnected(); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -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(); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,69 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2026 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.service.security.auth.oauth2; |
||||
|
|
||||
|
import jakarta.servlet.http.Cookie; |
||||
|
import jakarta.servlet.http.HttpServletRequest; |
||||
|
import jakarta.servlet.http.HttpServletResponse; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.junit.jupiter.params.ParameterizedTest; |
||||
|
import org.junit.jupiter.params.provider.NullAndEmptySource; |
||||
|
import org.junit.jupiter.params.provider.ValueSource; |
||||
|
import org.mockito.ArgumentCaptor; |
||||
|
import org.springframework.security.oauth2.core.endpoint.OAuth2AuthorizationRequest; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.mockito.Mockito.atLeastOnce; |
||||
|
import static org.mockito.Mockito.mock; |
||||
|
import static org.mockito.Mockito.verify; |
||||
|
import static org.mockito.Mockito.when; |
||||
|
import static org.thingsboard.server.service.security.auth.oauth2.HttpCookieOAuth2AuthorizationRequestRepository.PREV_URI_COOKIE_NAME; |
||||
|
import static org.thingsboard.server.service.security.auth.oauth2.HttpCookieOAuth2AuthorizationRequestRepository.PREV_URI_PARAMETER; |
||||
|
|
||||
|
public class HttpCookieOAuth2AuthorizationRequestRepositoryTest { |
||||
|
|
||||
|
private final HttpCookieOAuth2AuthorizationRequestRepository repository = new HttpCookieOAuth2AuthorizationRequestRepository(); |
||||
|
|
||||
|
@Test |
||||
|
public void testPrevUriSavedForInAppPath() { |
||||
|
assertThat(savePrevUri("/dashboards/3fa13530-6597-11ed-bd76-8bd591f0ec3e?state=someState")) |
||||
|
.isEqualTo("/dashboards/3fa13530-6597-11ed-bd76-8bd591f0ec3e?state=someState"); |
||||
|
} |
||||
|
|
||||
|
@ParameterizedTest |
||||
|
@NullAndEmptySource |
||||
|
@ValueSource(strings = {"@evil.com/"}) |
||||
|
public void testPrevUriNotSavedForInvalidValue(String prevUri) { |
||||
|
assertThat(savePrevUri(prevUri)).isNull(); |
||||
|
} |
||||
|
|
||||
|
private String savePrevUri(String prevUri) { |
||||
|
HttpServletRequest request = mock(HttpServletRequest.class); |
||||
|
HttpServletResponse response = mock(HttpServletResponse.class); |
||||
|
when(request.getParameter(PREV_URI_PARAMETER)).thenReturn(prevUri); |
||||
|
|
||||
|
repository.saveAuthorizationRequest(OAuth2AuthorizationRequest.authorizationCode() |
||||
|
.authorizationUri("testUri").clientId("testId").build(), request, response); |
||||
|
|
||||
|
ArgumentCaptor<Cookie> cookieCaptor = ArgumentCaptor.forClass(Cookie.class); |
||||
|
verify(response, atLeastOnce()).addCookie(cookieCaptor.capture()); |
||||
|
return cookieCaptor.getAllValues().stream() |
||||
|
.filter(cookie -> PREV_URI_COOKIE_NAME.equals(cookie.getName())) |
||||
|
.map(Cookie::getValue) |
||||
|
.findAny().orElse(null); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -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(); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,103 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2026 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.service.security.auth.oauth2; |
||||
|
|
||||
|
import jakarta.servlet.http.HttpServletRequest; |
||||
|
import jakarta.servlet.http.HttpServletResponse; |
||||
|
import org.junit.jupiter.api.BeforeEach; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.junit.jupiter.params.ParameterizedTest; |
||||
|
import org.junit.jupiter.params.provider.ValueSource; |
||||
|
import org.mockito.ArgumentCaptor; |
||||
|
import org.springframework.security.authentication.AuthenticationServiceException; |
||||
|
import org.springframework.security.oauth2.core.endpoint.OAuth2AuthorizationRequest; |
||||
|
import org.thingsboard.server.common.data.id.CustomerId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.service.security.system.SystemSecurityService; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.mockito.ArgumentMatchers.any; |
||||
|
import static org.mockito.ArgumentMatchers.anyString; |
||||
|
import static org.mockito.Mockito.mock; |
||||
|
import static org.mockito.Mockito.verify; |
||||
|
import static org.mockito.Mockito.when; |
||||
|
|
||||
|
public class Oauth2AuthenticationFailureHandlerTest { |
||||
|
|
||||
|
private static final String BASE_URL = "https://thingsboard.example.com"; |
||||
|
|
||||
|
private final HttpCookieOAuth2AuthorizationRequestRepository authorizationRequestRepository = |
||||
|
mock(HttpCookieOAuth2AuthorizationRequestRepository.class); |
||||
|
private final SystemSecurityService systemSecurityService = mock(SystemSecurityService.class); |
||||
|
private final Oauth2AuthenticationFailureHandler failureHandler = |
||||
|
new Oauth2AuthenticationFailureHandler(authorizationRequestRepository, systemSecurityService); |
||||
|
|
||||
|
private HttpServletRequest request; |
||||
|
private HttpServletResponse response; |
||||
|
|
||||
|
@BeforeEach |
||||
|
public void before() { |
||||
|
request = mock(HttpServletRequest.class); |
||||
|
response = mock(HttpServletResponse.class); |
||||
|
when(request.getContextPath()).thenReturn(""); |
||||
|
when(response.encodeRedirectURL(anyString())).thenAnswer(invocation -> invocation.getArgument(0)); |
||||
|
when(systemSecurityService.getBaseUrl(any(TenantId.class), any(CustomerId.class), any(HttpServletRequest.class))).thenReturn(BASE_URL); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testErrorIsSentToMobileAppScheme() throws Exception { |
||||
|
givenCallbackUrlScheme("tbmobile"); |
||||
|
assertThat(sendFailure()).isEqualTo("tbmobile:/?error=someError"); |
||||
|
} |
||||
|
|
||||
|
/** |
||||
|
* The scheme is restored from the oauth2_auth_request cookie, so a forged one must not turn the error redirect |
||||
|
* into a link to another host. |
||||
|
*/ |
||||
|
@ParameterizedTest |
||||
|
@ValueSource(strings = {"https://evil.com", "javascript"}) |
||||
|
public void testForgedCallbackUrlSchemeFallsBackToLoginPage(String callbackUrlScheme) throws Exception { |
||||
|
givenCallbackUrlScheme(callbackUrlScheme); |
||||
|
assertThat(sendFailure()).isEqualTo(BASE_URL + "/login?loginError=someError"); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testErrorIsSentToLoginPageWithoutCallbackUrlScheme() throws Exception { |
||||
|
givenCallbackUrlScheme(null); |
||||
|
assertThat(sendFailure()).isEqualTo(BASE_URL + "/login?loginError=someError"); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testErrorIsSentToLoginPageWithoutAuthorizationRequest() throws Exception { |
||||
|
assertThat(sendFailure()).isEqualTo(BASE_URL + "/login?loginError=someError"); |
||||
|
} |
||||
|
|
||||
|
private void givenCallbackUrlScheme(String callbackUrlScheme) { |
||||
|
when(authorizationRequestRepository.loadAuthorizationRequest(request)).thenReturn( |
||||
|
OAuth2AuthorizationRequest.authorizationCode().authorizationUri("testUri").clientId("testId") |
||||
|
.attributes(attributes -> attributes.put(TbOAuth2ParameterNames.CALLBACK_URL_SCHEME, callbackUrlScheme)) |
||||
|
.build()); |
||||
|
} |
||||
|
|
||||
|
private String sendFailure() throws Exception { |
||||
|
failureHandler.onAuthenticationFailure(request, response, new AuthenticationServiceException("someError")); |
||||
|
|
||||
|
ArgumentCaptor<String> redirectCaptor = ArgumentCaptor.forClass(String.class); |
||||
|
verify(response).sendRedirect(redirectCaptor.capture()); |
||||
|
return redirectCaptor.getValue(); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -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(); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -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)); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,168 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2026 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.edge.rpc; |
||||
|
|
||||
|
import io.grpc.Server; |
||||
|
import io.grpc.netty.shaded.io.grpc.netty.NettyServerBuilder; |
||||
|
import io.grpc.netty.shaded.io.netty.buffer.PooledByteBufAllocator; |
||||
|
import io.grpc.stub.StreamObserver; |
||||
|
import org.junit.jupiter.api.AfterEach; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.springframework.test.util.ReflectionTestUtils; |
||||
|
import org.thingsboard.server.gen.edge.v1.ConnectResponseCode; |
||||
|
import org.thingsboard.server.gen.edge.v1.ConnectResponseMsg; |
||||
|
import org.thingsboard.server.gen.edge.v1.EdgeRpcServiceGrpc; |
||||
|
import org.thingsboard.server.gen.edge.v1.RequestMsg; |
||||
|
import org.thingsboard.server.gen.edge.v1.RequestMsgType; |
||||
|
import org.thingsboard.server.gen.edge.v1.ResponseMsg; |
||||
|
import org.thingsboard.server.gen.edge.v1.UplinkMsg; |
||||
|
|
||||
|
import java.lang.reflect.Method; |
||||
|
import java.util.concurrent.TimeUnit; |
||||
|
import java.util.function.BooleanSupplier; |
||||
|
|
||||
|
import static org.junit.jupiter.api.Assertions.fail; |
||||
|
|
||||
|
class EdgeGrpcClientLeakTest { |
||||
|
|
||||
|
// The window in which the default shared event loop group would be destroyed after the channel
|
||||
|
// terminates (SharedResourceHolder delays destruction by 1 second). If EdgeGrpcClient ever goes
|
||||
|
// back to the shared group, writes after this window hit a terminated executor and every buffer
|
||||
|
// committed to grpc-netty's WriteQueue is pinned forever (4112 bytes per message, silent after
|
||||
|
// the first RejectedExecutionException).
|
||||
|
private static final long SHARED_GROUP_DEATH_WINDOW_MS = 3000; |
||||
|
private static final int MSG_COUNT = 50; |
||||
|
private static final long AWAIT_TIMEOUT_MS = 15_000; |
||||
|
|
||||
|
private Server server; |
||||
|
private EdgeGrpcClient client; |
||||
|
|
||||
|
@Test |
||||
|
void uplinksSentAfterTransportDeathDoNotPinPooledBuffers() throws Exception { |
||||
|
server = NettyServerBuilder.forPort(0) |
||||
|
.addService(new EdgeRpcServiceGrpc.EdgeRpcServiceImplBase() { |
||||
|
@Override |
||||
|
public StreamObserver<RequestMsg> handleMsgs(StreamObserver<ResponseMsg> outputStream) { |
||||
|
return new StreamObserver<>() { |
||||
|
@Override |
||||
|
public void onNext(RequestMsg requestMsg) { |
||||
|
if (requestMsg.hasConnectRequestMsg()) { |
||||
|
outputStream.onNext(ResponseMsg.newBuilder() |
||||
|
.setConnectResponseMsg(ConnectResponseMsg.newBuilder() |
||||
|
.setResponseCode(ConnectResponseCode.ACCEPTED) |
||||
|
.build()) |
||||
|
.build()); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void onError(Throwable t) { |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void onCompleted() { |
||||
|
} |
||||
|
}; |
||||
|
} |
||||
|
}) |
||||
|
.build() |
||||
|
.start(); |
||||
|
|
||||
|
client = new EdgeGrpcClient(); |
||||
|
ReflectionTestUtils.setField(client, "rpcHost", "localhost"); |
||||
|
ReflectionTestUtils.setField(client, "rpcPort", server.getPort()); |
||||
|
ReflectionTestUtils.setField(client, "timeoutSecs", 1); |
||||
|
ReflectionTestUtils.setField(client, "keepAliveTimeSec", 10); |
||||
|
ReflectionTestUtils.setField(client, "keepAliveTimeoutSec", 5); |
||||
|
ReflectionTestUtils.setField(client, "maxInboundMessageSize", 4194304); |
||||
|
|
||||
|
client.connect("leakTest", "leakTest", msg -> {}, cfg -> {}, msg -> {}, e -> {}); |
||||
|
await("client to connect", () -> client.isConnected()); |
||||
|
|
||||
|
server.shutdownNow(); |
||||
|
server.awaitTermination(10, TimeUnit.SECONDS); |
||||
|
await("client to observe the transport death", () -> !client.isConnected()); |
||||
|
Thread.sleep(SHARED_GROUP_DEATH_WINDOW_MS); |
||||
|
|
||||
|
long baseline = pinnedBytes(); |
||||
|
@SuppressWarnings("unchecked") |
||||
|
StreamObserver<RequestMsg> inputStream = (StreamObserver<RequestMsg>) ReflectionTestUtils.getField(client, "inputStream"); |
||||
|
RequestMsg uplink = RequestMsg.newBuilder() |
||||
|
.setMsgType(RequestMsgType.UPLINK_RPC_MESSAGE) |
||||
|
.setUplinkMsg(UplinkMsg.newBuilder().setUplinkMsgId(1).build()) |
||||
|
.build(); |
||||
|
// Bypasses the connected gate on purpose: this models the check-then-act straggler (and the
|
||||
|
// pre-gate retry loop) writing to a stream whose transport is already gone. Exceptions are
|
||||
|
// swallowed the same way the production retry loop survives them.
|
||||
|
for (int i = 0; i < MSG_COUNT; i++) { |
||||
|
try { |
||||
|
inputStream.onNext(uplink); |
||||
|
} catch (RuntimeException ignored) { |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
long deadline = System.currentTimeMillis() + AWAIT_TIMEOUT_MS; |
||||
|
while (pinnedBytes() > baseline) { |
||||
|
if (System.currentTimeMillis() > deadline) { |
||||
|
fail("Pinned pooled memory did not return to baseline: " + (pinnedBytes() - baseline) |
||||
|
+ " bytes retained after " + MSG_COUNT + " uplinks to a dead stream"); |
||||
|
} |
||||
|
Thread.sleep(50); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@AfterEach |
||||
|
void tearDown() throws Exception { |
||||
|
if (client != null) { |
||||
|
client.disconnect(true); |
||||
|
client.destroy(); |
||||
|
} |
||||
|
if (server != null) { |
||||
|
server.shutdownNow(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private void await(String what, BooleanSupplier condition) throws InterruptedException { |
||||
|
long deadline = System.currentTimeMillis() + AWAIT_TIMEOUT_MS; |
||||
|
while (!condition.getAsBoolean()) { |
||||
|
if (System.currentTimeMillis() > deadline) { |
||||
|
fail("Timed out waiting for " + what); |
||||
|
} |
||||
|
Thread.sleep(50); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
// grpc-netty builds its own PooledByteBufAllocator instances instead of using
|
||||
|
// PooledByteBufAllocator.DEFAULT, and the factory that owns them is package private - so they
|
||||
|
// have to be pulled out reflectively. Both variants are checked so the assertion holds no matter
|
||||
|
// which one the transport picks on this platform/version.
|
||||
|
private static long pinnedBytes() { |
||||
|
try { |
||||
|
Class<?> utils = Class.forName("io.grpc.netty.shaded.io.grpc.netty.Utils"); |
||||
|
Method getByteBufAllocator = utils.getDeclaredMethod("getByteBufAllocator", boolean.class); |
||||
|
getByteBufAllocator.setAccessible(true); |
||||
|
long total = 0; |
||||
|
for (boolean forceHeapBuffer : new boolean[]{false, true}) { |
||||
|
PooledByteBufAllocator pooled = (PooledByteBufAllocator) getByteBufAllocator.invoke(null, forceHeapBuffer); |
||||
|
total += pooled.pinnedDirectMemory() + pooled.pinnedHeapMemory(); |
||||
|
} |
||||
|
return total; |
||||
|
} catch (Exception e) { |
||||
|
throw new IllegalStateException("Failed to read the gRPC allocator metrics", e); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -1,74 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2026 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.client.tools.migrator; |
|
||||
|
|
||||
import org.apache.commons.io.FileUtils; |
|
||||
import org.apache.commons.io.LineIterator; |
|
||||
import org.thingsboard.server.common.data.StringUtils; |
|
||||
|
|
||||
import java.io.File; |
|
||||
import java.io.IOException; |
|
||||
import java.util.HashMap; |
|
||||
import java.util.Map; |
|
||||
|
|
||||
public class DictionaryParser { |
|
||||
private Map<String, String> dictionaryParsed = new HashMap<>(); |
|
||||
|
|
||||
public DictionaryParser(File sourceFile) throws IOException { |
|
||||
parseDictionaryDump(FileUtils.lineIterator(sourceFile)); |
|
||||
} |
|
||||
|
|
||||
public String getKeyByKeyId(String keyId) { |
|
||||
return dictionaryParsed.get(keyId); |
|
||||
} |
|
||||
|
|
||||
private boolean isBlockFinished(String line) { |
|
||||
return StringUtils.isBlank(line) || line.equals("\\."); |
|
||||
} |
|
||||
|
|
||||
private boolean isBlockStarted(String line) { |
|
||||
return line.startsWith("COPY public.key_dictionary ("); |
|
||||
} |
|
||||
|
|
||||
private void parseDictionaryDump(LineIterator iterator) throws IOException { |
|
||||
try { |
|
||||
String tempLine; |
|
||||
while (iterator.hasNext()) { |
|
||||
tempLine = iterator.nextLine(); |
|
||||
|
|
||||
if (isBlockStarted(tempLine)) { |
|
||||
processBlock(iterator); |
|
||||
} |
|
||||
} |
|
||||
} finally { |
|
||||
iterator.close(); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void processBlock(LineIterator lineIterator) { |
|
||||
String tempLine; |
|
||||
String[] lineSplited; |
|
||||
while(lineIterator.hasNext()) { |
|
||||
tempLine = lineIterator.nextLine(); |
|
||||
if(isBlockFinished(tempLine)) { |
|
||||
return; |
|
||||
} |
|
||||
|
|
||||
lineSplited = tempLine.split("\t"); |
|
||||
dictionaryParsed.put(lineSplited[1], lineSplited[0]); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
@ -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; |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -1,282 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2026 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.client.tools.migrator; |
|
||||
|
|
||||
import com.google.common.collect.Lists; |
|
||||
import org.apache.cassandra.io.sstable.CQLSSTableWriter; |
|
||||
import org.apache.commons.io.FileUtils; |
|
||||
import org.apache.commons.io.LineIterator; |
|
||||
import org.apache.commons.lang3.math.NumberUtils; |
|
||||
import org.thingsboard.server.common.data.StringUtils; |
|
||||
|
|
||||
import java.io.File; |
|
||||
import java.io.IOException; |
|
||||
import java.time.Instant; |
|
||||
import java.time.LocalDateTime; |
|
||||
import java.time.ZoneOffset; |
|
||||
import java.time.temporal.ChronoUnit; |
|
||||
import java.util.ArrayList; |
|
||||
import java.util.Arrays; |
|
||||
import java.util.Date; |
|
||||
import java.util.HashSet; |
|
||||
import java.util.List; |
|
||||
import java.util.Set; |
|
||||
import java.util.UUID; |
|
||||
import java.util.function.Function; |
|
||||
import java.util.stream.Collectors; |
|
||||
|
|
||||
public class PgCaMigrator { |
|
||||
|
|
||||
private final long LOG_BATCH = 1000000; |
|
||||
private final long rowPerFile = 1000000; |
|
||||
|
|
||||
private long linesTsMigrated = 0; |
|
||||
private long linesLatestMigrated = 0; |
|
||||
private long castErrors = 0; |
|
||||
private long castedOk = 0; |
|
||||
|
|
||||
private long currentWriterCount = 1; |
|
||||
|
|
||||
private final File sourceFile; |
|
||||
private final boolean castStringIfPossible; |
|
||||
|
|
||||
private final RelatedEntitiesParser entityIdsAndTypes; |
|
||||
private final DictionaryParser keyParser; |
|
||||
private CQLSSTableWriter currentTsWriter; |
|
||||
private CQLSSTableWriter currentPartitionsWriter; |
|
||||
private CQLSSTableWriter currentTsLatestWriter; |
|
||||
private final Set<String> partitions = new HashSet<>(); |
|
||||
|
|
||||
private File outTsDir; |
|
||||
private File outTsLatestDir; |
|
||||
|
|
||||
public PgCaMigrator(File sourceFile, |
|
||||
File ourTsDir, |
|
||||
File outTsPartitionDir, |
|
||||
File outTsLatestDir, |
|
||||
RelatedEntitiesParser allEntityIdsAndTypes, |
|
||||
DictionaryParser dictionaryParser, |
|
||||
boolean castStringsIfPossible) { |
|
||||
this.sourceFile = sourceFile; |
|
||||
this.entityIdsAndTypes = allEntityIdsAndTypes; |
|
||||
this.keyParser = dictionaryParser; |
|
||||
this.castStringIfPossible = castStringsIfPossible; |
|
||||
if(outTsLatestDir != null) { |
|
||||
this.currentTsLatestWriter = WriterBuilder.getLatestWriter(outTsLatestDir); |
|
||||
this.outTsLatestDir = outTsLatestDir; |
|
||||
} |
|
||||
if(ourTsDir != null) { |
|
||||
this.currentTsWriter = WriterBuilder.getTsWriter(ourTsDir); |
|
||||
this.currentPartitionsWriter = WriterBuilder.getPartitionWriter(outTsPartitionDir); |
|
||||
this.outTsDir = ourTsDir; |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
public void migrate() throws IOException { |
|
||||
boolean isTsDone = false; |
|
||||
boolean isLatestDone = false; |
|
||||
String line; |
|
||||
LineIterator iterator = FileUtils.lineIterator(this.sourceFile); |
|
||||
|
|
||||
try { |
|
||||
while(iterator.hasNext()) { |
|
||||
line = iterator.nextLine(); |
|
||||
if(!isLatestDone && isBlockLatestStarted(line)) { |
|
||||
System.out.println("START TO MIGRATE LATEST"); |
|
||||
long start = System.currentTimeMillis(); |
|
||||
processBlock(iterator, currentTsLatestWriter, outTsLatestDir, this::toValuesLatest); |
|
||||
System.out.println("TOTAL LINES MIGRATED: " + linesLatestMigrated + ", FORMING OF SSL FOR LATEST TS FINISHED WITH TIME: " + (System.currentTimeMillis() - start) + " ms."); |
|
||||
isLatestDone = true; |
|
||||
} |
|
||||
|
|
||||
if(!isTsDone && isBlockTsStarted(line)) { |
|
||||
System.out.println("START TO MIGRATE TS"); |
|
||||
long start = System.currentTimeMillis(); |
|
||||
processBlock(iterator, currentTsWriter, outTsDir, this::toValuesTs); |
|
||||
System.out.println("TOTAL LINES MIGRATED: " + linesTsMigrated + ", FORMING OF SSL FOR TS FINISHED WITH TIME: " + (System.currentTimeMillis() - start) + " ms."); |
|
||||
isTsDone = true; |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
System.out.println("Partitions collected " + partitions.size()); |
|
||||
long startTs = System.currentTimeMillis(); |
|
||||
for (String partition : partitions) { |
|
||||
String[] split = partition.split("\\|"); |
|
||||
List<Object> values = Lists.newArrayList(); |
|
||||
values.add(split[0]); |
|
||||
values.add(UUID.fromString(split[1])); |
|
||||
values.add(split[2]); |
|
||||
values.add(Long.parseLong(split[3])); |
|
||||
currentPartitionsWriter.addRow(values); |
|
||||
} |
|
||||
|
|
||||
System.out.println(new Date() + " Migrated partitions " + partitions.size() + " in " + (System.currentTimeMillis() - startTs)); |
|
||||
|
|
||||
System.out.println(); |
|
||||
System.out.println("Finished migrate Telemetry"); |
|
||||
|
|
||||
} finally { |
|
||||
iterator.close(); |
|
||||
currentTsLatestWriter.close(); |
|
||||
currentTsWriter.close(); |
|
||||
currentPartitionsWriter.close(); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void logLinesProcessed(long lines) { |
|
||||
if (lines % LOG_BATCH == 0) { |
|
||||
System.out.println(new Date() + " lines processed = " + lines + " in, castOk " + castedOk + " castErr " + castErrors); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void logLinesMigrated(long lines) { |
|
||||
if(lines % LOG_BATCH == 0) { |
|
||||
System.out.println(new Date() + " lines migrated = " + lines + " in, castOk " + castedOk + " castErr " + castErrors); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void addTypeIdKey(List<Object> result, List<String> raw) { |
|
||||
result.add(entityIdsAndTypes.getEntityType(raw.get(0))); |
|
||||
result.add(UUID.fromString(raw.get(0))); |
|
||||
result.add(keyParser.getKeyByKeyId(raw.get(1))); |
|
||||
} |
|
||||
|
|
||||
private void addPartitions(List<Object> result, List<String> raw) { |
|
||||
long ts = Long.parseLong(raw.get(2)); |
|
||||
long partition = toPartitionTs(ts); |
|
||||
result.add(partition); |
|
||||
result.add(ts); |
|
||||
} |
|
||||
|
|
||||
private void addTimeseries(List<Object> result, List<String> raw) { |
|
||||
result.add(Long.parseLong(raw.get(2))); |
|
||||
} |
|
||||
|
|
||||
private void addValues(List<Object> result, List<String> raw) { |
|
||||
result.add(raw.get(3).equals("\\N") ? null : raw.get(3).equals("t") ? Boolean.TRUE : Boolean.FALSE); |
|
||||
result.add(raw.get(4).equals("\\N") ? null : raw.get(4)); |
|
||||
result.add(raw.get(5).equals("\\N") ? null : Long.parseLong(raw.get(5))); |
|
||||
result.add(raw.get(6).equals("\\N") ? null : Double.parseDouble(raw.get(6))); |
|
||||
result.add(raw.get(7).equals("\\N") ? null : raw.get(7)); |
|
||||
} |
|
||||
|
|
||||
private List<Object> toValuesTs(List<String> raw) { |
|
||||
|
|
||||
logLinesMigrated(linesTsMigrated++); |
|
||||
|
|
||||
List<Object> result = new ArrayList<>(); |
|
||||
|
|
||||
addTypeIdKey(result, raw); |
|
||||
addPartitions(result, raw); |
|
||||
addValues(result, raw); |
|
||||
|
|
||||
processPartitions(result); |
|
||||
|
|
||||
return result; |
|
||||
} |
|
||||
|
|
||||
private List<Object> toValuesLatest(List<String> raw) { |
|
||||
logLinesMigrated(linesLatestMigrated++); |
|
||||
List<Object> result = new ArrayList<>(); |
|
||||
|
|
||||
addTypeIdKey(result, raw); |
|
||||
addTimeseries(result, raw); |
|
||||
addValues(result, raw); |
|
||||
|
|
||||
return result; |
|
||||
} |
|
||||
|
|
||||
private long toPartitionTs(long ts) { |
|
||||
LocalDateTime time = LocalDateTime.ofInstant(Instant.ofEpochMilli(ts), ZoneOffset.UTC); |
|
||||
return time.truncatedTo(ChronoUnit.DAYS).withDayOfMonth(1).toInstant(ZoneOffset.UTC).toEpochMilli(); |
|
||||
} |
|
||||
|
|
||||
private void processPartitions(List<Object> values) { |
|
||||
String key = values.get(0) + "|" + values.get(1) + "|" + values.get(2) + "|" + values.get(3); |
|
||||
partitions.add(key); |
|
||||
} |
|
||||
|
|
||||
private void processBlock(LineIterator iterator, CQLSSTableWriter writer, File outDir, Function<List<String>, List<Object>> function) { |
|
||||
String currentLine; |
|
||||
long linesProcessed = 0; |
|
||||
while(iterator.hasNext()) { |
|
||||
logLinesProcessed(linesProcessed++); |
|
||||
currentLine = iterator.nextLine(); |
|
||||
if(isBlockFinished(currentLine)) { |
|
||||
return; |
|
||||
} |
|
||||
|
|
||||
try { |
|
||||
List<String> raw = Arrays.stream(currentLine.trim().split("\t")) |
|
||||
.map(String::trim) |
|
||||
.collect(Collectors.toList()); |
|
||||
List<Object> values = function.apply(raw); |
|
||||
|
|
||||
if (this.currentWriterCount == 0) { |
|
||||
System.out.println(new Date() + " close writer " + new Date()); |
|
||||
writer.close(); |
|
||||
writer = WriterBuilder.getLatestWriter(outDir); |
|
||||
} |
|
||||
|
|
||||
if (this.castStringIfPossible) { |
|
||||
writer.addRow(castToNumericIfPossible(values)); |
|
||||
} else { |
|
||||
writer.addRow(values); |
|
||||
} |
|
||||
|
|
||||
currentWriterCount++; |
|
||||
if (currentWriterCount >= rowPerFile) { |
|
||||
currentWriterCount = 0; |
|
||||
} |
|
||||
} catch (Exception ex) { |
|
||||
System.out.println(ex.getMessage() + " -> " + currentLine); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private List<Object> castToNumericIfPossible(List<Object> values) { |
|
||||
try { |
|
||||
if (values.get(6) != null && NumberUtils.isCreatable(values.get(6).toString())) { |
|
||||
Double casted = NumberUtils.createDouble(values.get(6).toString()); |
|
||||
List<Object> numeric = Lists.newArrayList(); |
|
||||
numeric.addAll(values); |
|
||||
numeric.set(6, null); |
|
||||
numeric.set(8, casted); |
|
||||
castedOk++; |
|
||||
return numeric; |
|
||||
} |
|
||||
} catch (Throwable th) { |
|
||||
castErrors++; |
|
||||
} |
|
||||
|
|
||||
processPartitions(values); |
|
||||
|
|
||||
return values; |
|
||||
} |
|
||||
|
|
||||
private boolean isBlockFinished(String line) { |
|
||||
return StringUtils.isBlank(line) || line.equals("\\."); |
|
||||
} |
|
||||
|
|
||||
private boolean isBlockTsStarted(String line) { |
|
||||
return line.startsWith("COPY public.ts_kv ("); |
|
||||
} |
|
||||
|
|
||||
private boolean isBlockLatestStarted(String line) { |
|
||||
return line.startsWith("COPY public.ts_kv_latest ("); |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -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 |
|
||||
@ -1,88 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2026 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.client.tools.migrator; |
|
||||
|
|
||||
import org.apache.commons.io.FileUtils; |
|
||||
import org.apache.commons.io.LineIterator; |
|
||||
import org.thingsboard.server.common.data.EntityType; |
|
||||
import org.thingsboard.server.common.data.StringUtils; |
|
||||
|
|
||||
import java.io.File; |
|
||||
import java.io.IOException; |
|
||||
import java.util.HashMap; |
|
||||
import java.util.Map; |
|
||||
|
|
||||
public class RelatedEntitiesParser { |
|
||||
private final Map<String, String> allEntityIdsAndTypes = new HashMap<>(); |
|
||||
|
|
||||
private final Map<String, EntityType> tableNameAndEntityType = Map.ofEntries( |
|
||||
Map.entry("COPY public.alarm ", EntityType.ALARM), |
|
||||
Map.entry("COPY public.asset ", EntityType.ASSET), |
|
||||
Map.entry("COPY public.customer ", EntityType.CUSTOMER), |
|
||||
Map.entry("COPY public.dashboard ", EntityType.DASHBOARD), |
|
||||
Map.entry("COPY public.device ", EntityType.DEVICE), |
|
||||
Map.entry("COPY public.rule_chain ", EntityType.RULE_CHAIN), |
|
||||
Map.entry("COPY public.rule_node ", EntityType.RULE_NODE), |
|
||||
Map.entry("COPY public.tenant ", EntityType.TENANT), |
|
||||
Map.entry("COPY public.tb_user ", EntityType.USER), |
|
||||
Map.entry("COPY public.entity_view ", EntityType.ENTITY_VIEW), |
|
||||
Map.entry("COPY public.widgets_bundle ", EntityType.WIDGETS_BUNDLE), |
|
||||
Map.entry("COPY public.widget_type ", EntityType.WIDGET_TYPE), |
|
||||
Map.entry("COPY public.tenant_profile ", EntityType.TENANT_PROFILE), |
|
||||
Map.entry("COPY public.device_profile ", EntityType.DEVICE_PROFILE), |
|
||||
Map.entry("COPY public.asset_profile ", EntityType.ASSET_PROFILE), |
|
||||
Map.entry("COPY public.api_usage_state ", EntityType.API_USAGE_STATE) |
|
||||
); |
|
||||
|
|
||||
public RelatedEntitiesParser(File source) throws IOException { |
|
||||
processAllTables(FileUtils.lineIterator(source)); |
|
||||
} |
|
||||
|
|
||||
public String getEntityType(String uuid) { |
|
||||
return this.allEntityIdsAndTypes.get(uuid); |
|
||||
} |
|
||||
|
|
||||
private boolean isBlockFinished(String line) { |
|
||||
return StringUtils.isBlank(line) || line.equals("\\."); |
|
||||
} |
|
||||
|
|
||||
private void processAllTables(LineIterator lineIterator) throws IOException { |
|
||||
String currentLine; |
|
||||
try { |
|
||||
while (lineIterator.hasNext()) { |
|
||||
currentLine = lineIterator.nextLine(); |
|
||||
for(Map.Entry<String, EntityType> entry : tableNameAndEntityType.entrySet()) { |
|
||||
if(currentLine.startsWith(entry.getKey())) { |
|
||||
processBlock(lineIterator, entry.getValue()); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
} finally { |
|
||||
lineIterator.close(); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void processBlock(LineIterator lineIterator, EntityType entityType) { |
|
||||
String currentLine; |
|
||||
while(lineIterator.hasNext()) { |
|
||||
currentLine = lineIterator.nextLine(); |
|
||||
if(isBlockFinished(currentLine)) { |
|
||||
return; |
|
||||
} |
|
||||
allEntityIdsAndTypes.put(currentLine.split("\t")[0], entityType.name()); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
@ -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(); |
|
||||
} |
|
||||
} |
|
||||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue