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 404b4f59c4..0a85e5e993 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 @@ -52,6 +52,8 @@ 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.security.permission.Operation; +import org.thingsboard.server.service.telemetry.DefaultTelemetryWebSocketService; import org.thingsboard.server.service.telemetry.TelemetryWebSocketService; import org.thingsboard.server.service.telemetry.TelemetryWebSocketSessionRef; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataCmd; @@ -211,21 +213,25 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } else { historyFuture = Futures.immediateFuture(ctx); } - if (cmd.getLatestCmd() != null) { - Futures.addCallback(historyFuture, new FutureCallback() { - @Override - public void onSuccess(@Nullable TbEntityDataSubCtx theCtx) { + Futures.addCallback(historyFuture, new FutureCallback() { + @Override + public void onSuccess(@Nullable TbEntityDataSubCtx theCtx) { + if (cmd.getLatestCmd() != null) { handleLatestCmd(theCtx, cmd.getLatestCmd()); + } else if (cmd.getTsCmd() != null) { + handleTimeSeriesCmd(theCtx, cmd.getTsCmd()); + } else if (!theCtx.isInitialDataSent()) { + EntityDataUpdate update = new EntityDataUpdate(theCtx.getCmdId(), theCtx.getData(), null); + wsService.sendWsMsg(theCtx.getSessionId(), update); + theCtx.setInitialDataSent(true); } + } - @Override - public void onFailure(Throwable t) { - log.warn("[{}][{}] Failed to process command", session.getSessionId(), cmd.getCmdId()); - } - }, wsCallBackExecutor); - } else if (cmd.getTsCmd() != null) { - handleTimeseriesCmd(ctx, cmd.getTsCmd()); - } + @Override + public void onFailure(Throwable t) { + log.warn("[{}][{}] Failed to process command", session.getSessionId(), cmd.getCmdId()); + } + }, wsCallBackExecutor); } private TbEntityDataSubCtx createSubCtx(TelemetryWebSocketSessionRef sessionRef, EntityDataCmd cmd) { @@ -245,7 +251,47 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } } - private void handleTimeseriesCmd(TbEntityDataSubCtx ctx, TimeSeriesCmd tsCmd) { + private void handleTimeSeriesCmd(TbEntityDataSubCtx ctx, TimeSeriesCmd cmd) { + List keys = cmd.getKeys(); + log.debug("[{}][{}] Fetching time-series data for last {} ms for keys: ({})", ctx.getSessionId(), ctx.getCmdId(), cmd.getTimeWindow(), cmd.getKeys()); + long startTs = cmd.getStartTs(); + long endTs = cmd.getStartTs() + cmd.getTimeWindow(); + + Map>>> tsFutures = new HashMap<>(); + for (EntityData entityData : ctx.getData().getData()) { + List queries = keys.stream().map(key -> new BaseReadTsKvQuery(key, startTs, endTs, cmd.getInterval(), + getLimit(cmd.getLimit()), DefaultTelemetryWebSocketService.getAggregation(cmd.getAgg()))).collect(Collectors.toList()); + ListenableFuture> tsDataFutures = tsService.findAll(ctx.getTenantId(), entityData.getEntityId(), queries); + tsFutures.put(entityData, Futures.transform(tsDataFutures, this::toTsValues, MoreExecutors.directExecutor())); + } + Futures.addCallback(Futures.allAsList(tsFutures.values()), new FutureCallback>>>() { + @Override + public void onSuccess(@Nullable List>> result) { + tsFutures.forEach((key, value) -> { + try { + value.get().forEach((k, v) -> key.getTimeseries().put(k, v.toArray(new TsValue[v.size()]))); + } catch (InterruptedException | ExecutionException e) { + log.warn("[{}][{}] Failed to lookup time-series data: {}:{}", ctx.getSessionId(), ctx.getCmdId(), key.getEntityId(), keys, e); + } + }); + EntityDataUpdate update; + if (!ctx.isInitialDataSent()) { + update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null); + ctx.setInitialDataSent(true); + } else { + update = new EntityDataUpdate(ctx.getCmdId(), null, ctx.getData().getData()); + } + wsService.sendWsMsg(ctx.getSessionId(), update); + createSubscriptions(ctx, keys.stream().map(key -> new EntityKey(EntityKeyType.TIME_SERIES, key)).collect(Collectors.toList())); + } + + @Override + public void onFailure(Throwable t) { + log.warn("[{}][{}] Failed to process websocket command: {}:{}", ctx.getSessionId(), ctx.getCmdId(), ctx.getQuery(), cmd, t); + wsService.sendWsMsg(ctx.getSessionId(), + new EntityDataUpdate(ctx.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR.getCode(), "Failed to process websocket command!")); + } + }, wsCallBackExecutor); } private void handleLatestCmd(TbEntityDataSubCtx ctx, LatestValueCmd latestCmd) { @@ -257,7 +303,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc .filter(key -> key.getType().equals(EntityKeyType.TIME_SERIES)) .map(EntityKey::getKey).collect(Collectors.toList()); - Map>> missingTelemetryFurutes = new HashMap<>(); + Map>> missingTelemetryFutures = new HashMap<>(); for (EntityData entityData : ctx.getData().getData()) { Map> latestEntityData = entityData.getLatest(); Map tsEntityData = latestEntityData.get(EntityKeyType.TIME_SERIES); @@ -270,12 +316,12 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } ListenableFuture> missingTsData = tsService.findLatest(ctx.getTenantId(), entityData.getEntityId(), missingTsKeys); - missingTelemetryFurutes.put(entityData, Futures.transform(missingTsData, this::toTsValue, MoreExecutors.directExecutor())); + missingTelemetryFutures.put(entityData, Futures.transform(missingTsData, this::toTsValue, MoreExecutors.directExecutor())); } - Futures.addCallback(Futures.allAsList(missingTelemetryFurutes.values()), new FutureCallback>>() { + Futures.addCallback(Futures.allAsList(missingTelemetryFutures.values()), new FutureCallback>>() { @Override public void onSuccess(@Nullable List> result) { - missingTelemetryFurutes.forEach((key, value) -> { + missingTelemetryFutures.forEach((key, value) -> { try { key.getLatest().get(EntityKeyType.TIME_SERIES).putAll(value.get()); } catch (InterruptedException | ExecutionException e) { @@ -285,11 +331,12 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc EntityDataUpdate update; if (!ctx.isInitialDataSent()) { update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null); + ctx.setInitialDataSent(true); } else { update = new EntityDataUpdate(ctx.getCmdId(), null, ctx.getData().getData()); } wsService.sendWsMsg(ctx.getSessionId(), update); - createLatestSubscriptions(ctx, latestCmd); + createSubscriptions(ctx, latestCmd.getKeys()); } @Override @@ -303,14 +350,15 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc if (!ctx.isInitialDataSent()) { EntityDataUpdate update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null); wsService.sendWsMsg(ctx.getSessionId(), update); + ctx.setInitialDataSent(true); } - createLatestSubscriptions(ctx, latestCmd); + createSubscriptions(ctx, latestCmd.getKeys()); } } - private void createLatestSubscriptions(TbEntityDataSubCtx ctx, LatestValueCmd latestCmd) { + private void createSubscriptions(TbEntityDataSubCtx ctx, List keys) { //TODO: create context for this (session, cmdId) that contains query, latestCmd and update. Subscribe + periodic updates. - List tbSubs = ctx.createSubscriptions(latestCmd.getKeys()); + List tbSubs = ctx.createSubscriptions(keys); tbSubs.forEach(sub -> localSubscriptionService.addSubscription(sub)); } @@ -318,6 +366,14 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc return data.stream().collect(Collectors.toMap(TsKvEntry::getKey, value -> new TsValue(value.getTs(), value.getValueAsString()))); } + private Map> toTsValues(List data) { + Map> results = new HashMap<>(); + for (TsKvEntry tsKvEntry : data) { + results.computeIfAbsent(tsKvEntry.getKey(), k -> new ArrayList<>()).add(new TsValue(tsKvEntry.getTs(), tsKvEntry.getValueAsString())); + } + return results; + } + private ListenableFuture handleHistoryCmd(TbEntityDataSubCtx ctx, EntityHistoryCmd historyCmd) { List tsKvQueryList = historyCmd.getKeys().stream().map(key -> new BaseReadTsKvQuery( key, historyCmd.getStartTs(), historyCmd.getEndTs(), historyCmd.getInterval(), getLimit(historyCmd.getLimit()), historyCmd.getAgg() 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 ef08531761..b9f4ab022b 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 @@ -657,11 +657,6 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi "Query is empty!"); sendWsMsg(sessionRef, update); return false; - } else if (cmd.getHistoryCmd() == null && cmd.getLatestCmd() == null && cmd.getTsCmd() == null) { - SubscriptionUpdate update = new SubscriptionUpdate(cmd.getCmdId(), SubscriptionErrorCode.BAD_REQUEST, - "No history, latest or timeseries command present!"); - sendWsMsg(sessionRef, update); - return false; } return true; } @@ -823,7 +818,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi } - private static Aggregation getAggregation(String agg) { + public static Aggregation getAggregation(String agg) { return StringUtils.isEmpty(agg) ? DEFAULT_AGGREGATION : Aggregation.valueOf(agg); } 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 ec8ca198bb..e49dfc3264 100644 --- a/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java @@ -21,7 +21,6 @@ import org.checkerframework.checker.nullness.qual.Nullable; import org.junit.After; import org.junit.Assert; import org.junit.Before; -import org.junit.Ignore; import org.junit.Test; import org.springframework.beans.factory.annotation.Autowired; import org.thingsboard.server.common.data.Device; @@ -42,7 +41,6 @@ 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.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; @@ -210,12 +208,12 @@ public class BaseWebsocketApiTest extends AbstractWebsocketTest { 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.TIME_SERIES)); - TsValue tsValue = pageData.getData().get(0).getLatest().get(EntityKeyType.TIME_SERIES).get("temperature"); + 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(dataPoint1.getTs(), dataPoint1.getValueAsString()), tsValue); now = System.currentTimeMillis(); @@ -299,12 +297,12 @@ public class BaseWebsocketApiTest extends AbstractWebsocketTest { 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"); + 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.SERVER_ATTRIBUTE)); + TsValue tsValue = listData.get(0).getLatest().get(EntityKeyType.SERVER_ATTRIBUTE).get("serverAttributeKey"); Assert.assertEquals(new TsValue(dataPoint1.getLastUpdateTs(), dataPoint1.getValueAsString()), tsValue); now = System.currentTimeMillis(); @@ -386,7 +384,6 @@ public class BaseWebsocketApiTest extends AbstractWebsocketTest { 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); diff --git a/application/src/test/java/org/thingsboard/server/controller/ControllerSqlTestSuite.java b/application/src/test/java/org/thingsboard/server/controller/ControllerSqlTestSuite.java index 15da972cf5..0a5dae47b7 100644 --- a/application/src/test/java/org/thingsboard/server/controller/ControllerSqlTestSuite.java +++ b/application/src/test/java/org/thingsboard/server/controller/ControllerSqlTestSuite.java @@ -26,9 +26,9 @@ import java.util.Arrays; @RunWith(ClasspathSuite.class) @ClasspathSuite.ClassnameFilters({ -// "org.thingsboard.server.controller.sql.WebsocketApiSqlTest", + "org.thingsboard.server.controller.sql.WebsocketApiSqlTest", // "org.thingsboard.server.controller.sql.EntityQueryControllerSqlTest", - "org.thingsboard.server.controller.sql.*Test", +// "org.thingsboard.server.controller.sql.*Test", }) public class ControllerSqlTestSuite {