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..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) { + 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) { + if (!sub.isLatestValues() || kv.getTs() >= stateTs || kv.isDeletedEntry()) { 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 eab878854d..6ec0f3ad64 100644 --- a/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java @@ -734,6 +734,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..595e1aa26b 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvEntry.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvEntry.java @@ -37,4 +37,9 @@ public interface TsKvEntry extends KvEntry, HasVersion { return new TsValue(getTs(), getValueAsString()); } + @JsonIgnore + default boolean isDeletedEntry() { + return getTs() == 0 && (getValue() == null || getValueAsString().isEmpty()); + } + }