tsData) throws InterruptedException {
CountDownLatch latch = new CountDownLatch(1);
tsService.saveTimeseries(TimeseriesSaveRequest.builder()
diff --git a/application/src/test/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSslTest.java b/application/src/test/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSslTest.java
new file mode 100644
index 0000000000..04a1b87833
--- /dev/null
+++ b/application/src/test/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSslTest.java
@@ -0,0 +1,272 @@
+/**
+ * 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 io.grpc.ManagedChannel;
+import io.grpc.Server;
+import io.grpc.netty.shaded.io.grpc.netty.GrpcSslContexts;
+import io.grpc.netty.shaded.io.grpc.netty.NettyChannelBuilder;
+import io.grpc.netty.shaded.io.grpc.netty.NettyServerBuilder;
+import org.bouncycastle.asn1.x500.X500Name;
+import org.bouncycastle.cert.jcajce.JcaX509CertificateConverter;
+import org.bouncycastle.cert.jcajce.JcaX509v3CertificateBuilder;
+import org.bouncycastle.jce.provider.BouncyCastleProvider;
+import org.bouncycastle.openssl.jcajce.JcaPEMWriter;
+import org.bouncycastle.openssl.jcajce.JcePEMEncryptorBuilder;
+import org.bouncycastle.operator.jcajce.JcaContentSignerBuilder;
+import org.bouncycastle.util.io.pem.PemObject;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.EnumSource;
+import org.springframework.test.util.ReflectionTestUtils;
+import org.thingsboard.server.controller.AbstractWebTest;
+import org.thingsboard.server.gen.edge.v1.EdgeRpcServiceGrpc;
+
+import java.io.ByteArrayInputStream;
+import java.math.BigInteger;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.security.KeyPair;
+import java.security.KeyPairGenerator;
+import java.security.PrivateKey;
+import java.security.Security;
+import java.security.cert.X509Certificate;
+import java.security.spec.ECGenParameterSpec;
+import java.util.ArrayList;
+import java.util.Date;
+import java.util.List;
+import java.util.concurrent.TimeUnit;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+import static org.awaitility.Awaitility.await;
+
+/**
+ * Tests for Edge gRPC SSL setup using the production {@link EdgeGrpcService#setupSsl} method.
+ *
+ * Covers:
+ * 1. Separate cert and key PEM inputs
+ * 2. Combined PEM (cert + key in one file)
+ * 3. Encrypted private key with password
+ * 4. Missing key in combined PEM → error
+ *
+ * Each scenario is parameterized across key types: RSA-2048, RSA-4096, EC P-256, EC P-384.
+ */
+class EdgeGrpcSslTest {
+
+ static {
+ if (Security.getProvider(BouncyCastleProvider.PROVIDER_NAME) == null) {
+ Security.addProvider(new BouncyCastleProvider());
+ }
+ }
+
+ enum KeyType {
+ RSA_2048("RSA", 2048, null, "SHA256withRSA"),
+ RSA_4096("RSA", 4096, null, "SHA256withRSA"),
+ EC_P256("EC", 256, "secp256r1", "SHA256withECDSA"),
+ EC_P384("EC", 384, "secp384r1", "SHA384withECDSA");
+
+ final String algorithm;
+ final int size;
+ final String curve;
+ final String sigAlg;
+
+ KeyType(String algorithm, int size, String curve, String sigAlg) {
+ this.algorithm = algorithm;
+ this.size = size;
+ this.curve = curve;
+ this.sigAlg = sigAlg;
+ }
+
+ KeyPair generateKeyPair() throws Exception {
+ KeyPairGenerator kpg = KeyPairGenerator.getInstance(algorithm);
+ if (curve != null) {
+ kpg.initialize(new ECGenParameterSpec(curve));
+ } else {
+ kpg.initialize(size);
+ }
+ return kpg.generateKeyPair();
+ }
+ }
+
+ private final List tempFiles = new ArrayList<>();
+ private Server server;
+ private ManagedChannel channel;
+
+ @AfterEach
+ void cleanup() throws Exception {
+ if (channel != null) {
+ channel.shutdownNow().awaitTermination(2, TimeUnit.SECONDS);
+ }
+ if (server != null) {
+ server.shutdownNow().awaitTermination(2, TimeUnit.SECONDS);
+ }
+ for (Path p : tempFiles) {
+ Files.deleteIfExists(p);
+ }
+ }
+
+ @ParameterizedTest(name = "separateCertAndKey_{0}")
+ @EnumSource(KeyType.class)
+ void separateCertAndKey(KeyType keyType) throws Exception {
+ KeyPair kp = keyType.generateKeyPair();
+ X509Certificate cert = generateSelfSignedCert(kp, keyType.sigAlg);
+
+ Path certFile = writeTempPem("cert", cert);
+ Path keyFile = writeTempPem("key", kp.getPrivate());
+
+ server = startServer(certFile.toString(), keyFile.toString(), null);
+ assertTlsConnectivity(cert);
+ }
+
+ @ParameterizedTest(name = "combinedPemWithCertAndKey_{0}")
+ @EnumSource(KeyType.class)
+ void combinedPemWithCertAndKey(KeyType keyType) throws Exception {
+ KeyPair kp = keyType.generateKeyPair();
+ X509Certificate cert = generateSelfSignedCert(kp, keyType.sigAlg);
+
+ Path combinedFile = writeTempPem("combined", cert, kp.getPrivate());
+
+ server = startServer(combinedFile.toString(), "", null);
+ assertTlsConnectivity(cert);
+ }
+
+ // RSA-only: BouncyCastle writes encrypted EC keys in traditional PEM format (BEGIN EC PRIVATE KEY),
+ // which after decryption produces a PEMKeyPair without public key info — causing PemSslCredentials
+ // to fail with "Cannot invoke SubjectPublicKeyInfo.getEncoded() because getPublicKeyInfo() is null".
+ @ParameterizedTest(name = "encryptedPrivateKey_{0}")
+ @EnumSource(value = KeyType.class, names = {"RSA_2048", "RSA_4096"})
+ void encryptedPrivateKey(KeyType keyType) throws Exception {
+ KeyPair kp = keyType.generateKeyPair();
+ X509Certificate cert = generateSelfSignedCert(kp, keyType.sigAlg);
+ String password = "test-password";
+
+ Path combinedFile = writeTempPemEncrypted("enc-combined", password, cert, kp.getPrivate());
+
+ server = startServer(combinedFile.toString(), "", password);
+ assertTlsConnectivity(cert);
+ }
+
+ @ParameterizedTest(name = "combinedPemWithCertOnly_throwsException_{0}")
+ @EnumSource(KeyType.class)
+ void combinedPemWithCertOnly_throwsException(KeyType keyType) throws Exception {
+ KeyPair kp = keyType.generateKeyPair();
+ X509Certificate cert = generateSelfSignedCert(kp, keyType.sigAlg);
+
+ Path certOnlyFile = writeTempPem("cert-only", cert);
+
+ assertThatThrownBy(() -> startServer(certOnlyFile.toString(), "", null))
+ .isInstanceOf(IllegalArgumentException.class);
+ }
+
+ // --- Server startup using production EdgeGrpcService.setupSsl() ---
+
+ private Server startServer(String certFileResource, String privateKeyResource, String keyPassword) throws Exception {
+ EdgeGrpcService edgeGrpcService = new EdgeGrpcService();
+ ReflectionTestUtils.setField(edgeGrpcService, "certFileResource", certFileResource);
+ ReflectionTestUtils.setField(edgeGrpcService, "privateKeyResource", privateKeyResource);
+ ReflectionTestUtils.setField(edgeGrpcService, "keyPassword", keyPassword != null ? keyPassword : "");
+
+ NettyServerBuilder builder = NettyServerBuilder.forPort(0)
+ .addService(new EdgeRpcServiceGrpc.EdgeRpcServiceImplBase() {});
+
+ edgeGrpcService.setupSsl(builder);
+
+ return builder.build().start();
+ }
+
+ private void assertTlsConnectivity(X509Certificate trustedCert) throws Exception {
+ String certPem = toPem(trustedCert);
+ var clientSsl = GrpcSslContexts.forClient()
+ .trustManager(new ByteArrayInputStream(certPem.getBytes(StandardCharsets.UTF_8)))
+ .build();
+
+ channel = NettyChannelBuilder.forAddress("localhost", server.getPort())
+ .sslContext(clientSsl)
+ .build();
+
+ channel.getState(true); // trigger connection attempt
+ await().atMost(AbstractWebTest.TIMEOUT, TimeUnit.SECONDS)
+ .pollInterval(50, TimeUnit.MILLISECONDS)
+ .untilAsserted(() -> {
+ var state = channel.getState(false);
+ if (state == io.grpc.ConnectivityState.TRANSIENT_FAILURE) {
+ throw new AssertionError("TLS handshake failed: channel in TRANSIENT_FAILURE");
+ }
+ assertThat(state).isEqualTo(io.grpc.ConnectivityState.READY);
+ });
+ }
+
+ // --- Cert/key generation ---
+
+ private X509Certificate generateSelfSignedCert(KeyPair kp, String sigAlg) throws Exception {
+ X500Name subject = new X500Name("CN=localhost");
+ Date now = new Date();
+ return new JcaX509CertificateConverter().getCertificate(
+ new JcaX509v3CertificateBuilder(
+ subject, BigInteger.ONE, now,
+ new Date(now.getTime() + TimeUnit.DAYS.toMillis(1)),
+ subject, kp.getPublic())
+ .build(new JcaContentSignerBuilder(sigAlg).build(kp.getPrivate())));
+ }
+
+ // --- PEM file helpers ---
+
+ private String toPem(Object obj) throws Exception {
+ java.io.StringWriter sw = new java.io.StringWriter();
+ try (JcaPEMWriter w = new JcaPEMWriter(sw)) {
+ w.writeObject(obj);
+ }
+ return sw.toString();
+ }
+
+ private Path writeTempPem(String prefix, Object... objects) throws Exception {
+ Path p = Files.createTempFile(prefix + "-", ".pem");
+ tempFiles.add(p);
+ try (JcaPEMWriter w = new JcaPEMWriter(Files.newBufferedWriter(p))) {
+ for (Object o : objects) {
+ w.writeObject(toPkcs8IfKey(o));
+ }
+ }
+ return p;
+ }
+
+ private Path writeTempPemEncrypted(String prefix, String password, Object... objects) throws Exception {
+ Path p = Files.createTempFile(prefix + "-", ".pem");
+ tempFiles.add(p);
+ var encryptor = new JcePEMEncryptorBuilder("AES-256-CBC")
+ .setProvider(BouncyCastleProvider.PROVIDER_NAME)
+ .build(password.toCharArray());
+ try (JcaPEMWriter w = new JcaPEMWriter(Files.newBufferedWriter(p))) {
+ for (Object o : objects) {
+ if (o instanceof PrivateKey) {
+ w.writeObject(o, encryptor);
+ } else {
+ w.writeObject(o);
+ }
+ }
+ }
+ return p;
+ }
+
+ private Object toPkcs8IfKey(Object o) {
+ if (o instanceof PrivateKey pk) {
+ return new PemObject("PRIVATE KEY", pk.getEncoded());
+ }
+ return o;
+ }
+}
diff --git a/application/src/test/java/org/thingsboard/server/service/entitiy/EntityServiceTest.java b/application/src/test/java/org/thingsboard/server/service/entitiy/EntityServiceTest.java
index c54ffe4a07..05ab1165f6 100644
--- a/application/src/test/java/org/thingsboard/server/service/entitiy/EntityServiceTest.java
+++ b/application/src/test/java/org/thingsboard/server/service/entitiy/EntityServiceTest.java
@@ -115,6 +115,8 @@ import java.util.Map;
import java.util.Random;
import java.util.UUID;
import java.util.concurrent.ExecutionException;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
import java.util.stream.Collectors;
import java.util.stream.Stream;
@@ -1749,13 +1751,13 @@ public class EntityServiceTest extends AbstractControllerTest {
}
@Test
- public void testFindTenantTelemetry() {
+ public void testFindTenantTelemetry() throws ExecutionException, InterruptedException, TimeoutException {
// save timeseries by sys admin
BasicTsKvEntry timeseries = new BasicTsKvEntry(42L, new DoubleDataEntry("temperature", 45.5));
- timeseriesService.save(TenantId.SYS_TENANT_ID, tenantId, timeseries);
+ timeseriesService.save(TenantId.SYS_TENANT_ID, tenantId, timeseries).get(TIMEOUT, TimeUnit.SECONDS);
AttributeKvEntry attr = new BaseAttributeKvEntry(new LongDataEntry("attr", 10L), 42L);
- attributesService.save(TenantId.SYS_TENANT_ID, tenantId, SERVER_SCOPE, List.of(attr));
+ attributesService.save(TenantId.SYS_TENANT_ID, tenantId, SERVER_SCOPE, List.of(attr)).get(TIMEOUT, TimeUnit.SECONDS);
SingleEntityFilter singleEntityFilter = new SingleEntityFilter();
singleEntityFilter.setSingleEntity(AliasEntityId.fromEntityId(tenantId));
diff --git a/application/src/test/java/org/thingsboard/server/service/housekeeper/HousekeeperServiceTest.java b/application/src/test/java/org/thingsboard/server/service/housekeeper/HousekeeperServiceTest.java
index cf37afad5c..0c12767b60 100644
--- a/application/src/test/java/org/thingsboard/server/service/housekeeper/HousekeeperServiceTest.java
+++ b/application/src/test/java/org/thingsboard/server/service/housekeeper/HousekeeperServiceTest.java
@@ -23,8 +23,8 @@ import org.junit.Test;
import org.mockito.ArgumentMatcher;
import org.mockito.Mockito;
import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.boot.test.mock.mockito.SpyBean;
import org.springframework.test.context.TestPropertySource;
+import org.springframework.test.context.bean.override.mockito.MockitoSpyBean;
import org.testcontainers.shaded.org.apache.commons.lang3.RandomStringUtils;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.metadata.TbGetAttributesNode;
@@ -127,10 +127,12 @@ import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.
})
public class HousekeeperServiceTest extends AbstractControllerTest {
- @SpyBean
+ @MockitoSpyBean
private HousekeeperService housekeeperService;
- @SpyBean
+ @MockitoSpyBean
private HousekeeperReprocessingService housekeeperReprocessingService;
+ @MockitoSpyBean
+ private TsHistoryDeletionTaskProcessor tsHistoryDeletionTaskProcessor;
@Autowired
private EventService eventService;
@Autowired
@@ -153,8 +155,6 @@ public class HousekeeperServiceTest extends AbstractControllerTest {
private CustomerService customerService;
@Autowired
private DashboardService dashboardService;
- @SpyBean
- private TsHistoryDeletionTaskProcessor tsHistoryDeletionTaskProcessor;
private TenantId tenantId;
diff --git a/application/src/test/java/org/thingsboard/server/service/ttl/NotificationsCleanUpServiceTest.java b/application/src/test/java/org/thingsboard/server/service/ttl/NotificationsCleanUpServiceTest.java
new file mode 100644
index 0000000000..d74652325a
--- /dev/null
+++ b/application/src/test/java/org/thingsboard/server/service/ttl/NotificationsCleanUpServiceTest.java
@@ -0,0 +1,144 @@
+/**
+ * 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.ttl;
+
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.Mock;
+import org.mockito.junit.jupiter.MockitoExtension;
+import org.springframework.test.util.ReflectionTestUtils;
+import org.thingsboard.server.common.data.id.TenantId;
+import org.thingsboard.server.common.data.page.PageData;
+import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
+import org.thingsboard.server.dao.notification.NotificationRequestDao;
+import org.thingsboard.server.dao.sqlts.insert.sql.SqlPartitioningRepository;
+import org.thingsboard.server.dao.tenant.TenantService;
+import org.thingsboard.server.queue.discovery.PartitionService;
+
+import java.util.List;
+import java.util.UUID;
+
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyInt;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+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 NotificationsCleanUpServiceTest {
+
+ @Mock
+ private PartitionService partitionService;
+ @Mock
+ private SqlPartitioningRepository partitioningRepository;
+ @Mock
+ private NotificationRequestDao notificationRequestDao;
+ @Mock
+ private TenantService tenantService;
+
+ private NotificationsCleanUpService cleanUpService;
+
+ private static final int BATCH_SIZE = 3;
+
+ @BeforeEach
+ public void setUp() {
+ cleanUpService = new NotificationsCleanUpService(partitionService, partitioningRepository, notificationRequestDao, tenantService);
+ ReflectionTestUtils.setField(cleanUpService, "ttlInSec", 2592000L);
+ ReflectionTestUtils.setField(cleanUpService, "partitionSizeInHours", 168);
+ ReflectionTestUtils.setField(cleanUpService, "removalBatchSize", BATCH_SIZE);
+ }
+
+ @Test
+ public void testBatchLoopCallsDaoMultipleTimes() {
+ TopicPartitionInfo myPartition = TopicPartitionInfo.builder().topic("tb_core").myPartition(true).build();
+ when(partitionService.resolve(any(), any(), any())).thenReturn(myPartition);
+ when(partitioningRepository.dropPartitionsBefore(anyString(), anyLong(), anyLong()))
+ .thenReturn(System.currentTimeMillis());
+
+ TenantId tenantId = TenantId.fromUUID(UUID.randomUUID());
+ when(tenantService.findTenantsIds(any()))
+ .thenReturn(new PageData<>(List.of(tenantId), 1, 1, false));
+
+ // Sysadmin: returns 3 (full batch), then 1 (partial) -> 2 calls
+ when(notificationRequestDao.removeByTenantIdAndCreatedTimeBeforeBatch(eq(TenantId.SYS_TENANT_ID), anyLong(), eq(BATCH_SIZE)))
+ .thenReturn(BATCH_SIZE)
+ .thenReturn(1);
+ // Tenant: returns 3, 3, 0 -> 3 calls
+ when(notificationRequestDao.removeByTenantIdAndCreatedTimeBeforeBatch(eq(tenantId), anyLong(), eq(BATCH_SIZE)))
+ .thenReturn(BATCH_SIZE)
+ .thenReturn(BATCH_SIZE)
+ .thenReturn(0);
+
+ cleanUpService.cleanUp();
+
+ verify(notificationRequestDao, times(2))
+ .removeByTenantIdAndCreatedTimeBeforeBatch(eq(TenantId.SYS_TENANT_ID), anyLong(), eq(BATCH_SIZE));
+ verify(notificationRequestDao, times(3))
+ .removeByTenantIdAndCreatedTimeBeforeBatch(eq(tenantId), anyLong(), eq(BATCH_SIZE));
+ }
+
+ @Test
+ public void testSkipsTenantNotOnMyPartition() {
+ TopicPartitionInfo myPartition = TopicPartitionInfo.builder().topic("tb_core").myPartition(true).build();
+ TopicPartitionInfo notMyPartition = TopicPartitionInfo.builder().topic("tb_core").myPartition(false).build();
+ when(partitionService.resolve(any(), eq(TenantId.SYS_TENANT_ID), eq(TenantId.SYS_TENANT_ID)))
+ .thenReturn(myPartition);
+ when(partitioningRepository.dropPartitionsBefore(anyString(), anyLong(), anyLong()))
+ .thenReturn(System.currentTimeMillis());
+
+ // Sysadmin: no records
+ when(notificationRequestDao.removeByTenantIdAndCreatedTimeBeforeBatch(eq(TenantId.SYS_TENANT_ID), anyLong(), eq(BATCH_SIZE)))
+ .thenReturn(0);
+
+ TenantId myTenant = TenantId.fromUUID(UUID.randomUUID());
+ TenantId otherTenant = TenantId.fromUUID(UUID.randomUUID());
+ when(tenantService.findTenantsIds(any()))
+ .thenReturn(new PageData<>(List.of(myTenant, otherTenant), 2, 1, false));
+ when(partitionService.resolve(any(), eq(myTenant), eq(myTenant))).thenReturn(myPartition);
+ when(partitionService.resolve(any(), eq(otherTenant), eq(otherTenant))).thenReturn(notMyPartition);
+
+ when(notificationRequestDao.removeByTenantIdAndCreatedTimeBeforeBatch(eq(myTenant), anyLong(), eq(BATCH_SIZE)))
+ .thenReturn(0);
+
+ cleanUpService.cleanUp();
+
+ verify(notificationRequestDao).removeByTenantIdAndCreatedTimeBeforeBatch(eq(myTenant), anyLong(), eq(BATCH_SIZE));
+ verify(notificationRequestDao, never()).removeByTenantIdAndCreatedTimeBeforeBatch(eq(otherTenant), anyLong(), anyInt());
+ }
+
+ @Test
+ public void testNoPartitionsDropped_stillCleansUpRequests() {
+ TopicPartitionInfo myPartition = TopicPartitionInfo.builder().topic("tb_core").myPartition(true).build();
+ when(partitionService.resolve(any(), any(), any())).thenReturn(myPartition);
+ when(partitioningRepository.dropPartitionsBefore(anyString(), anyLong(), anyLong()))
+ .thenReturn(0L);
+
+ when(notificationRequestDao.removeByTenantIdAndCreatedTimeBeforeBatch(eq(TenantId.SYS_TENANT_ID), anyLong(), eq(BATCH_SIZE)))
+ .thenReturn(0);
+ when(tenantService.findTenantsIds(any()))
+ .thenReturn(new PageData<>(List.of(), 0, 0, false));
+
+ cleanUpService.cleanUp();
+
+ verify(notificationRequestDao).removeByTenantIdAndCreatedTimeBeforeBatch(eq(TenantId.SYS_TENANT_ID), anyLong(), eq(BATCH_SIZE));
+ }
+
+}
diff --git a/application/src/test/java/org/thingsboard/server/service/ttl/rpc/RpcCleanUpServiceTest.java b/application/src/test/java/org/thingsboard/server/service/ttl/rpc/RpcCleanUpServiceTest.java
new file mode 100644
index 0000000000..22cc1be1dd
--- /dev/null
+++ b/application/src/test/java/org/thingsboard/server/service/ttl/rpc/RpcCleanUpServiceTest.java
@@ -0,0 +1,136 @@
+/**
+ * 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.ttl.rpc;
+
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.Mock;
+import org.mockito.junit.jupiter.MockitoExtension;
+import org.springframework.test.util.ReflectionTestUtils;
+import org.thingsboard.server.common.data.TenantProfile;
+import org.thingsboard.server.common.data.id.TenantId;
+import org.thingsboard.server.common.data.page.PageData;
+import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
+import org.thingsboard.server.common.data.tenant.profile.TenantProfileData;
+import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
+import org.thingsboard.server.dao.rpc.RpcDao;
+import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
+import org.thingsboard.server.dao.tenant.TenantService;
+import org.thingsboard.server.queue.discovery.PartitionService;
+
+import java.util.List;
+import java.util.UUID;
+
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyInt;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.eq;
+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 RpcCleanUpServiceTest {
+
+ @Mock
+ private PartitionService partitionService;
+ @Mock
+ private RpcDao rpcDao;
+ @Mock
+ private TenantService tenantService;
+ @Mock
+ private TbTenantProfileCache tenantProfileCache;
+
+ private RpcCleanUpService cleanUpService;
+
+ private static final int BATCH_SIZE = 3;
+
+ @BeforeEach
+ public void setUp() {
+ cleanUpService = new RpcCleanUpService(tenantService, partitionService, tenantProfileCache, rpcDao);
+ ReflectionTestUtils.setField(cleanUpService, "removalBatchSize", BATCH_SIZE);
+ }
+
+ @Test
+ public void testBatchLoopCallsDaoMultipleTimes() {
+ TenantId tenantId = TenantId.fromUUID(UUID.randomUUID());
+ setupTenant(tenantId, 7);
+
+ // Returns 3 (full batch), 3 (full batch), 1 (partial) -> 3 calls
+ when(rpcDao.deleteOutdatedRpcByTenantIdBatch(eq(tenantId), anyLong(), eq(BATCH_SIZE)))
+ .thenReturn(BATCH_SIZE)
+ .thenReturn(BATCH_SIZE)
+ .thenReturn(1);
+
+ cleanUpService.cleanUp();
+
+ verify(rpcDao, times(3)).deleteOutdatedRpcByTenantIdBatch(eq(tenantId), anyLong(), eq(BATCH_SIZE));
+ }
+
+ @Test
+ public void testSkipsTenantNotOnMyPartition() {
+ TenantId myTenant = TenantId.fromUUID(UUID.randomUUID());
+ TenantId otherTenant = TenantId.fromUUID(UUID.randomUUID());
+
+ TopicPartitionInfo myPartition = TopicPartitionInfo.builder().topic("tb_core").myPartition(true).build();
+ TopicPartitionInfo notMyPartition = TopicPartitionInfo.builder().topic("tb_core").myPartition(false).build();
+
+ when(tenantService.findTenantsIds(any()))
+ .thenReturn(new PageData<>(List.of(myTenant, otherTenant), 2, 1, false));
+ when(partitionService.resolve(any(), eq(myTenant), eq(myTenant))).thenReturn(myPartition);
+ when(partitionService.resolve(any(), eq(otherTenant), eq(otherTenant))).thenReturn(notMyPartition);
+
+ setupTenantProfile(myTenant, 7);
+ when(rpcDao.deleteOutdatedRpcByTenantIdBatch(eq(myTenant), anyLong(), eq(BATCH_SIZE)))
+ .thenReturn(0);
+
+ cleanUpService.cleanUp();
+
+ verify(rpcDao).deleteOutdatedRpcByTenantIdBatch(eq(myTenant), anyLong(), eq(BATCH_SIZE));
+ verify(rpcDao, never()).deleteOutdatedRpcByTenantIdBatch(eq(otherTenant), anyLong(), anyInt());
+ }
+
+ @Test
+ public void testSkipsTenantWithZeroTtl() {
+ TenantId tenantId = TenantId.fromUUID(UUID.randomUUID());
+ setupTenant(tenantId, 0);
+
+ cleanUpService.cleanUp();
+
+ verify(rpcDao, never()).deleteOutdatedRpcByTenantIdBatch(any(), anyLong(), anyInt());
+ }
+
+ private void setupTenant(TenantId tenantId, int rpcTtlDays) {
+ TopicPartitionInfo myPartition = TopicPartitionInfo.builder().topic("tb_core").myPartition(true).build();
+ when(partitionService.resolve(any(), eq(tenantId), eq(tenantId))).thenReturn(myPartition);
+ when(tenantService.findTenantsIds(any()))
+ .thenReturn(new PageData<>(List.of(tenantId), 1, 1, false));
+ setupTenantProfile(tenantId, rpcTtlDays);
+ }
+
+ private void setupTenantProfile(TenantId tenantId, int rpcTtlDays) {
+ TenantProfile profile = new TenantProfile();
+ TenantProfileData profileData = new TenantProfileData();
+ DefaultTenantProfileConfiguration config = new DefaultTenantProfileConfiguration();
+ config.setRpcTtlDays(rpcTtlDays);
+ profileData.setConfiguration(config);
+ profile.setProfileData(profileData);
+ when(tenantProfileCache.get(tenantId)).thenReturn(profile);
+ }
+
+}
diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/FwLwM2MDevice.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/FwLwM2MDevice.java
index 9bed9bc483..35e12691c9 100644
--- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/FwLwM2MDevice.java
+++ b/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/FwLwM2MDevice.java
@@ -171,7 +171,14 @@ public class FwLwM2MDevice extends BaseInstanceEnabler implements Destroyable {
if (this.leshanClient != null) {
log.info("Stop/reboot LwM2M client {}", this.leshanClient.getEndpoint(identity));
- this.leshanClient.stop(false);
+ try {
+ this.leshanClient.stop(false);
+ } catch (Exception stopEx) {
+ // Leshan may throw NPE during CoAP observe-relation cleanup when the server
+ // reference is null (race condition in NotificationDataStore.toKey()).
+ // The client is still considered stopped at this point — proceed with restart.
+ log.warn("Exception during LwM2M client stop, proceeding with restart: {}", stopEx.getMessage());
+ }
log.info("Start after update fw LwM2M client {}", this.leshanClient.getEndpoint(identity));
this.leshanClient.start();
@@ -193,7 +200,7 @@ public class FwLwM2MDevice extends BaseInstanceEnabler implements Destroyable {
} catch (Exception e) {
log.error("Error during firmware update", e);
}
- }, 0, TimeUnit.SECONDS); // start immediately, without further delay
+ }, 1, TimeUnit.SECONDS); // delay 1 sec to allow CoAP Execute response to be delivered before client stops
}
protected void setLeshanClient(LeshanClient leshanClient) {
diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationReadCollectedValueTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationReadCollectedValueTest.java
index 432fac6d1a..c61ca0b123 100644
--- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationReadCollectedValueTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationReadCollectedValueTest.java
@@ -21,6 +21,7 @@ import com.fasterxml.jackson.databind.node.ObjectNode;
import lombok.extern.slf4j.Slf4j;
import org.junit.Before;
import org.junit.Test;
+import org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper;
import org.thingsboard.server.transport.lwm2m.rpc.AbstractRpcLwM2MIntegrationTest;
import java.util.concurrent.atomic.AtomicReference;
import static java.util.concurrent.TimeUnit.SECONDS;
@@ -37,6 +38,12 @@ import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID
@Slf4j
public class RpcLwm2mIntegrationReadCollectedValueTest extends AbstractRpcLwM2MIntegrationTest {
+ @Before
+ public void resetCollectedValueTimestamps() {
+ Lwm2mTestHelper.RESOURCE_ID_3303_12_5700_TS_0 = 0;
+ Lwm2mTestHelper.RESOURCE_ID_3303_12_5700_TS_1 = 0;
+ }
+
/**
* Read {"id":"/3303/12/5700"}
* Trigger a Send operation from the client with multiple values for the same resource as a payload
diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/AbstractMqttAttributesIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/AbstractMqttAttributesIntegrationTest.java
index a7307b2308..5f43aa7e00 100644
--- a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/AbstractMqttAttributesIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/AbstractMqttAttributesIntegrationTest.java
@@ -420,6 +420,7 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt
Awaitility.await()
.atMost(10, TimeUnit.SECONDS)
+ .ignoreExceptions()
.until(() -> {
List