From e276c9a936c7f0e9d16b7145ebf3367c6bd06c52 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Mon, 22 Jun 2020 19:36:56 +0300 Subject: [PATCH] Initial WebSocker API --- ...efaultTbEntityDataSubscriptionService.java | 26 ++++- .../DefaultTbLocalSubscriptionService.java | 7 +- .../subscription/TbAttributeSubscription.java | 5 +- .../subscription/TbEntityDataSubCtx.java | 97 ++++++++++++++++++- .../service/subscription/TbSubscription.java | 3 + .../TbTimeseriesSubscription.java | 19 ++-- .../DefaultTelemetryWebSocketService.java | 9 +- .../controller/BaseWebsocketApiTest.java | 59 +++++++++-- .../controller/TbTestWebSocketClient.java | 47 +++++++-- .../dao/sql/query/EntityKeyMapping.java | 2 +- 10 files changed, 238 insertions(+), 36 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 91db84a980..d9e8c043c5 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 @@ -49,6 +49,7 @@ import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.queue.discovery.ClusterTopologyChangeEvent; import org.thingsboard.server.queue.discovery.PartitionChangeEvent; import org.thingsboard.server.queue.discovery.PartitionService; +import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.queue.TbClusterService; import org.thingsboard.server.service.telemetry.TelemetryWebSocketService; @@ -105,19 +106,29 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc @Lazy private SubscriptionManagerService subscriptionManagerService; + @Autowired + @Lazy + private TbLocalSubscriptionService localSubscriptionService; + @Autowired private TimeseriesService tsService; + @Autowired + private TbServiceInfoProvider serviceInfoProvider; + @Value("${database.ts.type}") private String databaseTsType; private ExecutorService wsCallBackExecutor; private boolean tsInSqlDB; + private String serviceId; @PostConstruct public void initExecutor() { + serviceId = serviceInfoProvider.getServiceId(); wsCallBackExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("ws-entity-sub-callback")); tsInSqlDB = databaseTsType.equalsIgnoreCase("sql") || databaseTsType.equalsIgnoreCase("timescale"); + } @PreDestroy @@ -158,6 +169,9 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc TbEntityDataSubCtx ctx = getSubCtx(session.getSessionId(), cmd.getCmdId()); if (ctx != null) { log.debug("[{}][{}] Updating existing subscriptions using: {}", session.getSessionId(), cmd.getCmdId(), cmd); + if (cmd.getLatestCmd() != null || cmd.getTsCmd() != null) { + ctx.clearSubscriptions(); + } //TODO: cleanup old subscription; } else { log.debug("[{}][{}] Creating new subscription using: {}", session.getSessionId(), cmd.getCmdId(), cmd); @@ -209,7 +223,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc private TbEntityDataSubCtx createSubCtx(TelemetryWebSocketSessionRef sessionRef, EntityDataCmd cmd) { Map sessionSubs = subscriptionsBySessionId.computeIfAbsent(sessionRef.getSessionId(), k -> new HashMap<>()); - TbEntityDataSubCtx ctx = new TbEntityDataSubCtx(sessionRef, cmd.getCmdId()); + TbEntityDataSubCtx ctx = new TbEntityDataSubCtx(serviceId, wsService, sessionRef, cmd.getCmdId()); ctx.setQuery(cmd.getQuery()); sessionSubs.put(cmd.getCmdId(), ctx); return ctx; @@ -266,7 +280,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc update = new EntityDataUpdate(ctx.getCmdId(), null, ctx.getData().getData()); } wsService.sendWsMsg(ctx.getSessionId(), update); - //TODO: create context for this (session, cmdId) that contains query, latestCmd and update. Subscribe + periodic updates. + createLatestSubscriptions(ctx, latestCmd); } @Override @@ -281,10 +295,16 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc EntityDataUpdate update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null); wsService.sendWsMsg(ctx.getSessionId(), update); } - //TODO: create context for this (session, cmdId) that contains query, latestCmd and update. Subscribe + periodic updates. + createLatestSubscriptions(ctx, latestCmd); } } + private void createLatestSubscriptions(TbEntityDataSubCtx ctx, LatestValueCmd latestCmd) { + //TODO: create context for this (session, cmdId) that contains query, latestCmd and update. Subscribe + periodic updates. + List tbSubs = ctx.createSubscriptions(latestCmd.getKeys()); + tbSubs.forEach(sub -> localSubscriptionService.addSubscription(sub)); + } + private Map toTsValue(List data) { return data.stream().collect(Collectors.toMap(TsKvEntry::getKey, value -> new TsValue(value.getTs(), value.getValueAsString()))); } 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 d57067fa28..ee30fefe9a 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 @@ -5,7 +5,7 @@ * 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 + * 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, @@ -58,9 +58,6 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer private final Set currentPartitions = ConcurrentHashMap.newKeySet(); private final Map> subscriptionsBySessionId = new ConcurrentHashMap<>(); - @Autowired - private TelemetryWebSocketService wsService; - @Autowired private EntityViewService entityViewService; @@ -155,7 +152,7 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer update.getLatestValues().forEach((key, value) -> attrSub.getKeyStates().put(key, value)); break; } - wsService.sendWsMsg(sessionId, update); + subscription.getUpdateConsumer().accept(sessionId, update); } callback.onSuccess(); } diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbAttributeSubscription.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbAttributeSubscription.java index 83a86efefc..6c33e36ca2 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbAttributeSubscription.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbAttributeSubscription.java @@ -20,8 +20,10 @@ import lombok.Data; import lombok.Getter; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.service.telemetry.sub.SubscriptionUpdate; import java.util.Map; +import java.util.function.BiConsumer; public class TbAttributeSubscription extends TbSubscription { @@ -31,8 +33,9 @@ public class TbAttributeSubscription extends TbSubscription { @Builder public TbAttributeSubscription(String serviceId, String sessionId, int subscriptionId, TenantId tenantId, EntityId entityId, + BiConsumer updateConsumer, boolean allKeys, Map keyStates, TbAttributeSubscriptionScope scope) { - super(serviceId, sessionId, subscriptionId, tenantId, entityId, TbSubscriptionType.ATTRIBUTES); + super(serviceId, sessionId, subscriptionId, tenantId, entityId, TbSubscriptionType.ATTRIBUTES, updateConsumer); this.allKeys = allKeys; this.keyStates = keyStates; this.scope = scope; 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 36ee86a116..5f695d14de 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 @@ -2,17 +2,35 @@ package org.thingsboard.server.service.subscription; import lombok.Data; import org.thingsboard.server.common.data.id.CustomerId; +import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.query.EntityData; import org.thingsboard.server.common.data.query.EntityDataQuery; +import org.thingsboard.server.common.data.query.EntityKey; +import org.thingsboard.server.common.data.query.EntityKeyType; +import org.thingsboard.server.common.data.query.TsValue; +import org.thingsboard.server.service.telemetry.TelemetryWebSocketService; import org.thingsboard.server.service.telemetry.TelemetryWebSocketSessionRef; +import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUpdate; import org.thingsboard.server.service.telemetry.cmd.v2.LatestValueCmd; import org.thingsboard.server.service.telemetry.cmd.v2.TimeSeriesCmd; +import org.thingsboard.server.service.telemetry.sub.SubscriptionUpdate; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; @Data public class TbEntityDataSubCtx { + public static final int MAX_SUBS_PER_CMD = 1024 * 8; + private final String serviceId; + private final TelemetryWebSocketService wsService; private final TelemetryWebSocketSessionRef sessionRef; private final int cmdId; private EntityDataQuery query; @@ -20,8 +38,13 @@ public class TbEntityDataSubCtx { private TimeSeriesCmd tsCmd; private PageData data; private boolean initialDataSent; + private List tbSubs; + private int internalSubIdx; + private Map subToEntityIdMap; - public TbEntityDataSubCtx(TelemetryWebSocketSessionRef sessionRef, int cmdId) { + public TbEntityDataSubCtx(String serviceId, TelemetryWebSocketService wsService, TelemetryWebSocketSessionRef sessionRef, int cmdId) { + this.serviceId = serviceId; + this.wsService = wsService; this.sessionRef = sessionRef; this.cmdId = cmdId; } @@ -38,9 +61,79 @@ public class TbEntityDataSubCtx { return sessionRef.getSecurityCtx().getCustomerId(); } - public void setData(PageData data) { this.data = data; } + 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); + } + } + for (EntityData entityData : data.getData()) { + if (!tsSubKeys.isEmpty()) { + tbSubs.add(createTsSub(entityData, tsSubKeys)); + } + } + 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())); + } + } + if (entityData.getTimeseries() != null) { + entityData.getTimeseries().forEach((k, v) -> keyStates.put(k, Arrays.stream(v).map(TsValue::getTs).max(Long::compareTo).orElse(0L))); + } + + return TbTimeseriesSubscription.builder() + .serviceId(serviceId) + .sessionId(sessionRef.getSessionId()) + .subscriptionId(subIdx) + .tenantId(sessionRef.getSecurityCtx().getTenantId()) + .entityId(entityData.getEntityId()) + .updateConsumer(this::sendTsWsMsg) + .allKeys(false) + .keyStates(keyStates).build(); + } + + + private void sendTsWsMsg(String sessionId, SubscriptionUpdate subscriptionUpdate) { + EntityId entityId = subToEntityIdMap.get(subscriptionUpdate.getSubscriptionId()); + if (entityId != null) { + Map latest = new HashMap<>(); + subscriptionUpdate.getData().forEach((k, v) -> { + Object[] data = (Object[]) v.get(0); + latest.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))); + } + + } + + public void clearSubscriptions() { + subToEntityIdMap.clear(); + } } diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbSubscription.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbSubscription.java index 22b37ff690..68cbfbccb3 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbSubscription.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbSubscription.java @@ -19,8 +19,10 @@ import lombok.AllArgsConstructor; import lombok.Data; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.service.telemetry.sub.SubscriptionUpdate; import java.util.Objects; +import java.util.function.BiConsumer; @Data @AllArgsConstructor @@ -32,6 +34,7 @@ public abstract class TbSubscription { private final TenantId tenantId; private final EntityId entityId; private final TbSubscriptionType type; + private final BiConsumer updateConsumer; @Override public boolean equals(Object o) { diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbTimeseriesSubscription.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbTimeseriesSubscription.java index 0be63f7b65..ee4bd2c8a0 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbTimeseriesSubscription.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbTimeseriesSubscription.java @@ -5,7 +5,7 @@ * 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 + * 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, @@ -19,20 +19,27 @@ import lombok.Builder; import lombok.Getter; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.service.telemetry.sub.SubscriptionUpdate; import java.util.Map; +import java.util.function.BiConsumer; public class TbTimeseriesSubscription extends TbSubscription { - @Getter private final boolean allKeys; - @Getter private final Map keyStates; - @Getter private final long startTime; - @Getter private final long endTime; + @Getter + private final boolean allKeys; + @Getter + private final Map keyStates; + @Getter + private final long startTime; + @Getter + private final long endTime; @Builder public TbTimeseriesSubscription(String serviceId, String sessionId, int subscriptionId, TenantId tenantId, EntityId entityId, + BiConsumer updateConsumer, boolean allKeys, Map keyStates, long startTime, long endTime) { - super(serviceId, sessionId, subscriptionId, tenantId, entityId, TbSubscriptionType.TIMESERIES); + super(serviceId, sessionId, subscriptionId, tenantId, entityId, TbSubscriptionType.TIMESERIES, updateConsumer); this.allKeys = allKeys; this.keyStates = keyStates; this.startTime = startTime; diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java index 1a065c07fb..ef08531761 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java @@ -86,6 +86,7 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; +import java.util.function.BiConsumer; import java.util.function.Consumer; import java.util.stream.Collectors; @@ -129,6 +130,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi @Autowired private TbServiceInfoProvider serviceInfoProvider; + @Value("${server.ws.limits.max_subscriptions_per_tenant:0}") private int maxSubscriptionsPerTenant; @Value("${server.ws.limits.max_subscriptions_per_customer:0}") @@ -398,7 +400,9 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi .entityId(entityId) .allKeys(false) .keyStates(subState) - .scope(scope).build(); + .scope(scope) + .updateConsumer(DefaultTelemetryWebSocketService.this::sendWsMsg) + .build(); oldSubService.addSubscription(sub); } @@ -495,6 +499,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi .entityId(entityId) .allKeys(true) .keyStates(subState) + .updateConsumer(DefaultTelemetryWebSocketService.this::sendWsMsg) .scope(scope).build(); oldSubService.addSubscription(sub); } @@ -575,6 +580,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi .subscriptionId(cmd.getCmdId()) .tenantId(sessionRef.getSecurityCtx().getTenantId()) .entityId(entityId) + .updateConsumer(DefaultTelemetryWebSocketService.this::sendWsMsg) .allKeys(true) .keyStates(subState).build(); oldSubService.addSubscription(sub); @@ -612,6 +618,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi .subscriptionId(cmd.getCmdId()) .tenantId(sessionRef.getSecurityCtx().getTenantId()) .entityId(entityId) + .updateConsumer(DefaultTelemetryWebSocketService.this::sendWsMsg) .allKeys(false) .keyStates(subState).build(); oldSubService.addSubscription(sub); 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 41d9ed278a..939ea18ab1 100644 --- a/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java @@ -5,7 +5,7 @@ * 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 + * 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, @@ -15,6 +15,9 @@ */ package org.thingsboard.server.controller; +import com.google.common.util.concurrent.FutureCallback; +import lombok.extern.slf4j.Slf4j; +import org.checkerframework.checker.nullness.qual.Nullable; import org.junit.After; import org.junit.Assert; import org.junit.Before; @@ -38,6 +41,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.telemetry.TelemetrySubscriptionService; import org.thingsboard.server.service.telemetry.cmd.TelemetryPluginCmdsWrapper; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataCmd; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUpdate; @@ -46,10 +50,13 @@ import org.thingsboard.server.service.telemetry.cmd.v2.LatestValueCmd; import java.util.Arrays; import java.util.Collections; +import java.util.List; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; +@Slf4j public class BaseWebsocketApiTest extends AbstractWebsocketTest { private Tenant savedTenant; @@ -57,7 +64,7 @@ public class BaseWebsocketApiTest extends AbstractWebsocketTest { private TbTestWebSocketClient wsClient; @Autowired - private TimeseriesService tsService; + private TelemetrySubscriptionService tsService; @Before public void beforeTest() throws Exception { @@ -129,7 +136,10 @@ public class BaseWebsocketApiTest extends AbstractWebsocketTest { TsKvEntry dataPoint1 = new BasicTsKvEntry(now - TimeUnit.MINUTES.toMillis(1), new LongDataEntry("temperature", 42L)); TsKvEntry dataPoint2 = new BasicTsKvEntry(now - TimeUnit.MINUTES.toMillis(2), new LongDataEntry("temperature", 42L)); TsKvEntry dataPoint3 = new BasicTsKvEntry(now - TimeUnit.MINUTES.toMillis(3), new LongDataEntry("temperature", 42L)); - tsService.save(device.getTenantId(), device.getId(), Arrays.asList(dataPoint1, dataPoint2, dataPoint3), 0).get(); + List tsData = Arrays.asList(dataPoint1, dataPoint2, dataPoint3); + + sendTelemetry(device, tsData); + Thread.sleep(1000); wsClient.send(mapper.writeValueAsString(wrapper)); msg = wsClient.waitForReply(); @@ -146,6 +156,22 @@ 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 { @@ -177,12 +203,15 @@ public class BaseWebsocketApiTest extends AbstractWebsocketTest { Assert.assertNotNull(pageData); Assert.assertEquals(1, pageData.getData().size()); Assert.assertEquals(device.getId(), pageData.getData().get(0).getEntityId()); - Assert.assertNull(pageData.getData().get(0).getLatest().get(EntityKeyType.TIME_SERIES).get("temperature")); + Assert.assertNotNull(pageData.getData().get(0).getLatest().get(EntityKeyType.TIME_SERIES).get("temperature")); + Assert.assertEquals(0, pageData.getData().get(0).getLatest().get(EntityKeyType.TIME_SERIES).get("temperature").getTs()); + Assert.assertEquals("", pageData.getData().get(0).getLatest().get(EntityKeyType.TIME_SERIES).get("temperature").getValue()); TsKvEntry dataPoint1 = new BasicTsKvEntry(now - TimeUnit.MINUTES.toMillis(1), new LongDataEntry("temperature", 42L)); - tsService.save(device.getTenantId(), device.getId(), Arrays.asList(dataPoint1), 0).get(); + List tsData = Arrays.asList(dataPoint1); + sendTelemetry(device, tsData); - cmd = new EntityDataCmd(2, edq, null, latestCmd, null); + cmd = new EntityDataCmd(1, edq, null, latestCmd, null); wrapper = new TelemetryPluginCmdsWrapper(); wrapper.setEntityDataCmds(Collections.singletonList(cmd)); @@ -190,7 +219,7 @@ public class BaseWebsocketApiTest extends AbstractWebsocketTest { wsClient.send(mapper.writeValueAsString(wrapper)); msg = wsClient.waitForReply(); update = mapper.readValue(msg, EntityDataUpdate.class); - Assert.assertEquals(2, update.getCmdId()); + Assert.assertEquals(1, update.getCmdId()); pageData = update.getData(); Assert.assertNotNull(pageData); Assert.assertEquals(1, pageData.getData().size()); @@ -198,6 +227,22 @@ public class BaseWebsocketApiTest extends AbstractWebsocketTest { Assert.assertNotNull(pageData.getData().get(0).getLatest().get(EntityKeyType.TIME_SERIES)); 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)); + sendTelemetry(device, Arrays.asList(dataPoint2)); + + 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.TIME_SERIES)); + tsValue = eData.get(0).getLatest().get(EntityKeyType.TIME_SERIES).get("temperature"); + Assert.assertEquals(new TsValue(dataPoint2.getTs(), dataPoint2.getValueAsString()), tsValue); } } 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 29946a4f39..62dfdbb107 100644 --- a/application/src/test/java/org/thingsboard/server/controller/TbTestWebSocketClient.java +++ b/application/src/test/java/org/thingsboard/server/controller/TbTestWebSocketClient.java @@ -5,7 +5,7 @@ * 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 + * 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, @@ -27,9 +27,11 @@ import java.util.concurrent.TimeUnit; @Slf4j public class TbTestWebSocketClient extends WebSocketClient { - private volatile String lastMsg; + private volatile String lastReply; + private volatile String lastUpdate; private volatile boolean replyReceived; private CountDownLatch reply; + private CountDownLatch update; public TbTestWebSocketClient(URI serverUri) { super(serverUri); @@ -42,11 +44,22 @@ public class TbTestWebSocketClient extends WebSocketClient { @Override public void onMessage(String s) { - if (!replyReceived) { - replyReceived = true; - lastMsg = s; - if (reply != null) { - reply.countDown(); + 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(); } } } @@ -63,17 +76,31 @@ public class TbTestWebSocketClient extends WebSocketClient { @Override public void send(String text) throws NotYetConnectedException { - reply = new CountDownLatch(1); - replyReceived = false; + synchronized (this) { + reply = new CountDownLatch(1); + replyReceived = false; + } super.send(text); } + public String waitForUpdate() { + synchronized (this) { + update = new CountDownLatch(1); + } + try { + update.await(3, TimeUnit.SECONDS); + } catch (InterruptedException e) { + log.warn("Failed to await reply", e); + } + return lastUpdate; + } + public String waitForReply() { try { reply.await(3, TimeUnit.SECONDS); } catch (InterruptedException e) { log.warn("Failed to await reply", e); } - return lastMsg; + return lastReply; } } 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 08eb668ce0..6fa9daf01f 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 @@ -121,7 +121,7 @@ public class EntityKeyMapping { String join = hasFilter() ? "left join" : "left outer join"; ctx.addStringParameter(alias + "_key_id", entityKey.getKey()); if (entityKey.getType().equals(EntityKeyType.TIME_SERIES)) { - return String.format("%s ts_kv_latest %s ON %s.entity_id=to_uuid(entities.id) AND %s.key = (select key_id from ts_kv_dictionary where key = :%s_key_id)", + return String.format("%s ts_kv_latest %s ON %s.entity_id=entities.id AND %s.key = (select key_id from ts_kv_dictionary where key = :%s_key_id)", join, alias, alias, alias, alias); } else { String query = String.format("%s attribute_kv %s ON %s.entity_id=entities.id AND %s.entity_type=%s AND %s.attribute_key=:%s_key_id",