Browse Source

Merge branch 'develop/3.4' of github.com:thingsboard/thingsboard into feature/mqtt-tests-spead-up

pull/6537/head
ShvaykaD 4 years ago
parent
commit
f4aec3d731
  1. 73
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java
  2. 16
      application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractSubCtx.java
  3. 6
      application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java
  4. 7
      application/src/main/java/org/thingsboard/server/service/subscription/TbEntityCountSubCtx.java
  5. 8
      application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java
  6. 6
      application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java
  7. 53
      application/src/test/java/org/thingsboard/server/controller/BaseEntityViewControllerTest.java
  8. 14
      application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java
  9. 4
      application/src/test/java/org/thingsboard/server/controller/TbTestWebSocketClient.java
  10. 10
      ui-ngx/src/app/core/api/alias-controller.ts
  11. 4
      ui-ngx/src/app/core/api/widget-api.models.ts
  12. 8
      ui-ngx/src/app/core/api/widget-subscription.ts
  13. 6
      ui-ngx/src/app/core/http/entity.service.ts
  14. 8
      ui-ngx/src/app/modules/home/components/widget/lib/maps/map-widget2.ts
  15. 1
      ui-ngx/src/app/modules/home/components/widget/lib/maps/providers/image-map.ts
  16. 6
      ui-ngx/src/app/modules/home/components/widget/lib/markdown-widget.component.ts
  17. 8
      ui-ngx/src/app/modules/home/components/widget/widget-config.component.html
  18. 2
      ui-ngx/src/app/modules/home/components/widget/widget-config.component.ts
  19. 3
      ui-ngx/src/app/modules/home/components/widget/widget.component.ts
  20. 1
      ui-ngx/src/app/shared/models/widget.models.ts
  21. 1
      ui-ngx/src/assets/locale/locale.constant-en_US.json

73
application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java

@ -228,7 +228,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
}
} else if (!theCtx.isInitialDataSent()) {
EntityDataUpdate update = new EntityDataUpdate(theCtx.getCmdId(), theCtx.getData(), null, theCtx.getMaxEntitiesPerDataSubscription());
wsService.sendWsMsg(theCtx.getSessionId(), update);
theCtx.sendWsMsg(update);
theCtx.setInitialDataSent(true);
}
} catch (RuntimeException e) {
@ -287,7 +287,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
ctx.clearEntitySubscriptions();
if (entities.isEmpty()) {
AlarmDataUpdate update = new AlarmDataUpdate(cmd.getCmdId(), new PageData<>(), null, 0, 0);
wsService.sendWsMsg(ctx.getSessionId(), update);
ctx.sendWsMsg(update);
} else {
ctx.fetchAlarms();
ctx.createLatestValuesSubscriptions(cmd.getQuery().getLatestValues());
@ -420,22 +420,26 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
}
} catch (InterruptedException | ExecutionException e) {
log.warn("[{}][{}][{}] Failed to fetch historical data", ctx.getSessionId(), ctx.getCmdId(), entityData.getEntityId(), e);
wsService.sendWsMsg(ctx.getSessionId(),
new EntityDataUpdate(ctx.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR.getCode(), "Failed to fetch historical data!"));
ctx.sendWsMsg(new EntityDataUpdate(ctx.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR.getCode(), "Failed to fetch historical data!"));
}
});
EntityDataUpdate update;
if (!ctx.isInitialDataSent()) {
update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null, ctx.getMaxEntitiesPerDataSubscription());
ctx.setInitialDataSent(true);
} else {
update = new EntityDataUpdate(ctx.getCmdId(), null, ctx.getData().getData(), ctx.getMaxEntitiesPerDataSubscription());
}
wsService.sendWsMsg(ctx.getSessionId(), update);
if (subscribe) {
ctx.createTimeseriesSubscriptions(keys.stream().map(key -> new EntityKey(EntityKeyType.TIME_SERIES, key)).collect(Collectors.toList()), cmd.getStartTs(), cmd.getEndTs());
ctx.getWsLock().lock();
try {
EntityDataUpdate update;
if (!ctx.isInitialDataSent()) {
update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null, ctx.getMaxEntitiesPerDataSubscription());
ctx.setInitialDataSent(true);
} else {
update = new EntityDataUpdate(ctx.getCmdId(), null, ctx.getData().getData(), ctx.getMaxEntitiesPerDataSubscription());
}
if (subscribe) {
ctx.createTimeseriesSubscriptions(keys.stream().map(key -> new EntityKey(EntityKeyType.TIME_SERIES, key)).collect(Collectors.toList()), cmd.getStartTs(), cmd.getEndTs());
}
ctx.sendWsMsg(update);
ctx.getData().getData().forEach(ed -> ed.getTimeseries().clear());
} finally {
ctx.getWsLock().unlock();
}
ctx.getData().getData().forEach(ed -> ed.getTimeseries().clear());
return ctx;
}, wsCallBackExecutor);
}
@ -464,7 +468,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
ListenableFuture<List<TsKvEntry>> missingTsData = tsService.findLatest(ctx.getTenantId(), entityData.getEntityId(), missingTsKeys);
missingTelemetryFutures.put(entityData, Futures.transform(missingTsData, this::toTsValue, MoreExecutors.directExecutor()));
}
Futures.addCallback(Futures.allAsList(missingTelemetryFutures.values()), new FutureCallback<List<Map<String, TsValue>>>() {
Futures.addCallback(Futures.allAsList(missingTelemetryFutures.values()), new FutureCallback<>() {
@Override
public void onSuccess(@Nullable List<Map<String, TsValue>> result) {
missingTelemetryFutures.forEach((key, value) -> {
@ -475,30 +479,39 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
}
});
EntityDataUpdate update;
if (!ctx.isInitialDataSent()) {
update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null, ctx.getMaxEntitiesPerDataSubscription());
ctx.setInitialDataSent(true);
} else {
update = new EntityDataUpdate(ctx.getCmdId(), null, ctx.getData().getData(), ctx.getMaxEntitiesPerDataSubscription());
ctx.getWsLock().lock();
try {
ctx.createLatestValuesSubscriptions(latestCmd.getKeys());
if (!ctx.isInitialDataSent()) {
update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null, ctx.getMaxEntitiesPerDataSubscription());
ctx.setInitialDataSent(true);
} else {
update = new EntityDataUpdate(ctx.getCmdId(), null, ctx.getData().getData(), ctx.getMaxEntitiesPerDataSubscription());
}
ctx.sendWsMsg(update);
} finally {
ctx.getWsLock().unlock();
}
wsService.sendWsMsg(ctx.getSessionId(), update);
ctx.createLatestValuesSubscriptions(latestCmd.getKeys());
}
@Override
public void onFailure(Throwable t) {
log.warn("[{}][{}] Failed to process websocket command: {}:{}", ctx.getSessionId(), ctx.getCmdId(), ctx.getQuery(), latestCmd, t);
wsService.sendWsMsg(ctx.getSessionId(),
new EntityDataUpdate(ctx.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR.getCode(), "Failed to process websocket command!"));
ctx.sendWsMsg(new EntityDataUpdate(ctx.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR.getCode(), "Failed to process websocket command!"));
}
}, wsCallBackExecutor);
} else {
if (!ctx.isInitialDataSent()) {
EntityDataUpdate update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null, ctx.getMaxEntitiesPerDataSubscription());
wsService.sendWsMsg(ctx.getSessionId(), update);
ctx.setInitialDataSent(true);
ctx.getWsLock().lock();
try {
ctx.createLatestValuesSubscriptions(latestCmd.getKeys());
if (!ctx.isInitialDataSent()) {
EntityDataUpdate update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null, ctx.getMaxEntitiesPerDataSubscription());
ctx.sendWsMsg(update);
ctx.setInitialDataSent(true);
}
} finally {
ctx.getWsLock().unlock();
}
ctx.createLatestValuesSubscriptions(latestCmd.getKeys());
}
}

16
application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractSubCtx.java

@ -41,6 +41,7 @@ import org.thingsboard.server.dao.attributes.AttributesService;
import org.thingsboard.server.dao.entity.EntityService;
import org.thingsboard.server.service.telemetry.TelemetryWebSocketService;
import org.thingsboard.server.service.telemetry.TelemetryWebSocketSessionRef;
import org.thingsboard.server.service.telemetry.cmd.v2.CmdUpdate;
import org.thingsboard.server.service.telemetry.sub.TelemetrySubscriptionUpdate;
import java.util.ArrayList;
@ -52,14 +53,18 @@ import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
@Slf4j
@Data
public abstract class TbAbstractSubCtx<T extends EntityCountQuery> {
@Getter
protected final Lock wsLock = new ReentrantLock(true);
protected final String serviceId;
protected final SubscriptionServiceStatistics stats;
protected final TelemetryWebSocketService wsService;
private final TelemetryWebSocketService wsService;
protected final EntityService entityService;
protected final TbLocalSubscriptionService localSubscriptionService;
protected final AttributesService attributesService;
@ -314,4 +319,13 @@ public abstract class TbAbstractSubCtx<T extends EntityCountQuery> {
private final String sourceAttribute;
}
public void sendWsMsg(CmdUpdate update) {
wsLock.lock();
try {
wsService.sendWsMsg(sessionRef.getSessionId(), update);
} finally {
wsLock.unlock();
}
}
}

6
application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java

@ -117,7 +117,7 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx<AlarmDataQuery> {
} else {
update = new AlarmDataUpdate(cmdId, new PageData<>(), null, maxEntitiesPerAlarmSubscription, data.getTotalElements());
}
wsService.sendWsMsg(getSessionId(), update);
sendWsMsg(update);
}
public void fetchData() {
@ -198,7 +198,7 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx<AlarmDataQuery> {
return alarm;
}).collect(Collectors.toList());
if (!update.isEmpty()) {
wsService.sendWsMsg(sessionId, new AlarmDataUpdate(cmdId, null, update, maxEntitiesPerAlarmSubscription, data.getTotalElements()));
sendWsMsg(new AlarmDataUpdate(cmdId, null, update, maxEntitiesPerAlarmSubscription, data.getTotalElements()));
}
} else {
log.trace("[{}][{}][{}][{}] Received stale subscription update: {}", sessionId, cmdId, subscriptionUpdate.getSubscriptionId(), keyType, subscriptionUpdate);
@ -222,7 +222,7 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx<AlarmDataQuery> {
AlarmData updated = new AlarmData(alarm, current.getOriginatorName(), current.getEntityId());
updated.getLatest().putAll(current.getLatest());
alarmsMap.put(alarmId, updated);
wsService.sendWsMsg(sessionId, new AlarmDataUpdate(cmdId, null, Collections.singletonList(updated), maxEntitiesPerAlarmSubscription, data.getTotalElements()));
sendWsMsg(new AlarmDataUpdate(cmdId, null, Collections.singletonList(updated), maxEntitiesPerAlarmSubscription, data.getTotalElements()));
} else {
fetchAlarms();
}

7
application/src/main/java/org/thingsboard/server/service/subscription/TbEntityCountSubCtx.java

@ -17,14 +17,11 @@ package org.thingsboard.server.service.subscription;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.server.common.data.query.EntityCountQuery;
import org.thingsboard.server.common.data.query.EntityKeyType;
import org.thingsboard.server.dao.attributes.AttributesService;
import org.thingsboard.server.dao.entity.EntityService;
import org.thingsboard.server.service.telemetry.TelemetryWebSocketService;
import org.thingsboard.server.service.telemetry.TelemetryWebSocketSessionRef;
import org.thingsboard.server.service.telemetry.cmd.v2.EntityCountUpdate;
import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUpdate;
import org.thingsboard.server.service.telemetry.sub.TelemetrySubscriptionUpdate;
@Slf4j
public class TbEntityCountSubCtx extends TbAbstractSubCtx<EntityCountQuery> {
@ -40,7 +37,7 @@ public class TbEntityCountSubCtx extends TbAbstractSubCtx<EntityCountQuery> {
@Override
public void fetchData() {
result = (int) entityService.countEntitiesByQuery(getTenantId(), getCustomerId(), query);
wsService.sendWsMsg(sessionRef.getSessionId(), new EntityCountUpdate(cmdId, result));
sendWsMsg(new EntityCountUpdate(cmdId, result));
}
@Override
@ -48,7 +45,7 @@ public class TbEntityCountSubCtx extends TbAbstractSubCtx<EntityCountQuery> {
int newCount = (int) entityService.countEntitiesByQuery(getTenantId(), getCustomerId(), query);
if (newCount != result) {
result = newCount;
wsService.sendWsMsg(sessionRef.getSessionId(), new EntityCountUpdate(cmdId, result));
sendWsMsg(new EntityCountUpdate(cmdId, result));
}
}

8
application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java

@ -51,7 +51,7 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx<EntityDataQuery> {
@Getter
@Setter
private boolean initialDataSent;
private volatile boolean initialDataSent;
private TimeSeriesCmd curTsCmd;
private LatestValueCmd latestValueCmd;
@Getter
@ -121,7 +121,7 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx<EntityDataQuery> {
if (!latestUpdate.isEmpty()) {
Map<EntityKeyType, Map<String, TsValue>> latestMap = Collections.singletonMap(keyType, latestUpdate);
entityData = new EntityData(entityId, latestMap, null);
wsService.sendWsMsg(sessionId, new EntityDataUpdate(cmdId, null, Collections.singletonList(entityData), maxEntitiesPerDataSubscription));
sendWsMsg(new EntityDataUpdate(cmdId, null, Collections.singletonList(entityData), maxEntitiesPerDataSubscription));
}
}
@ -162,7 +162,7 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx<EntityDataQuery> {
Map<String, TsValue[]> tsMap = new HashMap<>();
tsUpdate.forEach((key, tsValue) -> tsMap.put(key, tsValue.toArray(new TsValue[tsValue.size()])));
EntityData entityData = new EntityData(entityId, null, tsMap);
wsService.sendWsMsg(sessionId, new EntityDataUpdate(cmdId, null, Collections.singletonList(entityData), maxEntitiesPerDataSubscription));
sendWsMsg(new EntityDataUpdate(cmdId, null, Collections.singletonList(entityData), maxEntitiesPerDataSubscription));
}
}
@ -219,9 +219,9 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx<EntityDataQuery> {
}
}
}
wsService.sendWsMsg(sessionRef.getSessionId(), new EntityDataUpdate(cmdId, data, null, maxEntitiesPerDataSubscription));
subIdsToCancel.forEach(subId -> localSubscriptionService.cancelSubscription(getSessionId(), subId));
subsToAdd.forEach(localSubscriptionService::addSubscription);
sendWsMsg(new EntityDataUpdate(cmdId, data, null, maxEntitiesPerDataSubscription));
}
public void setCurrentCmd(EntityDataCmd cmd) {

6
application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java

@ -442,7 +442,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi
private void handleWsAttributesSubscriptionByKeys(TelemetryWebSocketSessionRef sessionRef,
AttributesSubscriptionCmd cmd, String sessionId, EntityId entityId,
List<String> keys) {
FutureCallback<List<AttributeKvEntry>> callback = new FutureCallback<List<AttributeKvEntry>>() {
FutureCallback<List<AttributeKvEntry>> callback = new FutureCallback<>() {
@Override
public void onSuccess(List<AttributeKvEntry> data) {
List<TsKvEntry> attributesData = data.stream().map(d -> new BasicTsKvEntry(d.getLastUpdateTs(), d)).collect(Collectors.toList());
@ -542,7 +542,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi
private void handleWsAttributesSubscription(TelemetryWebSocketSessionRef sessionRef,
AttributesSubscriptionCmd cmd, String sessionId, EntityId entityId) {
FutureCallback<List<AttributeKvEntry>> callback = new FutureCallback<List<AttributeKvEntry>>() {
FutureCallback<List<AttributeKvEntry>> callback = new FutureCallback<>() {
@Override
public void onSuccess(List<AttributeKvEntry> data) {
List<TsKvEntry> attributesData = data.stream().map(d -> new BasicTsKvEntry(d.getLastUpdateTs(), d)).collect(Collectors.toList());
@ -666,7 +666,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi
}
private FutureCallback<List<TsKvEntry>> getSubscriptionCallback(final TelemetryWebSocketSessionRef sessionRef, final TimeseriesSubscriptionCmd cmd, final String sessionId, final EntityId entityId, final long startTs, final List<String> keys) {
return new FutureCallback<List<TsKvEntry>>() {
return new FutureCallback<>() {
@Override
public void onSuccess(List<TsKvEntry> data) {
sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), data));

53
application/src/test/java/org/thingsboard/server/controller/BaseEntityViewControllerTest.java

@ -62,6 +62,7 @@ import java.util.Set;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
import static java.util.concurrent.TimeUnit.HOURS;
import static java.util.concurrent.TimeUnit.MILLISECONDS;
import static java.util.concurrent.TimeUnit.SECONDS;
import static org.assertj.core.api.Assertions.assertThat;
@ -368,7 +369,7 @@ public abstract class BaseEntityViewControllerTest extends AbstractControllerTes
public void testTheCopyOfAttrsIntoTSForTheView() throws Exception {
Set<String> expectedActualAttributesSet = Set.of("caKey1", "caKey2", "caKey3", "caKey4");
Set<String> actualAttributesSet =
getAttributesByKeys("{\"caKey1\":\"value1\", \"caKey2\":true, \"caKey3\":42.0, \"caKey4\":73}", expectedActualAttributesSet);
putAttributesAndWait("{\"caKey1\":\"value1\", \"caKey2\":true, \"caKey3\":42.0, \"caKey4\":73}", expectedActualAttributesSet);
log.debug("got correct actualAttributesSet, saving new entity view...");
EntityView savedView = getNewSavedEntityView("Test entity view");
@ -389,13 +390,15 @@ public abstract class BaseEntityViewControllerTest extends AbstractControllerTes
@Test
public void testTheCopyOfAttrsOutOfTSForTheView() throws Exception {
long now = System.currentTimeMillis();
Set<String> expectedActualAttributesSet = Set.of("caKey1", "caKey2", "caKey3", "caKey4");
Set<String> actualAttributesSet =
getAttributesByKeys("{\"caKey1\":\"value1\", \"caKey2\":true, \"caKey3\":42.0, \"caKey4\":73}", expectedActualAttributesSet);
putAttributesAndWait("{\"caKey1\":\"value1\", \"caKey2\":true, \"caKey3\":42.0, \"caKey4\":73}", expectedActualAttributesSet);
List<Map<String, Object>> valueTelemetryOfDevices = doGetAsyncTyped("/api/plugins/telemetry/DEVICE/" + testDevice.getId().getId().toString() +
"/values/attributes?keys=" + String.join(",", actualAttributesSet), new TypeReference<>() {
List<Map<String, Object>> values = doGetAsyncTyped("/api/plugins/telemetry/DEVICE/" + testDevice.getId() +
"/values/attributes?keys=" + String.join(",", expectedActualAttributesSet), new TypeReference<>() {
});
assertEquals(expectedActualAttributesSet.size(), values.size());
EntityView view = new EntityView();
view.setEntityId(testDevice.getId());
@ -403,12 +406,12 @@ public abstract class BaseEntityViewControllerTest extends AbstractControllerTes
view.setName("Test entity view");
view.setType("default");
view.setKeys(telemetry);
view.setStartTimeMs((long) getValue(valueTelemetryOfDevices, "lastActivityTime") * 10);
view.setEndTimeMs((long) getValue(valueTelemetryOfDevices, "lastActivityTime") / 10);
view.setStartTimeMs(now - HOURS.toMillis(1));
view.setEndTimeMs(now - 1);
EntityView savedView = doPost("/api/entityView", view, EntityView.class);
List<Map<String, Object>> values = doGetAsyncTyped("/api/plugins/telemetry/ENTITY_VIEW/" + savedView.getId().getId().toString() +
"/values/attributes?keys=" + String.join(",", actualAttributesSet), new TypeReference<>() {
values = doGetAsyncTyped("/api/plugins/telemetry/ENTITY_VIEW/" + savedView.getId().getId().toString() +
"/values/attributes?keys=" + String.join(",", expectedActualAttributesSet), new TypeReference<>() {
});
assertEquals(0, values.size());
}
@ -431,19 +434,19 @@ public abstract class BaseEntityViewControllerTest extends AbstractControllerTes
uploadTelemetry("{\"tsKey1\":\"value1\", \"tsKey2\":true, \"tsKey3\":40.0}", accessToken);
getWsClient().waitForUpdate();
long startTimeMs = System.currentTimeMillis();
long startTimeMs = getCurTsButNotPrevTs(now);
getWsClient().registerWaitForUpdate();
uploadTelemetry("{\"tsKey1\":\"value2\", \"tsKey2\":false, \"tsKey3\":80.0}", accessToken);
getWsClient().waitForUpdate();
Thread.sleep(3);
long middleOfTestMs = getCurTsButNotPrevTs(startTimeMs);
getWsClient().registerWaitForUpdate();
uploadTelemetry("{\"tsKey1\":\"value3\", \"tsKey2\":false, \"tsKey3\":120.0}", accessToken);
getWsClient().waitForUpdate();
long endTimeMs = System.currentTimeMillis();
long endTimeMs = getCurTsButNotPrevTs(middleOfTestMs);
getWsClient().registerWaitForUpdate();
uploadTelemetry("{\"tsKey1\":\"value4\", \"tsKey2\":true, \"tsKey3\":160.0}", accessToken);
getWsClient().waitForUpdate();
@ -455,15 +458,25 @@ public abstract class BaseEntityViewControllerTest extends AbstractControllerTes
EntityView savedView = doPost("/api/entityView", view, EntityView.class);
String entityViewId = savedView.getId().getId().toString();
Map<String, List<Map<String, String>>> expectedValues = getTelemetryValues("DEVICE", deviceId, keys, 0L, (startTimeMs + endTimeMs) / 2);
Assert.assertEquals(2, expectedValues.get("tsKey1").size());
Assert.assertEquals(2, expectedValues.get("tsKey2").size());
Assert.assertEquals(2, expectedValues.get("tsKey3").size());
Map<String, List<Map<String, String>>> actualDeviceValues = getTelemetryValues("DEVICE", deviceId, keys, 0L, middleOfTestMs);
Assert.assertEquals(2, actualDeviceValues.get("tsKey1").size());
Assert.assertEquals(2, actualDeviceValues.get("tsKey2").size());
Assert.assertEquals(2, actualDeviceValues.get("tsKey3").size());
Map<String, List<Map<String, String>>> actualEntityViewValues = getTelemetryValues("ENTITY_VIEW", entityViewId, keys, 0L, middleOfTestMs);
Assert.assertEquals(1, actualEntityViewValues.get("tsKey1").size());
Assert.assertEquals(1, actualEntityViewValues.get("tsKey2").size());
Assert.assertEquals(1, actualEntityViewValues.get("tsKey3").size());
}
Map<String, List<Map<String, String>>> actualValues = getTelemetryValues("ENTITY_VIEW", entityViewId, keys, 0L, (startTimeMs + endTimeMs) / 2);
Assert.assertEquals(1, actualValues.get("tsKey1").size());
Assert.assertEquals(1, actualValues.get("tsKey2").size());
Assert.assertEquals(1, actualValues.get("tsKey3").size());
private static long getCurTsButNotPrevTs(long prevTs) throws InterruptedException {
long result = System.currentTimeMillis();
if (prevTs == result) {
Thread.sleep(1);
return getCurTsButNotPrevTs(prevTs);
} else {
return result;
}
}
private void uploadTelemetry(String strKvs, String accessToken) throws Exception {
@ -504,7 +517,7 @@ public abstract class BaseEntityViewControllerTest extends AbstractControllerTes
});
}
private Set<String> getAttributesByKeys(String stringKV, Set<String> expectedKeySet) throws Exception {
private Set<String> putAttributesAndWait(String stringKV, Set<String> expectedKeySet) throws Exception {
DeviceTypeFilter dtf = new DeviceTypeFilter(testDevice.getType(), testDevice.getName());
List<EntityKey> keysToSubscribe = expectedKeySet.stream()
.map(key -> new EntityKey(EntityKeyType.CLIENT_ATTRIBUTE, key))

14
application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java

@ -102,7 +102,6 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest {
List<TsKvEntry> tsData = Arrays.asList(dataPoint1, dataPoint2, dataPoint3);
sendTelemetry(device, tsData);
Thread.sleep(100);
update = getWsClient().sendHistoryCmd(keys, now, TimeUnit.HOURS.toMillis(1), dtf);
@ -136,7 +135,6 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest {
List<TsKvEntry> tsData = Arrays.asList(dataPoint1, dataPoint2, dataPoint3);
sendTelemetry(device, tsData);
Thread.sleep(100);
update = getWsClient().subscribeTsUpdate(List.of("temperature"), now, TimeUnit.HOURS.toMillis(1));
Assert.assertEquals(1, update.getCmdId());
@ -153,7 +151,6 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest {
now = System.currentTimeMillis();
TsKvEntry dataPoint4 = new BasicTsKvEntry(now, new LongDataEntry("temperature", 45L));
getWsClient().registerWaitForUpdate();
Thread.sleep(100);
sendTelemetry(device, Arrays.asList(dataPoint4));
String msg = getWsClient().waitForUpdate();
@ -309,13 +306,12 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest {
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());
getWsClient().registerWaitForUpdate();
TsKvEntry dataPoint1 = new BasicTsKvEntry(now - TimeUnit.MINUTES.toMillis(1), new LongDataEntry("temperature", 42L));
List<TsKvEntry> tsData = Arrays.asList(dataPoint1);
sendTelemetry(device, tsData);
Thread.sleep(100);
update = getWsClient().subscribeLatestUpdate(keys, dtf);
update = getWsClient().parseDataReply(getWsClient().waitForUpdate());
Assert.assertEquals(1, update.getCmdId());
@ -329,7 +325,6 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest {
now = System.currentTimeMillis();
TsKvEntry dataPoint2 = new BasicTsKvEntry(now, new LongDataEntry("temperature", 52L));
getWsClient().registerWaitForUpdate();
sendTelemetry(device, Arrays.asList(dataPoint2));
update = getWsClient().parseDataReply(getWsClient().waitForUpdate());
@ -371,7 +366,6 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest {
Assert.assertEquals("", pageData.getData().get(0).getLatest().get(EntityKeyType.SERVER_ATTRIBUTE).get("serverAttributeKey").getValue());
getWsClient().registerWaitForUpdate();
Thread.sleep(500);
AttributeKvEntry dataPoint1 = new BaseAttributeKvEntry(now - TimeUnit.MINUTES.toMillis(1), new LongDataEntry("serverAttributeKey", 42L));
List<AttributeKvEntry> tsData = Arrays.asList(dataPoint1);
@ -394,7 +388,6 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest {
AttributeKvEntry dataPoint2 = new BaseAttributeKvEntry(now, new LongDataEntry("serverAttributeKey", 52L));
getWsClient().registerWaitForUpdate();
Thread.sleep(500);
sendAttributes(device, TbAttributeSubscriptionScope.SERVER_SCOPE, Arrays.asList(dataPoint2));
msg = getWsClient().waitForUpdate();
Assert.assertNotNull(msg);
@ -411,14 +404,12 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest {
//Sending update from the past, while latest value has new timestamp;
getWsClient().registerWaitForUpdate();
Thread.sleep(500);
sendAttributes(device, TbAttributeSubscriptionScope.SERVER_SCOPE, Arrays.asList(dataPoint1));
msg = getWsClient().waitForUpdate(TimeUnit.SECONDS.toMillis(1));
Assert.assertNull(msg);
//Sending duplicate update again
getWsClient().registerWaitForUpdate();
Thread.sleep(500);
sendAttributes(device, TbAttributeSubscriptionScope.SERVER_SCOPE, Arrays.asList(dataPoint2));
msg = getWsClient().waitForUpdate(TimeUnit.SECONDS.toMillis(1));
Assert.assertNull(msg);
@ -456,7 +447,6 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest {
getWsClient().registerWaitForUpdate();
AttributeKvEntry dataPoint1 = new BaseAttributeKvEntry(now - TimeUnit.MINUTES.toMillis(1), new LongDataEntry("serverAttributeKey", 42L));
List<AttributeKvEntry> tsData = Arrays.asList(dataPoint1);
Thread.sleep(100);
sendAttributes(device, TbAttributeSubscriptionScope.SERVER_SCOPE, tsData);

4
application/src/test/java/org/thingsboard/server/controller/TbTestWebSocketClient.java

@ -44,8 +44,8 @@ import java.util.concurrent.TimeUnit;
public class TbTestWebSocketClient extends WebSocketClient {
private volatile String lastMsg;
private CountDownLatch reply;
private CountDownLatch update;
private volatile CountDownLatch reply;
private volatile CountDownLatch update;
public TbTestWebSocketClient(URI serverUri) {
super(serverUri);

10
ui-ngx/src/app/core/api/alias-controller.ts

@ -17,14 +17,15 @@
import { AliasInfo, IAliasController, StateControllerHolder, StateEntityInfo } from '@core/api/widget-api.models';
import { forkJoin, Observable, of, ReplaySubject, Subject } from 'rxjs';
import { Datasource, DatasourceType, datasourceTypeTranslationMap } from '@app/shared/models/widget.models';
import { deepClone, isEqual } from '@core/utils';
import { deepClone, isDefinedAndNotNull, isEqual } from '@core/utils';
import { EntityService } from '@core/http/entity.service';
import { UtilsService } from '@core/services/utils.service';
import { AliasFilterType, EntityAliases, SingleEntityFilter } from '@shared/models/alias.models';
import { EntityInfo } from '@shared/models/entity.models';
import { map, mergeMap } from 'rxjs/operators';
import {
defaultEntityDataPageLink, Filter, FilterInfo, filterInfoToKeyFilters, Filters, KeyFilter, singleEntityDataPageLink,
createDefaultEntityDataPageLink,
Filter, FilterInfo, filterInfoToKeyFilters, Filters, KeyFilter, singleEntityDataPageLink,
updateDatasourceFromEntityInfo
} from '@shared/models/query/query.models';
import { TranslateService } from '@ngx-translate/core';
@ -322,7 +323,7 @@ export class AliasController implements IAliasController {
);
}
resolveDatasources(datasources: Array<Datasource>, singleEntity?: boolean): Observable<Array<Datasource>> {
resolveDatasources(datasources: Array<Datasource>, singleEntity?: boolean, pageSize = 1024): Observable<Array<Datasource>> {
if (!datasources || !datasources.length) {
return of([]);
}
@ -360,7 +361,8 @@ export class AliasController implements IAliasController {
if (singleEntity) {
datasource.pageLink = deepClone(singleEntityDataPageLink);
} else if (!datasource.pageLink) {
datasource.pageLink = deepClone(defaultEntityDataPageLink);
pageSize = isDefinedAndNotNull(pageSize) && pageSize > 0 ? pageSize : 1024;
datasource.pageLink = createDefaultEntityDataPageLink(pageSize);
}
}
});

4
ui-ngx/src/app/core/api/widget-api.models.ts

@ -119,7 +119,7 @@ export interface IAliasController {
getEntityAliasId(aliasName: string): string;
getInstantAliasInfo(aliasId: string): AliasInfo;
resolveSingleEntityInfo(aliasId: string): Observable<EntityInfo>;
resolveDatasources(datasources: Array<Datasource>, singleEntity?: boolean): Observable<Array<Datasource>>;
resolveDatasources(datasources: Array<Datasource>, singleEntity?: boolean, pageSize?: number): Observable<Array<Datasource>>;
resolveAlarmSource(alarmSource: Datasource): Observable<Datasource>;
getEntityAliases(): EntityAliases;
getFilters(): Filters;
@ -184,6 +184,7 @@ export interface SubscriptionInfo {
deviceName?: string;
deviceNamePrefix?: string;
deviceIds?: Array<string>;
pageSize?: number;
}
export class WidgetSubscriptionContext {
@ -242,6 +243,7 @@ export interface WidgetSubscriptionOptions {
datasourcesOptional?: boolean;
hasDataPageLink?: boolean;
singleEntity?: boolean;
pageSize?: number;
warnOnPageDataOverflow?: boolean;
ignoreDataUpdateOnIntervalTick?: boolean;
targetDeviceAliasIds?: Array<string>;

8
ui-ngx/src/app/core/api/widget-subscription.ts

@ -97,6 +97,7 @@ export class WidgetSubscription implements IWidgetSubscription {
hasDataPageLink: boolean;
singleEntity: boolean;
pageSize: number;
warnOnPageDataOverflow: boolean;
ignoreDataUpdateOnIntervalTick: boolean;
@ -229,6 +230,7 @@ export class WidgetSubscription implements IWidgetSubscription {
this.entityDataListeners = [];
this.hasDataPageLink = options.hasDataPageLink;
this.singleEntity = options.singleEntity;
this.pageSize = options.pageSize;
this.warnOnPageDataOverflow = options.warnOnPageDataOverflow;
this.ignoreDataUpdateOnIntervalTick = options.ignoreDataUpdateOnIntervalTick;
this.datasourcePages = [];
@ -387,7 +389,7 @@ export class WidgetSubscription implements IWidgetSubscription {
}
);
} else {
this.ctx.aliasController.resolveDatasources(this.configuredDatasources, this.singleEntity).subscribe(
this.ctx.aliasController.resolveDatasources(this.configuredDatasources, this.singleEntity, this.pageSize).subscribe(
(datasources) => {
this.configuredDatasources = datasources;
this.prepareDataSubscriptions().subscribe(
@ -1132,7 +1134,7 @@ export class WidgetSubscription implements IWidgetSubscription {
}
);
} else {
this.ctx.aliasController.resolveDatasources(this.configuredDatasources, this.singleEntity).subscribe(
this.ctx.aliasController.resolveDatasources(this.configuredDatasources, this.singleEntity, this.pageSize).subscribe(
(datasources) => {
this.configuredDatasources = datasources;
this.prepareDataSubscriptions().subscribe(
@ -1271,7 +1273,7 @@ export class WidgetSubscription implements IWidgetSubscription {
totalPages: pageData.totalPages
};
if (datasource.type === DatasourceType.entity &&
pageData.hasNext && pageLink.pageSize > 1) {
pageData.hasNext && !this.singleEntity) {
if (this.warnOnPageDataOverflow) {
const message = this.ctx.translate.instant('widget.data-overflow',
{count: pageData.data.length, total: pageData.totalElements});

6
ui-ngx/src/app/core/http/entity.service.ts

@ -1319,7 +1319,8 @@ export class EntityService {
pageLink = deepClone(singleEntityDataPageLink);
} else {
nameFilter = subscriptionInfo.entityNamePrefix;
pageLink = deepClone(defaultEntityDataPageLink);
const pageSize = isDefinedAndNotNull(subscriptionInfo.pageSize) && subscriptionInfo.pageSize > 0 ? subscriptionInfo.pageSize : 1024;
pageLink = createDefaultEntityDataPageLink(pageSize);
}
datasource.entityFilter = {
type: AliasFilterType.entityName,
@ -1333,7 +1334,8 @@ export class EntityService {
entityType: subscriptionInfo.entityType,
entityList: subscriptionInfo.entityIds
};
datasource.pageLink = deepClone(defaultEntityDataPageLink);
const pageSize = isDefinedAndNotNull(subscriptionInfo.pageSize) && subscriptionInfo.pageSize > 0 ? subscriptionInfo.pageSize : 1024;
datasource.pageLink = createDefaultEntityDataPageLink(pageSize);
}
}

8
ui-ngx/src/app/modules/home/components/widget/lib/maps/map-widget2.ts

@ -45,7 +45,7 @@ import { TranslateService } from '@ngx-translate/core';
import { UtilsService } from '@core/services/utils.service';
import { EntityDataPageLink } from '@shared/models/query/query.models';
import { providerClass } from '@home/components/widget/lib/maps/providers';
import { isDefined, parseFunction } from '@core/utils';
import { isDefined, isDefinedAndNotNull, parseFunction } from '@core/utils';
import L from 'leaflet';
import { forkJoin, Observable, of } from 'rxjs';
import { AttributeService } from '@core/http/attribute.service';
@ -87,9 +87,13 @@ export class MapWidgetController implements MapWidgetInterface {
this.map.saveMarkerLocation = this.setMarkerLocation.bind(this);
this.map.savePolygonLocation = this.savePolygonLocation.bind(this);
this.map.saveLocation = this.saveLocation.bind(this);
let pageSize = this.settings.mapPageSize;
if (isDefinedAndNotNull(this.ctx.widgetConfig.pageSize)) {
pageSize = Math.max(pageSize, this.ctx.widgetConfig.pageSize);
}
this.pageLink = {
page: 0,
pageSize: this.settings.mapPageSize,
pageSize,
textSearch: null,
dynamic: true
};

1
ui-ngx/src/app/modules/home/components/widget/lib/maps/providers/image-map.ts

@ -88,6 +88,7 @@ export class ImageMap extends LeafletMap {
const imageUrlSubscriptionOptions: WidgetSubscriptionOptions = {
datasources,
hasDataPageLink: true,
singleEntity: true,
useDashboardTimewindow: false,
type: widgetType.latest,
callbacks: {

6
ui-ngx/src/app/modules/home/components/widget/lib/markdown-widget.component.ts

@ -26,7 +26,7 @@ import {
fillDataPattern,
flatFormattedData,
formattedDataFormDatasourceData,
hashCode,
hashCode, isDefinedAndNotNull,
isNotEmptyStr,
parseFunction, processDataPattern,
safeExecute
@ -83,9 +83,11 @@ export class MarkdownWidgetComponent extends PageComponent implements OnInit {
cssParser.cssPreviewNamespace = this.markdownClass;
cssParser.createStyleElement(this.markdownClass, cssString);
}
const pageSize = isDefinedAndNotNull(this.ctx.widgetConfig.pageSize) &&
this.ctx.widgetConfig.pageSize > 0 ? this.ctx.widgetConfig.pageSize : 16384;
const pageLink: EntityDataPageLink = {
page: 0,
pageSize: 16384,
pageSize,
textSearch: null,
dynamic: true
};

8
ui-ngx/src/app/modules/home/components/widget/widget-config.component.html

@ -317,6 +317,14 @@
<mat-panel-title translate>widget-config.data-settings</mat-panel-title>
</mat-expansion-panel-header>
<ng-template matExpansionPanelContent>
<div fxLayout="row" *ngIf="widgetType !== widgetTypes.rpc &&
widgetType !== widgetTypes.alarm &&
modelValue?.isDataEnabled && !modelValue?.typeParameters?.singleEntity">
<mat-form-field fxFlex>
<mat-label translate>widget-config.data-page-size</mat-label>
<input matInput formControlName="pageSize" type="number" min="1" step="1">
</mat-form-field>
</div>
<div fxLayout.xs="column" fxLayout="row" fxLayoutGap="8px">
<mat-form-field fxFlex>
<mat-label translate>widget-config.units</mat-label>

2
ui-ngx/src/app/modules/home/components/widget/widget-config.component.ts

@ -213,6 +213,7 @@ export class WidgetConfigComponent extends PageComponent implements OnInit, Cont
widgetStyle: [null, []],
widgetCss: [null, []],
titleStyle: [null, []],
pageSize: [1024, [Validators.min(1), Validators.pattern(/^\d*$/)]],
units: [null, []],
decimals: [null, [Validators.min(0), Validators.max(15), Validators.pattern(/^\d*$/)]],
noDataDisplayMessage: [null, []],
@ -420,6 +421,7 @@ export class WidgetConfigComponent extends PageComponent implements OnInit, Cont
fontSize: '16px',
fontWeight: 400
},
pageSize: isDefined(config.pageSize) ? config.pageSize : 1024,
units: config.units,
decimals: config.decimals,
noDataDisplayMessage: isDefined(config.noDataDisplayMessage) ? config.noDataDisplayMessage : '',

3
ui-ngx/src/app/modules/home/components/widget/widget.component.ts

@ -977,7 +977,8 @@ export class WidgetComponent extends PageComponent implements OnInit, AfterViewI
ignoreDataUpdateOnIntervalTick: this.typeParameters.ignoreDataUpdateOnIntervalTick,
comparisonEnabled: comparisonSettings.comparisonEnabled,
timeForComparison: comparisonSettings.timeForComparison,
comparisonCustomIntervalValue: comparisonSettings.comparisonCustomIntervalValue
comparisonCustomIntervalValue: comparisonSettings.comparisonCustomIntervalValue,
pageSize: this.widget.config.pageSize
};
if (this.widget.type === widgetType.alarm) {
options.alarmSource = deepClone(this.widget.config.alarmSource);

1
ui-ngx/src/app/shared/models/widget.models.ts

@ -554,6 +554,7 @@ export interface WidgetConfig {
units?: string;
decimals?: number;
noDataDisplayMessage?: string;
pageSize?: number;
actions?: {[actionSourceId: string]: Array<WidgetActionDescriptor>};
settings?: WidgetSettings;
alarmSource?: Datasource;

1
ui-ngx/src/assets/locale/locale.constant-en_US.json

@ -3246,6 +3246,7 @@
"advanced-settings": "Advanced settings",
"data-settings": "Data settings",
"no-data-display-message": "\"No data to display\" alternative message",
"data-page-size": "Maximum entities per datasource",
"settings-component-not-found": "Settings form component not found for selector '{{selector}}'"
},
"widget-type": {

Loading…
Cancel
Save