From 0468451f3ba02983e540e26f968305dd7a2ed04a Mon Sep 17 00:00:00 2001 From: dashevchenko Date: Mon, 12 Jan 2026 10:52:29 +0200 Subject: [PATCH 001/106] added ws update on telemetry deletion --- .../DefaultTbLocalSubscriptionService.java | 4 +-- .../server/controller/WebsocketApiTest.java | 35 +++++++++++++++++++ .../server/common/data/kv/TsKvEntry.java | 5 +++ 3 files changed, 42 insertions(+), 2 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java index 51663eea9c..df28da2dad 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java @@ -348,7 +348,7 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer if (sub.isLatestValues()) { for (TsKvEntry kv : data) { Long stateTs = keyStates.get(kv.getKey()); - if (stateTs == null || kv.getTs() >= stateTs) { + if (stateTs == null || kv.getTs() >= stateTs || kv.isDeletedEntryMarker()) { if (updateData == null) { updateData = new ArrayList<>(); } @@ -362,7 +362,7 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer for (TsKvEntry kv : data) { Long stateTs = keyStates.get(kv.getKey()); if (stateTs != null) { - if (!sub.isLatestValues() || kv.getTs() >= stateTs) { + if (!sub.isLatestValues() || kv.getTs() >= stateTs || kv.isDeletedEntryMarker()) { if (updateData == null) { updateData = new ArrayList<>(); } diff --git a/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java b/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java index 87ba0ec3e8..e88808ad2a 100644 --- a/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java @@ -726,6 +726,41 @@ public class WebsocketApiTest extends AbstractControllerTest { Assert.assertNull(msg); } + @Test + public void testShouldSendWsUpdateMessageWhenTelemetryWasDeleted() throws Exception { + long now = System.currentTimeMillis() - 100; + TsKvEntry dataPoint = new BasicTsKvEntry(now, new LongDataEntry("temperature", 42L)); + List tsData = List.of(dataPoint); + sendTelemetry(device, tsData); + + List keys = List.of(new EntityKey(EntityKeyType.TIME_SERIES, "temperature")); + EntityDataUpdate update = getWsClient().subscribeLatestUpdate(keys, dtf); + + Assert.assertEquals(1, update.getCmdId()); + PageData pageData = update.getData(); + Assert.assertNotNull(pageData); + Assert.assertEquals(1, pageData.getData().size()); + Assert.assertEquals(device.getId(), pageData.getData().get(0).getEntityId()); + Assert.assertNotNull(pageData.getData().get(0).getLatest().get(EntityKeyType.TIME_SERIES).get("temperature")); + Assert.assertEquals(now, pageData.getData().get(0).getLatest().get(EntityKeyType.TIME_SERIES).get("temperature").getTs()); + Assert.assertEquals("42", pageData.getData().get(0).getLatest().get(EntityKeyType.TIME_SERIES).get("temperature").getValue()); + + // delete telemetry + getWsClient().registerWaitForUpdate(); + doDeleteAsync("/api/plugins/telemetry/DEVICE/" + device.getId() + "/timeseries/delete?keys=temperature&deleteAllDataForKeys=true", String.class); + update = getWsClient().parseDataReply(getWsClient().waitForUpdate()); + + Assert.assertEquals(1, update.getCmdId()); + + List listData = update.getUpdate(); + Assert.assertNotNull(listData); + Assert.assertEquals(1, listData.size()); + Assert.assertEquals(device.getId(), listData.get(0).getEntityId()); + Assert.assertNotNull(listData.get(0).getLatest().get(EntityKeyType.TIME_SERIES)); + TsValue tsValue = listData.get(0).getLatest().get(EntityKeyType.TIME_SERIES).get("temperature"); + Assert.assertEquals(new TsValue(0, ""), tsValue); + } + @Test public void testEntityDataLatestAttrWsCmd() throws Exception { long now = System.currentTimeMillis(); 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..eca0609704 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 isDeletedEntryMarker() { + return getTs() == 0 && (getValue() == null || getValueAsString().isEmpty()); + } + } From 97d68dda90433a441eac439bb1d377a9b58bd978 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Wed, 28 Jan 2026 12:00:25 +0200 Subject: [PATCH 002/106] Ota package unlink data object --- .../housekeeper/HousekeeperServiceTest.java | 10 ++-- .../AlarmsDeletionHousekeeperTask.java | 4 ++ .../AlarmsUnassignHousekeeperTask.java | 4 ++ .../EntitiesDeletionHousekeeperTask.java | 4 ++ .../data/housekeeper/HousekeeperTask.java | 4 ++ .../LatestTsDeletionHousekeeperTask.java | 5 ++ ...TenantEntitiesDeletionHousekeeperTask.java | 5 ++ .../TsHistoryDeletionHousekeeperTask.java | 5 ++ .../server/dao/ota/BaseOtaPackageService.java | 48 ++++++++++++------- .../server/dao/ota/OtaPackageDao.java | 7 ++- .../server/dao/sql/ota/JpaOtaPackageDao.java | 10 ++++ .../dao/sql/ota/OtaPackageRepository.java | 9 +++- .../dao/service/OtaPackageServiceTest.java | 37 ++++++++++++++ 13 files changed, 128 insertions(+), 24 deletions(-) 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/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/dao/src/main/java/org/thingsboard/server/dao/ota/BaseOtaPackageService.java b/dao/src/main/java/org/thingsboard/server/dao/ota/BaseOtaPackageService.java index 95515bd15e..aac4157706 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 = otaPackageDao.getDataOidById(otaPackageId.getId()); 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,20 @@ 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/sql/ota/JpaOtaPackageDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/ota/JpaOtaPackageDao.java index d67110d53d..76bf5ca26e 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/ota/JpaOtaPackageDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/ota/JpaOtaPackageDao.java @@ -77,6 +77,16 @@ public class JpaOtaPackageDao extends JpaAbstractDao findIdsByTenantId(@Param("tenantId") UUID tenantId, Pageable pageable); + @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/test/java/org/thingsboard/server/dao/service/OtaPackageServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/OtaPackageServiceTest.java index a38499c82c..e79a7bad07 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 @@ -37,6 +37,7 @@ import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileCon import org.thingsboard.server.dao.device.DeviceProfileService; import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.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; @@ -77,6 +78,8 @@ public class OtaPackageServiceTest extends AbstractServiceTest { @Autowired TenantProfileService tenantProfileService; @Autowired + OtaPackageDao otaPackageDao; + @Autowired TbTenantProfileCache tenantProfileCache; @Before @@ -533,6 +536,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<>(); @@ -726,4 +762,5 @@ public class OtaPackageServiceTest extends AbstractServiceTest { firmware.setDataSize(DATA_SIZE); return firmware; } + } From cc6af84a928711c5d1d9f20e03e1661419969b7d Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Wed, 28 Jan 2026 12:04:42 +0200 Subject: [PATCH 003/106] Small revert --- .../server/dao/ota/BaseOtaPackageService.java | 12 +++++++++++- 1 file changed, 11 insertions(+), 1 deletion(-) 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 aac4157706..54577ff9eb 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,7 @@ public class BaseOtaPackageService extends AbstractCachedEntityService INCORRECT_OTA_PACKAGE_ID + id); try { - Long oid = otaPackageDao.getDataOidById(otaPackageId.getId()); + Long oid = getDataOidById(tenantId, otaPackageId); otaPackageDao.removeById(tenantId, otaPackageId.getId()); unlinkDataIfPresent(tenantId, otaPackageId, oid); publishEvictEvent(new OtaPackageCacheEvictEvent(otaPackageId)); @@ -217,6 +217,16 @@ public class BaseOtaPackageService extends AbstractCachedEntityService Date: Wed, 4 Mar 2026 14:44:54 +0100 Subject: [PATCH 004/106] Fix LwM2M Redis stores using separate connections for SCAN and GET MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Using the same connection for both SCAN cursor iteration and GET value fetches causes Jedis 5.x response-ordering corruption: the SCAN response parser receives a GET response (byte[]) where it expects a List, and vice versa, resulting in ClassCastException on startup. Fix: open two connections per getAll() call — one dedicated to the scan cursor and one for value fetches — eliminating any interleaving. Affected: TbRedisLwM2MClientStore, TbRedisLwM2MModelConfigStore, TbLwM2mRedisRegistrationStore. Co-Authored-By: Claude Sonnet 4.6 --- .../store/TbLwM2mRedisRegistrationStore.java | 18 ++- .../server/store/TbRedisLwM2MClientStore.java | 14 +- .../store/TbRedisLwM2MModelConfigStore.java | 18 ++- .../store/TbRedisLwM2MClientStoreTest.java | 137 ++++++++++++++++++ 4 files changed, 164 insertions(+), 23 deletions(-) create mode 100644 common/transport/lwm2m/src/test/java/org/thingsboard/server/transport/lwm2m/server/store/TbRedisLwM2MClientStoreTest.java diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mRedisRegistrationStore.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mRedisRegistrationStore.java index 4b2ef07994..8158e82a81 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mRedisRegistrationStore.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mRedisRegistrationStore.java @@ -323,22 +323,24 @@ public class TbLwM2mRedisRegistrationStore implements RegistrationStore, Startab @Override public Iterator getAllRegistrations() { - try (var connection = connectionFactory.getConnection()) { + try (var scanConnection = connectionFactory.getConnection(); + var getConnection = connectionFactory.getConnection()) { Collection list = new LinkedList<>(); ScanOptions scanOptions = ScanOptions.scanOptions().count(100).match(REG_EP + "*").build(); List> scans = new ArrayList<>(); - if (connection instanceof RedisClusterConnection) { - ((RedisClusterConnection) connection).clusterGetNodes().forEach(node -> { - scans.add(((RedisClusterConnection) connection).scan(node, scanOptions)); - }); + if (scanConnection instanceof RedisClusterConnection clusterConnection) { + clusterConnection.clusterGetNodes().forEach(node -> + scans.add(clusterConnection.scan(node, scanOptions))); } else { - scans.add(connection.scan(scanOptions)); + scans.add(scanConnection.scan(scanOptions)); } scans.forEach(scan -> { scan.forEachRemaining(key -> { - byte[] element = connection.get(key); - list.add(deserializeReg(element)); + byte[] element = getConnection.get(key); + if (element != null) { + list.add(deserializeReg(element)); + } }); }); return list.iterator(); diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbRedisLwM2MClientStore.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbRedisLwM2MClientStore.java index 3293cd8b53..4beefb4896 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbRedisLwM2MClientStore.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbRedisLwM2MClientStore.java @@ -61,21 +61,21 @@ public class TbRedisLwM2MClientStore implements TbLwM2MClientStore { @Override public Set getAll() { - try (var connection = connectionFactory.getConnection()) { + try (var scanConnection = connectionFactory.getConnection(); + var getConnection = connectionFactory.getConnection()) { Set clients = new HashSet<>(); ScanOptions scanOptions = ScanOptions.scanOptions().count(100).match(CLIENT_EP + "*").build(); List> scans = new ArrayList<>(); - if (connection instanceof RedisClusterConnection) { - ((RedisClusterConnection) connection).clusterGetNodes().forEach(node -> { - scans.add(((RedisClusterConnection) connection).scan(node, scanOptions)); - }); + if (scanConnection instanceof RedisClusterConnection clusterConnection) { + clusterConnection.clusterGetNodes().forEach(node -> + scans.add(clusterConnection.scan(node, scanOptions))); } else { - scans.add(connection.scan(scanOptions)); + scans.add(scanConnection.scan(scanOptions)); } scans.forEach(scan -> { scan.forEachRemaining(key -> { - byte[] element = connection.get(key); + byte[] element = getConnection.get(key); if (element != null) { try { clients.add(deserialize(element)); diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbRedisLwM2MModelConfigStore.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbRedisLwM2MModelConfigStore.java index 73b6f3c8df..31a78234d0 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbRedisLwM2MModelConfigStore.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbRedisLwM2MModelConfigStore.java @@ -35,22 +35,24 @@ public class TbRedisLwM2MModelConfigStore implements TbLwM2MModelConfigStore { @Override public List getAll() { - try (var connection = connectionFactory.getConnection()) { + try (var scanConnection = connectionFactory.getConnection(); + var getConnection = connectionFactory.getConnection()) { List configs = new ArrayList<>(); ScanOptions scanOptions = ScanOptions.scanOptions().count(100).match(MODEL_EP + "*").build(); List> scans = new ArrayList<>(); - if (connection instanceof RedisClusterConnection) { - ((RedisClusterConnection) connection).clusterGetNodes().forEach(node -> { - scans.add(((RedisClusterConnection) connection).scan(node, scanOptions)); - }); + if (scanConnection instanceof RedisClusterConnection clusterConnection) { + clusterConnection.clusterGetNodes().forEach(node -> + scans.add(clusterConnection.scan(node, scanOptions))); } else { - scans.add(connection.scan(scanOptions)); + scans.add(scanConnection.scan(scanOptions)); } scans.forEach(scan -> { scan.forEachRemaining(key -> { - byte[] element = connection.get(key); - configs.add(JacksonUtil.fromBytes(element, LwM2MModelConfig.class)); + byte[] element = getConnection.get(key); + if (element != null) { + configs.add(JacksonUtil.fromBytes(element, LwM2MModelConfig.class)); + } }); }); return configs; diff --git a/common/transport/lwm2m/src/test/java/org/thingsboard/server/transport/lwm2m/server/store/TbRedisLwM2MClientStoreTest.java b/common/transport/lwm2m/src/test/java/org/thingsboard/server/transport/lwm2m/server/store/TbRedisLwM2MClientStoreTest.java new file mode 100644 index 0000000000..9fe29c1188 --- /dev/null +++ b/common/transport/lwm2m/src/test/java/org/thingsboard/server/transport/lwm2m/server/store/TbRedisLwM2MClientStoreTest.java @@ -0,0 +1,137 @@ +/** + * 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.transport.lwm2m.server.store; + +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.data.redis.connection.RedisConnection; +import org.springframework.data.redis.connection.RedisConnectionFactory; +import org.springframework.data.redis.core.Cursor; +import org.springframework.data.redis.core.ScanOptions; +import org.thingsboard.server.transport.lwm2m.server.client.LwM2MClientState; +import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; + +import java.util.List; +import java.util.Set; +import java.util.function.Consumer; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; +import static org.thingsboard.server.transport.lwm2m.server.store.util.LwM2MClientSerDes.serialize; + +/** + * Verifies that {@link TbRedisLwM2MClientStore#getAll()} uses separate connections for + * SCAN and GET operations to prevent Jedis 5.x response-ordering corruption that occurs + * when both commands share the same connection. + */ +@ExtendWith(MockitoExtension.class) +class TbRedisLwM2MClientStoreTest { + + @Mock + RedisConnectionFactory connectionFactory; + + @Mock + RedisConnection scanConnection; + + @Mock + RedisConnection getConnection; + + TbRedisLwM2MClientStore store; + + @BeforeEach + void setUp() { + // First getConnection() call → scanConnection, second → getConnection + when(connectionFactory.getConnection()) + .thenReturn(scanConnection) + .thenReturn(getConnection); + store = new TbRedisLwM2MClientStore(connectionFactory); + } + + @Test + void getAll_returnsSingleClient() { + LwM2mClient client = new LwM2mClient("nodeId", "testEndpoint"); + client.setState(LwM2MClientState.REGISTERED); + byte[] key = "CLIENT#EP#testEndpoint".getBytes(); + byte[] value = serialize(client); + + // Cursor created before thenReturn to avoid Mockito unfinished-stubbing error + Cursor cursor = cursorOf(key); + when(scanConnection.scan(any(ScanOptions.class))).thenReturn(cursor); + when(getConnection.get(key)).thenReturn(value); + + Set result = store.getAll(); + + assertThat(result).hasSize(1); + assertThat(result.iterator().next().getEndpoint()).isEqualTo("testEndpoint"); + } + + @Test + void getAll_getIsNeverCalledOnScanConnection() { + Cursor cursor = cursorOf(); + when(scanConnection.scan(any(ScanOptions.class))).thenReturn(cursor); + + store.getAll(); + + verify(scanConnection, never()).get(any(byte[].class)); + } + + @Test + void getAll_scanIsNeverCalledOnGetConnection() { + Cursor cursor = cursorOf(); + when(scanConnection.scan(any(ScanOptions.class))).thenReturn(cursor); + + store.getAll(); + + verify(getConnection, never()).scan(any(ScanOptions.class)); + } + + @Test + void getAll_skipsKeyWhenValueIsNull() { + byte[] key = "CLIENT#EP#gone".getBytes(); + Cursor cursor = cursorOf(key); + when(scanConnection.scan(any(ScanOptions.class))).thenReturn(cursor); + // getConnection.get(key) returns null by default — no stubbing needed + + Set result = store.getAll(); + + assertThat(result).isEmpty(); + } + + /** + * Creates a mock {@link Cursor} that iterates over the given keys via {@code forEachRemaining}. + * The cursor is created separately (not inside a {@code thenReturn()} argument) to avoid + * Mockito's "unfinished stubbing" error caused by nested {@code when()} calls. + */ + @SuppressWarnings("unchecked") + private static Cursor cursorOf(byte[]... keys) { + Cursor cursor = mock(Cursor.class); + List keyList = List.of(keys); + doAnswer(inv -> { + Consumer action = inv.getArgument(0); + keyList.forEach(action); + return null; + }).when(cursor).forEachRemaining(any(Consumer.class)); + return cursor; + } +} From 455f62eaefcbf7f996af54f1d7da33238f797601 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 5 Mar 2026 16:19:59 +0100 Subject: [PATCH 005/106] fix: delay LwM2M client stop in FW update to prevent Execute response race When Execute RPC is sent to FW Update resource (/5/0/2), the test client's startUpdating() scheduled the client stop with 0 delay. This caused a race where the client stopped before the CoAP Execute response (2.04 Changed) was delivered to the server, resulting in RequestCanceledException and INTERNAL_SERVER_ERROR instead of CHANGED. Adding a 1-second delay before leshanClient.stop() ensures the CoAP response is transmitted and received before the client disconnects, fixing the flaky testExecuteUpdateFWById_Result_CHANGED test. Co-Authored-By: Claude Sonnet 4.6 --- .../server/transport/lwm2m/client/FwLwM2MDevice.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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..6a7d631eb1 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 @@ -193,7 +193,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) { From 30b046fa30ce6d6c3b490733a78d3bb6f769d779 Mon Sep 17 00:00:00 2001 From: dashevchenko Date: Thu, 5 Mar 2026 17:53:27 +0200 Subject: [PATCH 006/106] sending ws error when the telemetry queries exceed limit --- ...efaultTbEntityDataSubscriptionService.java | 19 +++++++- .../server/controller/WebsocketApiTest.java | 46 +++++++++++++++++++ 2 files changed, 63 insertions(+), 2 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java index f0f2e1eaa7..71df830f8d 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java @@ -34,6 +34,7 @@ import org.springframework.web.socket.CloseStatus; import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.dao.nosql.ResultSetSizeLimitExceededException; import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult; @@ -242,7 +243,10 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc @Override public void onFailure(Throwable t) { - log.warn("[{}][{}] Failed to process command", finalCtx.getSessionId(), finalCtx.getCmdId()); + log.warn("[{}][{}] Failed to process command", finalCtx.getSessionId(), finalCtx.getCmdId(), t); + if (t instanceof ResultSetSizeLimitExceededException) { + finalCtx.sendWsMsg(new EntityDataUpdate(finalCtx.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR.getCode(), t.getMessage())); + } } }, wsCallBackExecutor); } @@ -258,7 +262,18 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc handleLatestCmd(ctx, cmd.getLatestCmd()); } if (cmd.getTsCmd() != null) { - handleTimeSeriesCmd(ctx, cmd.getTsCmd()); + Futures.addCallback(handleTimeSeriesCmd(ctx, cmd.getTsCmd()), new FutureCallback<>() { + @Override + public void onSuccess(TbEntityDataSubCtx result) {} + + @Override + public void onFailure(Throwable t) { + log.warn("[{}][{}] Failed to process timeseries command", ctx.getSessionId(), ctx.getCmdId(), t); + if (t instanceof ResultSetSizeLimitExceededException) { + ctx.sendWsMsg(new EntityDataUpdate(ctx.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR.getCode(), t.getMessage())); + } + } + }, wsCallBackExecutor); } } else { checkAndSendInitialData(ctx); diff --git a/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java b/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java index 87ba0ec3e8..aa7ffaeaf5 100644 --- a/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java @@ -19,13 +19,16 @@ import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.node.ArrayNode; import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.FutureCallback; +import com.google.common.util.concurrent.Futures; import lombok.extern.slf4j.Slf4j; import org.checkerframework.checker.nullness.qual.Nullable; import org.junit.After; import org.junit.Assert; import org.junit.Before; import org.junit.Test; +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.testcontainers.shaded.org.apache.commons.lang3.RandomStringUtils; import org.thingsboard.common.util.JacksonUtil; @@ -60,7 +63,9 @@ import org.thingsboard.server.common.data.query.NumericFilterPredicate; import org.thingsboard.server.common.data.query.SingleEntityFilter; import org.thingsboard.server.common.data.query.TsValue; import org.thingsboard.server.common.data.relation.EntityRelation; +import org.thingsboard.server.dao.nosql.ResultSetSizeLimitExceededException; import org.thingsboard.server.dao.service.DaoSqlTest; +import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.service.subscription.SubscriptionErrorCode; import org.thingsboard.server.service.subscription.TbAttributeSubscriptionScope; import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; @@ -95,6 +100,9 @@ public class WebsocketApiTest extends AbstractControllerTest { @Autowired private TelemetrySubscriptionService tsService; + @SpyBean + private TimeseriesService timeseriesService; + Device device; DeviceTypeFilter dtf; @@ -965,6 +973,44 @@ public class WebsocketApiTest extends AbstractControllerTest { } + @Test + public void testHistoryCmdSendsWsErrorOnResultSetSizeLimitExceeded() throws Exception { + ResultSetSizeLimitExceededException exception = new ResultSetSizeLimitExceededException(100L, 200L); + Mockito.doReturn(Futures.immediateFailedFuture(exception)) + .when(timeseriesService).findAllByQueries(Mockito.any(), Mockito.any(), Mockito.any()); + + List keys = List.of("temperature"); + long now = System.currentTimeMillis(); + + // Register for 2 messages: initial entity page data + error + getWsClient().registerWaitForUpdate(2); + getWsClient().sendHistoryCmd(keys, now, TimeUnit.HOURS.toMillis(1), dtf); + getWsClient().waitForUpdate(); + + EntityDataUpdate errorUpdate = JacksonUtil.fromString(getWsClient().getLastMsg(), EntityDataUpdate.class); + assertThat(errorUpdate.getErrorCode()).isEqualTo(SubscriptionErrorCode.INTERNAL_ERROR.getCode()); + assertThat(errorUpdate.getErrorMsg()).isEqualTo(exception.getMessage()); + } + + @Test + public void testTimeSeriesCmdSendsWsErrorOnResultSetSizeLimitExceeded() throws Exception { + ResultSetSizeLimitExceededException exception = new ResultSetSizeLimitExceededException(100L, 200L); + Mockito.doReturn(Futures.immediateFailedFuture(exception)) + .when(timeseriesService).findAllByQueries(Mockito.any(), Mockito.any(), Mockito.any()); + + List keys = List.of("temperature"); + long now = System.currentTimeMillis(); + + // Register for 2 messages: initial entity page data + error + getWsClient().registerWaitForUpdate(2); + getWsClient().subscribeTsUpdate(keys, now, TimeUnit.HOURS.toMillis(1), dtf); + getWsClient().waitForUpdate(); + + EntityDataUpdate errorUpdate = JacksonUtil.fromString(getWsClient().getLastMsg(), EntityDataUpdate.class); + assertThat(errorUpdate.getErrorCode()).isEqualTo(SubscriptionErrorCode.INTERNAL_ERROR.getCode()); + assertThat(errorUpdate.getErrorMsg()).isEqualTo(exception.getMessage()); + } + private void sendTelemetry(Device device, List tsData) throws InterruptedException { CountDownLatch latch = new CountDownLatch(1); tsService.saveTimeseries(TimeseriesSaveRequest.builder() From bc39695fabe5441ce1a18861664bcd3d86aead74 Mon Sep 17 00:00:00 2001 From: dashevchenko Date: Thu, 5 Mar 2026 18:24:39 +0200 Subject: [PATCH 007/106] test fixes --- .../server/controller/WebsocketApiTest.java | 14 ++------------ 1 file changed, 2 insertions(+), 12 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java b/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java index aa7ffaeaf5..eab878854d 100644 --- a/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java @@ -982,12 +982,7 @@ public class WebsocketApiTest extends AbstractControllerTest { List keys = List.of("temperature"); long now = System.currentTimeMillis(); - // Register for 2 messages: initial entity page data + error - getWsClient().registerWaitForUpdate(2); - getWsClient().sendHistoryCmd(keys, now, TimeUnit.HOURS.toMillis(1), dtf); - getWsClient().waitForUpdate(); - - EntityDataUpdate errorUpdate = JacksonUtil.fromString(getWsClient().getLastMsg(), EntityDataUpdate.class); + EntityDataUpdate errorUpdate = getWsClient().sendHistoryCmd(keys, now, TimeUnit.HOURS.toMillis(1), dtf); assertThat(errorUpdate.getErrorCode()).isEqualTo(SubscriptionErrorCode.INTERNAL_ERROR.getCode()); assertThat(errorUpdate.getErrorMsg()).isEqualTo(exception.getMessage()); } @@ -1001,12 +996,7 @@ public class WebsocketApiTest extends AbstractControllerTest { List keys = List.of("temperature"); long now = System.currentTimeMillis(); - // Register for 2 messages: initial entity page data + error - getWsClient().registerWaitForUpdate(2); - getWsClient().subscribeTsUpdate(keys, now, TimeUnit.HOURS.toMillis(1), dtf); - getWsClient().waitForUpdate(); - - EntityDataUpdate errorUpdate = JacksonUtil.fromString(getWsClient().getLastMsg(), EntityDataUpdate.class); + EntityDataUpdate errorUpdate = getWsClient().subscribeTsUpdate(keys, now, TimeUnit.HOURS.toMillis(1), dtf); assertThat(errorUpdate.getErrorCode()).isEqualTo(SubscriptionErrorCode.INTERNAL_ERROR.getCode()); assertThat(errorUpdate.getErrorMsg()).isEqualTo(exception.getMessage()); } From 2cd0a746a3e45496338af4d686f89ead181641b8 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Fri, 6 Mar 2026 11:15:59 +0200 Subject: [PATCH 008/106] Refactoring --- .../server/dao/ota/BaseOtaPackageService.java | 4 ++-- .../server/dao/sql/ota/OtaPackageRepository.java | 3 ++- .../server/dao/service/OtaPackageServiceTest.java | 13 +++++-------- 3 files changed, 9 insertions(+), 11 deletions(-) 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 54577ff9eb..c43274cc26 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 @@ -234,10 +234,10 @@ public class BaseOtaPackageService extends AbstractCachedEntityService 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); 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 e79a7bad07..2f296d38f2 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 @@ -44,7 +44,6 @@ import org.thingsboard.server.dao.tenant.TenantProfileService; import java.nio.ByteBuffer; import java.util.ArrayList; -import java.util.Collections; import java.util.List; import static org.assertj.core.api.Assertions.assertThat; @@ -121,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)); } @@ -603,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); @@ -658,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); From 47369026fb05552725f1fb572e364b6d227ea311 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Fri, 6 Mar 2026 14:25:16 +0200 Subject: [PATCH 009/106] CE: Fixed notification requests and RPC cleanup timeout on large datasets --- .../service/ttl/AbstractCleanUpService.java | 4 + .../ttl/NotificationsCleanUpService.java | 63 +++++++- .../service/ttl/rpc/RpcCleanUpService.java | 94 +++++++----- .../src/main/resources/thingsboard.yml | 4 +- .../ttl/NotificationsCleanUpServiceTest.java | 139 ++++++++++++++++++ .../ttl/rpc/RpcCleanUpServiceTest.java | 136 +++++++++++++++++ .../notification/NotificationRequestDao.java | 2 +- .../thingsboard/server/dao/rpc/RpcDao.java | 3 +- .../JpaNotificationRequestDao.java | 4 +- .../NotificationRequestRepository.java | 8 +- .../server/dao/sql/rpc/JpaRpcDao.java | 4 +- .../server/dao/sql/rpc/RpcRepository.java | 9 +- .../JpaNotificationRequestDaoTest.java | 133 +++++++++++++++++ .../server/dao/sql/rpc/JpaRpcDaoTest.java | 7 +- 14 files changed, 554 insertions(+), 56 deletions(-) create mode 100644 application/src/test/java/org/thingsboard/server/service/ttl/NotificationsCleanUpServiceTest.java create mode 100644 application/src/test/java/org/thingsboard/server/service/ttl/rpc/RpcCleanUpServiceTest.java create mode 100644 dao/src/test/java/org/thingsboard/server/dao/sql/notification/JpaNotificationRequestDaoTest.java diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/AbstractCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/AbstractCleanUpService.java index 5865d9e39d..596f5a8754 100644 --- a/application/src/main/java/org/thingsboard/server/service/ttl/AbstractCleanUpService.java +++ b/application/src/main/java/org/thingsboard/server/service/ttl/AbstractCleanUpService.java @@ -32,4 +32,8 @@ public abstract class AbstractCleanUpService { return partitionService.resolve(ServiceType.TB_CORE, TenantId.SYS_TENANT_ID, TenantId.SYS_TENANT_ID).isMyPartition(); } + protected boolean isTenantPartitionMine(TenantId tenantId) { + return partitionService.resolve(ServiceType.TB_CORE, tenantId, tenantId).isMyPartition(); + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/NotificationsCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/NotificationsCleanUpService.java index ddc95b74e9..5e327a72db 100644 --- a/application/src/main/java/org/thingsboard/server/service/ttl/NotificationsCleanUpService.java +++ b/application/src/main/java/org/thingsboard/server/service/ttl/NotificationsCleanUpService.java @@ -20,33 +20,41 @@ import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.notification.NotificationRequestConfig; +import org.thingsboard.server.common.data.page.PageDataIterable; 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.time.Instant; import java.util.concurrent.TimeUnit; import static org.thingsboard.server.dao.model.ModelConstants.NOTIFICATION_TABLE_NAME; +@Slf4j @Service @ConditionalOnExpression("${sql.ttl.notifications.enabled:true} && ${sql.ttl.notifications.ttl:0} > 0") -@Slf4j public class NotificationsCleanUpService extends AbstractCleanUpService { private final SqlPartitioningRepository partitioningRepository; private final NotificationRequestDao notificationRequestDao; + private final TenantService tenantService; @Value("${sql.ttl.notifications.ttl:2592000}") private long ttlInSec; @Value("${sql.notifications.partition_size:168}") private int partitionSizeInHours; + @Value("${sql.ttl.notifications.removal_batch_size:10000}") + private int removalBatchSize; public NotificationsCleanUpService(PartitionService partitionService, SqlPartitioningRepository partitioningRepository, - NotificationRequestDao notificationRequestDao) { + NotificationRequestDao notificationRequestDao, TenantService tenantService) { super(partitionService); this.partitioningRepository = partitioningRepository; this.notificationRequestDao = notificationRequestDao; + this.tenantService = tenantService; } @Scheduled(initialDelayString = "#{T(org.apache.commons.lang3.RandomUtils).nextLong(0, ${sql.ttl.notifications.checking_interval_ms:86400000})}", @@ -63,9 +71,56 @@ public class NotificationsCleanUpService extends AbstractCleanUpService { if (lastRemovedNotificationTs > 0) { long gap = TimeUnit.MINUTES.toMillis(10); long requestExpTime = lastRemovedNotificationTs - TimeUnit.SECONDS.toMillis(NotificationRequestConfig.MAX_SENDING_DELAY) - gap; - int removed = notificationRequestDao.removeAllByCreatedTimeBefore(requestExpTime); - log.info("Removed {} outdated notification requests older than {}", removed, requestExpTime); + cleanUpNotificationRequests(requestExpTime); + } + } + + private void cleanUpNotificationRequests(long expirationTime) { + log.info("Starting notification requests cleanup for records older than {}", Instant.ofEpochMilli(expirationTime)); + int totalRemoved = 0; + int tenantsProcessed = 0; + + // Clean up SYSADMIN's notification requests: + try { + totalRemoved += cleanUpByTenant(TenantId.SYS_TENANT_ID, expirationTime); + } catch (Exception e) { + log.warn("Failed to clean up notification requests for sysadmin {}", TenantId.SYS_TENANT_ID, e); + } + // Clean up notification requests for tenants + PageDataIterable tenants = new PageDataIterable<>(tenantService::findTenantsIds, 10_000); + for (TenantId tenantId : tenants) { + try { + if (!isTenantPartitionMine(tenantId)) { + continue; + } + int tenantRemoved = cleanUpByTenant(tenantId, expirationTime); + totalRemoved += tenantRemoved; + tenantsProcessed++; + if (tenantRemoved > 0) { + log.trace("Removed {} notification requests for tenant {}", tenantRemoved, tenantId); + } + } catch (Exception e) { + log.warn("Failed to clean up notification requests for tenant {}", tenantId, e); + } } + + log.info("Notification requests cleanup completed. Processed {} tenants, removed {} total records older than {}", tenantsProcessed, totalRemoved, Instant.ofEpochMilli(expirationTime)); + } + + private int cleanUpByTenant(TenantId tenantId, long expirationTime) { + int totalRemoved = 0; + int batchRemoved; + + do { + batchRemoved = notificationRequestDao.removeByTenantIdAndCreatedTimeBeforeBatch(tenantId, expirationTime, removalBatchSize); + totalRemoved += batchRemoved; + + if (batchRemoved > 0) { + log.trace("Removed {} notification requests in batch for tenant {}", batchRemoved, tenantId); + } + } while (batchRemoved >= removalBatchSize); + + return totalRemoved; } } diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/rpc/RpcCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/rpc/RpcCleanUpService.java index 9404a0ae72..a945bfd8a5 100644 --- a/application/src/main/java/org/thingsboard/server/service/ttl/rpc/RpcCleanUpService.java +++ b/application/src/main/java/org/thingsboard/server/service/ttl/rpc/RpcCleanUpService.java @@ -15,69 +15,87 @@ */ package org.thingsboard.server.service.ttl.rpc; -import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; +import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.data.page.PageData; -import org.thingsboard.server.common.data.page.PageLink; +import org.thingsboard.server.common.data.page.PageDataIterable; import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; -import org.thingsboard.server.common.msg.queue.ServiceType; 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 org.thingsboard.server.queue.util.TbCoreComponent; +import org.thingsboard.server.service.ttl.AbstractCleanUpService; -import java.util.Date; +import java.time.Instant; import java.util.Optional; import java.util.concurrent.TimeUnit; -@TbCoreComponent -@Service @Slf4j -@RequiredArgsConstructor -public class RpcCleanUpService { - @Value("${sql.ttl.rpc.enabled}") - private boolean ttlTaskExecutionEnabled; +@Service +@TbCoreComponent +@ConditionalOnExpression("${sql.ttl.rpc.enabled:true}") +public class RpcCleanUpService extends AbstractCleanUpService { + + @Value("${sql.ttl.rpc.removal_batch_size:10000}") + private int removalBatchSize; + private final RpcDao rpcDao; private final TenantService tenantService; - private final PartitionService partitionService; private final TbTenantProfileCache tenantProfileCache; - private final RpcDao rpcDao; + + public RpcCleanUpService(TenantService tenantService, PartitionService partitionService, TbTenantProfileCache tenantProfileCache, RpcDao rpcDao) { + super(partitionService); + this.tenantService = tenantService; + this.tenantProfileCache = tenantProfileCache; + this.rpcDao = rpcDao; + } @Scheduled(initialDelayString = "#{T(org.apache.commons.lang3.RandomUtils).nextLong(0, ${sql.ttl.rpc.checking_interval})}", fixedDelayString = "${sql.ttl.rpc.checking_interval}") public void cleanUp() { - if (ttlTaskExecutionEnabled) { - PageLink tenantsBatchRequest = new PageLink(10_000, 0); - PageData tenantsIds; - do { - tenantsIds = tenantService.findTenantsIds(tenantsBatchRequest); - for (TenantId tenantId : tenantsIds.getData()) { - if (!partitionService.resolve(ServiceType.TB_CORE, tenantId, tenantId).isMyPartition()) { - continue; - } - - Optional tenantProfileConfiguration = tenantProfileCache.get(tenantId).getProfileConfiguration(); - if (tenantProfileConfiguration.isEmpty() || tenantProfileConfiguration.get().getRpcTtlDays() == 0) { - continue; - } - - long ttl = TimeUnit.DAYS.toMillis(tenantProfileConfiguration.get().getRpcTtlDays()); - long expirationTime = System.currentTimeMillis() - ttl; - - int totalRemoved = rpcDao.deleteOutdatedRpcByTenantId(tenantId, expirationTime); - - if (totalRemoved > 0) { - log.info("Removed {} outdated rpc(s) for tenant {} older than {}", totalRemoved, tenantId, new Date(expirationTime)); - } + PageDataIterable tenants = new PageDataIterable<>(tenantService::findTenantsIds, 10_000); + for (TenantId tenantId : tenants) { + try { + if (!isTenantPartitionMine(tenantId)) { + continue; } - tenantsBatchRequest = tenantsBatchRequest.nextPageLink(); - } while (tenantsIds.hasNext()); + Optional tenantProfileConfiguration = tenantProfileCache.get(tenantId).getProfileConfiguration(); + if (tenantProfileConfiguration.isEmpty() || tenantProfileConfiguration.get().getRpcTtlDays() == 0) { + continue; + } + + long ttl = TimeUnit.DAYS.toMillis(tenantProfileConfiguration.get().getRpcTtlDays()); + long expirationTime = System.currentTimeMillis() - ttl; + + int totalRemoved = cleanUpByTenant(tenantId, expirationTime); + + if (totalRemoved > 0) { + log.info("Removed {} outdated rpc(s) for tenant {} older than {}", totalRemoved, tenantId, Instant.ofEpochMilli(expirationTime)); + } + } catch (Exception e) { + log.warn("Failed to clean up rpc by ttl for tenant {}", tenantId, e); + } } } + private int cleanUpByTenant(TenantId tenantId, long expirationTime) { + int totalRemoved = 0; + int batchRemoved; + + do { + batchRemoved = rpcDao.deleteOutdatedRpcByTenantIdBatch(tenantId, expirationTime, removalBatchSize); + totalRemoved += batchRemoved; + + if (batchRemoved > 0) { + log.trace("Removed {} rpc in batch for tenant {}", batchRemoved, tenantId); + } + } while (batchRemoved >= removalBatchSize); + + return totalRemoved; + } + } diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 751d1c0e0b..8afdbb9ca9 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -421,10 +421,11 @@ sql: edge_events_ttl: "${SQL_TTL_EDGE_EVENTS_TTL:2628000}" # Number of seconds. The current value corresponds to one month alarms: checking_interval: "${SQL_ALARMS_TTL_CHECKING_INTERVAL:7200000}" # Number of milliseconds. The current value corresponds to two hours - removal_batch_size: "${SQL_ALARMS_TTL_REMOVAL_BATCH_SIZE:3000}" # To delete outdated alarms not all at once but in batches + removal_batch_size: "${SQL_ALARMS_TTL_REMOVAL_BATCH_SIZE:3000}" # Batch size for records removal rpc: enabled: "${SQL_TTL_RPC_ENABLED:true}" # Enable/disable TTL (Time To Live) for rpc call records checking_interval: "${SQL_RPC_TTL_CHECKING_INTERVAL:7200000}" # Number of milliseconds. The current value corresponds to two hours + removal_batch_size: "${SQL_RPC_TTL_REMOVAL_BATCH_SIZE:10000}" # Batch size for records removal audit_logs: enabled: "${SQL_TTL_AUDIT_LOGS_ENABLED:true}" # Enable/disable TTL (Time To Live) for audit log records ttl: "${SQL_TTL_AUDIT_LOGS_SECS:0}" # Disabled by default. The accuracy of the cleanup depends on the sql.audit_logs.partition_size @@ -433,6 +434,7 @@ sql: enabled: "${SQL_TTL_NOTIFICATIONS_ENABLED:true}" # Enable/disable TTL (Time To Live) for notification center records ttl: "${SQL_TTL_NOTIFICATIONS_SECS:2592000}" # Default value - 30 days checking_interval_ms: "${SQL_TTL_NOTIFICATIONS_CHECKING_INTERVAL_MS:86400000}" # Default value - 1 day + removal_batch_size: "${SQL_TTL_NOTIFICATIONS_REMOVAL_BATCH_SIZE:10000}" # Batch size for records removal relations: max_level: "${SQL_RELATIONS_MAX_LEVEL:50}" # This value has to be reasonably small to prevent infinite recursion as early as possible pool_size: "${SQL_RELATIONS_POOL_SIZE:4}" # This value has to be reasonably small to prevent the relation query from blocking all other DB calls 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..645eb600a4 --- /dev/null +++ b/application/src/test/java/org/thingsboard/server/service/ttl/NotificationsCleanUpServiceTest.java @@ -0,0 +1,139 @@ +/** + * 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_skipsRequestCleanup() { + 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); + + cleanUpService.cleanUp(); + + verify(notificationRequestDao, never()).removeByTenantIdAndCreatedTimeBeforeBatch(any(), anyLong(), anyInt()); + } + +} 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/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/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 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..7be106219b 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 @@ -27,6 +27,7 @@ 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); @@ -34,7 +35,11 @@ public interface RpcRepository extends JpaRepository { Page findAllByTenantId(UUID tenantId, Pageable pageable); @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/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); } } From 96b742189d25d832e47f95fdd056550629fc6db9 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 5 Mar 2026 14:11:33 +0100 Subject: [PATCH 010/106] Fix flaky Sparkplug connection test: handle 404 during device await doGet(url, Class) asserts HTTP 200 internally, so when the Sparkplug device hasn't been created yet the method throws AssertionError instead of returning null. Awaitility propagates Error immediately rather than continuing to poll, causing the test to fail after ~3 s instead of retrying for up to 200 s. Add .ignoreExceptions() to both await() calls in connectClientWithCorrectAccessTokenWithNDEATHCreatedDevices and connectClientWithCorrectAccessTokenWithNDEATHWithAliasCreatedDevices so that a transient 404 is treated as "condition not yet met" and polling continues as intended. Co-Authored-By: Claude Sonnet 4.6 --- .../mqtt/sparkplug/AbstractMqttV5ClientSparkplugTest.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/AbstractMqttV5ClientSparkplugTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/AbstractMqttV5ClientSparkplugTest.java index 4e4d0debf9..a9bdbb55db 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/AbstractMqttV5ClientSparkplugTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/AbstractMqttV5ClientSparkplugTest.java @@ -191,6 +191,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte AtomicReference device = new AtomicReference<>(); await(alias + "find device [" + deviceName + "] after created") .atMost(200, TimeUnit.SECONDS) + .ignoreExceptions() .until(() -> { device.set(doGet("/api/tenant/devices?deviceName=" + deviceName, Device.class)); return device.get() != null; @@ -236,6 +237,7 @@ public abstract class AbstractMqttV5ClientSparkplugTest extends AbstractMqttInte AtomicReference device = new AtomicReference<>(); await(alias + "find device [" + deviceName + "] after created") .atMost(200, TimeUnit.SECONDS) + .ignoreExceptions() .until(() -> { device.set(doGet("/api/tenant/devices?deviceName=" + deviceName, Device.class)); return device.get() != null; From 0e3381bca36a644fa97342bd752ebcfb816585e8 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 5 Mar 2026 15:26:12 +0100 Subject: [PATCH 011/106] Fix similar flaky await patterns in MQTT transport tests doGet/doGetAsyncTyped assert HTTP 200 internally, so any non-200 response throws AssertionError which Awaitility re-throws immediately instead of continuing to poll. Add .ignoreExceptions() to three additional await() polling loops that call HTTP helpers: - AbstractMqttV5ClientSparkplugAttributesTest: two doGetAsyncTyped calls polling for attribute keys after NBIRTH/DBIRTH - AbstractMqttAttributesIntegrationTest: doGetAsyncTyped polling for attribute values after client publish Co-Authored-By: Claude Sonnet 4.6 --- .../attributes/AbstractMqttAttributesIntegrationTest.java | 1 + .../attributes/AbstractMqttV5ClientSparkplugAttributesTest.java | 2 ++ 2 files changed, 3 insertions(+) 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> attributes = doGetAsyncTyped(attributeValuesUrl, new TypeReference<>() { }); diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/AbstractMqttV5ClientSparkplugAttributesTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/AbstractMqttV5ClientSparkplugAttributesTest.java index 756c8e603c..d42adfbcf0 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/AbstractMqttV5ClientSparkplugAttributesTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/AbstractMqttV5ClientSparkplugAttributesTest.java @@ -468,6 +468,7 @@ public abstract class AbstractMqttV5ClientSparkplugAttributesTest extends Abstra AtomicReference> actualKeys = new AtomicReference<>(); await(alias + SparkplugMessageType.NBIRTH.name()) .atMost(40, TimeUnit.SECONDS) + .ignoreExceptions() .until(() -> { actualKeys.set(doGetAsyncTyped(urlTemplate, new TypeReference<>() { })); @@ -483,6 +484,7 @@ public abstract class AbstractMqttV5ClientSparkplugAttributesTest extends Abstra AtomicReference> actualKeys = new AtomicReference<>(); await(alias + SparkplugMessageType.DBIRTH.name()) .atMost(40, TimeUnit.SECONDS) + .ignoreExceptions() .until(() -> { actualKeys.set(doGetAsyncTyped(urlTemplate, new TypeReference<>() { })); From 148dd17ce1e3c1a50bd5ee4d8e4aeb53f4fcccd8 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Fri, 6 Mar 2026 09:07:06 +0100 Subject: [PATCH 012/106] Fix flaky TenantControllerTest by draining housekeeper before teardown MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Tests like testFindTenantsByTitle create 261 tenants and delete them via deleteEntitiesAsync, which only waits for HTTP responses — not housekeeper completion. Each tenant deletion submits ~30 TenantEntitiesDeletionHousekeeper tasks (~7800 tasks total), which cascade further. The teardown's deleteTenant then waits for lag==0 with a 90s timeout, which is insufficient for this backlog and causes ConditionTimeoutException. Fix: add awaitHousekeeperDrained() (5-min timeout) called at the start of teardownWebTest so any pending housekeeper work from the test body drains before per-tenant teardown deletions begin. Co-Authored-By: Claude Sonnet 4.6 --- .../thingsboard/server/controller/AbstractWebTest.java | 9 +++++++++ .../src/test/resources/application-test.properties | 1 + 2 files changed, 10 insertions(+) diff --git a/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java b/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java index 3b7286cf01..5f4b7e952f 100644 --- a/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java @@ -405,6 +405,10 @@ public abstract class AbstractWebTest extends AbstractInMemoryStorageTest { public void teardownWebTest() throws Exception { log.debug("Executing web test teardown"); + // Drain any pending housekeeper work left by the test body (e.g., bulk tenant deletes) + // before proceeding with teardown deletions, to avoid 90s per-tenant wait timing out. + awaitHousekeeperDrained(); + loginSysAdmin(); deleteTenant(tenantId); deleteDifferentTenant(); @@ -436,6 +440,11 @@ public abstract class AbstractWebTest extends AbstractInMemoryStorageTest { .until(() -> storage.getLag("tb_housekeeper") == 0); } + protected void awaitHousekeeperDrained() { + Awaitility.await("housekeeper drained").atMost(5, TimeUnit.MINUTES).during(300, TimeUnit.MILLISECONDS) + .until(() -> storage.getLag("tb_housekeeper") == 0); + } + private List getAllTenants() throws Exception { List loadedTenants = new ArrayList<>(); PageLink pageLink = new PageLink(10); diff --git a/application/src/test/resources/application-test.properties b/application/src/test/resources/application-test.properties index e79289340c..7f0ab964d6 100644 --- a/application/src/test/resources/application-test.properties +++ b/application/src/test/resources/application-test.properties @@ -44,6 +44,7 @@ queue.transport_api.response_poll_interval=5 queue.transport.poll_interval=5 queue.core.poll-interval=5 queue.core.partitions=2 +queue.core.housekeeper.task-reprocessing-delay-ms=0 queue.rule-engine.poll-interval=5 queue.rule-engine.stats.enabled=true From f73abcdebd1bcfd2b910f97169cd4a6d3bf3b60c Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Mon, 9 Mar 2026 09:30:31 +0100 Subject: [PATCH 013/106] bump frontend-maven-plugin version to 2.0.0 --- pom.xml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pom.xml b/pom.xml index f00718c322..cd58b11f54 100755 --- a/pom.xml +++ b/pom.xml @@ -655,7 +655,7 @@ com.github.eirslett frontend-maven-plugin - 1.12.0 + 2.0.0 org.apache.maven.plugins From e82861b3e9ae7ce6398e0544bd742dd23627e7af Mon Sep 17 00:00:00 2001 From: dashevchenko Date: Mon, 9 Mar 2026 10:50:30 +0200 Subject: [PATCH 014/106] refactoring --- .../subscription/DefaultTbLocalSubscriptionService.java | 4 ++-- .../java/org/thingsboard/server/common/data/kv/TsKvEntry.java | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java index df28da2dad..2b1a810924 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java @@ -348,7 +348,7 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer if (sub.isLatestValues()) { for (TsKvEntry kv : data) { Long stateTs = keyStates.get(kv.getKey()); - if (stateTs == null || kv.getTs() >= stateTs || kv.isDeletedEntryMarker()) { + if (stateTs == null || kv.getTs() >= stateTs || kv.isDeletedEntry()) { if (updateData == null) { updateData = new ArrayList<>(); } @@ -362,7 +362,7 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer for (TsKvEntry kv : data) { Long stateTs = keyStates.get(kv.getKey()); if (stateTs != null) { - if (!sub.isLatestValues() || kv.getTs() >= stateTs || kv.isDeletedEntryMarker()) { + if (!sub.isLatestValues() || kv.getTs() >= stateTs || kv.isDeletedEntry()) { if (updateData == null) { updateData = new ArrayList<>(); } 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 eca0609704..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 @@ -38,7 +38,7 @@ public interface TsKvEntry extends KvEntry, HasVersion { } @JsonIgnore - default boolean isDeletedEntryMarker() { + default boolean isDeletedEntry() { return getTs() == 0 && (getValue() == null || getValueAsString().isEmpty()); } From 0587adb789a9f0dfbb47b176fe240f82f781d9b6 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Mon, 9 Mar 2026 10:53:07 +0200 Subject: [PATCH 015/106] Refactoring after review --- .../ttl/NotificationsCleanUpService.java | 28 +++++++++---------- .../server/dao/sql/rpc/RpcRepository.java | 2 ++ 2 files changed, 16 insertions(+), 14 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/NotificationsCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/NotificationsCleanUpService.java index 5e327a72db..83855bc8f8 100644 --- a/application/src/main/java/org/thingsboard/server/service/ttl/NotificationsCleanUpService.java +++ b/application/src/main/java/org/thingsboard/server/service/ttl/NotificationsCleanUpService.java @@ -62,17 +62,15 @@ public class NotificationsCleanUpService extends AbstractCleanUpService { public void cleanUp() { long expTime = System.currentTimeMillis() - TimeUnit.SECONDS.toMillis(ttlInSec); long partitionDurationMs = TimeUnit.HOURS.toMillis(partitionSizeInHours); - if (!isSystemTenantPartitionMine()) { + if (isSystemTenantPartitionMine()) { + partitioningRepository.dropPartitionsBefore(NOTIFICATION_TABLE_NAME, expTime, partitionDurationMs); + } else { partitioningRepository.cleanupPartitionsCache(NOTIFICATION_TABLE_NAME, expTime, partitionDurationMs); - return; } - long lastRemovedNotificationTs = partitioningRepository.dropPartitionsBefore(NOTIFICATION_TABLE_NAME, expTime, partitionDurationMs); - if (lastRemovedNotificationTs > 0) { - long gap = TimeUnit.MINUTES.toMillis(10); - long requestExpTime = lastRemovedNotificationTs - TimeUnit.SECONDS.toMillis(NotificationRequestConfig.MAX_SENDING_DELAY) - gap; - cleanUpNotificationRequests(requestExpTime); - } + long gap = TimeUnit.MINUTES.toMillis(10); + long requestExpTime = expTime - TimeUnit.SECONDS.toMillis(NotificationRequestConfig.MAX_SENDING_DELAY) - gap; + cleanUpNotificationRequests(requestExpTime); } private void cleanUpNotificationRequests(long expirationTime) { @@ -80,13 +78,15 @@ public class NotificationsCleanUpService extends AbstractCleanUpService { int totalRemoved = 0; int tenantsProcessed = 0; - // Clean up SYSADMIN's notification requests: - try { - totalRemoved += cleanUpByTenant(TenantId.SYS_TENANT_ID, expirationTime); - } catch (Exception e) { - log.warn("Failed to clean up notification requests for sysadmin {}", TenantId.SYS_TENANT_ID, e); + // Clean up SYSADMIN's notification requests on the system node only + if (isSystemTenantPartitionMine()) { + try { + totalRemoved += cleanUpByTenant(TenantId.SYS_TENANT_ID, expirationTime); + } catch (Exception e) { + log.warn("Failed to clean up notification requests for sysadmin {}", TenantId.SYS_TENANT_ID, e); + } } - // Clean up notification requests for tenants + // Each node cleans up notification requests for its own tenants PageDataIterable tenants = new PageDataIterable<>(tenantService::findTenantsIds, 10_000); for (TenantId tenantId : tenants) { try { 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 7be106219b..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,6 +21,7 @@ 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; @@ -34,6 +35,7 @@ public interface RpcRepository extends JpaRepository { Page findAllByTenantId(UUID tenantId, Pageable pageable); + @Transactional @Modifying @Query(value = "DELETE FROM rpc WHERE id IN " + "(SELECT id FROM rpc WHERE tenant_id = :tenantId AND created_time < :expirationTime LIMIT :batchSize)", From dd9fdb4181166084a21faa76f91e9708bb98f308 Mon Sep 17 00:00:00 2001 From: dashevchenko Date: Mon, 9 Mar 2026 11:07:26 +0200 Subject: [PATCH 016/106] refactoring --- .../DefaultTbEntityDataSubscriptionService.java | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java index 71df830f8d..fc8c5c3be2 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java @@ -245,7 +245,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc public void onFailure(Throwable t) { log.warn("[{}][{}] Failed to process command", finalCtx.getSessionId(), finalCtx.getCmdId(), t); if (t instanceof ResultSetSizeLimitExceededException) { - finalCtx.sendWsMsg(new EntityDataUpdate(finalCtx.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR.getCode(), t.getMessage())); + sendError(finalCtx, t); } } }, wsCallBackExecutor); @@ -270,7 +270,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc public void onFailure(Throwable t) { log.warn("[{}][{}] Failed to process timeseries command", ctx.getSessionId(), ctx.getCmdId(), t); if (t instanceof ResultSetSizeLimitExceededException) { - ctx.sendWsMsg(new EntityDataUpdate(ctx.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR.getCode(), t.getMessage())); + sendError(ctx, t); } } }, wsCallBackExecutor); @@ -283,6 +283,10 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } } + private void sendError(TbEntityDataSubCtx ctx, Throwable t) { + ctx.sendWsMsg(new EntityDataUpdate(ctx.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR.getCode(), t.getMessage())); + } + private void checkAndSendInitialData(@Nullable TbEntityDataSubCtx theCtx) { if (!theCtx.isInitialDataSent()) { EntityDataUpdate update = new EntityDataUpdate(theCtx.getCmdId(), theCtx.getData(), null, theCtx.getMaxEntitiesPerDataSubscription()); From 7cce918f90158a10fadc6d070b6f5e07576186e1 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Mon, 9 Mar 2026 11:13:38 +0200 Subject: [PATCH 017/106] Change log level to warn --- .../org/thingsboard/server/dao/ota/BaseOtaPackageService.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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 c43274cc26..54577ff9eb 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 @@ -234,10 +234,10 @@ public class BaseOtaPackageService extends AbstractCachedEntityService Date: Sat, 30 Aug 2025 08:56:43 +0200 Subject: [PATCH 018/106] MQTTS metrics --- .../transport/mqtt/MqttTransportContext.java | 22 ++++++++++++++----- .../transport/mqtt/MqttTransportHandler.java | 8 +++++-- .../common/transport/TransportService.java | 2 +- .../service/DefaultTransportService.java | 18 ++++++++++++--- 4 files changed, 38 insertions(+), 12 deletions(-) diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java index 3d05e999e0..8a60168154 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java @@ -88,20 +88,30 @@ public class MqttTransportContext extends TransportContext { @Value("${transport.mqtt.proxy_enabled:false}") private boolean proxyEnabled; - private final AtomicInteger connectionsCounter = new AtomicInteger(); + private final AtomicInteger connectionsActiveCounterMQTT = new AtomicInteger(); + private final AtomicInteger connectionsActiveCounterMQTTS = new AtomicInteger(); @PostConstruct public void init() { super.init(); - transportService.createGaugeStats("openConnections", connectionsCounter); + transportService.createGaugeStats("connections_active", connectionsActiveCounterMQTT, "protocol", "MQTT"); + transportService.createGaugeStats("connections_active", connectionsActiveCounterMQTTS, "protocol", "MQTTS"); } - public void channelRegistered() { - connectionsCounter.incrementAndGet(); + public void channelRegistered(boolean isSSL) { + if (isSSL) { + connectionsActiveCounterMQTTS.incrementAndGet(); + } else { + connectionsActiveCounterMQTT.incrementAndGet(); + } } - public void channelUnregistered() { - connectionsCounter.decrementAndGet(); + public void channelUnregistered(boolean isSSL) { + if (isSSL) { + connectionsActiveCounterMQTTS.decrementAndGet(); + } else { + connectionsActiveCounterMQTT.decrementAndGet(); + } } public boolean checkAddress(InetSocketAddress address) { diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index f753579292..bf623d67d8 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -180,16 +180,20 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement this.rpcAwaitingAck = new ConcurrentHashMap<>(); } + boolean isSSL() { + return sslHandler != null; + } + @Override public void channelRegistered(ChannelHandlerContext ctx) throws Exception { super.channelRegistered(ctx); - context.channelRegistered(); + context.channelRegistered(isSSL()); } @Override public void channelUnregistered(ChannelHandlerContext ctx) throws Exception { super.channelUnregistered(ctx); - context.channelUnregistered(); + context.channelUnregistered(isSSL()); } @Override diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java index 5c7d552855..612f331757 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java @@ -161,5 +161,5 @@ public interface TransportService { boolean hasSession(SessionInfoProto sessionInfo); - void createGaugeStats(String openConnections, AtomicInteger connectionsCounter); + void createGaugeStats(String openConnections, AtomicInteger connectionsCounter, String... tags); } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java index 260957c990..47e341547e 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java @@ -1257,9 +1257,21 @@ public class DefaultTransportService extends TransportActivityManager implements } @Override - public void createGaugeStats(String statsName, AtomicInteger number) { - statsFactory.createGauge(StatsType.TRANSPORT + "." + statsName, number); - statsMap.put(statsName, number); + public void createGaugeStats(String statsName, AtomicInteger number, String... tags) { + String key = "thingsboard" + "." + StatsType.TRANSPORT.getName() + "." + statsName; + statsFactory.createGauge(key, number, tags); + statsMap.put(statsName + TagsKey(tags), number); + } + + String TagsKey(String... tags) { + if (tags == null || tags.length < 2) return ""; + StringBuilder sb = new StringBuilder("["); + for (int i = 0; i < tags.length; i += 2) { + if (i > 0) sb.append(','); + sb.append(tags[i]).append('=').append(i + 1 < tags.length ? tags[i + 1] : ""); + } + sb.append(']'); + return sb.toString(); } @Scheduled(fixedDelayString = "${transport.stats.print-interval-ms:60000}") From 2a926cfbdfc5a7c09363a6dec078a211af8cf8d6 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Sat, 30 Aug 2025 09:39:04 +0200 Subject: [PATCH 019/106] MQTT transport: client remote address logged on exceptionCaught --- .../transport/mqtt/MqttTransportHandler.java | 17 ++++++++++++++--- 1 file changed, 14 insertions(+), 3 deletions(-) diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index bf623d67d8..aadeff0dee 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -1152,21 +1152,32 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement @Override public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { + String clientAddr = null; + try { + InetSocketAddress remote = getAddress(ctx); + clientAddr = remote != null ? (remote.getAddress() != null ? remote.getAddress().getHostAddress() : remote.getHostString()) + ":" + remote.getPort() : "IPunknown"; + } catch (Exception ignored) { + } + if (cause instanceof IOException) { if (log.isDebugEnabled()) { - log.debug("[{}][{}][{}] IOException: {}", sessionId, + log.debug("[{}][{}][{}][{}] {}: {}", sessionId, Optional.ofNullable(this.deviceSessionCtx.getDeviceInfo()).map(TransportDeviceInfo::getDeviceId).orElse(null), Optional.ofNullable(this.deviceSessionCtx.getDeviceInfo()).map(TransportDeviceInfo::getDeviceName).orElse(""), + clientAddr, + cause.getClass().getSimpleName(), cause.getMessage(), cause); } else if (log.isInfoEnabled()) { - log.info("[{}][{}][{}] IOException: {}", sessionId, + log.info("[{}][{}][{}][{}] {}: {}", sessionId, Optional.ofNullable(this.deviceSessionCtx.getDeviceInfo()).map(TransportDeviceInfo::getDeviceId).orElse(null), Optional.ofNullable(this.deviceSessionCtx.getDeviceInfo()).map(TransportDeviceInfo::getDeviceName).orElse(""), + clientAddr, + cause.getClass().getSimpleName(), cause.getMessage()); } } else { - log.error("[{}] Unexpected Exception", sessionId, cause); + log.error("[{}][{}] Unexpected Exception", sessionId, clientAddr, cause); } closeCtx(ctx, MqttReasonCodes.Disconnect.SERVER_SHUTTING_DOWN); From a643a340d8136e4da657488e34a5e840750b1791 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 5 Mar 2026 13:41:05 +0100 Subject: [PATCH 020/106] Address code review: fix naming, visibility, and lazy address resolution MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Rename TagsKey → tagsKey and make it private (Java naming convention) - Make isSSL() private in MqttTransportHandler (internal use only) - Fix double space in if (isSSL) in MqttTransportContext - Extract getClientAddr() helper and move clientAddr computation inside logging guards so address resolution is skipped when logging is disabled Co-Authored-By: Claude Sonnet 4.6 --- .../transport/mqtt/MqttTransportContext.java | 4 ++-- .../transport/mqtt/MqttTransportHandler.java | 24 ++++++++++++------- .../service/DefaultTransportService.java | 4 ++-- 3 files changed, 19 insertions(+), 13 deletions(-) diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java index 8a60168154..10245c5a24 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java @@ -99,7 +99,7 @@ public class MqttTransportContext extends TransportContext { } public void channelRegistered(boolean isSSL) { - if (isSSL) { + if (isSSL) { connectionsActiveCounterMQTTS.incrementAndGet(); } else { connectionsActiveCounterMQTT.incrementAndGet(); @@ -107,7 +107,7 @@ public class MqttTransportContext extends TransportContext { } public void channelUnregistered(boolean isSSL) { - if (isSSL) { + if (isSSL) { connectionsActiveCounterMQTTS.decrementAndGet(); } else { connectionsActiveCounterMQTT.decrementAndGet(); diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index aadeff0dee..d427c9b78e 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -180,7 +180,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement this.rpcAwaitingAck = new ConcurrentHashMap<>(); } - boolean isSSL() { + private boolean isSSL() { return sslHandler != null; } @@ -258,6 +258,17 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } } + private String getClientAddr(ChannelHandlerContext ctx) { + try { + InetSocketAddress remote = getAddress(ctx); + if (remote == null) return "unknown"; + String host = remote.getAddress() != null ? remote.getAddress().getHostAddress() : remote.getHostString(); + return host + ":" + remote.getPort(); + } catch (Exception ignored) { + return "unknown"; + } + } + InetSocketAddress getAddress(ChannelHandlerContext ctx) { var address = ctx.channel().attr(MqttTransportService.ADDRESS).get(); if (address == null) { @@ -1152,15 +1163,9 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement @Override public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { - String clientAddr = null; - try { - InetSocketAddress remote = getAddress(ctx); - clientAddr = remote != null ? (remote.getAddress() != null ? remote.getAddress().getHostAddress() : remote.getHostString()) + ":" + remote.getPort() : "IPunknown"; - } catch (Exception ignored) { - } - if (cause instanceof IOException) { if (log.isDebugEnabled()) { + String clientAddr = getClientAddr(ctx); log.debug("[{}][{}][{}][{}] {}: {}", sessionId, Optional.ofNullable(this.deviceSessionCtx.getDeviceInfo()).map(TransportDeviceInfo::getDeviceId).orElse(null), Optional.ofNullable(this.deviceSessionCtx.getDeviceInfo()).map(TransportDeviceInfo::getDeviceName).orElse(""), @@ -1169,6 +1174,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement cause.getMessage(), cause); } else if (log.isInfoEnabled()) { + String clientAddr = getClientAddr(ctx); log.info("[{}][{}][{}][{}] {}: {}", sessionId, Optional.ofNullable(this.deviceSessionCtx.getDeviceInfo()).map(TransportDeviceInfo::getDeviceId).orElse(null), Optional.ofNullable(this.deviceSessionCtx.getDeviceInfo()).map(TransportDeviceInfo::getDeviceName).orElse(""), @@ -1177,7 +1183,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement cause.getMessage()); } } else { - log.error("[{}][{}] Unexpected Exception", sessionId, clientAddr, cause); + log.error("[{}][{}] Unexpected Exception", sessionId, getClientAddr(ctx), cause); } closeCtx(ctx, MqttReasonCodes.Disconnect.SERVER_SHUTTING_DOWN); diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java index 47e341547e..e8ade498db 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java @@ -1260,10 +1260,10 @@ public class DefaultTransportService extends TransportActivityManager implements public void createGaugeStats(String statsName, AtomicInteger number, String... tags) { String key = "thingsboard" + "." + StatsType.TRANSPORT.getName() + "." + statsName; statsFactory.createGauge(key, number, tags); - statsMap.put(statsName + TagsKey(tags), number); + statsMap.put(statsName + tagsKey(tags), number); } - String TagsKey(String... tags) { + private String tagsKey(String... tags) { if (tags == null || tags.length < 2) return ""; StringBuilder sb = new StringBuilder("["); for (int i = 0; i < tags.length; i += 2) { From a2393f368c7bcaae76a0d98426987fff6a724180 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 5 Mar 2026 13:42:49 +0100 Subject: [PATCH 021/106] Keep isSSL() and getClientAddr() package-private for test accessibility Methods stubbed via Mockito spy in same-package tests must remain package-private. Revert isSSL() to package-private; getClientAddr() follows the same convention as getAddress(). Co-Authored-By: Claude Sonnet 4.6 --- .../server/transport/mqtt/MqttTransportHandler.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index d427c9b78e..48fbad3032 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -180,7 +180,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement this.rpcAwaitingAck = new ConcurrentHashMap<>(); } - private boolean isSSL() { + boolean isSSL() { return sslHandler != null; } @@ -258,7 +258,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } } - private String getClientAddr(ChannelHandlerContext ctx) { + String getClientAddr(ChannelHandlerContext ctx) { try { InetSocketAddress remote = getAddress(ctx); if (remote == null) return "unknown"; From 2e612899e250b5313d0268c7214dab12475a266b Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Tue, 10 Mar 2026 10:16:42 +0100 Subject: [PATCH 022/106] Fix NotificationRuleRecipientsConfig serialization with Jackson 2.18.x In Jackson 2.18.x, EXISTING_PROPERTY type info combined with the no-arg @JsonIgnoreProperties causes the triggerType discriminator field to be silently excluded from the serialized JSON. When the server then tries to deserialize the POST body for /api/notification/rule, Jackson cannot find triggerType and throws "missing type id property 'triggerType'", resulting in a 500 for NotificationEdgeTest.testNotificationRule. Fix by: 1. Adding @JsonProperty("triggerType") to force the field into normal bean serialization, overriding any suppression by the type info machinery. 2. Replacing the no-arg @JsonIgnoreProperties with @JsonIgnoreProperties( ignoreUnknown = true) so unknown properties are ignored rather than causing errors (e.g. for forward compatibility). Co-Authored-By: Claude Sonnet 4.6 --- .../notification/rule/NotificationRuleRecipientsConfig.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) 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 From 95016447ccdaf9b120c0885292ed346748166b56 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Tue, 10 Mar 2026 11:10:04 +0100 Subject: [PATCH 023/106] Fix EntityViewControllerTest MQTT port collision with other test contexts EntityViewControllerTest was importing MQTT_PORT from AbstractMqttIntegrationTest, a static final field initialized once per JVM. When running in the same Surefire fork alongside other test classes that also use this constant (e.g. MqttGatewayRateLimitsTest, DeviceEdgeTest), each class gets a different Spring context key but all try to bind MqttTransportService to the same port, causing BindException. Fix: define a private static MQTT_PORT/MQTT_URL directly in EntityViewControllerTest so its Spring context gets its own independently allocated port. Co-Authored-By: Claude Sonnet 4.6 --- .../server/controller/EntityViewControllerTest.java | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/controller/EntityViewControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/EntityViewControllerTest.java index 421f3785ca..311ef54fa4 100644 --- a/application/src/test/java/org/thingsboard/server/controller/EntityViewControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/EntityViewControllerTest.java @@ -41,6 +41,7 @@ import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.DynamicPropertyRegistry; import org.springframework.test.context.DynamicPropertySource; import org.springframework.test.context.TestPropertySource; +import org.springframework.test.util.TestSocketUtils; import org.springframework.test.web.servlet.ResultActions; import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.server.common.data.Customer; @@ -87,8 +88,6 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; import static org.thingsboard.server.dao.model.ModelConstants.NULL_UUID; -import static org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest.MQTT_PORT; -import static org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest.MQTT_URL; @TestPropertySource(properties = { "transport.mqtt.enabled=true", @@ -98,6 +97,9 @@ import static org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest. @ContextConfiguration(classes = {EntityViewControllerTest.Config.class}) @DaoSqlTest public class EntityViewControllerTest extends AbstractControllerTest { + static final int MQTT_PORT = TestSocketUtils.findAvailableTcpPort(); + static final String MQTT_URL = "tcp://localhost:" + MQTT_PORT; + @DynamicPropertySource static void props(DynamicPropertyRegistry registry) { log.warn("transport.mqtt.bind_port = {}", MQTT_PORT); From 42179bddf8dadf78a242c8943cc81574f9628245 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Tue, 10 Mar 2026 11:24:53 +0100 Subject: [PATCH 024/106] Add comment explaining why EntityViewControllerTest owns its MQTT port Co-Authored-By: Claude Sonnet 4.6 --- .../server/controller/EntityViewControllerTest.java | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/application/src/test/java/org/thingsboard/server/controller/EntityViewControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/EntityViewControllerTest.java index 311ef54fa4..b65d4178f0 100644 --- a/application/src/test/java/org/thingsboard/server/controller/EntityViewControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/EntityViewControllerTest.java @@ -97,6 +97,12 @@ import static org.thingsboard.server.dao.model.ModelConstants.NULL_UUID; @ContextConfiguration(classes = {EntityViewControllerTest.Config.class}) @DaoSqlTest public class EntityViewControllerTest extends AbstractControllerTest { + // Must NOT be imported from AbstractMqttIntegrationTest. That field is a static final initialized + // once per JVM. Other test classes (e.g. MqttGatewayRateLimitsTest, DeviceEdgeTest) share the same + // constant but produce a different Spring context cache key, so Spring creates a separate + // ApplicationContext for each of them. Every context starts its own MqttTransportService and tries + // to bind the same port, causing BindException when tests run in the same Surefire JVM fork. + // Declaring the port here gives this context its own independently allocated port. static final int MQTT_PORT = TestSocketUtils.findAvailableTcpPort(); static final String MQTT_URL = "tcp://localhost:" + MQTT_PORT; From bd181f115ed8ea536a551bfc3fd984f3527a5548 Mon Sep 17 00:00:00 2001 From: Vladyslav_Prykhodko Date: Tue, 10 Mar 2026 17:26:15 +0200 Subject: [PATCH 025/106] UI: Hidden show on widgets button in Sys Admin users --- .../components/attribute/attribute-table.component.html | 3 ++- .../home/components/attribute/attribute-table.component.ts | 7 +++++-- 2 files changed, 7 insertions(+), 3 deletions(-) diff --git a/ui-ngx/src/app/modules/home/components/attribute/attribute-table.component.html b/ui-ngx/src/app/modules/home/components/attribute/attribute-table.component.html index 6a797a37aa..6d758b8c7b 100644 --- a/ui-ngx/src/app/modules/home/components/attribute/attribute-table.component.html +++ b/ui-ngx/src/app/modules/home/components/attribute/attribute-table.component.html @@ -93,7 +93,8 @@ (click)="deleteTelemetry($event)"> delete -