From 0468451f3ba02983e540e26f968305dd7a2ed04a Mon Sep 17 00:00:00 2001 From: dashevchenko Date: Mon, 12 Jan 2026 10:52:29 +0200 Subject: [PATCH 1/2] 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 e82861b3e9ae7ce6398e0544bd742dd23627e7af Mon Sep 17 00:00:00 2001 From: dashevchenko Date: Mon, 9 Mar 2026 10:50:30 +0200 Subject: [PATCH 2/2] 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()); }