1207 changed files with 16796 additions and 11370 deletions
@ -0,0 +1,31 @@ |
|||
|
|||
## Running tests in parallel with a reasonable memory usage |
|||
|
|||
```bash |
|||
export MAVEN_OPTS="-Xmx1024m" |
|||
export NODE_OPTIONS="--max_old_space_size=4096" |
|||
export SUREFIRE_JAVA_OPTS="-Xmx1200m -Xss256k -XX:+ExitOnOutOfMemoryError" |
|||
|
|||
mvn clean install -T6 -DskipTests |
|||
mvn test -pl='!application,!dao,!ui-ngx,!msa/js-executor,!msa/web-ui' -T4 |
|||
mvn test -pl dao -Dparallel=packages -DforkCount=4 |
|||
|
|||
mvn test -pl application -Dsurefire.excludes='**/nosql/*Test.java' -Dtest='org.thingsboard.server.controller.**' -DforkCount=6 -Dparallel=classes -Dsurefire.rerunFailingTestsCount=2 -Dsurefire.failOnFlakeCount=5 |
|||
mvn test -pl application -Dsurefire.excludes='**/nosql/*Test.java' -Dtest='org.thingsboard.server.edge.**' -DforkCount=4 -Dparallel=packages -Dsurefire.rerunFailingTestsCount=2 -Dsurefire.failOnFlakeCount=5 |
|||
mvn test -pl application -Dsurefire.excludes='**/nosql/*Test.java' -Dtest='org.thingsboard.server.service.**' -DforkCount=6 -Dparallel=packages -Dsurefire.rerunFailingTestsCount=2 -Dsurefire.failOnFlakeCount=5 |
|||
mvn test -pl application -Dsurefire.excludes='**/nosql/*Test.java' -Dtest='org.thingsboard.server.transport.mqtt.**' -DforkCount=6 -Dparallel=classes -Dsurefire.rerunFailingTestsCount=2 -Dsurefire.failOnFlakeCount=5 |
|||
mvn test -pl application -Dsurefire.excludes='**/nosql/*Test.java' -Dtest='org.thingsboard.server.transport.coap.**' -DforkCount=6 -Dparallel=classes -Dsurefire.rerunFailingTestsCount=2 -Dsurefire.failOnFlakeCount=5 |
|||
mvn test -pl application -Dsurefire.excludes='**/nosql/*Test.java' -Dtest='org.thingsboard.server.transport.lwm2m.**' -DforkCount=6 -Dparallel=packages -Dsurefire.rerunFailingTestsCount=2 -Dsurefire.failOnFlakeCount=5 |
|||
mvn test -pl application -Dsurefire.excludes='**/nosql/*Test.java' -Dtest='**/*TestSuite.java' -DforkCount=4 -Dparallel=classes -Dsurefire.rerunFailingTestsCount=2 -Dsurefire.failOnFlakeCount=5 |
|||
|
|||
#the rest of application tests |
|||
mvn test -pl application -Dtest=' |
|||
!**/nosql/*Test.java, |
|||
!org.thingsboard.server.controller.**, |
|||
!org.thingsboard.server.edge.**, |
|||
!org.thingsboard.server.service.**, |
|||
!org.thingsboard.server.transport.mqtt.**, |
|||
!org.thingsboard.server.transport.coap.**, |
|||
!org.thingsboard.server.transport.lwm2m.** |
|||
' -DforkCount=6 -Dparallel=packages -Dsurefire.rerunFailingTestsCount=2 -Dsurefire.failOnFlakeCount=5 |
|||
``` |
|||
@ -0,0 +1,41 @@ |
|||
/** |
|||
* 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.Assert; |
|||
import org.junit.Test; |
|||
import org.thingsboard.edge.rpc.EdgeVersionComparator; |
|||
import org.thingsboard.server.gen.edge.v1.EdgeVersion; |
|||
|
|||
public class EdgeLatestVersionTest { |
|||
|
|||
@Test |
|||
public void edgeLatestVersionIsSynchronizedTest() { |
|||
EdgeVersion currentHighestEdgeVersion = EdgeVersionComparator.getNewestEdgeVersion(); |
|||
|
|||
String projectVersion = EdgeLatestVersionTest.class.getPackage().getImplementationVersion(); |
|||
if (projectVersion == null || projectVersion.isBlank()) { |
|||
projectVersion = System.getProperty("project.version", "UNKNOWN"); |
|||
} |
|||
|
|||
String projectVersionDigits = projectVersion.replaceAll("\\D", ""); |
|||
String currentHighestEdgeVersionDigits = currentHighestEdgeVersion.name().replaceAll("\\D", ""); |
|||
|
|||
String msg = "EdgeVersion enum in edge.proto is out of sync. Please add respective " + projectVersionDigits + " to EdgeVersion"; |
|||
Assert.assertEquals(msg, projectVersionDigits, currentHighestEdgeVersionDigits); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,204 @@ |
|||
/** |
|||
* 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.system; |
|||
|
|||
import org.junit.Before; |
|||
import org.junit.Test; |
|||
import org.junit.runner.RunWith; |
|||
import org.mockito.Mock; |
|||
import org.mockito.junit.MockitoJUnitRunner; |
|||
import org.springframework.security.authentication.BadCredentialsException; |
|||
import org.springframework.security.authentication.DisabledException; |
|||
import org.springframework.security.authentication.LockedException; |
|||
import org.springframework.security.crypto.bcrypt.BCryptPasswordEncoder; |
|||
import org.thingsboard.rule.engine.api.MailService; |
|||
import org.thingsboard.server.common.data.exception.ThingsboardException; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.id.UserId; |
|||
import org.thingsboard.server.common.data.security.UserCredentials; |
|||
import org.thingsboard.server.common.data.security.model.SecuritySettings; |
|||
import org.thingsboard.server.common.data.security.model.UserPasswordPolicy; |
|||
import org.thingsboard.server.dao.audit.AuditLogService; |
|||
import org.thingsboard.server.dao.settings.AdminSettingsService; |
|||
import org.thingsboard.server.dao.settings.SecuritySettingsService; |
|||
import org.thingsboard.server.dao.user.UserService; |
|||
|
|||
import java.util.UUID; |
|||
|
|||
import static org.assertj.core.api.Assertions.assertThatThrownBy; |
|||
import static org.mockito.ArgumentMatchers.any; |
|||
import static org.mockito.ArgumentMatchers.eq; |
|||
import static org.mockito.Mockito.never; |
|||
import static org.mockito.Mockito.verify; |
|||
import static org.mockito.Mockito.when; |
|||
|
|||
@RunWith(MockitoJUnitRunner.class) |
|||
public class DefaultSystemSecurityServiceTest { |
|||
|
|||
@Mock |
|||
private AdminSettingsService adminSettingsService; |
|||
@Mock |
|||
private BCryptPasswordEncoder encoder; |
|||
@Mock |
|||
private UserService userService; |
|||
@Mock |
|||
private MailService mailService; |
|||
@Mock |
|||
private AuditLogService auditLogService; |
|||
@Mock |
|||
private SecuritySettingsService securitySettingsService; |
|||
|
|||
private DefaultSystemSecurityService systemSecurityService; |
|||
|
|||
private TenantId tenantId; |
|||
private UserId userId; |
|||
private UserCredentials userCredentials; |
|||
private SecuritySettings securitySettings; |
|||
private String username; |
|||
private String password; |
|||
private String encodedPassword; |
|||
|
|||
@Before |
|||
public void setUp() { |
|||
systemSecurityService = new DefaultSystemSecurityService(adminSettingsService, encoder, userService, mailService, auditLogService, securitySettingsService); |
|||
|
|||
tenantId = TenantId.fromUUID(UUID.randomUUID()); |
|||
userId = new UserId(UUID.randomUUID()); |
|||
username = "tenant@example.com"; |
|||
password = "correctPassword"; |
|||
encodedPassword = "$2a$10$encodedPasswordHash"; |
|||
|
|||
userCredentials = new UserCredentials(); |
|||
userCredentials.setUserId(userId); |
|||
userCredentials.setEnabled(true); |
|||
userCredentials.setPassword(encodedPassword); |
|||
userCredentials.setCreatedTime(System.currentTimeMillis()); |
|||
|
|||
securitySettings = new SecuritySettings(); |
|||
securitySettings.setMaxFailedLoginAttempts(5); |
|||
securitySettings.setPasswordPolicy(new UserPasswordPolicy()); |
|||
} |
|||
|
|||
@Test |
|||
public void testValidateUserCredentials_successfulLogin() { |
|||
when(encoder.matches(password, encodedPassword)).thenReturn(true); |
|||
when(securitySettingsService.getSecuritySettings()).thenReturn(securitySettings); |
|||
|
|||
systemSecurityService.validateUserCredentials(tenantId, userCredentials, username, password); |
|||
|
|||
verify(encoder).matches(password, encodedPassword); |
|||
verify(userService).resetFailedLoginAttempts(tenantId, userId); |
|||
verify(userService, never()).increaseFailedLoginAttempts(any(), any()); |
|||
} |
|||
|
|||
@Test |
|||
public void testValidateUserCredentials_wrongPassword_incrementsFailedAttempts() { |
|||
when(encoder.matches(password, encodedPassword)).thenReturn(false); |
|||
when(userService.increaseFailedLoginAttempts(tenantId, userId)).thenReturn(3); |
|||
when(securitySettingsService.getSecuritySettings()).thenReturn(securitySettings); |
|||
|
|||
assertThatThrownBy(() -> systemSecurityService.validateUserCredentials(tenantId, userCredentials, username, password)) |
|||
.isInstanceOf(BadCredentialsException.class) |
|||
.hasMessageContaining("Authentication Failed"); |
|||
|
|||
verify(userService).increaseFailedLoginAttempts(tenantId, userId); |
|||
verify(userService, never()).setUserCredentialsEnabled(any(), any(), eq(false)); |
|||
} |
|||
|
|||
@Test |
|||
public void testValidateUserCredentials_wrongPassword_accountLocked() { |
|||
when(encoder.matches(password, encodedPassword)).thenReturn(false); |
|||
when(userService.increaseFailedLoginAttempts(tenantId, userId)).thenReturn(6); |
|||
when(securitySettingsService.getSecuritySettings()).thenReturn(securitySettings); |
|||
|
|||
assertThatThrownBy(() -> systemSecurityService.validateUserCredentials(tenantId, userCredentials, username, password)) |
|||
.isInstanceOf(LockedException.class) |
|||
.hasMessageContaining("locked due to security policy"); |
|||
|
|||
verify(userService).increaseFailedLoginAttempts(tenantId, userId); |
|||
verify(userService).setUserCredentialsEnabled(TenantId.SYS_TENANT_ID, userId, false); |
|||
} |
|||
|
|||
@Test |
|||
public void testValidateUserCredentials_wrongPassword_exactlyAtThreshold_noLock() { |
|||
when(encoder.matches(password, encodedPassword)).thenReturn(false); |
|||
when(userService.increaseFailedLoginAttempts(tenantId, userId)).thenReturn(5); |
|||
when(securitySettingsService.getSecuritySettings()).thenReturn(securitySettings); |
|||
|
|||
assertThatThrownBy(() -> systemSecurityService.validateUserCredentials(tenantId, userCredentials, username, password)) |
|||
.isInstanceOf(BadCredentialsException.class) |
|||
.hasMessageContaining("Authentication Failed"); |
|||
|
|||
verify(userService).increaseFailedLoginAttempts(tenantId, userId); |
|||
verify(userService, never()).setUserCredentialsEnabled(any(), any(), eq(false)); |
|||
} |
|||
|
|||
@Test |
|||
public void testValidateUserCredentials_correctPassword_disabledUser_throwsDisabledException() { |
|||
userCredentials.setEnabled(false); |
|||
|
|||
assertThatThrownBy(() -> systemSecurityService.validateUserCredentials(tenantId, userCredentials, username, password)) |
|||
.isInstanceOf(DisabledException.class) |
|||
.hasMessage("User is not active"); |
|||
|
|||
verify(encoder, never()).matches(any(), any()); |
|||
verify(userService, never()).increaseFailedLoginAttempts(any(), any()); |
|||
} |
|||
|
|||
@Test |
|||
public void testValidateUserCredentials_wrongPassword_maxAttemptsDisabled() { |
|||
securitySettings.setMaxFailedLoginAttempts(null); |
|||
when(encoder.matches(password, encodedPassword)).thenReturn(false); |
|||
when(userService.increaseFailedLoginAttempts(tenantId, userId)).thenReturn(100); |
|||
when(securitySettingsService.getSecuritySettings()).thenReturn(securitySettings); |
|||
|
|||
assertThatThrownBy(() -> systemSecurityService.validateUserCredentials(tenantId, userCredentials, username, password)) |
|||
.isInstanceOf(BadCredentialsException.class); |
|||
|
|||
verify(userService).increaseFailedLoginAttempts(tenantId, userId); |
|||
verify(userService, never()).setUserCredentialsEnabled(any(), any(), eq(false)); |
|||
} |
|||
|
|||
@Test |
|||
public void testValidateUserCredentials_wrongPassword_maxAttemptsSetToZero() { |
|||
securitySettings.setMaxFailedLoginAttempts(0); |
|||
when(encoder.matches(password, encodedPassword)).thenReturn(false); |
|||
when(userService.increaseFailedLoginAttempts(tenantId, userId)).thenReturn(100); |
|||
when(securitySettingsService.getSecuritySettings()).thenReturn(securitySettings); |
|||
|
|||
assertThatThrownBy(() -> systemSecurityService.validateUserCredentials(tenantId, userCredentials, username, password)) |
|||
.isInstanceOf(BadCredentialsException.class); |
|||
|
|||
verify(userService).increaseFailedLoginAttempts(tenantId, userId); |
|||
verify(userService, never()).setUserCredentialsEnabled(any(), any(), eq(false)); |
|||
} |
|||
|
|||
@Test |
|||
public void testValidateUserCredentials_wrongPassword_withNotificationEmail() throws ThingsboardException { |
|||
String notificationEmail = "admin@example.com"; |
|||
securitySettings.setUserLockoutNotificationEmail(notificationEmail); |
|||
when(encoder.matches(password, encodedPassword)).thenReturn(false); |
|||
when(userService.increaseFailedLoginAttempts(tenantId, userId)).thenReturn(6); |
|||
when(securitySettingsService.getSecuritySettings()).thenReturn(securitySettings); |
|||
|
|||
assertThatThrownBy(() -> systemSecurityService.validateUserCredentials(tenantId, userCredentials, username, password)) |
|||
.isInstanceOf(LockedException.class); |
|||
|
|||
verify(userService).setUserCredentialsEnabled(TenantId.SYS_TENANT_ID, userId, false); |
|||
verify(mailService).sendAccountLockoutEmail(eq(username), eq(notificationEmail), eq(5)); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,42 @@ |
|||
/** |
|||
* 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.utils; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
|
|||
import java.io.IOException; |
|||
import java.net.DatagramSocket; |
|||
import java.net.SocketException; |
|||
|
|||
@Slf4j |
|||
public class PortFinder { |
|||
public static int findAvailableUdpPort() { |
|||
try (DatagramSocket socket = new DatagramSocket(0)) { |
|||
return socket.getLocalPort(); |
|||
} catch (SocketException e) { |
|||
throw new IllegalStateException("No available UDP ports found", e); |
|||
} |
|||
} |
|||
|
|||
public static boolean isUDPPortAvailable(int port) { |
|||
try (DatagramSocket socket = new DatagramSocket(port)) { |
|||
return true; |
|||
} catch (IOException e) { |
|||
log.debug("Failed to open UDP port {}", port, e); |
|||
return false; |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,86 @@ |
|||
/** |
|||
* 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 org.thingsboard.server.gen.edge.v1.EdgeVersion; |
|||
|
|||
import java.util.Comparator; |
|||
|
|||
public class EdgeVersionComparator implements Comparator<EdgeVersion> { |
|||
|
|||
public static final EdgeVersionComparator INSTANCE = new EdgeVersionComparator(); |
|||
|
|||
@Override |
|||
public int compare(EdgeVersion v1, EdgeVersion v2) { |
|||
if (v1 == v2) { |
|||
return 0; |
|||
} |
|||
// UNRECOGNIZED is less than any other version
|
|||
if (v1 == EdgeVersion.UNRECOGNIZED) { |
|||
return -1; |
|||
} |
|||
if (v2 == EdgeVersion.UNRECOGNIZED) { |
|||
return 1; |
|||
} |
|||
// V_LATEST is treated as the newest version
|
|||
if (v1 == EdgeVersion.V_LATEST) { |
|||
v1 = getNewestEdgeVersion(); |
|||
} |
|||
if (v2 == EdgeVersion.V_LATEST) { |
|||
v2 = getNewestEdgeVersion(); |
|||
} |
|||
return compareVersionParts(parseVersionParts(v1), parseVersionParts(v2)); |
|||
} |
|||
|
|||
public static EdgeVersion getNewestEdgeVersion() { |
|||
EdgeVersion newest = null; |
|||
for (EdgeVersion v : EdgeVersion.values()) { |
|||
if (v == EdgeVersion.V_LATEST || v == EdgeVersion.UNRECOGNIZED) { |
|||
continue; |
|||
} |
|||
if (newest == null || INSTANCE.compare(v, newest) > 0) { |
|||
newest = v; |
|||
} |
|||
} |
|||
return newest; |
|||
} |
|||
|
|||
private static int[] parseVersionParts(EdgeVersion version) { |
|||
String name = version.name(); |
|||
if (name.startsWith("V_")) { |
|||
name = name.substring(2); |
|||
} |
|||
String[] parts = name.split("_"); |
|||
int[] result = new int[parts.length]; |
|||
for (int i = 0; i < parts.length; i++) { |
|||
result[i] = Integer.parseInt(parts[i]); |
|||
} |
|||
return result; |
|||
} |
|||
|
|||
private static int compareVersionParts(int[] a, int[] b) { |
|||
int maxLen = Math.max(a.length, b.length); |
|||
for (int i = 0; i < maxLen; i++) { |
|||
int partA = i < a.length ? a[i] : 0; |
|||
int partB = i < b.length ? b[i] : 0; |
|||
if (partA != partB) { |
|||
return Integer.compare(partA, partB); |
|||
} |
|||
} |
|||
return 0; |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,102 @@ |
|||
/** |
|||
* 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 org.junit.jupiter.api.Test; |
|||
import org.thingsboard.server.gen.edge.v1.EdgeVersion; |
|||
|
|||
import static org.assertj.core.api.Assertions.assertThat; |
|||
|
|||
class EdgeVersionComparatorTest { |
|||
|
|||
@Test |
|||
void compare_sameVersion_returnsZero() { |
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.V_3_3_0, EdgeVersion.V_3_3_0)).isEqualTo(0); |
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.V_4_0_0, EdgeVersion.V_4_0_0)).isEqualTo(0); |
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.V_4_2_1_2, EdgeVersion.V_4_2_1_2)).isEqualTo(0); |
|||
} |
|||
|
|||
@Test |
|||
void compare_majorVersionDifference_returnsCorrectOrder() { |
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.V_3_3_0, EdgeVersion.V_4_0_0)).isLessThan(0); |
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.V_4_0_0, EdgeVersion.V_3_3_0)).isGreaterThan(0); |
|||
} |
|||
|
|||
@Test |
|||
void compare_minorVersionDifference_returnsCorrectOrder() { |
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.V_3_3_0, EdgeVersion.V_3_6_0)).isLessThan(0); |
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.V_3_6_0, EdgeVersion.V_3_3_0)).isGreaterThan(0); |
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.V_4_0_0, EdgeVersion.V_4_1_0)).isLessThan(0); |
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.V_4_1_0, EdgeVersion.V_4_0_0)).isGreaterThan(0); |
|||
} |
|||
|
|||
@Test |
|||
void compare_patchVersionDifference_returnsCorrectOrder() { |
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.V_3_6_0, EdgeVersion.V_3_6_1)).isLessThan(0); |
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.V_3_6_1, EdgeVersion.V_3_6_0)).isGreaterThan(0); |
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.V_3_6_1, EdgeVersion.V_3_6_2)).isLessThan(0); |
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.V_3_6_2, EdgeVersion.V_3_6_4)).isLessThan(0); |
|||
} |
|||
|
|||
@Test |
|||
void compare_fourPartVersion_returnsCorrectOrder() { |
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.V_4_2_0, EdgeVersion.V_4_2_1_2)).isLessThan(0); |
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.V_4_2_1_2, EdgeVersion.V_4_2_0)).isGreaterThan(0); |
|||
} |
|||
|
|||
@Test |
|||
void compare_threePartVsFourPart_treatsImplicitZero() { |
|||
// V_4_2_0 should be less than V_4_2_1_2 (4.2.0.0 < 4.2.1.2)
|
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.V_4_2_0, EdgeVersion.V_4_2_1_2)).isLessThan(0); |
|||
} |
|||
|
|||
@Test |
|||
void getNewestEdgeVersion_excludesLatestAndUnrecognized() { |
|||
EdgeVersion newest = EdgeVersionComparator.getNewestEdgeVersion(); |
|||
assertThat(newest).isNotNull(); |
|||
assertThat(newest).isNotEqualTo(EdgeVersion.V_LATEST); |
|||
assertThat(newest).isNotEqualTo(EdgeVersion.UNRECOGNIZED); |
|||
} |
|||
|
|||
@Test |
|||
void compare_vLatest_treatedAsNewestVersion() { |
|||
EdgeVersion newest = EdgeVersionComparator.getNewestEdgeVersion(); |
|||
// V_LATEST equals the newest version
|
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.V_LATEST, newest)).isEqualTo(0); |
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(newest, EdgeVersion.V_LATEST)).isEqualTo(0); |
|||
// V_LATEST is greater than older versions
|
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.V_LATEST, EdgeVersion.V_3_3_0)).isGreaterThan(0); |
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.V_3_3_0, EdgeVersion.V_LATEST)).isLessThan(0); |
|||
} |
|||
|
|||
@Test |
|||
void compare_vLatest_withItself_returnsZero() { |
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.V_LATEST, EdgeVersion.V_LATEST)).isEqualTo(0); |
|||
} |
|||
|
|||
@Test |
|||
void compare_unrecognized_isLessThanAnyVersion() { |
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.UNRECOGNIZED, EdgeVersion.V_3_3_0)).isLessThan(0); |
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.UNRECOGNIZED, EdgeVersion.V_LATEST)).isLessThan(0); |
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.V_3_3_0, EdgeVersion.UNRECOGNIZED)).isGreaterThan(0); |
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.V_LATEST, EdgeVersion.UNRECOGNIZED)).isGreaterThan(0); |
|||
} |
|||
|
|||
@Test |
|||
void compare_unrecognized_withItself_returnsZero() { |
|||
assertThat(EdgeVersionComparator.INSTANCE.compare(EdgeVersion.UNRECOGNIZED, EdgeVersion.UNRECOGNIZED)).isEqualTo(0); |
|||
} |
|||
} |
|||
@ -1,154 +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.server.dao.service.timeseries.sql; |
|||
|
|||
import com.google.common.util.concurrent.Futures; |
|||
import com.google.common.util.concurrent.ListenableFuture; |
|||
import com.google.common.util.concurrent.ListeningExecutorService; |
|||
import com.google.common.util.concurrent.MoreExecutors; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.apache.commons.lang3.RandomStringUtils; |
|||
import org.junit.After; |
|||
import org.junit.Assert; |
|||
import org.junit.Before; |
|||
import org.junit.Test; |
|||
import org.springframework.beans.factory.annotation.Autowired; |
|||
import org.thingsboard.common.util.ThingsBoardThreadFactory; |
|||
import org.thingsboard.server.common.data.Tenant; |
|||
import org.thingsboard.server.common.data.id.DeviceId; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.kv.BasicTsKvEntry; |
|||
import org.thingsboard.server.common.data.kv.BooleanDataEntry; |
|||
import org.thingsboard.server.common.data.kv.DoubleDataEntry; |
|||
import org.thingsboard.server.common.data.kv.LongDataEntry; |
|||
import org.thingsboard.server.common.data.kv.StringDataEntry; |
|||
import org.thingsboard.server.common.data.kv.TsKvEntry; |
|||
import org.thingsboard.server.dao.service.AbstractServiceTest; |
|||
import org.thingsboard.server.dao.service.DaoSqlTest; |
|||
import org.thingsboard.server.dao.timeseries.TimeseriesLatestDao; |
|||
|
|||
import java.util.ArrayList; |
|||
import java.util.List; |
|||
import java.util.Random; |
|||
import java.util.UUID; |
|||
import java.util.concurrent.Executors; |
|||
import java.util.concurrent.TimeUnit; |
|||
import java.util.concurrent.atomic.AtomicLong; |
|||
|
|||
@DaoSqlTest |
|||
@Slf4j |
|||
public class LatestTimeseriesPerformanceTest extends AbstractServiceTest { |
|||
|
|||
private static final String STRING_KEY = "stringKey"; |
|||
private static final String LONG_KEY = "longKey"; |
|||
private static final String DOUBLE_KEY = "doubleKey"; |
|||
private static final String BOOLEAN_KEY = "booleanKey"; |
|||
private static final int AMOUNT_OF_UNIQ_KEY = 10000; |
|||
private static final int TIMEOUT = 100; |
|||
|
|||
private final Random random = new Random(); |
|||
|
|||
@Autowired |
|||
private TimeseriesLatestDao timeseriesLatestDao; |
|||
|
|||
private ListeningExecutorService testExecutor; |
|||
|
|||
private EntityId entityId; |
|||
|
|||
private AtomicLong saveCounter; |
|||
|
|||
@Before |
|||
public void before() { |
|||
Tenant tenant = new Tenant(); |
|||
tenant.setTitle("My tenant"); |
|||
Tenant savedTenant = tenantService.saveTenant(tenant); |
|||
Assert.assertNotNull(savedTenant); |
|||
tenantId = savedTenant.getId(); |
|||
entityId = new DeviceId(UUID.randomUUID()); |
|||
saveCounter = new AtomicLong(0); |
|||
testExecutor = MoreExecutors.listeningDecorator(Executors.newFixedThreadPool(200, ThingsBoardThreadFactory.forName(getClass().getSimpleName() + "-test-scope"))); |
|||
} |
|||
|
|||
@After |
|||
public void after() { |
|||
tenantService.deleteTenant(tenantId); |
|||
if (testExecutor != null) { |
|||
testExecutor.shutdownNow(); |
|||
} |
|||
} |
|||
|
|||
@Test |
|||
public void test_save_latest_timeseries() throws Exception { |
|||
warmup(); |
|||
saveCounter.set(0); |
|||
|
|||
long startTime = System.currentTimeMillis(); |
|||
List<ListenableFuture<?>> futures = new ArrayList<>(); |
|||
for (int i = 0; i < 25_000; i++) { |
|||
futures.add(save(generateStrEntry(getRandomKey()))); |
|||
futures.add(save(generateLngEntry(getRandomKey()))); |
|||
futures.add(save(generateDblEntry(getRandomKey()))); |
|||
futures.add(save(generateBoolEntry(getRandomKey()))); |
|||
} |
|||
Futures.allAsList(futures).get(TIMEOUT, TimeUnit.SECONDS); |
|||
long endTime = System.currentTimeMillis(); |
|||
|
|||
long totalTime = endTime - startTime; |
|||
|
|||
log.info("Total time: {}", totalTime); |
|||
log.info("Saved count: {}", saveCounter.get()); |
|||
log.warn("Saved per 1 sec: {}", saveCounter.get() * 1000 / totalTime); |
|||
} |
|||
|
|||
private void warmup() throws Exception { |
|||
List<ListenableFuture<?>> futures = new ArrayList<>(); |
|||
for (int i = 0; i < AMOUNT_OF_UNIQ_KEY; i++) { |
|||
futures.add(save(generateStrEntry(i))); |
|||
futures.add(save(generateLngEntry(i))); |
|||
futures.add(save(generateDblEntry(i))); |
|||
futures.add(save(generateBoolEntry(i))); |
|||
} |
|||
Futures.allAsList(futures).get(TIMEOUT, TimeUnit.SECONDS); |
|||
} |
|||
|
|||
private ListenableFuture<?> save(TsKvEntry tsKvEntry) { |
|||
return Futures.transformAsync(testExecutor.submit(() -> timeseriesLatestDao.saveLatest(tenantId, entityId, tsKvEntry)), result -> { |
|||
saveCounter.incrementAndGet(); |
|||
return result; |
|||
}, testExecutor); |
|||
} |
|||
|
|||
private TsKvEntry generateStrEntry(int keyIndex) { |
|||
return new BasicTsKvEntry(System.currentTimeMillis(), new StringDataEntry(STRING_KEY + keyIndex, RandomStringUtils.random(10))); |
|||
} |
|||
|
|||
private TsKvEntry generateLngEntry(int keyIndex) { |
|||
return new BasicTsKvEntry(System.currentTimeMillis(), new LongDataEntry(LONG_KEY + keyIndex, random.nextLong())); |
|||
} |
|||
|
|||
private TsKvEntry generateDblEntry(int keyIndex) { |
|||
return new BasicTsKvEntry(System.currentTimeMillis(), new DoubleDataEntry(DOUBLE_KEY + keyIndex, random.nextDouble())); |
|||
} |
|||
|
|||
private TsKvEntry generateBoolEntry(int keyIndex) { |
|||
return new BasicTsKvEntry(System.currentTimeMillis(), new BooleanDataEntry(BOOLEAN_KEY + keyIndex, random.nextBoolean())); |
|||
} |
|||
|
|||
private int getRandomKey() { |
|||
return random.nextInt(AMOUNT_OF_UNIQ_KEY); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,42 @@ |
|||
/** |
|||
* 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.msa; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
|
|||
import java.io.IOException; |
|||
import java.net.DatagramSocket; |
|||
import java.net.SocketException; |
|||
|
|||
@Slf4j |
|||
public class PortFinder { |
|||
public static int findAvailableUdpPort() { |
|||
try (DatagramSocket socket = new DatagramSocket(0)) { |
|||
return socket.getLocalPort(); |
|||
} catch (SocketException e) { |
|||
throw new IllegalStateException("No available UDP ports found", e); |
|||
} |
|||
} |
|||
|
|||
public static boolean isUDPPortAvailable(int port) { |
|||
try (DatagramSocket socket = new DatagramSocket(port)) { |
|||
return true; |
|||
} catch (IOException e) { |
|||
log.debug("Failed to open UDP port {}", port, e); |
|||
return false; |
|||
} |
|||
} |
|||
} |
|||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue