Browse Source

Merge pull request #14781 from dashevchenko/noWebSocketOnTsDeletion

Added WS update on telemetry deletion
pull/15195/merge
Viacheslav Klimov 5 months ago
committed by GitHub
parent
commit
055c498ca3
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 4
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java
  2. 35
      application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java
  3. 5
      common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvEntry.java

4
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<>();
}

35
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<TsKvEntry> tsData = List.of(dataPoint);
sendTelemetry(device, tsData);
List<EntityKey> keys = List.of(new EntityKey(EntityKeyType.TIME_SERIES, "temperature"));
EntityDataUpdate update = getWsClient().subscribeLatestUpdate(keys, dtf);
Assert.assertEquals(1, update.getCmdId());
PageData<EntityData> 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<EntityData> 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();

5
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());
}
}

Loading…
Cancel
Save