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 4fd46195d0..404b4f59c4 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 @@ -66,6 +66,7 @@ import org.thingsboard.server.service.telemetry.sub.SubscriptionUpdate; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; import java.util.ArrayList; +import java.util.Collection; import java.util.HashMap; import java.util.LinkedHashMap; import java.util.LinkedHashSet; @@ -169,7 +170,8 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc if (ctx != null) { log.debug("[{}][{}] Updating existing subscriptions using: {}", session.getSessionId(), cmd.getCmdId(), cmd); if (cmd.getLatestCmd() != null || cmd.getTsCmd() != null) { - ctx.clearSubscriptions(); + Collection oldSubIds = ctx.clearSubscriptions(); + oldSubIds.forEach(subId -> localSubscriptionService.cancelSubscription(serviceId, subId)); } //TODO: cleanup old subscription; } else { @@ -195,10 +197,16 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc }); } PageData data = entityService.findEntityDataByQuery(tenantId, customerId, ctx.getQuery()); + if (log.isTraceEnabled()) { + data.getData().forEach(ed -> { + log.trace("[{}][{}] EntityData: {}", session.getSessionId(), cmd.getCmdId(), ed); + }); + } ctx.setData(data); } ListenableFuture historyFuture; if (cmd.getHistoryCmd() != null) { + log.trace("[{}][{}] Going to process history command: {}", session.getSessionId(), cmd.getCmdId(), cmd.getHistoryCmd()); historyFuture = handleHistoryCmd(ctx, cmd.getHistoryCmd()); } else { historyFuture = Futures.immediateFuture(ctx); @@ -241,8 +249,10 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } private void handleLatestCmd(TbEntityDataSubCtx ctx, LatestValueCmd latestCmd) { + log.trace("[{}][{}] Going to process latest command: {}", ctx.getSessionId(), ctx.getCmdId(), latestCmd); //Fetch the latest values for telemetry keys (in case they are not copied from NoSQL to SQL DB in hybrid mode. if (!tsInSqlDB) { + log.trace("[{}][{}] Going to fetch missing latest values: {}", ctx.getSessionId(), ctx.getCmdId(), latestCmd); List allTsKeys = latestCmd.getKeys().stream() .filter(key -> key.getType().equals(EntityKeyType.TIME_SERIES)) .map(EntityKey::getKey).collect(Collectors.toList()); diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java index 98ba5c954f..f00cd19c0e 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java @@ -16,6 +16,7 @@ package org.thingsboard.server.service.subscription; import lombok.Data; +import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -34,12 +35,13 @@ import org.thingsboard.server.service.telemetry.sub.SubscriptionUpdate; import java.util.ArrayList; import java.util.Arrays; +import java.util.Collection; import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; -import java.util.stream.Collectors; +@Slf4j @Data public class TbEntityDataSubCtx { @@ -54,7 +56,6 @@ public class TbEntityDataSubCtx { private PageData data; private boolean initialDataSent; private List tbSubs; - private int internalSubIdx; private Map subToEntityIdMap; public TbEntityDataSubCtx(String serviceId, TelemetryWebSocketService wsService, TelemetryWebSocketSessionRef sessionRef, int cmdId) { @@ -82,45 +83,61 @@ public class TbEntityDataSubCtx { public List createSubscriptions(List keys) { this.subToEntityIdMap = new HashMap<>(); - this.internalSubIdx = cmdId * MAX_SUBS_PER_CMD; tbSubs = new ArrayList<>(); - List attrSubKeys = new ArrayList<>(); - List tsSubKeys = new ArrayList<>(); - for (EntityKey key : keys) { - switch (key.getType()) { - case TIME_SERIES: - tsSubKeys.add(key); - break; - case ATTRIBUTE: - case CLIENT_ATTRIBUTE: - case SHARED_ATTRIBUTE: - case SERVER_ATTRIBUTE: - attrSubKeys.add(key); - } - } + Map> keysByType = new HashMap<>(); + keys.forEach(key -> keysByType.computeIfAbsent(key.getType(), k -> new ArrayList<>()).add(key)); for (EntityData entityData : data.getData()) { - if (!tsSubKeys.isEmpty()) { - tbSubs.add(createTsSub(entityData, tsSubKeys)); - } + keysByType.forEach((keysType, keysList) -> { + int subIdx = sessionRef.getSessionSubIdSeq().incrementAndGet(); + subToEntityIdMap.put(subIdx, entityData.getEntityId()); + switch (keysType) { + case TIME_SERIES: + tbSubs.add(createTsSub(entityData, subIdx, keysList)); + break; + case CLIENT_ATTRIBUTE: + tbSubs.add(createAttrSub(entityData, subIdx, keysType, TbAttributeSubscriptionScope.CLIENT_SCOPE, keysList)); + break; + case SHARED_ATTRIBUTE: + tbSubs.add(createAttrSub(entityData, subIdx, keysType, TbAttributeSubscriptionScope.SHARED_SCOPE, keysList)); + break; + case SERVER_ATTRIBUTE: + tbSubs.add(createAttrSub(entityData, subIdx, keysType, TbAttributeSubscriptionScope.SERVER_SCOPE, keysList)); + break; + case ATTRIBUTE: + tbSubs.add(createAttrSub(entityData, subIdx, keysType, TbAttributeSubscriptionScope.ANY_SCOPE, keysList)); + break; + } + }); } return tbSubs; } - private TbSubscription createTsSub(EntityData entityData, List tsSubKeys) { - int subIdx = internalSubIdx++; - subToEntityIdMap.put(subIdx, entityData.getEntityId()); - Map keyStates = new HashMap<>(); - tsSubKeys.forEach(key -> keyStates.put(key.getKey(), 0L)); - if (entityData.getLatest() != null) { - Map currentValues = entityData.getLatest().get(EntityKeyType.TIME_SERIES); - if (currentValues != null) { - currentValues.forEach((k, v) -> keyStates.put(k, v.getTs())); - } - } + private TbSubscription createAttrSub(EntityData entityData, int subIdx, EntityKeyType keysType, TbAttributeSubscriptionScope scope, List subKeys) { + Map keyStates = buildKeyStats(entityData, keysType, subKeys); + log.trace("[{}][{}][{}] Creating attributes subscription with keys: {}", serviceId, cmdId, subIdx, keyStates); + return TbAttributeSubscription.builder() + .serviceId(serviceId) + .sessionId(sessionRef.getSessionId()) + .subscriptionId(subIdx) + .tenantId(sessionRef.getSecurityCtx().getTenantId()) + .entityId(entityData.getEntityId()) + .updateConsumer((s, subscriptionUpdate) -> sendWsMsg(s, subscriptionUpdate, keysType)) + .allKeys(false) + .keyStates(keyStates) + .scope(scope) + .build(); + } + + private TbSubscription createTsSub(EntityData entityData, int subIdx, List subKeys) { + Map keyStates = buildKeyStats(entityData, EntityKeyType.TIME_SERIES, subKeys); if (entityData.getTimeseries() != null) { - entityData.getTimeseries().forEach((k, v) -> keyStates.put(k, Arrays.stream(v).map(TsValue::getTs).max(Long::compareTo).orElse(0L))); + entityData.getTimeseries().forEach((k, v) -> { + long ts = Arrays.stream(v).map(TsValue::getTs).max(Long::compareTo).orElse(0L); + log.trace("[{}][{}] Updating key: {} with ts: {}", serviceId, cmdId, k, ts); + keyStates.put(k, ts); + }); } - + log.trace("[{}][{}][{}] Creating time-series subscription with keys: {}", serviceId, cmdId, subIdx, keyStates); return TbTimeseriesSubscription.builder() .serviceId(serviceId) .sessionId(sessionRef.getSessionId()) @@ -129,26 +146,74 @@ public class TbEntityDataSubCtx { .entityId(entityData.getEntityId()) .updateConsumer(this::sendTsWsMsg) .allKeys(false) - .keyStates(keyStates).build(); + .keyStates(keyStates) + .build(); } + private Map buildKeyStats(EntityData entityData, EntityKeyType keysType, List subKeys) { + Map keyStates = new HashMap<>(); + subKeys.forEach(key -> keyStates.put(key.getKey(), 0L)); + if (entityData.getLatest() != null) { + Map currentValues = entityData.getLatest().get(keysType); + if (currentValues != null) { + currentValues.forEach((k, v) -> { + log.trace("[{}][{}] Updating key: {} with ts: {}", serviceId, cmdId, k, v.getTs()); + keyStates.put(k, v.getTs()); + }); + } + } + return keyStates; + } private void sendTsWsMsg(String sessionId, SubscriptionUpdate subscriptionUpdate) { + sendWsMsg(sessionId, subscriptionUpdate, EntityKeyType.TIME_SERIES); + } + + private void sendWsMsg(String sessionId, SubscriptionUpdate subscriptionUpdate, EntityKeyType keyType) { EntityId entityId = subToEntityIdMap.get(subscriptionUpdate.getSubscriptionId()); if (entityId != null) { - Map latest = new HashMap<>(); + log.trace("[{}][{}][{}][{}] Received subscription update: {}", sessionId, cmdId, subscriptionUpdate.getSubscriptionId(), keyType, subscriptionUpdate); + Map latestUpdate = new HashMap<>(); subscriptionUpdate.getData().forEach((k, v) -> { Object[] data = (Object[]) v.get(0); - latest.put(k, new TsValue((Long) data[0], (String) data[1])); + latestUpdate.put(k, new TsValue((Long) data[0], (String) data[1])); }); - Map> latestMap = Collections.singletonMap(EntityKeyType.TIME_SERIES, latest); - EntityData entityData = new EntityData(entityId, latestMap, null); - wsService.sendWsMsg(sessionId, new EntityDataUpdate(cmdId, null, Collections.singletonList(entityData))); + EntityData entityData = getDataForEntity(entityId); + if (entityData != null && entityData.getLatest() != null) { + Map latestCtxValues = entityData.getLatest().get(keyType); + log.trace("[{}][{}][{}] Going to compare update with {}", sessionId, cmdId, subscriptionUpdate.getSubscriptionId(), latestCtxValues); + if (latestCtxValues != null) { + latestCtxValues.forEach((k, v) -> { + TsValue update = latestUpdate.get(k); + if (update.getTs() < v.getTs()) { + log.trace("[{}][{}][{}] Removed stale update for key: {} and ts: {}", sessionId, cmdId, subscriptionUpdate.getSubscriptionId(), k, update.getTs()); + latestUpdate.remove(k); + } else if ((update.getTs() == v.getTs() && update.getValue().equals(v.getValue()))) { + log.trace("[{}][{}][{}] Removed duplicate update for key: {} and ts: {}", sessionId, cmdId, subscriptionUpdate.getSubscriptionId(), k, update.getTs()); + latestUpdate.remove(k); + } + }); + //Setting new values + latestUpdate.forEach(latestCtxValues::put); + } + } + if (!latestUpdate.isEmpty()) { + Map> latestMap = Collections.singletonMap(keyType, latestUpdate); + entityData = new EntityData(entityId, latestMap, null); + wsService.sendWsMsg(sessionId, new EntityDataUpdate(cmdId, null, Collections.singletonList(entityData))); + } + } else { + log.trace("[{}][{}][{}][{}] Received stale subscription update: {}", sessionId, cmdId, subscriptionUpdate.getSubscriptionId(), keyType, subscriptionUpdate); } + } + private EntityData getDataForEntity(EntityId entityId) { + return data.getData().stream().filter(item -> item.getEntityId().equals(entityId)).findFirst().orElse(null); } - public void clearSubscriptions() { + public Collection clearSubscriptions() { + List oldSubIds = new ArrayList<>(subToEntityIdMap.keySet()); subToEntityIdMap.clear(); + return oldSubIds; } } diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/TelemetryWebSocketSessionRef.java b/application/src/main/java/org/thingsboard/server/service/telemetry/TelemetryWebSocketSessionRef.java index cd70f5af4f..0c2a7cedbb 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/TelemetryWebSocketSessionRef.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/TelemetryWebSocketSessionRef.java @@ -20,6 +20,7 @@ import org.thingsboard.server.service.security.model.SecurityUser; import java.net.InetSocketAddress; import java.util.Objects; +import java.util.concurrent.atomic.AtomicInteger; /** * Created by ashvayka on 27.03.18. @@ -36,12 +37,15 @@ public class TelemetryWebSocketSessionRef { private final InetSocketAddress localAddress; @Getter private final InetSocketAddress remoteAddress; + @Getter + private final AtomicInteger sessionSubIdSeq; public TelemetryWebSocketSessionRef(String sessionId, SecurityUser securityCtx, InetSocketAddress localAddress, InetSocketAddress remoteAddress) { this.sessionId = sessionId; this.securityCtx = securityCtx; this.localAddress = localAddress; this.remoteAddress = remoteAddress; + this.sessionSubIdSeq = new AtomicInteger(); } @Override diff --git a/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java b/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java index 8fb6f59fd1..ec8ca198bb 100644 --- a/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java @@ -28,6 +28,8 @@ import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.kv.Aggregation; +import org.thingsboard.server.common.data.kv.AttributeKvEntry; +import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; @@ -41,6 +43,7 @@ import org.thingsboard.server.common.data.query.EntityKeyType; import org.thingsboard.server.common.data.query.TsValue; import org.thingsboard.server.common.data.security.Authority; import org.thingsboard.server.dao.timeseries.TimeseriesService; +import org.thingsboard.server.service.subscription.TbAttributeSubscriptionScope; import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; import org.thingsboard.server.service.telemetry.cmd.TelemetryPluginCmdsWrapper; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataCmd; @@ -48,6 +51,7 @@ import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUpdate; import org.thingsboard.server.service.telemetry.cmd.v2.EntityHistoryCmd; import org.thingsboard.server.service.telemetry.cmd.v2.LatestValueCmd; +import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; import java.util.List; @@ -139,7 +143,7 @@ public class BaseWebsocketApiTest extends AbstractWebsocketTest { List tsData = Arrays.asList(dataPoint1, dataPoint2, dataPoint3); sendTelemetry(device, tsData); - Thread.sleep(1000); + Thread.sleep(100); wsClient.send(mapper.writeValueAsString(wrapper)); msg = wsClient.waitForReply(); @@ -156,25 +160,8 @@ public class BaseWebsocketApiTest extends AbstractWebsocketTest { Assert.assertEquals(new TsValue(dataPoint3.getTs(), dataPoint3.getValueAsString()), tsArray[2]); } - private void sendTelemetry(Device device, List tsData) throws InterruptedException { - CountDownLatch latch = new CountDownLatch(1); - tsService.saveAndNotify(device.getTenantId(), device.getId(), tsData, 0, new FutureCallback() { - @Override - public void onSuccess(@Nullable Void result) { - latch.countDown(); - } - - @Override - public void onFailure(Throwable t) { - latch.countDown(); - } - }); - latch.await(3, TimeUnit.SECONDS); - } - @Test - @Ignore - public void testEntityDataLatestWsCmd() throws Exception { + public void testEntityDataLatestTsWsCmd() throws Exception { Device device = new Device(); device.setName("Device"); device.setType("default"); @@ -211,9 +198,9 @@ public class BaseWebsocketApiTest extends AbstractWebsocketTest { List tsData = Arrays.asList(dataPoint1); sendTelemetry(device, tsData); - cmd = new EntityDataCmd(1, edq, null, latestCmd, null); - + Thread.sleep(100); + cmd = new EntityDataCmd(1, edq, null, latestCmd, null); wrapper = new TelemetryPluginCmdsWrapper(); wrapper.setEntityDataCmds(Collections.singletonList(cmd)); @@ -231,11 +218,12 @@ public class BaseWebsocketApiTest extends AbstractWebsocketTest { TsValue tsValue = pageData.getData().get(0).getLatest().get(EntityKeyType.TIME_SERIES).get("temperature"); Assert.assertEquals(new TsValue(dataPoint1.getTs(), dataPoint1.getValueAsString()), tsValue); - log.error("GOING TO LISTEN FOR UPDATES"); - msg = wsClient.waitForUpdate(); now = System.currentTimeMillis(); TsKvEntry dataPoint2 = new BasicTsKvEntry(now, new LongDataEntry("temperature", 52L)); + + wsClient.registerWaitForUpdate(); sendTelemetry(device, Arrays.asList(dataPoint2)); + msg = wsClient.waitForUpdate(); update = mapper.readValue(msg, EntityDataUpdate.class); Assert.assertEquals(1, update.getCmdId()); @@ -247,6 +235,280 @@ public class BaseWebsocketApiTest extends AbstractWebsocketTest { tsValue = eData.get(0).getLatest().get(EntityKeyType.TIME_SERIES).get("temperature"); Assert.assertEquals(new TsValue(dataPoint2.getTs(), dataPoint2.getValueAsString()), tsValue); + //Sending update from the past, while latest value has new timestamp; + wsClient.registerWaitForUpdate(); + sendTelemetry(device, Arrays.asList(dataPoint1)); + msg = wsClient.waitForUpdate(TimeUnit.SECONDS.toMillis(1)); + Assert.assertNull(msg); + + //Sending duplicate update again + wsClient.registerWaitForUpdate(); + sendTelemetry(device, Arrays.asList(dataPoint2)); + msg = wsClient.waitForUpdate(TimeUnit.SECONDS.toMillis(1)); + Assert.assertNull(msg); + } + + @Test + public void testEntityDataLatestAttrWsCmd() throws Exception { + Device device = new Device(); + device.setName("Device"); + device.setType("default"); + device.setLabel("testLabel" + (int) (Math.random() * 1000)); + device = doPost("/api/device", device, Device.class); + + long now = System.currentTimeMillis(); + + DeviceTypeFilter dtf = new DeviceTypeFilter(); + dtf.setDeviceNameFilter("D"); + dtf.setDeviceType("default"); + EntityDataQuery edq = new EntityDataQuery(dtf, new EntityDataPageLink(1, 0, null, null), + Collections.emptyList(), Collections.emptyList(), Collections.emptyList()); + + LatestValueCmd latestCmd = new LatestValueCmd(); + latestCmd.setKeys(Collections.singletonList(new EntityKey(EntityKeyType.SERVER_ATTRIBUTE, "serverAttributeKey"))); + EntityDataCmd cmd = new EntityDataCmd(1, edq, null, latestCmd, null); + + TelemetryPluginCmdsWrapper wrapper = new TelemetryPluginCmdsWrapper(); + wrapper.setEntityDataCmds(Collections.singletonList(cmd)); + + wsClient.send(mapper.writeValueAsString(wrapper)); + String msg = wsClient.waitForReply(); + EntityDataUpdate update = mapper.readValue(msg, EntityDataUpdate.class); + 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.SERVER_ATTRIBUTE).get("serverAttributeKey")); + Assert.assertEquals(0, pageData.getData().get(0).getLatest().get(EntityKeyType.SERVER_ATTRIBUTE).get("serverAttributeKey").getTs()); + Assert.assertEquals("", pageData.getData().get(0).getLatest().get(EntityKeyType.SERVER_ATTRIBUTE).get("serverAttributeKey").getValue()); + + AttributeKvEntry dataPoint1 = new BaseAttributeKvEntry(now - TimeUnit.MINUTES.toMillis(1), new LongDataEntry("serverAttributeKey", 42L)); + List tsData = Arrays.asList(dataPoint1); + sendAttributes(device, TbAttributeSubscriptionScope.SERVER_SCOPE, tsData); + + Thread.sleep(100); + + cmd = new EntityDataCmd(1, edq, null, latestCmd, null); + wrapper = new TelemetryPluginCmdsWrapper(); + wrapper.setEntityDataCmds(Collections.singletonList(cmd)); + + wsClient.send(mapper.writeValueAsString(wrapper)); + msg = wsClient.waitForReply(); + update = mapper.readValue(msg, EntityDataUpdate.class); + + Assert.assertEquals(1, update.getCmdId()); + + 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.SERVER_ATTRIBUTE)); + TsValue tsValue = pageData.getData().get(0).getLatest().get(EntityKeyType.SERVER_ATTRIBUTE).get("serverAttributeKey"); + Assert.assertEquals(new TsValue(dataPoint1.getLastUpdateTs(), dataPoint1.getValueAsString()), tsValue); + + now = System.currentTimeMillis(); + AttributeKvEntry dataPoint2 = new BaseAttributeKvEntry(now, new LongDataEntry("serverAttributeKey", 52L)); + + wsClient.registerWaitForUpdate(); + sendAttributes(device, TbAttributeSubscriptionScope.SERVER_SCOPE, Arrays.asList(dataPoint2)); + msg = wsClient.waitForUpdate(); + + update = mapper.readValue(msg, EntityDataUpdate.class); + Assert.assertEquals(1, update.getCmdId()); + List eData = update.getUpdate(); + Assert.assertNotNull(eData); + Assert.assertEquals(1, eData.size()); + Assert.assertEquals(device.getId(), eData.get(0).getEntityId()); + Assert.assertNotNull(eData.get(0).getLatest().get(EntityKeyType.SERVER_ATTRIBUTE)); + tsValue = eData.get(0).getLatest().get(EntityKeyType.SERVER_ATTRIBUTE).get("serverAttributeKey"); + Assert.assertEquals(new TsValue(dataPoint2.getLastUpdateTs(), dataPoint2.getValueAsString()), tsValue); + + //Sending update from the past, while latest value has new timestamp; + wsClient.registerWaitForUpdate(); + sendAttributes(device, TbAttributeSubscriptionScope.SERVER_SCOPE, Arrays.asList(dataPoint1)); + msg = wsClient.waitForUpdate(TimeUnit.SECONDS.toMillis(1)); + Assert.assertNull(msg); + + //Sending duplicate update again + wsClient.registerWaitForUpdate(); + sendAttributes(device, TbAttributeSubscriptionScope.SERVER_SCOPE, Arrays.asList(dataPoint2)); + msg = wsClient.waitForUpdate(TimeUnit.SECONDS.toMillis(1)); + Assert.assertNull(msg); + } + + @Test + public void testEntityDataLatestAttrTypesWsCmd() throws Exception { + Device device = new Device(); + device.setName("Device"); + device.setType("default"); + device.setLabel("testLabel" + (int) (Math.random() * 1000)); + device = doPost("/api/device", device, Device.class); + + long now = System.currentTimeMillis(); + + DeviceTypeFilter dtf = new DeviceTypeFilter(); + dtf.setDeviceNameFilter("D"); + dtf.setDeviceType("default"); + EntityDataQuery edq = new EntityDataQuery(dtf, new EntityDataPageLink(1, 0, null, null), + Collections.emptyList(), Collections.emptyList(), Collections.emptyList()); + + LatestValueCmd latestCmd = new LatestValueCmd(); + List keys = new ArrayList<>(); + keys.add(new EntityKey(EntityKeyType.SERVER_ATTRIBUTE, "serverAttributeKey")); + keys.add(new EntityKey(EntityKeyType.CLIENT_ATTRIBUTE, "clientAttributeKey")); + keys.add(new EntityKey(EntityKeyType.SHARED_ATTRIBUTE, "sharedAttributeKey")); + keys.add(new EntityKey(EntityKeyType.ATTRIBUTE, "anyAttributeKey")); + latestCmd.setKeys(keys); + EntityDataCmd cmd = new EntityDataCmd(1, edq, null, latestCmd, null); + + TelemetryPluginCmdsWrapper wrapper = new TelemetryPluginCmdsWrapper(); + wrapper.setEntityDataCmds(Collections.singletonList(cmd)); + + wsClient.send(mapper.writeValueAsString(wrapper)); + String msg = wsClient.waitForReply(); + EntityDataUpdate update = mapper.readValue(msg, EntityDataUpdate.class); + 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.SERVER_ATTRIBUTE).get("serverAttributeKey")); + Assert.assertEquals(0, pageData.getData().get(0).getLatest().get(EntityKeyType.SERVER_ATTRIBUTE).get("serverAttributeKey").getTs()); + Assert.assertEquals("", pageData.getData().get(0).getLatest().get(EntityKeyType.SERVER_ATTRIBUTE).get("serverAttributeKey").getValue()); + Assert.assertNotNull(pageData.getData().get(0).getLatest().get(EntityKeyType.CLIENT_ATTRIBUTE).get("clientAttributeKey")); + Assert.assertEquals(0, pageData.getData().get(0).getLatest().get(EntityKeyType.CLIENT_ATTRIBUTE).get("clientAttributeKey").getTs()); + Assert.assertEquals("", pageData.getData().get(0).getLatest().get(EntityKeyType.CLIENT_ATTRIBUTE).get("clientAttributeKey").getValue()); + Assert.assertNotNull(pageData.getData().get(0).getLatest().get(EntityKeyType.SHARED_ATTRIBUTE).get("sharedAttributeKey")); + Assert.assertEquals(0, pageData.getData().get(0).getLatest().get(EntityKeyType.SHARED_ATTRIBUTE).get("sharedAttributeKey").getTs()); + Assert.assertEquals("", pageData.getData().get(0).getLatest().get(EntityKeyType.SHARED_ATTRIBUTE).get("sharedAttributeKey").getValue()); + Assert.assertNotNull(pageData.getData().get(0).getLatest().get(EntityKeyType.ATTRIBUTE).get("anyAttributeKey")); + Assert.assertEquals(0, pageData.getData().get(0).getLatest().get(EntityKeyType.ATTRIBUTE).get("anyAttributeKey").getTs()); + Assert.assertEquals("", pageData.getData().get(0).getLatest().get(EntityKeyType.ATTRIBUTE).get("anyAttributeKey").getValue()); + + + wsClient.registerWaitForUpdate(); + AttributeKvEntry dataPoint1 = new BaseAttributeKvEntry(now - TimeUnit.MINUTES.toMillis(1), new LongDataEntry("serverAttributeKey", 42L)); + List tsData = Arrays.asList(dataPoint1); + sendAttributes(device, TbAttributeSubscriptionScope.SERVER_SCOPE, tsData); + + Thread.sleep(100); + + cmd = new EntityDataCmd(1, edq, null, latestCmd, null); + wrapper = new TelemetryPluginCmdsWrapper(); + wrapper.setEntityDataCmds(Collections.singletonList(cmd)); + + msg = wsClient.waitForUpdate(); + + update = mapper.readValue(msg, EntityDataUpdate.class); + Assert.assertEquals(1, update.getCmdId()); + List eData = update.getUpdate(); + Assert.assertNotNull(eData); + Assert.assertEquals(1, eData.size()); + Assert.assertEquals(device.getId(), eData.get(0).getEntityId()); + Assert.assertNotNull(eData.get(0).getLatest().get(EntityKeyType.SERVER_ATTRIBUTE)); + TsValue attrValue = eData.get(0).getLatest().get(EntityKeyType.SERVER_ATTRIBUTE).get("serverAttributeKey"); + Assert.assertEquals(new TsValue(dataPoint1.getLastUpdateTs(), dataPoint1.getValueAsString()), attrValue); + + //Sending update from the past, while latest value has new timestamp; + wsClient.registerWaitForUpdate(); + sendAttributes(device, TbAttributeSubscriptionScope.SHARED_SCOPE, Arrays.asList(dataPoint1)); + msg = wsClient.waitForUpdate(TimeUnit.SECONDS.toMillis(1)); + Assert.assertNull(msg); + + //Sending duplicate update again + wsClient.registerWaitForUpdate(); + sendAttributes(device, TbAttributeSubscriptionScope.CLIENT_SCOPE, Arrays.asList(dataPoint1)); + msg = wsClient.waitForUpdate(TimeUnit.SECONDS.toMillis(1)); + Assert.assertNull(msg); + + //Sending update from the past, while latest value has new timestamp; + wsClient.registerWaitForUpdate(); + AttributeKvEntry dataPoint2 = new BaseAttributeKvEntry(now, new LongDataEntry("sharedAttributeKey", 42L)); + sendAttributes(device, TbAttributeSubscriptionScope.SHARED_SCOPE, Arrays.asList(dataPoint2)); + msg = wsClient.waitForUpdate(TimeUnit.SECONDS.toMillis(1)); + update = mapper.readValue(msg, EntityDataUpdate.class); + Assert.assertEquals(1, update.getCmdId()); + eData = update.getUpdate(); + Assert.assertNotNull(eData); + Assert.assertEquals(1, eData.size()); + Assert.assertEquals(device.getId(), eData.get(0).getEntityId()); + Assert.assertNotNull(eData.get(0).getLatest().get(EntityKeyType.SHARED_ATTRIBUTE)); + attrValue = eData.get(0).getLatest().get(EntityKeyType.SHARED_ATTRIBUTE).get("sharedAttributeKey"); + Assert.assertEquals(new TsValue(dataPoint2.getLastUpdateTs(), dataPoint2.getValueAsString()), attrValue); + + wsClient.registerWaitForUpdate(); + AttributeKvEntry dataPoint3 = new BaseAttributeKvEntry(now, new LongDataEntry("clientAttributeKey", 42L)); + sendAttributes(device, TbAttributeSubscriptionScope.CLIENT_SCOPE, Arrays.asList(dataPoint3)); + msg = wsClient.waitForUpdate(TimeUnit.SECONDS.toMillis(1)); + update = mapper.readValue(msg, EntityDataUpdate.class); + Assert.assertEquals(1, update.getCmdId()); + eData = update.getUpdate(); + Assert.assertNotNull(eData); + Assert.assertEquals(1, eData.size()); + Assert.assertEquals(device.getId(), eData.get(0).getEntityId()); + Assert.assertNotNull(eData.get(0).getLatest().get(EntityKeyType.CLIENT_ATTRIBUTE)); + attrValue = eData.get(0).getLatest().get(EntityKeyType.CLIENT_ATTRIBUTE).get("clientAttributeKey"); + Assert.assertEquals(new TsValue(dataPoint3.getLastUpdateTs(), dataPoint3.getValueAsString()), attrValue); + + wsClient.registerWaitForUpdate(); + AttributeKvEntry dataPoint4 = new BaseAttributeKvEntry(now, new LongDataEntry("anyAttributeKey", 42L)); + sendAttributes(device, TbAttributeSubscriptionScope.CLIENT_SCOPE, Arrays.asList(dataPoint4)); + msg = wsClient.waitForUpdate(TimeUnit.SECONDS.toMillis(1)); + update = mapper.readValue(msg, EntityDataUpdate.class); + Assert.assertEquals(1, update.getCmdId()); + eData = update.getUpdate(); + Assert.assertNotNull(eData); + Assert.assertEquals(1, eData.size()); + Assert.assertEquals(device.getId(), eData.get(0).getEntityId()); + Assert.assertNotNull(eData.get(0).getLatest().get(EntityKeyType.ATTRIBUTE)); + attrValue = eData.get(0).getLatest().get(EntityKeyType.ATTRIBUTE).get("anyAttributeKey"); + Assert.assertEquals(new TsValue(dataPoint4.getLastUpdateTs(), dataPoint4.getValueAsString()), attrValue); + + wsClient.registerWaitForUpdate(); + AttributeKvEntry dataPoint5 = new BaseAttributeKvEntry(now, new LongDataEntry("anyAttributeKey", 43L)); + sendAttributes(device, TbAttributeSubscriptionScope.SERVER_SCOPE, Arrays.asList(dataPoint5)); + msg = wsClient.waitForUpdate(TimeUnit.SECONDS.toMillis(1)); + update = mapper.readValue(msg, EntityDataUpdate.class); + Assert.assertEquals(1, update.getCmdId()); + eData = update.getUpdate(); + Assert.assertNotNull(eData); + Assert.assertEquals(1, eData.size()); + Assert.assertEquals(device.getId(), eData.get(0).getEntityId()); + Assert.assertNotNull(eData.get(0).getLatest().get(EntityKeyType.ATTRIBUTE)); + attrValue = eData.get(0).getLatest().get(EntityKeyType.ATTRIBUTE).get("anyAttributeKey"); + Assert.assertEquals(new TsValue(dataPoint5.getLastUpdateTs(), dataPoint5.getValueAsString()), attrValue); + } + + private void sendTelemetry(Device device, List tsData) throws InterruptedException { + CountDownLatch latch = new CountDownLatch(1); + tsService.saveAndNotify(device.getTenantId(), device.getId(), tsData, 0, new FutureCallback() { + @Override + public void onSuccess(@Nullable Void result) { + latch.countDown(); + } + + @Override + public void onFailure(Throwable t) { + latch.countDown(); + } + }); + latch.await(3, TimeUnit.SECONDS); } + private void sendAttributes(Device device, TbAttributeSubscriptionScope scope, List attrData) throws InterruptedException { + CountDownLatch latch = new CountDownLatch(1); + tsService.saveAndNotify(device.getTenantId(), device.getId(), scope.name(), attrData, new FutureCallback() { + @Override + public void onSuccess(@Nullable Void result) { + latch.countDown(); + } + + @Override + public void onFailure(Throwable t) { + latch.countDown(); + } + }); + latch.await(3, TimeUnit.SECONDS); + } } diff --git a/application/src/test/java/org/thingsboard/server/controller/TbTestWebSocketClient.java b/application/src/test/java/org/thingsboard/server/controller/TbTestWebSocketClient.java index cf1c31c851..10f72a406b 100644 --- a/application/src/test/java/org/thingsboard/server/controller/TbTestWebSocketClient.java +++ b/application/src/test/java/org/thingsboard/server/controller/TbTestWebSocketClient.java @@ -27,9 +27,7 @@ import java.util.concurrent.TimeUnit; @Slf4j public class TbTestWebSocketClient extends WebSocketClient { - private volatile String lastReply; - private volatile String lastUpdate; - private volatile boolean replyReceived; + private volatile String lastMsg; private CountDownLatch reply; private CountDownLatch update; @@ -44,23 +42,13 @@ public class TbTestWebSocketClient extends WebSocketClient { @Override public void onMessage(String s) { - log.error("RECEIVED: {}", s); - synchronized (this) { - if (!replyReceived) { - replyReceived = true; - lastReply = s; - log.error("LAST REPLY: {}", s); - if (reply != null) { - reply.countDown(); - } - } else { - lastUpdate = s; - log.error("LAST UPDATE: {}", s); - if (update == null) { - update = new CountDownLatch(1); - } - update.countDown(); - } + log.info("RECEIVED: {}", s); + lastMsg = s; + if (reply != null) { + reply.countDown(); + } + if (update != null) { + update.countDown(); } } @@ -74,25 +62,28 @@ public class TbTestWebSocketClient extends WebSocketClient { } + public void registerWaitForUpdate() { + lastMsg = null; + update = new CountDownLatch(1); + } + @Override public void send(String text) throws NotYetConnectedException { - synchronized (this) { - reply = new CountDownLatch(1); - replyReceived = false; - } + reply = new CountDownLatch(1); super.send(text); } public String waitForUpdate() { - synchronized (this) { - update = new CountDownLatch(1); - } + return waitForUpdate(TimeUnit.SECONDS.toMillis(3)); + } + + public String waitForUpdate(long ms) { try { - update.await(3, TimeUnit.SECONDS); + update.await(ms, TimeUnit.MILLISECONDS); } catch (InterruptedException e) { log.warn("Failed to await reply", e); } - return lastUpdate; + return lastMsg; } public String waitForReply() { @@ -101,6 +92,6 @@ public class TbTestWebSocketClient extends WebSocketClient { } catch (InterruptedException e) { log.warn("Failed to await reply", e); } - return lastReply; + return lastMsg; } } diff --git a/application/src/test/resources/logback.xml b/application/src/test/resources/logback.xml index 4f083906e5..f991a40078 100644 --- a/application/src/test/resources/logback.xml +++ b/application/src/test/resources/logback.xml @@ -7,6 +7,8 @@ + + diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityDataQuery.java b/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityDataQuery.java index 01c695ceaf..e0c1d12699 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityDataQuery.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityDataQuery.java @@ -17,9 +17,11 @@ package org.thingsboard.server.common.data.query; import com.fasterxml.jackson.annotation.JsonIgnore; import lombok.Getter; +import lombok.ToString; import java.util.List; +@ToString public class EntityDataQuery extends EntityCountQuery { @Getter diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/query/EntityKeyMapping.java b/dao/src/main/java/org/thingsboard/server/dao/sql/query/EntityKeyMapping.java index b7cb331237..0b2c74a25c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/query/EntityKeyMapping.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/query/EntityKeyMapping.java @@ -136,7 +136,7 @@ public class EntityKeyMapping { } else { scope = DataConstants.SERVER_SCOPE; } - query = String.format("%s AND %s.attribute_type=%s", query, alias, scope); + query = String.format("%s AND %s.attribute_type='%s'", query, alias, scope); } return query; }