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/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/AlarmsDeletionHousekeeperTask.java b/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/AlarmsDeletionHousekeeperTask.java
index dea590295e..d66ec046de 100644
--- a/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/AlarmsDeletionHousekeeperTask.java
+++ b/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/AlarmsDeletionHousekeeperTask.java
@@ -23,6 +23,7 @@ import lombok.ToString;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
+import java.io.Serial;
import java.util.List;
import java.util.UUID;
@@ -32,6 +33,9 @@ import java.util.UUID;
@NoArgsConstructor(access = AccessLevel.PROTECTED)
public class AlarmsDeletionHousekeeperTask extends HousekeeperTask {
+ @Serial
+ private static final long serialVersionUID = 9214680001573764374L;
+
private List alarms;
public AlarmsDeletionHousekeeperTask(TenantId tenantId, EntityId entityId) {
diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/AlarmsUnassignHousekeeperTask.java b/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/AlarmsUnassignHousekeeperTask.java
index 0313190056..445850d387 100644
--- a/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/AlarmsUnassignHousekeeperTask.java
+++ b/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/AlarmsUnassignHousekeeperTask.java
@@ -24,6 +24,7 @@ import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.UserId;
+import java.io.Serial;
import java.util.List;
import java.util.UUID;
@@ -33,6 +34,9 @@ import java.util.UUID;
@NoArgsConstructor(access = AccessLevel.PROTECTED)
public class AlarmsUnassignHousekeeperTask extends HousekeeperTask {
+ @Serial
+ private static final long serialVersionUID = 9156667024462937756L;
+
private String userTitle;
private List alarms;
diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/EntitiesDeletionHousekeeperTask.java b/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/EntitiesDeletionHousekeeperTask.java
index fe25a98a1d..c1a233b1c4 100644
--- a/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/EntitiesDeletionHousekeeperTask.java
+++ b/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/EntitiesDeletionHousekeeperTask.java
@@ -23,6 +23,7 @@ import lombok.ToString;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.TenantId;
+import java.io.Serial;
import java.util.List;
import java.util.UUID;
@@ -32,6 +33,9 @@ import java.util.UUID;
@NoArgsConstructor
public class EntitiesDeletionHousekeeperTask extends HousekeeperTask {
+ @Serial
+ private static final long serialVersionUID = 9009068831061529286L;
+
private EntityType entityType;
private List entities;
diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/HousekeeperTask.java b/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/HousekeeperTask.java
index 875ef2765f..2df7cf4dd4 100644
--- a/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/HousekeeperTask.java
+++ b/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/HousekeeperTask.java
@@ -29,6 +29,7 @@ import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
+import java.io.Serial;
import java.io.Serializable;
@JsonIgnoreProperties(ignoreUnknown = true)
@@ -45,6 +46,9 @@ import java.io.Serializable;
@NoArgsConstructor(access = AccessLevel.PROTECTED)
public class HousekeeperTask implements Serializable {
+ @Serial
+ private static final long serialVersionUID = -2585974110832225152L;
+
private TenantId tenantId;
private EntityId entityId;
private HousekeeperTaskType taskType;
diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/LatestTsDeletionHousekeeperTask.java b/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/LatestTsDeletionHousekeeperTask.java
index cd3e94e5c6..931c2931cb 100644
--- a/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/LatestTsDeletionHousekeeperTask.java
+++ b/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/LatestTsDeletionHousekeeperTask.java
@@ -23,12 +23,17 @@ import lombok.ToString;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
+import java.io.Serial;
+
@Data
@ToString(callSuper = true)
@EqualsAndHashCode(callSuper = true)
@NoArgsConstructor(access = AccessLevel.PROTECTED)
public class LatestTsDeletionHousekeeperTask extends HousekeeperTask {
+ @Serial
+ private static final long serialVersionUID = 5193191938513490138L;
+
private String key;
public LatestTsDeletionHousekeeperTask(TenantId tenantId, EntityId entityId, String key) {
diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/TenantEntitiesDeletionHousekeeperTask.java b/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/TenantEntitiesDeletionHousekeeperTask.java
index 443d929d8b..be7ff6f7ec 100644
--- a/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/TenantEntitiesDeletionHousekeeperTask.java
+++ b/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/TenantEntitiesDeletionHousekeeperTask.java
@@ -23,12 +23,17 @@ import lombok.ToString;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.TenantId;
+import java.io.Serial;
+
@Data
@ToString(callSuper = true)
@EqualsAndHashCode(callSuper = true)
@NoArgsConstructor
public class TenantEntitiesDeletionHousekeeperTask extends HousekeeperTask {
+ @Serial
+ private static final long serialVersionUID = -8033108795318393447L;
+
private EntityType entityType;
public TenantEntitiesDeletionHousekeeperTask(TenantId tenantId, EntityType entityType) {
diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/TsHistoryDeletionHousekeeperTask.java b/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/TsHistoryDeletionHousekeeperTask.java
index d9315f0ff4..b520899ca4 100644
--- a/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/TsHistoryDeletionHousekeeperTask.java
+++ b/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/TsHistoryDeletionHousekeeperTask.java
@@ -23,12 +23,17 @@ import lombok.ToString;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
+import java.io.Serial;
+
@Data
@ToString(callSuper = true)
@EqualsAndHashCode(callSuper = true)
@NoArgsConstructor(access = AccessLevel.PROTECTED)
public class TsHistoryDeletionHousekeeperTask extends HousekeeperTask {
+ @Serial
+ private static final long serialVersionUID = 4573851542705079043L;
+
private String key;
public TsHistoryDeletionHousekeeperTask(TenantId tenantId, EntityId entityId, String key) {
diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvEntry.java b/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvEntry.java
index cb4092f433..595e1aa26b 100644
--- a/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvEntry.java
+++ b/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvEntry.java
@@ -37,4 +37,9 @@ public interface TsKvEntry extends KvEntry, HasVersion {
return new TsValue(getTs(), getValueAsString());
}
+ @JsonIgnore
+ default boolean isDeletedEntry() {
+ return getTs() == 0 && (getValue() == null || getValueAsString().isEmpty());
+ }
+
}
diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/NotificationRuleRecipientsConfig.java b/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/NotificationRuleRecipientsConfig.java
index fda73b8059..de4e05a2cb 100644
--- a/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/NotificationRuleRecipientsConfig.java
+++ b/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/NotificationRuleRecipientsConfig.java
@@ -17,6 +17,7 @@ package org.thingsboard.server.common.data.notification.rule;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
+import com.fasterxml.jackson.annotation.JsonProperty;
import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonSubTypes.Type;
import com.fasterxml.jackson.annotation.JsonTypeInfo;
@@ -29,7 +30,7 @@ import java.util.List;
import java.util.Map;
import java.util.UUID;
-@JsonIgnoreProperties
+@JsonIgnoreProperties(ignoreUnknown = true)
@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "triggerType", visible = true, include = JsonTypeInfo.As.EXISTING_PROPERTY, defaultImpl = DefaultNotificationRuleRecipientsConfig.class)
@JsonSubTypes({
@Type(name = "ALARM", value = EscalatedNotificationRuleRecipientsConfig.class),
@@ -38,6 +39,7 @@ import java.util.UUID;
public abstract class NotificationRuleRecipientsConfig implements Serializable {
@NotNull
+ @JsonProperty("triggerType")
private NotificationRuleTriggerType triggerType;
@JsonIgnore
diff --git a/common/edge-api/src/main/proto/edge.proto b/common/edge-api/src/main/proto/edge.proto
index ea45efaf7c..0f12446818 100644
--- a/common/edge-api/src/main/proto/edge.proto
+++ b/common/edge-api/src/main/proto/edge.proto
@@ -47,8 +47,10 @@ enum EdgeVersion {
V_4_3_0 = 13;
V_4_2_1_2 = 14;
V_4_2_2 = 4220;
+ V_4_2_2_1 = 4221;
V_4_3_0_1 = 15;
V_4_3_1 = 4310;
+ V_4_3_1_1 = 4311;
V_LATEST = 99999;
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestDao.java b/dao/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestDao.java
index 96a86073d4..8e1408f4fc 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/notification/NotificationRequestDao.java
@@ -50,7 +50,7 @@ public interface NotificationRequestDao extends Dao {
boolean existsByTenantIdAndStatusAndTemplateId(TenantId tenantId, NotificationRequestStatus status, NotificationTemplateId templateId);
- int removeAllByCreatedTimeBefore(long ts);
+ int removeByTenantIdAndCreatedTimeBeforeBatch(TenantId tenantId, long ts, int batchSize);
NotificationRequestInfo findInfoById(TenantId tenantId, NotificationRequestId id);
diff --git a/dao/src/main/java/org/thingsboard/server/dao/ota/BaseOtaPackageService.java b/dao/src/main/java/org/thingsboard/server/dao/ota/BaseOtaPackageService.java
index 484b3c49b3..2e30c0d598 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/ota/BaseOtaPackageService.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/ota/BaseOtaPackageService.java
@@ -196,7 +196,9 @@ public class BaseOtaPackageService extends AbstractCachedEntityService INCORRECT_OTA_PACKAGE_ID + id);
try {
+ Long oid = getDataOidById(tenantId, otaPackageId);
otaPackageDao.removeById(tenantId, otaPackageId.getId());
+ unlinkDataIfPresent(tenantId, otaPackageId, oid);
publishEvictEvent(new OtaPackageCacheEvictEvent(otaPackageId));
eventPublisher.publishEvent(DeleteEntityEvent.builder().tenantId(tenantId).entityId(otaPackageId).build());
} catch (Exception t) {
@@ -215,6 +217,30 @@ public class BaseOtaPackageService extends AbstractCachedEntityService tenantOtaPackageRemover =
- new PaginatedRemover<>() {
-
- @Override
- protected PageData findEntities(TenantId tenantId, TenantId id, PageLink pageLink) {
- return otaPackageInfoDao.findOtaPackageInfoByTenantId(id, pageLink);
- }
+ private final PaginatedRemover tenantOtaPackageRemover = new PaginatedRemover<>() {
+ @Override
+ protected PageData findEntities(TenantId tenantId, TenantId id, PageLink pageLink) {
+ return otaPackageInfoDao.findOtaPackageInfoByTenantId(id, pageLink);
+ }
- @Override
- protected void removeEntity(TenantId tenantId, OtaPackageInfo entity) {
- deleteOtaPackage(tenantId, entity.getId());
- }
- };
+ @Override
+ protected void removeEntity(TenantId tenantId, OtaPackageInfo entity) {
+ deleteOtaPackage(tenantId, entity.getId());
+ }
+ };
@Override
public Optional> findEntity(TenantId tenantId, EntityId entityId) {
diff --git a/dao/src/main/java/org/thingsboard/server/dao/ota/OtaPackageDao.java b/dao/src/main/java/org/thingsboard/server/dao/ota/OtaPackageDao.java
index c11a13cbe1..875aaea4ed 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/ota/OtaPackageDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/ota/OtaPackageDao.java
@@ -18,15 +18,20 @@ package org.thingsboard.server.dao.ota;
import org.thingsboard.server.common.data.OtaPackage;
import org.thingsboard.server.common.data.id.OtaPackageId;
import org.thingsboard.server.common.data.id.TenantId;
-import org.thingsboard.server.common.data.ota.OtaPackageType;
import org.thingsboard.server.dao.Dao;
import org.thingsboard.server.dao.ExportableEntityDao;
import org.thingsboard.server.dao.TenantEntityWithDataDao;
+import java.util.UUID;
+
public interface OtaPackageDao extends Dao, TenantEntityWithDataDao, ExportableEntityDao {
Long sumDataSizeByTenantId(TenantId tenantId);
OtaPackage findOtaPackageByTenantIdAndTitleAndVersion(TenantId tenantId, String title, String version);
+ Long getDataOidById(UUID id);
+
+ Integer unlinkLargeObject(Long dataOid);
+
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/rpc/RpcDao.java b/dao/src/main/java/org/thingsboard/server/dao/rpc/RpcDao.java
index f88b672a34..37fe2950b4 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/rpc/RpcDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/rpc/RpcDao.java
@@ -24,12 +24,13 @@ import org.thingsboard.server.common.data.rpc.RpcStatus;
import org.thingsboard.server.dao.Dao;
public interface RpcDao extends Dao {
+
PageData findAllByDeviceId(TenantId tenantId, DeviceId deviceId, PageLink pageLink);
PageData findAllByDeviceIdAndStatus(TenantId tenantId, DeviceId deviceId, RpcStatus rpcStatus, PageLink pageLink);
PageData findAllRpcByTenantId(TenantId tenantId, PageLink pageLink);
- int deleteOutdatedRpcByTenantId(TenantId tenantId, Long expirationTime);
+ int deleteOutdatedRpcByTenantIdBatch(TenantId tenantId, Long expirationTime, int batchSize);
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationRequestDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationRequestDao.java
index 9d32e91ca2..f4cc406a75 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationRequestDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationRequestDao.java
@@ -98,8 +98,8 @@ public class JpaNotificationRequestDao extends JpaAbstractDao findIdsByTenantId(@Param("tenantId") UUID tenantId, Pageable pageable);
+ // The 'data' column is of type OID (PostgreSQL large object reference), so it returns the OID as Long
+ @Query(value = "SELECT data FROM ota_package WHERE id = :id AND data IS NOT NULL", nativeQuery = true)
+ Long getDataOidById(@Param("id") UUID id);
+
+ @Transactional
+ @Query(value = "SELECT lo_unlink(:oid)", nativeQuery = true)
+ Integer unlinkLargeObject(@Param("oid") Long oid);
+
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDao.java
index 9daae69d2a..a929b6d1e4 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDao.java
@@ -71,8 +71,8 @@ public class JpaRpcDao extends JpaAbstractDao implements RpcDao,
@Transactional
@Override
- public int deleteOutdatedRpcByTenantId(TenantId tenantId, Long expirationTime) {
- return rpcRepository.deleteOutdatedRpcByTenantId(tenantId.getId(), expirationTime);
+ public int deleteOutdatedRpcByTenantIdBatch(TenantId tenantId, Long expirationTime, int batchSize) {
+ return rpcRepository.deleteOutdatedRpcByTenantIdBatch(tenantId.getId(), expirationTime, batchSize);
}
@Override
diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcRepository.java
index 8a9489333f..3f76170c84 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcRepository.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcRepository.java
@@ -21,20 +21,27 @@ import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.data.jpa.repository.Modifying;
import org.springframework.data.jpa.repository.Query;
import org.springframework.data.repository.query.Param;
+import org.springframework.transaction.annotation.Transactional;
import org.thingsboard.server.common.data.rpc.RpcStatus;
import org.thingsboard.server.dao.model.sql.RpcEntity;
import java.util.UUID;
public interface RpcRepository extends JpaRepository {
+
Page findAllByTenantIdAndDeviceId(UUID tenantId, UUID deviceId, Pageable pageable);
Page findAllByTenantIdAndDeviceIdAndStatus(UUID tenantId, UUID deviceId, RpcStatus status, Pageable pageable);
Page findAllByTenantId(UUID tenantId, Pageable pageable);
+ @Transactional
@Modifying
- @Query(value = "DELETE FROM rpc WHERE tenant_id = :tenantId AND created_time < :expirationTime",
+ @Query(value = "DELETE FROM rpc WHERE id IN " +
+ "(SELECT id FROM rpc WHERE tenant_id = :tenantId AND created_time < :expirationTime LIMIT :batchSize)",
nativeQuery = true)
- int deleteOutdatedRpcByTenantId(@Param("tenantId") UUID tenantId, @Param("expirationTime") Long expirationTime);
+ int deleteOutdatedRpcByTenantIdBatch(@Param("tenantId") UUID tenantId,
+ @Param("expirationTime") Long expirationTime,
+ @Param("batchSize") int batchSize);
+
}
diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/OtaPackageServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/OtaPackageServiceTest.java
index 002a5a7845..69fa3fbc5d 100644
--- a/dao/src/test/java/org/thingsboard/server/dao/service/OtaPackageServiceTest.java
+++ b/dao/src/test/java/org/thingsboard/server/dao/service/OtaPackageServiceTest.java
@@ -36,14 +36,14 @@ import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
import org.thingsboard.server.dao.device.DeviceProfileService;
import org.thingsboard.server.dao.device.DeviceService;
-import org.thingsboard.server.exception.DataValidationException;
+import org.thingsboard.server.dao.ota.OtaPackageDao;
import org.thingsboard.server.dao.ota.OtaPackageService;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.dao.tenant.TenantProfileService;
+import org.thingsboard.server.exception.DataValidationException;
import java.nio.ByteBuffer;
import java.util.ArrayList;
-import java.util.Collections;
import java.util.List;
import static org.assertj.core.api.Assertions.assertThat;
@@ -77,6 +77,8 @@ public class OtaPackageServiceTest extends AbstractServiceTest {
@Autowired
TenantProfileService tenantProfileService;
@Autowired
+ OtaPackageDao otaPackageDao;
+ @Autowired
TbTenantProfileCache tenantProfileCache;
@Before
@@ -118,10 +120,8 @@ public class OtaPackageServiceTest extends AbstractServiceTest {
Assert.assertEquals(1, otaPackageService.sumDataSizeByTenantId(tenantId));
int maxSumDataSize = 8;
- List packages = new ArrayList<>(maxSumDataSize);
-
for (int i = 2; i <= maxSumDataSize; i++) {
- packages.add(createAndSaveFirmware(tenantId, "0." + i));
+ createAndSaveFirmware(tenantId, "0." + i);
Assert.assertEquals(i, otaPackageService.sumDataSizeByTenantId(tenantId));
}
@@ -533,6 +533,39 @@ public class OtaPackageServiceTest extends AbstractServiceTest {
Assert.assertNull(foundFirmware);
}
+ @Test
+ public void testDeleteOtaPackageWithoutData() {
+ OtaPackageInfo firmwareInfo = new OtaPackageInfo();
+ firmwareInfo.setTenantId(tenantId);
+ firmwareInfo.setDeviceProfileId(deviceProfileId);
+ firmwareInfo.setType(FIRMWARE);
+ firmwareInfo.setTitle(TITLE);
+ firmwareInfo.setVersion(VERSION);
+ OtaPackageInfo savedFirmwareInfo = otaPackageService.saveOtaPackageInfo(firmwareInfo, false);
+
+ Assert.assertNotNull(savedFirmwareInfo);
+ Assert.assertNotNull(savedFirmwareInfo.getId());
+
+ // Should not throw NPE when deleting package without data (OID is null)
+ otaPackageService.deleteOtaPackage(tenantId, savedFirmwareInfo.getId());
+
+ OtaPackageInfo foundFirmware = otaPackageService.findOtaPackageInfoById(tenantId, savedFirmwareInfo.getId());
+ Assert.assertNull(foundFirmware);
+ }
+
+ @Test
+ public void testDeleteOtaPackageUnlinksLargeObject() {
+ OtaPackage savedFirmware = createAndSaveFirmware(tenantId, VERSION);
+
+ Long oid = otaPackageDao.getDataOidById(savedFirmware.getId().getId());
+ Assert.assertNotNull(oid);
+
+ otaPackageService.deleteOtaPackage(tenantId, savedFirmware.getId());
+
+ // Verify the large object was unlinked - PostgreSQL throws an exception when the object doesn't exist
+ assertThatThrownBy(() -> otaPackageDao.unlinkLargeObject(oid)).hasMessageContaining("large object " + oid + " does not exist");
+ }
+
@Test
public void testFindTenantFirmwaresByTenantId() {
List firmwares = new ArrayList<>();
@@ -567,8 +600,8 @@ public class OtaPackageServiceTest extends AbstractServiceTest {
}
} while (pageData.hasNext());
- Collections.sort(firmwares, idComparator);
- Collections.sort(loadedFirmwares, idComparator);
+ firmwares.sort(idComparator);
+ loadedFirmwares.sort(idComparator);
assertThat(firmwares).isEqualTo(loadedFirmwares);
@@ -622,8 +655,8 @@ public class OtaPackageServiceTest extends AbstractServiceTest {
}
} while (pageData.hasNext());
- Collections.sort(firmwares, idComparator);
- Collections.sort(loadedFirmwares, idComparator);
+ firmwares.sort(idComparator);
+ loadedFirmwares.sort(idComparator);
assertThat(firmwares).isEqualTo(loadedFirmwares);
@@ -726,4 +759,5 @@ public class OtaPackageServiceTest extends AbstractServiceTest {
firmware.setDataSize(DATA_SIZE);
return firmware;
}
+
}
diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/nosql/TimeseriesServiceNoSqlTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/nosql/TimeseriesServiceNoSqlTest.java
index aa43826734..b66cdb4363 100644
--- a/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/nosql/TimeseriesServiceNoSqlTest.java
+++ b/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/nosql/TimeseriesServiceNoSqlTest.java
@@ -65,7 +65,7 @@ public class TimeseriesServiceNoSqlTest extends BaseTimeseriesServiceTest {
new BasicTsKvEntry(TimeUnit.MINUTES.toMillis(5), new JsonDataEntry("test", "{\"test\":\"testValue\"}")));
DeviceId deviceId = new DeviceId(Uuids.timeBased());
- tsService.save(tenantId, deviceId, timeseries, ttlInSec);
+ tsService.save(tenantId, deviceId, timeseries, ttlInSec).get(MAX_TIMEOUT, TimeUnit.SECONDS);
List fullList = tsService.findAll(tenantId, deviceId, Collections.singletonList(new BaseReadTsKvQuery("test", 0L,
TimeUnit.MINUTES.toMillis(6), 1000, 10, Aggregation.NONE))).get(MAX_TIMEOUT, TimeUnit.SECONDS);
diff --git a/dao/src/test/java/org/thingsboard/server/dao/sql/notification/JpaNotificationRequestDaoTest.java b/dao/src/test/java/org/thingsboard/server/dao/sql/notification/JpaNotificationRequestDaoTest.java
new file mode 100644
index 0000000000..5e2fe2308b
--- /dev/null
+++ b/dao/src/test/java/org/thingsboard/server/dao/sql/notification/JpaNotificationRequestDaoTest.java
@@ -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.dao.sql.notification;
+
+import org.junit.After;
+import org.junit.Test;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.thingsboard.server.common.data.id.NotificationRequestId;
+import org.thingsboard.server.common.data.id.TenantId;
+import org.thingsboard.server.common.data.notification.NotificationRequest;
+import org.thingsboard.server.common.data.notification.NotificationRequestStatus;
+import org.thingsboard.server.dao.AbstractJpaDaoTest;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.UUID;
+import java.util.concurrent.TimeUnit;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+public class JpaNotificationRequestDaoTest extends AbstractJpaDaoTest {
+
+ @Autowired
+ JpaNotificationRequestDao notificationRequestDao;
+
+ private final List createdRequests = new ArrayList<>();
+
+ @After
+ public void tearDown() {
+ for (NotificationRequest request : createdRequests) {
+ notificationRequestDao.removeById(request.getTenantId(), request.getId().getId());
+ }
+ createdRequests.clear();
+ }
+
+ @Test
+ public void testBatchDeletion() {
+ TenantId sysTenantId = TenantId.SYS_TENANT_ID;
+ long now = System.currentTimeMillis();
+ long oldTimestamp = now - TimeUnit.DAYS.toMillis(30);
+
+ NotificationRequest oldRequest1 = createNotificationRequest(sysTenantId, oldTimestamp);
+ notificationRequestDao.save(sysTenantId, oldRequest1);
+
+ NotificationRequest oldRequest2 = createNotificationRequest(sysTenantId, oldTimestamp);
+ notificationRequestDao.save(sysTenantId, oldRequest2);
+
+ NotificationRequest freshRequest = createNotificationRequest(sysTenantId, now);
+ notificationRequestDao.save(sysTenantId, freshRequest);
+
+ TenantId tenant2Id = TenantId.fromUUID(UUID.fromString("3d193a7a-774b-4c05-84d5-f7fdcf7a37cf"));
+ NotificationRequest tenant2Request = createNotificationRequest(tenant2Id, oldTimestamp);
+ notificationRequestDao.save(tenant2Id, tenant2Request);
+
+ int batchSize = 10_000;
+
+ assertThat(notificationRequestDao.removeByTenantIdAndCreatedTimeBeforeBatch(sysTenantId, oldTimestamp - 1, batchSize)).isEqualTo(0);
+
+ long expirationTime = now - TimeUnit.DAYS.toMillis(15);
+ assertThat(notificationRequestDao.removeByTenantIdAndCreatedTimeBeforeBatch(sysTenantId, expirationTime, batchSize)).isEqualTo(2);
+
+ assertThat(notificationRequestDao.findById(sysTenantId, freshRequest.getId().getId())).isNotNull();
+ assertThat(notificationRequestDao.removeByTenantIdAndCreatedTimeBeforeBatch(tenant2Id, now + 1, batchSize)).isEqualTo(1);
+ }
+
+ @Test
+ public void testBatchDeletionWithSmallBatchSize() {
+ TenantId tenantId = TenantId.SYS_TENANT_ID;
+ long oldTimestamp = System.currentTimeMillis() - TimeUnit.DAYS.toMillis(30);
+
+ for (int i = 0; i < 10; i++) {
+ NotificationRequest request = createNotificationRequest(tenantId, oldTimestamp);
+ notificationRequestDao.save(tenantId, request);
+ }
+
+ int batchSize = 3;
+ long expirationTime = System.currentTimeMillis();
+
+ assertThat(notificationRequestDao.removeByTenantIdAndCreatedTimeBeforeBatch(tenantId, expirationTime, batchSize)).isEqualTo(3);
+ assertThat(notificationRequestDao.removeByTenantIdAndCreatedTimeBeforeBatch(tenantId, expirationTime, batchSize)).isEqualTo(3);
+ assertThat(notificationRequestDao.removeByTenantIdAndCreatedTimeBeforeBatch(tenantId, expirationTime, batchSize)).isEqualTo(3);
+ assertThat(notificationRequestDao.removeByTenantIdAndCreatedTimeBeforeBatch(tenantId, expirationTime, batchSize)).isEqualTo(1);
+ assertThat(notificationRequestDao.removeByTenantIdAndCreatedTimeBeforeBatch(tenantId, expirationTime, batchSize)).isEqualTo(0);
+ }
+
+ @Test
+ public void testBatchDeletionIsolationBetweenTenants() {
+ TenantId tenant1 = TenantId.SYS_TENANT_ID;
+ TenantId tenant2 = TenantId.fromUUID(UUID.fromString("3d193a7a-774b-4c05-84d5-f7fdcf7a37cf"));
+ long oldTimestamp = System.currentTimeMillis() - TimeUnit.DAYS.toMillis(30);
+
+ for (int i = 0; i < 5; i++) {
+ NotificationRequest request = createNotificationRequest(tenant1, oldTimestamp);
+ notificationRequestDao.save(tenant1, request);
+ }
+
+ for (int i = 0; i < 3; i++) {
+ NotificationRequest request = createNotificationRequest(tenant2, oldTimestamp);
+ notificationRequestDao.save(tenant2, request);
+ }
+
+ int batchSize = 10_000;
+ long expirationTime = System.currentTimeMillis();
+
+ assertThat(notificationRequestDao.removeByTenantIdAndCreatedTimeBeforeBatch(tenant1, expirationTime, batchSize)).isEqualTo(5);
+ assertThat(notificationRequestDao.removeByTenantIdAndCreatedTimeBeforeBatch(tenant2, expirationTime, batchSize)).isEqualTo(3);
+ }
+
+ private NotificationRequest createNotificationRequest(TenantId tenantId, long createdTime) {
+ NotificationRequest request = new NotificationRequest();
+ request.setId(new NotificationRequestId(UUID.randomUUID()));
+ request.setTenantId(tenantId);
+ request.setCreatedTime(createdTime);
+ request.setTargets(List.of(UUID.randomUUID()));
+ request.setStatus(NotificationRequestStatus.SENT);
+ createdRequests.add(request);
+ return request;
+ }
+
+}
diff --git a/dao/src/test/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDaoTest.java b/dao/src/test/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDaoTest.java
index 1629922685..921339a92b 100644
--- a/dao/src/test/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDaoTest.java
+++ b/dao/src/test/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDaoTest.java
@@ -51,9 +51,10 @@ public class JpaRpcDaoTest extends AbstractJpaDaoTest {
rpc.setDeviceId(new DeviceId(UUID.randomUUID()));
rpcDao.saveAndFlush(rpc.getTenantId(), rpc);
- assertThat(rpcDao.deleteOutdatedRpcByTenantId(TenantId.SYS_TENANT_ID, 0L)).isEqualTo(0);
- assertThat(rpcDao.deleteOutdatedRpcByTenantId(TenantId.SYS_TENANT_ID, Long.MAX_VALUE)).isEqualTo(2);
- assertThat(rpcDao.deleteOutdatedRpcByTenantId(tenantId, System.currentTimeMillis() + 1)).isEqualTo(1);
+ int batchSize = 10_000;
+ assertThat(rpcDao.deleteOutdatedRpcByTenantIdBatch(TenantId.SYS_TENANT_ID, 0L, batchSize)).isEqualTo(0);
+ assertThat(rpcDao.deleteOutdatedRpcByTenantIdBatch(TenantId.SYS_TENANT_ID, Long.MAX_VALUE, batchSize)).isEqualTo(2);
+ assertThat(rpcDao.deleteOutdatedRpcByTenantIdBatch(tenantId, System.currentTimeMillis() + 1, batchSize)).isEqualTo(1);
}
}