Browse Source

Merge branch 'master' into feature/log-telemetry-updated

pull/3602/head
Viacheslav Kukhtyn 6 years ago
parent
commit
9c44920fe7
  1. 26
      application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java
  2. 133
      application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java
  3. 6
      application/src/main/java/org/thingsboard/server/service/query/EntityQueryService.java
  4. 3
      application/src/main/resources/thingsboard.yml
  5. 3
      common/dao-api/src/main/java/org/thingsboard/server/dao/attributes/AttributesService.java
  6. 2
      common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java
  7. 37
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
  8. 3
      dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesDao.java
  9. 6
      dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesService.java
  10. 3
      dao/src/main/java/org/thingsboard/server/dao/sql/attributes/AttributeKvRepository.java
  11. 7
      dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java
  12. 6
      dao/src/main/java/org/thingsboard/server/dao/sqlts/SqlTimeseriesLatestDao.java
  13. 5
      dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/TsKvLatestRepository.java
  14. 5
      dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java
  15. 75
      dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java
  16. 5
      dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java
  17. 30
      dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraPartitionCacheKey.java
  18. 42
      dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraTsPartitionsCache.java
  19. 2
      dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesLatestDao.java
  20. 110
      dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java
  21. 2
      dao/src/test/resources/cassandra-test.properties
  22. 74
      ui-ngx/src/app/core/http/entity.service.ts
  23. 164
      ui-ngx/src/app/modules/home/components/profile/alarm/device-profile-alarm.component.html
  24. 38
      ui-ngx/src/app/modules/home/components/profile/device-profile.component.html
  25. 10
      ui-ngx/src/app/modules/home/components/profile/tenant-profile-data.component.html
  26. 34
      ui-ngx/src/app/modules/home/components/widget/data-key-config.component.ts
  27. 2
      ui-ngx/src/app/modules/home/components/widget/data-keys.component.models.ts
  28. 37
      ui-ngx/src/app/modules/home/components/widget/data-keys.component.ts
  29. 34
      ui-ngx/src/app/modules/home/components/widget/lib/maps/leaflet-map.ts
  30. 8
      ui-ngx/src/app/modules/home/components/widget/trip-animation/trip-animation.component.html
  31. 5
      ui-ngx/src/app/modules/home/components/widget/trip-animation/trip-animation.component.scss
  32. 80
      ui-ngx/src/app/modules/home/components/widget/trip-animation/trip-animation.component.ts
  33. 62
      ui-ngx/src/app/modules/home/components/widget/widget-config.component.ts
  34. 6
      ui-ngx/src/app/shared/models/entity.models.ts

26
application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java

@ -16,18 +16,23 @@
package org.thingsboard.server.controller; package org.thingsboard.server.controller;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.http.ResponseEntity;
import org.springframework.security.access.prepost.PreAuthorize; import org.springframework.security.access.prepost.PreAuthorize;
import org.springframework.web.bind.annotation.RequestBody; import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestMethod; import org.springframework.web.bind.annotation.RequestMethod;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.ResponseBody; import org.springframework.web.bind.annotation.ResponseBody;
import org.springframework.web.bind.annotation.RestController; import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.context.request.async.DeferredResult;
import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.exception.ThingsboardException;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.query.AlarmData; import org.thingsboard.server.common.data.query.AlarmData;
import org.thingsboard.server.common.data.query.AlarmDataQuery; import org.thingsboard.server.common.data.query.AlarmDataQuery;
import org.thingsboard.server.common.data.query.EntityCountQuery; import org.thingsboard.server.common.data.query.EntityCountQuery;
import org.thingsboard.server.common.data.query.EntityData; import org.thingsboard.server.common.data.query.EntityData;
import org.thingsboard.server.common.data.query.EntityDataPageLink;
import org.thingsboard.server.common.data.query.EntityDataQuery; import org.thingsboard.server.common.data.query.EntityDataQuery;
import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.query.EntityQueryService; import org.thingsboard.server.service.query.EntityQueryService;
@ -40,6 +45,7 @@ public class EntityQueryController extends BaseController {
@Autowired @Autowired
private EntityQueryService entityQueryService; private EntityQueryService entityQueryService;
private static final int MAX_PAGE_SIZE = 100;
@PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") @PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')")
@RequestMapping(value = "/entitiesQuery/count", method = RequestMethod.POST) @RequestMapping(value = "/entitiesQuery/count", method = RequestMethod.POST)
@ -76,4 +82,24 @@ public class EntityQueryController extends BaseController {
throw handleException(e); throw handleException(e);
} }
} }
@PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')")
@RequestMapping(value = "/entitiesQuery/find/keys", method = RequestMethod.POST)
@ResponseBody
public DeferredResult<ResponseEntity> findEntityTimeseriesAndAttributesKeysByQuery(@RequestBody EntityDataQuery query,
@RequestParam("timeseries") boolean isTimeseries,
@RequestParam("attributes") boolean isAttributes) throws ThingsboardException {
TenantId tenantId = getTenantId();
checkNotNull(query);
try {
EntityDataPageLink pageLink = query.getPageLink();
if (pageLink.getPageSize() > MAX_PAGE_SIZE) {
pageLink.setPageSize(MAX_PAGE_SIZE);
}
return entityQueryService.getKeysByQuery(getCurrentUser(), tenantId, query, isTimeseries, isAttributes);
} catch (Exception e) {
throw handleException(e);
}
}
} }

133
application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java

@ -15,11 +15,23 @@
*/ */
package org.thingsboard.server.service.query; package org.thingsboard.server.service.query;
import com.fasterxml.jackson.databind.node.ArrayNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.checkerframework.checker.nullness.qual.Nullable;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.http.HttpStatus;
import org.springframework.http.ResponseEntity;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.springframework.util.CollectionUtils;
import org.springframework.web.context.request.async.DeferredResult;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.EntityId; 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.page.PageData;
import org.thingsboard.server.common.data.query.AlarmData; import org.thingsboard.server.common.data.query.AlarmData;
import org.thingsboard.server.common.data.query.AlarmDataQuery; import org.thingsboard.server.common.data.query.AlarmDataQuery;
@ -31,12 +43,25 @@ import org.thingsboard.server.common.data.query.EntityDataSortOrder;
import org.thingsboard.server.common.data.query.EntityKey; import org.thingsboard.server.common.data.query.EntityKey;
import org.thingsboard.server.common.data.query.EntityKeyType; import org.thingsboard.server.common.data.query.EntityKeyType;
import org.thingsboard.server.dao.alarm.AlarmService; import org.thingsboard.server.dao.alarm.AlarmService;
import org.thingsboard.server.dao.attributes.AttributesService;
import org.thingsboard.server.dao.entity.EntityService; import org.thingsboard.server.dao.entity.EntityService;
import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.model.ModelConstants;
import org.thingsboard.server.dao.timeseries.TimeseriesService;
import org.thingsboard.server.dao.util.mapping.JacksonUtil;
import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.executors.DbCallbackExecutorService;
import org.thingsboard.server.service.security.AccessValidator;
import org.thingsboard.server.service.security.model.SecurityUser; import org.thingsboard.server.service.security.model.SecurityUser;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.LinkedHashMap; import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.function.Consumer;
import java.util.stream.Collectors;
@Service @Service
@Slf4j @Slf4j
@ -52,6 +77,15 @@ public class DefaultEntityQueryService implements EntityQueryService {
@Value("${server.ws.max_entities_per_alarm_subscription:1000}") @Value("${server.ws.max_entities_per_alarm_subscription:1000}")
private int maxEntitiesPerAlarmSubscription; private int maxEntitiesPerAlarmSubscription;
@Autowired
private DbCallbackExecutorService dbCallbackExecutor;
@Autowired
private TimeseriesService timeseriesService;
@Autowired
private AttributesService attributesService;
@Override @Override
public long countEntitiesByQuery(SecurityUser securityUser, EntityCountQuery query) { public long countEntitiesByQuery(SecurityUser securityUser, EntityCountQuery query) {
return entityService.countEntitiesByQuery(securityUser.getTenantId(), securityUser.getCustomerId(), query); return entityService.countEntitiesByQuery(securityUser.getTenantId(), securityUser.getCustomerId(), query);
@ -100,4 +134,103 @@ public class DefaultEntityQueryService implements EntityQueryService {
EntityDataPageLink edpl = new EntityDataPageLink(maxEntitiesPerAlarmSubscription, 0, null, entitiesSortOrder); EntityDataPageLink edpl = new EntityDataPageLink(maxEntitiesPerAlarmSubscription, 0, null, entitiesSortOrder);
return new EntityDataQuery(query.getEntityFilter(), edpl, query.getEntityFields(), query.getLatestValues(), query.getKeyFilters()); return new EntityDataQuery(query.getEntityFilter(), edpl, query.getEntityFields(), query.getLatestValues(), query.getKeyFilters());
} }
@Override
public DeferredResult<ResponseEntity> getKeysByQuery(SecurityUser securityUser, TenantId tenantId, EntityDataQuery query,
boolean isTimeseries, boolean isAttributes) {
final DeferredResult<ResponseEntity> response = new DeferredResult<>();
if (!isAttributes && !isTimeseries) {
replyWithEmptyResponse(response);
return response;
}
List<EntityId> ids = this.findEntityDataByQuery(securityUser, query).getData().stream()
.map(EntityData::getEntityId)
.collect(Collectors.toList());
if (ids.isEmpty()) {
replyWithEmptyResponse(response);
return response;
}
Set<EntityType> types = ids.stream().map(EntityId::getEntityType).collect(Collectors.toSet());
final ListenableFuture<List<String>> timeseriesKeysFuture;
final ListenableFuture<List<String>> attributesKeysFuture;
if (isTimeseries) {
timeseriesKeysFuture = dbCallbackExecutor.submit(() -> timeseriesService.findAllKeysByEntityIds(tenantId, ids));
} else {
timeseriesKeysFuture = null;
}
if (isAttributes) {
Map<EntityType, List<EntityId>> typesMap = ids.stream().collect(Collectors.groupingBy(EntityId::getEntityType));
List<ListenableFuture<List<String>>> futures = new ArrayList<>(typesMap.size());
typesMap.forEach((type, entityIds) -> futures.add(dbCallbackExecutor.submit(() -> attributesService.findAllKeysByEntityIds(tenantId, type, entityIds))));
attributesKeysFuture = Futures.transform(Futures.allAsList(futures), lists -> {
if (CollectionUtils.isEmpty(lists)) {
return Collections.emptyList();
}
return lists.stream().flatMap(List::stream).distinct().sorted().collect(Collectors.toList());
}, dbCallbackExecutor);
} else {
attributesKeysFuture = null;
}
if (isTimeseries && isAttributes) {
Futures.whenAllComplete(timeseriesKeysFuture, attributesKeysFuture).run(() -> {
try {
replyWithResponse(response, types, timeseriesKeysFuture.get(), attributesKeysFuture.get());
} catch (Exception e) {
log.error("Failed to fetch timeseries and attributes keys!", e);
AccessValidator.handleError(e, response, HttpStatus.INTERNAL_SERVER_ERROR);
}
}, dbCallbackExecutor);
} else if (isTimeseries) {
addCallback(timeseriesKeysFuture, keys -> replyWithResponse(response, types, keys, null),
error -> {
log.error("Failed to fetch timeseries keys!", error);
AccessValidator.handleError(error, response, HttpStatus.INTERNAL_SERVER_ERROR);
});
} else {
addCallback(attributesKeysFuture, keys -> replyWithResponse(response, types, null, keys),
error -> {
log.error("Failed to fetch attributes keys!", error);
AccessValidator.handleError(error, response, HttpStatus.INTERNAL_SERVER_ERROR);
});
}
return response;
}
private void replyWithResponse(DeferredResult<ResponseEntity> response, Set<EntityType> types, List<String> timeseriesKeys, List<String> attributesKeys) {
ObjectNode json = JacksonUtil.newObjectNode();
addItemsToArrayNode(json.putArray("entityTypes"), types);
addItemsToArrayNode(json.putArray("timeseries"), timeseriesKeys);
addItemsToArrayNode(json.putArray("attribute"), attributesKeys);
response.setResult(new ResponseEntity(json, HttpStatus.OK));
}
private void replyWithEmptyResponse(DeferredResult<ResponseEntity> response) {
replyWithResponse(response, Collections.emptySet(), Collections.emptyList(), Collections.emptyList());
}
private void addItemsToArrayNode(ArrayNode arrayNode, Collection<?> collection) {
if (!CollectionUtils.isEmpty(collection)) {
collection.forEach(item -> arrayNode.add(item.toString()));
}
}
private void addCallback(ListenableFuture<List<String>> future, Consumer<List<String>> success, Consumer<Throwable> error) {
Futures.addCallback(future, new FutureCallback<List<String>>() {
@Override
public void onSuccess(@Nullable List<String> keys) {
success.accept(keys);
}
@Override
public void onFailure(Throwable t) {
error.accept(t);
}
}, dbCallbackExecutor);
}
} }

6
application/src/main/java/org/thingsboard/server/service/query/EntityQueryService.java

@ -15,6 +15,9 @@
*/ */
package org.thingsboard.server.service.query; package org.thingsboard.server.service.query;
import org.springframework.http.ResponseEntity;
import org.springframework.web.context.request.async.DeferredResult;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.query.AlarmData; import org.thingsboard.server.common.data.query.AlarmData;
import org.thingsboard.server.common.data.query.AlarmDataQuery; import org.thingsboard.server.common.data.query.AlarmDataQuery;
@ -31,4 +34,7 @@ public interface EntityQueryService {
PageData<AlarmData> findAlarmDataByQuery(SecurityUser securityUser, AlarmDataQuery query); PageData<AlarmData> findAlarmDataByQuery(SecurityUser securityUser, AlarmDataQuery query);
DeferredResult<ResponseEntity> getKeysByQuery(SecurityUser securityUser, TenantId tenantId, EntityDataQuery query,
boolean isTimeseries, boolean isAttributes);
} }

3
application/src/main/resources/thingsboard.yml

@ -192,8 +192,9 @@ cassandra:
read_consistency_level: "${CASSANDRA_READ_CONSISTENCY_LEVEL:ONE}" read_consistency_level: "${CASSANDRA_READ_CONSISTENCY_LEVEL:ONE}"
write_consistency_level: "${CASSANDRA_WRITE_CONSISTENCY_LEVEL:ONE}" write_consistency_level: "${CASSANDRA_WRITE_CONSISTENCY_LEVEL:ONE}"
default_fetch_size: "${CASSANDRA_DEFAULT_FETCH_SIZE:2000}" default_fetch_size: "${CASSANDRA_DEFAULT_FETCH_SIZE:2000}"
# Specify partitioning size for timestamp key-value storage. Example: MINUTES, HOURS, DAYS, MONTHS,INDEFINITE # Specify partitioning size for timestamp key-value storage. Example: MINUTES, HOURS, DAYS, MONTHS, INDEFINITE
ts_key_value_partitioning: "${TS_KV_PARTITIONING:MONTHS}" ts_key_value_partitioning: "${TS_KV_PARTITIONING:MONTHS}"
ts_key_value_partitions_max_cache_size: "${TS_KV_PARTITIONS_MAX_CACHE_SIZE:100000}"
ts_key_value_ttl: "${TS_KV_TTL:0}" ts_key_value_ttl: "${TS_KV_TTL:0}"
events_ttl: "${TS_EVENTS_TTL:0}" events_ttl: "${TS_EVENTS_TTL:0}"
# Specify TTL of debug log in seconds. The current value corresponds to one week # Specify TTL of debug log in seconds. The current value corresponds to one week

3
common/dao-api/src/main/java/org/thingsboard/server/dao/attributes/AttributesService.java

@ -16,6 +16,7 @@
package org.thingsboard.server.dao.attributes; package org.thingsboard.server.dao.attributes;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.DeviceProfileId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
@ -42,4 +43,6 @@ public interface AttributesService {
List<String> findAllKeysByDeviceProfileId(TenantId tenantId, DeviceProfileId deviceProfileId); List<String> findAllKeysByDeviceProfileId(TenantId tenantId, DeviceProfileId deviceProfileId);
List<String> findAllKeysByEntityIds(TenantId tenantId, EntityType entityType, List<EntityId> entityIds);
} }

2
common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java

@ -50,4 +50,6 @@ public interface TimeseriesService {
ListenableFuture<Collection<String>> removeAllLatest(TenantId tenantId, EntityId entityId); ListenableFuture<Collection<String>> removeAllLatest(TenantId tenantId, EntityId entityId);
List<String> findAllKeysByDeviceProfileId(TenantId tenantId, DeviceProfileId deviceProfileId); List<String> findAllKeysByDeviceProfileId(TenantId tenantId, DeviceProfileId deviceProfileId);
List<String> findAllKeysByEntityIds(TenantId tenantId, List<EntityId> entityIds);
} }

37
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java

@ -462,26 +462,26 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
String clientId = msg.payload().clientIdentifier(); String clientId = msg.payload().clientIdentifier();
if (DataConstants.PROVISION.equals(userName) || DataConstants.PROVISION.equals(clientId)) { if (DataConstants.PROVISION.equals(userName) || DataConstants.PROVISION.equals(clientId)) {
deviceSessionCtx.setProvisionOnly(true); deviceSessionCtx.setProvisionOnly(true);
ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED)); ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED, msg));
} else { } else {
X509Certificate cert; X509Certificate cert;
if (sslHandler != null && (cert = getX509Certificate()) != null) { if (sslHandler != null && (cert = getX509Certificate()) != null) {
processX509CertConnect(ctx, cert); processX509CertConnect(ctx, cert, msg);
} else { } else {
processAuthTokenConnect(ctx, msg); processAuthTokenConnect(ctx, msg);
} }
} }
} }
private void processAuthTokenConnect(ChannelHandlerContext ctx, MqttConnectMessage msg) { private void processAuthTokenConnect(ChannelHandlerContext ctx, MqttConnectMessage connectMessage) {
String userName = msg.payload().userName(); String userName = connectMessage.payload().userName();
log.info("[{}] Processing connect msg for client with user name: {}!", sessionId, userName); log.info("[{}] Processing connect msg for client with user name: {}!", sessionId, userName);
TransportProtos.ValidateBasicMqttCredRequestMsg.Builder request = TransportProtos.ValidateBasicMqttCredRequestMsg.newBuilder() TransportProtos.ValidateBasicMqttCredRequestMsg.Builder request = TransportProtos.ValidateBasicMqttCredRequestMsg.newBuilder()
.setClientId(msg.payload().clientIdentifier()); .setClientId(connectMessage.payload().clientIdentifier());
if (userName != null) { if (userName != null) {
request.setUserName(userName); request.setUserName(userName);
} }
String password = msg.payload().password(); String password = connectMessage.payload().password();
if (password != null) { if (password != null) {
request.setPassword(password); request.setPassword(password);
} }
@ -489,19 +489,19 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
new TransportServiceCallback<ValidateDeviceCredentialsResponse>() { new TransportServiceCallback<ValidateDeviceCredentialsResponse>() {
@Override @Override
public void onSuccess(ValidateDeviceCredentialsResponse msg) { public void onSuccess(ValidateDeviceCredentialsResponse msg) {
onValidateDeviceResponse(msg, ctx); onValidateDeviceResponse(msg, ctx, connectMessage);
} }
@Override @Override
public void onError(Throwable e) { public void onError(Throwable e) {
log.trace("[{}] Failed to process credentials: {}", address, userName, e); log.trace("[{}] Failed to process credentials: {}", address, userName, e);
ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_SERVER_UNAVAILABLE)); ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_SERVER_UNAVAILABLE, connectMessage));
ctx.close(); ctx.close();
} }
}); });
} }
private void processX509CertConnect(ChannelHandlerContext ctx, X509Certificate cert) { private void processX509CertConnect(ChannelHandlerContext ctx, X509Certificate cert, MqttConnectMessage connectMessage) {
try { try {
if (!context.isSkipValidityCheckForClientCert()) { if (!context.isSkipValidityCheckForClientCert()) {
cert.checkValidity(); cert.checkValidity();
@ -512,18 +512,18 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
new TransportServiceCallback<ValidateDeviceCredentialsResponse>() { new TransportServiceCallback<ValidateDeviceCredentialsResponse>() {
@Override @Override
public void onSuccess(ValidateDeviceCredentialsResponse msg) { public void onSuccess(ValidateDeviceCredentialsResponse msg) {
onValidateDeviceResponse(msg, ctx); onValidateDeviceResponse(msg, ctx, connectMessage);
} }
@Override @Override
public void onError(Throwable e) { public void onError(Throwable e) {
log.trace("[{}] Failed to process credentials: {}", address, sha3Hash, e); log.trace("[{}] Failed to process credentials: {}", address, sha3Hash, e);
ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_SERVER_UNAVAILABLE)); ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_SERVER_UNAVAILABLE, connectMessage));
ctx.close(); ctx.close();
} }
}); });
} catch (Exception e) { } catch (Exception e) {
ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED)); ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED, connectMessage));
ctx.close(); ctx.close();
} }
} }
@ -547,11 +547,11 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
doDisconnect(); doDisconnect();
} }
private MqttConnAckMessage createMqttConnAckMsg(MqttConnectReturnCode returnCode) { private MqttConnAckMessage createMqttConnAckMsg(MqttConnectReturnCode returnCode, MqttConnectMessage msg) {
MqttFixedHeader mqttFixedHeader = MqttFixedHeader mqttFixedHeader =
new MqttFixedHeader(CONNACK, false, AT_MOST_ONCE, false, 0); new MqttFixedHeader(CONNACK, false, AT_MOST_ONCE, false, 0);
MqttConnAckVariableHeader mqttConnAckVariableHeader = MqttConnAckVariableHeader mqttConnAckVariableHeader =
new MqttConnAckVariableHeader(returnCode, true); new MqttConnAckVariableHeader(returnCode, !msg.variableHeader().isCleanSession());
return new MqttConnAckMessage(mqttFixedHeader, mqttConnAckVariableHeader); return new MqttConnAckMessage(mqttFixedHeader, mqttConnAckVariableHeader);
} }
@ -627,9 +627,10 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
} }
} }
private void onValidateDeviceResponse(ValidateDeviceCredentialsResponse msg, ChannelHandlerContext ctx) {
private void onValidateDeviceResponse(ValidateDeviceCredentialsResponse msg, ChannelHandlerContext ctx, MqttConnectMessage connectMessage) {
if (!msg.hasDeviceInfo()) { if (!msg.hasDeviceInfo()) {
ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED)); ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED, connectMessage));
ctx.close(); ctx.close();
} else { } else {
deviceSessionCtx.setDeviceInfo(msg.getDeviceInfo()); deviceSessionCtx.setDeviceInfo(msg.getDeviceInfo());
@ -640,7 +641,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
public void onSuccess(Void msg) { public void onSuccess(Void msg) {
transportService.registerAsyncSession(deviceSessionCtx.getSessionInfo(), MqttTransportHandler.this); transportService.registerAsyncSession(deviceSessionCtx.getSessionInfo(), MqttTransportHandler.this);
checkGatewaySession(); checkGatewaySession();
ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED)); ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED, connectMessage));
log.info("[{}] Client connected!", sessionId); log.info("[{}] Client connected!", sessionId);
} }
@ -651,7 +652,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
} else { } else {
log.warn("[{}] Failed to submit session event", sessionId, e); log.warn("[{}] Failed to submit session event", sessionId, e);
} }
ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_SERVER_UNAVAILABLE)); ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_SERVER_UNAVAILABLE, connectMessage));
ctx.close(); ctx.close();
} }
}); });

3
dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesDao.java

@ -16,6 +16,7 @@
package org.thingsboard.server.dao.attributes; package org.thingsboard.server.dao.attributes;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.DeviceProfileId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
@ -41,4 +42,6 @@ public interface AttributesDao {
ListenableFuture<List<Void>> removeAll(TenantId tenantId, EntityId entityId, String attributeType, List<String> keys); ListenableFuture<List<Void>> removeAll(TenantId tenantId, EntityId entityId, String attributeType, List<String> keys);
List<String> findAllKeysByDeviceProfileId(TenantId tenantId, DeviceProfileId deviceProfileId); List<String> findAllKeysByDeviceProfileId(TenantId tenantId, DeviceProfileId deviceProfileId);
List<String> findAllKeysByEntityIds(TenantId tenantId, EntityType entityType, List<EntityId> entityIds);
} }

6
dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesService.java

@ -20,6 +20,7 @@ import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.DeviceProfileId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
@ -65,6 +66,11 @@ public class BaseAttributesService implements AttributesService {
return attributesDao.findAllKeysByDeviceProfileId(tenantId, deviceProfileId); return attributesDao.findAllKeysByDeviceProfileId(tenantId, deviceProfileId);
} }
@Override
public List<String> findAllKeysByEntityIds(TenantId tenantId, EntityType entityType, List<EntityId> entityIds) {
return attributesDao.findAllKeysByEntityIds(tenantId, entityType, entityIds);
}
@Override @Override
public ListenableFuture<List<Void>> save(TenantId tenantId, EntityId entityId, String scope, List<AttributeKvEntry> attributes) { public ListenableFuture<List<Void>> save(TenantId tenantId, EntityId entityId, String scope, List<AttributeKvEntry> attributes) {
validate(entityId, scope); validate(entityId, scope);

3
dao/src/main/java/org/thingsboard/server/dao/sql/attributes/AttributeKvRepository.java

@ -56,5 +56,8 @@ public interface AttributeKvRepository extends CrudRepository<AttributeKvEntity,
"AND entity_id in (SELECT id FROM device WHERE tenant_id = :tenantId limit 100) ORDER BY attribute_key", nativeQuery = true) "AND entity_id in (SELECT id FROM device WHERE tenant_id = :tenantId limit 100) ORDER BY attribute_key", nativeQuery = true)
List<String> findAllKeysByTenantId(@Param("tenantId") UUID tenantId); List<String> findAllKeysByTenantId(@Param("tenantId") UUID tenantId);
@Query(value = "SELECT DISTINCT attribute_key FROM attribute_kv WHERE entity_type = :entityType " +
"AND entity_id in :entityIds ORDER BY attribute_key", nativeQuery = true)
List<String> findAllKeysByEntityIds(@Param("entityType") String entityType, @Param("entityIds") List<UUID> entityIds);
} }

7
dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java

@ -22,6 +22,7 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.DeviceProfileId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
@ -145,6 +146,12 @@ public class JpaAttributeDao extends JpaAbstractDaoListeningExecutorService impl
} }
} }
@Override
public List<String> findAllKeysByEntityIds(TenantId tenantId, EntityType entityType, List<EntityId> entityIds) {
return attributeKvRepository
.findAllKeysByEntityIds(entityType.name(), entityIds.stream().map(EntityId::getId).collect(Collectors.toList()));
}
@Override @Override
public ListenableFuture<Void> save(TenantId tenantId, EntityId entityId, String attributeType, AttributeKvEntry attribute) { public ListenableFuture<Void> save(TenantId tenantId, EntityId entityId, String attributeType, AttributeKvEntry attribute) {
AttributeKvEntity entity = new AttributeKvEntity(); AttributeKvEntity entity = new AttributeKvEntity();

6
dao/src/main/java/org/thingsboard/server/dao/sqlts/SqlTimeseriesLatestDao.java

@ -61,6 +61,7 @@ import java.util.Optional;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutionException;
import java.util.function.Function; import java.util.function.Function;
import java.util.stream.Collectors;
@Slf4j @Slf4j
@Component @Component
@ -169,6 +170,11 @@ public class SqlTimeseriesLatestDao extends BaseAbstractSqlTimeseriesDao impleme
} }
} }
@Override
public List<String> findAllKeysByEntityIds(TenantId tenantId, List<EntityId> entityIds) {
return tsKvLatestRepository.findAllKeysByEntityIds(entityIds.stream().map(EntityId::getId).collect(Collectors.toList()));
}
private ListenableFuture<Void> getNewLatestEntryFuture(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query) { private ListenableFuture<Void> getNewLatestEntryFuture(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query) {
ListenableFuture<List<TsKvEntry>> future = findNewLatestEntryFuture(tenantId, entityId, query); ListenableFuture<List<TsKvEntry>> future = findNewLatestEntryFuture(tenantId, entityId, query);
return Futures.transformAsync(future, entryList -> { return Futures.transformAsync(future, entryList -> {

5
dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/TsKvLatestRepository.java

@ -36,4 +36,9 @@ public interface TsKvLatestRepository extends CrudRepository<TsKvLatestEntity, T
"WHERE ts_kv_latest.entity_id IN (SELECT id FROM device WHERE tenant_id = :tenant_id limit 100) ORDER BY ts_kv_dictionary.key", nativeQuery = true) "WHERE ts_kv_latest.entity_id IN (SELECT id FROM device WHERE tenant_id = :tenant_id limit 100) ORDER BY ts_kv_dictionary.key", nativeQuery = true)
List<String> getKeysByTenantId(@Param("tenant_id") UUID tenantId); List<String> getKeysByTenantId(@Param("tenant_id") UUID tenantId);
@Query(value = "SELECT DISTINCT ts_kv_dictionary.key AS strKey FROM ts_kv_latest " +
"INNER JOIN ts_kv_dictionary ON ts_kv_latest.key = ts_kv_dictionary.key_id " +
"WHERE ts_kv_latest.entity_id IN :entityIds ORDER BY ts_kv_dictionary.key", nativeQuery = true)
List<String> findAllKeysByEntityIds(@Param("entityIds") List<UUID> entityIds);
} }

5
dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java

@ -121,6 +121,11 @@ public class BaseTimeseriesService implements TimeseriesService {
return timeseriesLatestDao.findAllKeysByDeviceProfileId(tenantId, deviceProfileId); return timeseriesLatestDao.findAllKeysByDeviceProfileId(tenantId, deviceProfileId);
} }
@Override
public List<String> findAllKeysByEntityIds(TenantId tenantId, List<EntityId> entityIds) {
return timeseriesLatestDao.findAllKeysByEntityIds(tenantId, entityIds);
}
@Override @Override
public ListenableFuture<Integer> save(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry) { public ListenableFuture<Integer> save(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry) {
validate(entityId); validate(entityId);

75
dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java

@ -79,12 +79,17 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
protected static List<Long> FIXED_PARTITION = Arrays.asList(new Long[]{0L}); protected static List<Long> FIXED_PARTITION = Arrays.asList(new Long[]{0L});
private CassandraTsPartitionsCache cassandraTsPartitionsCache;
@Autowired @Autowired
private Environment environment; private Environment environment;
@Value("${cassandra.query.ts_key_value_partitioning}") @Value("${cassandra.query.ts_key_value_partitioning}")
private String partitioning; private String partitioning;
@Value("${cassandra.query.ts_key_value_partitions_max_cache_size:100000}")
private long partitionsCacheSize;
@Value("${cassandra.query.ts_key_value_ttl}") @Value("${cassandra.query.ts_key_value_ttl}")
private long systemTtl; private long systemTtl;
@ -111,13 +116,16 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
super.startExecutor(); super.startExecutor();
if (!isInstall()) { if (!isInstall()) {
getFetchStmt(Aggregation.NONE, DESC_ORDER); getFetchStmt(Aggregation.NONE, DESC_ORDER);
} Optional<NoSqlTsPartitionDate> partition = NoSqlTsPartitionDate.parse(partitioning);
Optional<NoSqlTsPartitionDate> partition = NoSqlTsPartitionDate.parse(partitioning); if (partition.isPresent()) {
if (partition.isPresent()) { tsFormat = partition.get();
tsFormat = partition.get(); if (!isFixedPartitioning() && partitionsCacheSize > 0) {
} else { cassandraTsPartitionsCache = new CassandraTsPartitionsCache(partitionsCacheSize);
log.warn("Incorrect configuration of partitioning {}", partitioning); }
throw new RuntimeException("Failed to parse partitioning property: " + partitioning + "!"); } else {
log.warn("Incorrect configuration of partitioning {}", partitioning);
throw new RuntimeException("Failed to parse partitioning property: " + partitioning + "!");
}
} }
} }
@ -175,17 +183,18 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
} }
ttl = computeTtl(ttl); ttl = computeTtl(ttl);
long partition = toPartitionTs(tsKvEntryTs); long partition = toPartitionTs(tsKvEntryTs);
log.debug("Saving partition {} for the entity [{}-{}] and key {}", partition, entityId.getEntityType(), entityId.getId(), key); if (cassandraTsPartitionsCache == null) {
BoundStatementBuilder stmtBuilder = new BoundStatementBuilder((ttl == 0 ? getPartitionInsertStmt() : getPartitionInsertTtlStmt()).bind()); return doSavePartition(tenantId, entityId, key, ttl, partition);
stmtBuilder.setString(0, entityId.getEntityType().name()) } else {
.setUuid(1, entityId.getId()) CassandraPartitionCacheKey partitionSearchKey = new CassandraPartitionCacheKey(entityId, key, partition);
.setLong(2, partition) if (!cassandraTsPartitionsCache.has(partitionSearchKey)) {
.setString(3, key); ListenableFuture<Integer> result = doSavePartition(tenantId, entityId, key, ttl, partition);
if (ttl > 0) { Futures.addCallback(result, new CacheCallback<>(partitionSearchKey), MoreExecutors.directExecutor());
stmtBuilder.setInt(4, (int) ttl); return result;
} else {
return Futures.immediateFuture(0);
}
} }
BoundStatement stmt = stmtBuilder.build();
return getFuture(executeAsyncWrite(tenantId, stmt), rs -> 0);
} }
@Override @Override
@ -461,6 +470,38 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
return getFuture(executeAsyncWrite(tenantId, stmt), rs -> null); return getFuture(executeAsyncWrite(tenantId, stmt), rs -> null);
} }
private ListenableFuture<Integer> doSavePartition(TenantId tenantId, EntityId entityId, String key, long ttl, long partition) {
log.debug("Saving partition {} for the entity [{}-{}] and key {}", partition, entityId.getEntityType(), entityId.getId(), key);
PreparedStatement preparedStatement = ttl == 0 ? getPartitionInsertStmt() : getPartitionInsertTtlStmt();
BoundStatement stmt = preparedStatement.bind();
stmt = stmt.setString(0, entityId.getEntityType().name())
.setUuid(1, entityId.getId())
.setLong(2, partition)
.setString(3, key);
if (ttl > 0) {
stmt = stmt.setInt(4, (int) ttl);
}
return getFuture(executeAsyncWrite(tenantId, stmt), rs -> 0);
}
private class CacheCallback<Void> implements FutureCallback<Void> {
private final CassandraPartitionCacheKey key;
private CacheCallback(CassandraPartitionCacheKey key) {
this.key = key;
}
@Override
public void onSuccess(Void result) {
cassandraTsPartitionsCache.put(key);
}
@Override
public void onFailure(Throwable t) {
}
}
private long computeTtl(long ttl) { private long computeTtl(long ttl) {
if (systemTtl > 0) { if (systemTtl > 0) {
if (ttl == 0) { if (ttl == 0) {

5
dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java

@ -86,6 +86,11 @@ public class CassandraBaseTimeseriesLatestDao extends AbstractCassandraBaseTimes
return Collections.emptyList(); return Collections.emptyList();
} }
@Override
public List<String> findAllKeysByEntityIds(TenantId tenantId, List<EntityId> entityIds) {
return Collections.emptyList();
}
@Override @Override
public ListenableFuture<Void> saveLatest(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry) { public ListenableFuture<Void> saveLatest(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry) {
BoundStatementBuilder stmtBuilder = new BoundStatementBuilder(getLatestStmt().bind()); BoundStatementBuilder stmtBuilder = new BoundStatementBuilder(getLatestStmt().bind());

30
dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraPartitionCacheKey.java

@ -0,0 +1,30 @@
/**
* Copyright © 2016-2020 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* 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
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.dao.timeseries;
import lombok.AllArgsConstructor;
import lombok.Data;
import org.thingsboard.server.common.data.id.EntityId;
@Data
@AllArgsConstructor
public class CassandraPartitionCacheKey {
private EntityId entityId;
private String key;
private long partition;
}

42
dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraTsPartitionsCache.java

@ -0,0 +1,42 @@
/**
* Copyright © 2016-2020 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* 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
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.dao.timeseries;
import com.github.benmanes.caffeine.cache.AsyncLoadingCache;
import com.github.benmanes.caffeine.cache.Caffeine;
import java.util.concurrent.CompletableFuture;
public class CassandraTsPartitionsCache {
private AsyncLoadingCache<CassandraPartitionCacheKey, Boolean> partitionsCache;
public CassandraTsPartitionsCache(long maxCacheSize) {
this.partitionsCache = Caffeine.newBuilder()
.maximumSize(maxCacheSize)
.buildAsync(key -> {
throw new IllegalStateException("'get' methods calls are not supported!");
});
}
public boolean has(CassandraPartitionCacheKey key) {
return partitionsCache.getIfPresent(key) != null;
}
public void put(CassandraPartitionCacheKey key) {
partitionsCache.put(key, CompletableFuture.completedFuture(true));
}
}

2
dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesLatestDao.java

@ -35,4 +35,6 @@ public interface TimeseriesLatestDao {
ListenableFuture<Void> removeLatest(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query); ListenableFuture<Void> removeLatest(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query);
List<String> findAllKeysByDeviceProfileId(TenantId tenantId, DeviceProfileId deviceProfileId); List<String> findAllKeysByDeviceProfileId(TenantId tenantId, DeviceProfileId deviceProfileId);
List<String> findAllKeysByEntityIds(TenantId tenantId, List<EntityId> entityIds);
} }

110
dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java

@ -0,0 +1,110 @@
/**
* Copyright © 2016-2020 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* 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
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.dao.nosql;
import com.datastax.oss.driver.api.core.ConsistencyLevel;
import com.datastax.oss.driver.api.core.cql.BoundStatement;
import com.datastax.oss.driver.api.core.cql.PreparedStatement;
import com.datastax.oss.driver.api.core.cql.Statement;
import com.google.common.util.concurrent.Futures;
import org.junit.Before;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.mockito.Mock;
import org.mockito.Spy;
import org.mockito.runners.MockitoJUnitRunner;
import org.springframework.core.env.Environment;
import org.springframework.test.util.ReflectionTestUtils;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.dao.cassandra.CassandraCluster;
import org.thingsboard.server.dao.cassandra.guava.GuavaSession;
import org.thingsboard.server.dao.timeseries.CassandraBaseTimeseriesDao;
import java.util.UUID;
import static org.mockito.Matchers.any;
import static org.mockito.Matchers.anyInt;
import static org.mockito.Matchers.anyString;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@RunWith(MockitoJUnitRunner.class)
public class CassandraPartitionsCacheTest {
@Spy
private CassandraBaseTimeseriesDao cassandraBaseTimeseriesDao;
@Mock
private PreparedStatement preparedStatement;
@Mock
private BoundStatement boundStatement;
@Mock
private Environment environment;
@Mock
private CassandraBufferedRateExecutor rateLimiter;
@Mock
private CassandraCluster cluster;
@Mock
private GuavaSession session;
@Before
public void setUp() throws Exception {
ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "partitioning", "MONTHS");
ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "partitionsCacheSize", 100000);
ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "systemTtl", 0);
ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "setNullValuesEnabled", false);
ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "environment", environment);
ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "rateLimiter", rateLimiter);
ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "cluster", cluster);
when(cluster.getDefaultReadConsistencyLevel()).thenReturn(ConsistencyLevel.ONE);
when(cluster.getDefaultWriteConsistencyLevel()).thenReturn(ConsistencyLevel.ONE);
when(cluster.getSession()).thenReturn(session);
when(session.prepare(anyString())).thenReturn(preparedStatement);
when(preparedStatement.bind()).thenReturn(boundStatement);
when(boundStatement.setString(anyInt(), anyString())).thenReturn(boundStatement);
when(boundStatement.setUuid(anyInt(), any(UUID.class))).thenReturn(boundStatement);
when(boundStatement.setLong(anyInt(), any(Long.class))).thenReturn(boundStatement);
doReturn(Futures.immediateFuture(0)).when(cassandraBaseTimeseriesDao).getFuture(any(TbResultSetFuture.class), any());
}
@Test
public void testPartitionSave() throws Exception {
cassandraBaseTimeseriesDao.init();
UUID id = UUID.randomUUID();
TenantId tenantId = new TenantId(id);
long tsKvEntryTs = System.currentTimeMillis();
for (int i = 0; i < 50000; i++) {
cassandraBaseTimeseriesDao.savePartition(tenantId, tenantId, tsKvEntryTs, "test" + i, 0);
}
for (int i = 0; i < 60000; i++) {
cassandraBaseTimeseriesDao.savePartition(tenantId, tenantId, tsKvEntryTs, "test" + i, 0);
}
verify(cassandraBaseTimeseriesDao, times(60000)).executeAsyncWrite(any(TenantId.class), any(Statement.class));
}
}

2
dao/src/test/resources/cassandra-test.properties

@ -54,6 +54,8 @@ cassandra.query.default_fetch_size=2000
cassandra.query.ts_key_value_partitioning=HOURS cassandra.query.ts_key_value_partitioning=HOURS
cassandra.query.ts_key_value_partitions_max_cache_size=100000
cassandra.query.ts_key_value_ttl=0 cassandra.query.ts_key_value_ttl=0
cassandra.query.debug_events_ttl=604800 cassandra.query.debug_events_ttl=604800

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

@ -41,10 +41,16 @@ import { AttributeScope, DataKeyType } from '@shared/models/telemetry/telemetry.
import { defaultHttpOptionsFromConfig, RequestConfig } from '@core/http/http-utils'; import { defaultHttpOptionsFromConfig, RequestConfig } from '@core/http/http-utils';
import { RuleChainService } from '@core/http/rule-chain.service'; import { RuleChainService } from '@core/http/rule-chain.service';
import { AliasInfo, StateParams, SubscriptionInfo } from '@core/api/widget-api.models'; import { AliasInfo, StateParams, SubscriptionInfo } from '@core/api/widget-api.models';
import { Datasource, DatasourceType, KeyInfo } from '@app/shared/models/widget.models'; import { DataKey, Datasource, DatasourceType, KeyInfo } from '@app/shared/models/widget.models';
import { UtilsService } from '@core/services/utils.service'; import { UtilsService } from '@core/services/utils.service';
import { AliasFilterType, EntityAlias, EntityAliasFilter, EntityAliasFilterResult } from '@shared/models/alias.models'; import { AliasFilterType, EntityAlias, EntityAliasFilter, EntityAliasFilterResult } from '@shared/models/alias.models';
import { entityFields, EntityInfo, ImportEntitiesResultInfo, ImportEntityData } from '@shared/models/entity.models'; import {
EntitiesKeysByQuery,
entityFields,
EntityInfo,
ImportEntitiesResultInfo,
ImportEntityData
} from '@shared/models/entity.models';
import { EntityRelationService } from '@core/http/entity-relation.service'; import { EntityRelationService } from '@core/http/entity-relation.service';
import { deepClone, isDefined, isDefinedAndNotNull } from '@core/utils'; import { deepClone, isDefined, isDefinedAndNotNull } from '@core/utils';
import { Asset } from '@shared/models/asset.models'; import { Asset } from '@shared/models/asset.models';
@ -376,6 +382,13 @@ export class EntityService {
return this.http.post<PageData<EntityData>>('/api/entitiesQuery/find', query, defaultHttpOptionsFromConfig(config)); return this.http.post<PageData<EntityData>>('/api/entitiesQuery/find', query, defaultHttpOptionsFromConfig(config));
} }
public findEntityKeysByQuery(query: EntityDataQuery, attributes = true, timeseries = true,
config?: RequestConfig): Observable<EntitiesKeysByQuery> {
return this.http.post<EntitiesKeysByQuery>(
`/api/entitiesQuery/find/keys?attributes=${attributes}&timeseries=${timeseries}`,
query, defaultHttpOptionsFromConfig(config));
}
public findAlarmDataByQuery(query: AlarmDataQuery, config?: RequestConfig): Observable<PageData<AlarmData>> { public findAlarmDataByQuery(query: AlarmDataQuery, config?: RequestConfig): Observable<PageData<AlarmData>> {
return this.http.post<PageData<AlarmData>>('/api/alarmsQuery/find', query, defaultHttpOptionsFromConfig(config)); return this.http.post<PageData<AlarmData>>('/api/alarmsQuery/find', query, defaultHttpOptionsFromConfig(config));
} }
@ -595,7 +608,7 @@ export class EntityService {
return entityTypes; return entityTypes;
} }
private getEntityFieldKeys(entityType: EntityType, searchText: string): Array<string> { private getEntityFieldKeys(entityType: EntityType, searchText: string = ''): Array<string> {
const entityFieldKeys: string[] = [entityFields.createdTime.keyName]; const entityFieldKeys: string[] = [entityFields.createdTime.keyName];
const query = searchText.toLowerCase(); const query = searchText.toLowerCase();
switch (entityType) { switch (entityType) {
@ -637,7 +650,7 @@ export class EntityService {
return query ? entityFieldKeys.filter((entityField) => entityField.toLowerCase().indexOf(query) === 0) : entityFieldKeys; return query ? entityFieldKeys.filter((entityField) => entityField.toLowerCase().indexOf(query) === 0) : entityFieldKeys;
} }
private getAlarmKeys(searchText: string): Array<string> { private getAlarmKeys(searchText: string = ''): Array<string> {
const alarmKeys: string[] = Object.keys(alarmFields); const alarmKeys: string[] = Object.keys(alarmFields);
const query = searchText.toLowerCase(); const query = searchText.toLowerCase();
return query ? alarmKeys.filter((alarmField) => alarmField.toLowerCase().indexOf(query) === 0) : alarmKeys; return query ? alarmKeys.filter((alarmField) => alarmField.toLowerCase().indexOf(query) === 0) : alarmKeys;
@ -672,6 +685,59 @@ export class EntityService {
); );
} }
public getEntityKeysByEntityFilter(filter: EntityFilter, types: DataKeyType[], config?: RequestConfig): Observable<Array<DataKey>> {
if (!types.length) {
return of([]);
}
let entitiesKeysByQuery$: Observable<EntitiesKeysByQuery>;
if (filter !== null && types.some(type => [DataKeyType.timeseries, DataKeyType.attribute].includes(type))) {
const dataQuery = {
entityFilter: filter,
pageLink: createDefaultEntityDataPageLink(100),
};
entitiesKeysByQuery$ = this.findEntityKeysByQuery(dataQuery, types.includes(DataKeyType.attribute),
types.includes(DataKeyType.timeseries), config);
} else {
entitiesKeysByQuery$ = of({
attribute: [],
timeseries: [],
entityTypes: [],
});
}
return entitiesKeysByQuery$.pipe(
map((entitiesKeys) => {
const dataKeys: Array<DataKey> = [];
types.forEach(type => {
let keys: Array<string>;
switch (type) {
case DataKeyType.entityField:
if (entitiesKeys.entityTypes.length) {
const entitiesFields = [];
entitiesKeys.entityTypes.forEach(entityType => entitiesFields.push(...this.getEntityFieldKeys(entityType)));
keys = Array.from(new Set(entitiesFields));
}
break;
case DataKeyType.alarm:
keys = this.getAlarmKeys();
break;
case DataKeyType.attribute:
case DataKeyType.timeseries:
if (entitiesKeys[type].length) {
keys = entitiesKeys[type];
}
break;
}
if (keys) {
dataKeys.push(...keys.map(key => {
return {name: key, type};
}));
}
});
return dataKeys;
})
);
}
public createDatasourcesFromSubscriptionsInfo(subscriptionsInfo: Array<SubscriptionInfo>): Array<Datasource> { public createDatasourcesFromSubscriptionsInfo(subscriptionsInfo: Array<SubscriptionInfo>): Array<Datasource> {
const datasources = subscriptionsInfo.map(subscriptionInfo => this.createDatasourceFromSubscriptionInfo(subscriptionInfo)); const datasources = subscriptionsInfo.map(subscriptionInfo => this.createDatasourceFromSubscriptionInfo(subscriptionInfo));
this.utils.generateColors(datasources); this.utils.generateColors(datasources);

164
ui-ngx/src/app/modules/home/components/profile/alarm/device-profile-alarm.component.html

@ -33,87 +33,91 @@
</button> </button>
</div> </div>
</mat-expansion-panel-header> </mat-expansion-panel-header>
<div fxLayout="column" fxLayoutGap="0.5em"> <ng-template matExpansionPanelContent>
<mat-divider></mat-divider> <div fxLayout="column" fxLayoutGap="0.5em">
<mat-form-field fxFlex floatLabel="always"> <mat-divider></mat-divider>
<mat-label>{{'device-profile.alarm-type' | translate}}</mat-label> <mat-form-field fxFlex floatLabel="always">
<input required matInput formControlName="alarmType" placeholder="Enter alarm type"> <mat-label>{{'device-profile.alarm-type' | translate}}</mat-label>
<mat-error *ngIf="alarmFormGroup.get('alarmType').hasError('required')"> <input required matInput formControlName="alarmType" placeholder="Enter alarm type">
{{ 'device-profile.alarm-type-required' | translate }} <mat-error *ngIf="alarmFormGroup.get('alarmType').hasError('required')">
</mat-error> {{ 'device-profile.alarm-type-required' | translate }}
<mat-error *ngIf="alarmFormGroup.get('alarmType').hasError('unique')"> </mat-error>
{{ 'device-profile.alarm-type-unique' | translate }} <mat-error *ngIf="alarmFormGroup.get('alarmType').hasError('unique')">
</mat-error> {{ 'device-profile.alarm-type-unique' | translate }}
</mat-form-field> </mat-error>
</div>
<mat-expansion-panel class="advanced-settings" [expanded]="false">
<mat-expansion-panel-header>
<mat-panel-title>
<div fxFlex fxLayout="row" fxLayoutAlign="end center">
<div class="tb-small" translate>device-profile.advanced-settings</div>
</div>
</mat-panel-title>
</mat-expansion-panel-header>
<mat-checkbox formControlName="propagate" style="display: block; padding-bottom: 16px;">
{{ 'device-profile.propagate-alarm' | translate }}
</mat-checkbox>
<section *ngIf="alarmFormGroup.get('propagate').value === true" style="padding-bottom: 1em;">
<mat-form-field floatLabel="always" class="mat-block">
<mat-label translate>device-profile.alarm-rule-relation-types-list</mat-label>
<mat-chip-list #relationTypesChipList [disabled]="disabled">
<mat-chip
*ngFor="let key of alarmFormGroup.get('propagateRelationTypes').value;"
(removed)="removeRelationType(key)">
{{key}}
<mat-icon matChipRemove>close</mat-icon>
</mat-chip>
<input matInput type="text" placeholder="{{'device-profile.alarm-rule-relation-types-list' | translate}}"
style="max-width: 200px;"
[matChipInputFor]="relationTypesChipList"
[matChipInputSeparatorKeyCodes]="separatorKeysCodes"
(matChipInputTokenEnd)="addRelationType($event)"
[matChipInputAddOnBlur]="true">
</mat-chip-list>
<mat-hint innerHTML="{{ 'device-profile.alarm-rule-relation-types-list-hint' | translate }}"></mat-hint>
</mat-form-field> </mat-form-field>
</section>
</mat-expansion-panel>
<div fxFlex fxLayout="column">
<div translate class="tb-small" style="padding-bottom: 8px;">device-profile.create-alarm-rules</div>
<tb-create-alarm-rules formControlName="createRules"
style="padding-bottom: 16px;"
[deviceProfileId]="deviceProfileId">
</tb-create-alarm-rules>
<div translate class="tb-small" style="padding-bottom: 8px;">device-profile.clear-alarm-rule</div>
<div fxLayout="row" fxLayoutGap="8px;" fxLayoutAlign="start center"
[fxShow]="alarmFormGroup.get('clearRule').value"
style="padding-bottom: 8px;">
<div class="clear-alarm-rule" fxFlex fxLayout="row">
<tb-alarm-rule formControlName="clearRule" fxFlex [deviceProfileId]="deviceProfileId">
</tb-alarm-rule>
</div>
<button *ngIf="!disabled"
mat-icon-button color="primary" style="min-width: 40px;"
type="button"
(click)="removeClearAlarmRule()"
matTooltip="{{ 'action.remove' | translate }}"
matTooltipPosition="above">
<mat-icon>remove_circle_outline</mat-icon>
</button>
</div> </div>
<div *ngIf="disabled && !alarmFormGroup.get('clearRule').value"> <mat-expansion-panel class="advanced-settings" [expanded]="false">
<span translate fxLayoutAlign="center center" style="margin: 16px 0" <mat-expansion-panel-header>
class="tb-prompt">device-profile.no-clear-alarm-rule</span> <mat-panel-title>
</div> <div fxFlex fxLayout="row" fxLayoutAlign="end center">
<div *ngIf="!disabled" [fxShow]="!alarmFormGroup.get('clearRule').value"> <div class="tb-small" translate>device-profile.advanced-settings</div>
<button mat-stroked-button color="primary" </div>
type="button" </mat-panel-title>
(click)="addClearAlarmRule()" </mat-expansion-panel-header>
matTooltip="{{ 'device-profile.add-clear-alarm-rule' | translate }}" <ng-template matExpansionPanelContent>
matTooltipPosition="above"> <mat-checkbox formControlName="propagate" style="display: block; padding-bottom: 16px;">
<mat-icon class="button-icon">add_circle_outline</mat-icon> {{ 'device-profile.propagate-alarm' | translate }}
{{ 'device-profile.add-clear-alarm-rule' | translate }} </mat-checkbox>
</button> <section *ngIf="alarmFormGroup.get('propagate').value === true" style="padding-bottom: 1em;">
<mat-form-field floatLabel="always" class="mat-block">
<mat-label translate>device-profile.alarm-rule-relation-types-list</mat-label>
<mat-chip-list #relationTypesChipList [disabled]="disabled">
<mat-chip
*ngFor="let key of alarmFormGroup.get('propagateRelationTypes').value;"
(removed)="removeRelationType(key)">
{{key}}
<mat-icon matChipRemove>close</mat-icon>
</mat-chip>
<input matInput type="text" placeholder="{{'device-profile.alarm-rule-relation-types-list' | translate}}"
style="max-width: 200px;"
[matChipInputFor]="relationTypesChipList"
[matChipInputSeparatorKeyCodes]="separatorKeysCodes"
(matChipInputTokenEnd)="addRelationType($event)"
[matChipInputAddOnBlur]="true">
</mat-chip-list>
<mat-hint innerHTML="{{ 'device-profile.alarm-rule-relation-types-list-hint' | translate }}"></mat-hint>
</mat-form-field>
</section>
</ng-template>
</mat-expansion-panel>
<div fxFlex fxLayout="column">
<div translate class="tb-small" style="padding-bottom: 8px;">device-profile.create-alarm-rules</div>
<tb-create-alarm-rules formControlName="createRules"
style="padding-bottom: 16px;"
[deviceProfileId]="deviceProfileId">
</tb-create-alarm-rules>
<div translate class="tb-small" style="padding-bottom: 8px;">device-profile.clear-alarm-rule</div>
<div fxLayout="row" fxLayoutGap="8px;" fxLayoutAlign="start center"
[fxShow]="alarmFormGroup.get('clearRule').value"
style="padding-bottom: 8px;">
<div class="clear-alarm-rule" fxFlex fxLayout="row">
<tb-alarm-rule formControlName="clearRule" fxFlex [deviceProfileId]="deviceProfileId">
</tb-alarm-rule>
</div>
<button *ngIf="!disabled"
mat-icon-button color="primary" style="min-width: 40px;"
type="button"
(click)="removeClearAlarmRule()"
matTooltip="{{ 'action.remove' | translate }}"
matTooltipPosition="above">
<mat-icon>remove_circle_outline</mat-icon>
</button>
</div>
<div *ngIf="disabled && !alarmFormGroup.get('clearRule').value">
<span translate fxLayoutAlign="center center" style="margin: 16px 0"
class="tb-prompt">device-profile.no-clear-alarm-rule</span>
</div>
<div *ngIf="!disabled" [fxShow]="!alarmFormGroup.get('clearRule').value">
<button mat-stroked-button color="primary"
type="button"
(click)="addClearAlarmRule()"
matTooltip="{{ 'device-profile.add-clear-alarm-rule' | translate }}"
matTooltipPosition="above">
<mat-icon class="button-icon">add_circle_outline</mat-icon>
{{ 'device-profile.add-clear-alarm-rule' | translate }}
</button>
</div>
</div> </div>
</div> </ng-template>
</mat-expansion-panel> </mat-expansion-panel>

38
ui-ngx/src/app/modules/home/components/profile/device-profile.component.html

@ -91,10 +91,12 @@
<div translate>device-profile.profile-configuration</div> <div translate>device-profile.profile-configuration</div>
</mat-panel-title> </mat-panel-title>
</mat-expansion-panel-header> </mat-expansion-panel-header>
<tb-device-profile-configuration <ng-template matExpansionPanelContent>
formControlName="configuration" <tb-device-profile-configuration
required> formControlName="configuration"
</tb-device-profile-configuration> required>
</tb-device-profile-configuration>
</ng-template>
</mat-expansion-panel> </mat-expansion-panel>
<mat-expansion-panel *ngIf="displayTransportConfiguration" [expanded]="true"> <mat-expansion-panel *ngIf="displayTransportConfiguration" [expanded]="true">
<mat-expansion-panel-header> <mat-expansion-panel-header>
@ -102,10 +104,12 @@
<div translate>device-profile.transport-configuration</div> <div translate>device-profile.transport-configuration</div>
</mat-panel-title> </mat-panel-title>
</mat-expansion-panel-header> </mat-expansion-panel-header>
<tb-device-profile-transport-configuration <ng-template matExpansionPanelContent>
formControlName="transportConfiguration" <tb-device-profile-transport-configuration
required> formControlName="transportConfiguration"
</tb-device-profile-transport-configuration> required>
</tb-device-profile-transport-configuration>
</ng-template>
</mat-expansion-panel> </mat-expansion-panel>
<mat-expansion-panel [expanded]="false"> <mat-expansion-panel [expanded]="false">
<mat-expansion-panel-header> <mat-expansion-panel-header>
@ -115,10 +119,12 @@
entityForm.get('profileData.alarms').value.length : 0} }}</div> entityForm.get('profileData.alarms').value.length : 0} }}</div>
</mat-panel-title> </mat-panel-title>
</mat-expansion-panel-header> </mat-expansion-panel-header>
<tb-device-profile-alarms <ng-template matExpansionPanelContent>
formControlName="alarms" <tb-device-profile-alarms
[deviceProfileId]="deviceProfileId"> formControlName="alarms"
</tb-device-profile-alarms> [deviceProfileId]="deviceProfileId">
</tb-device-profile-alarms>
</ng-template>
</mat-expansion-panel> </mat-expansion-panel>
<mat-expansion-panel [expanded]="true"> <mat-expansion-panel [expanded]="true">
<mat-expansion-panel-header> <mat-expansion-panel-header>
@ -126,9 +132,11 @@
<div translate>device-profile.device-provisioning</div> <div translate>device-profile.device-provisioning</div>
</mat-panel-title> </mat-panel-title>
</mat-expansion-panel-header> </mat-expansion-panel-header>
<tb-device-profile-provision-configuration <ng-template matExpansionPanelContent>
formControlName="provisionConfiguration"> <tb-device-profile-provision-configuration
</tb-device-profile-provision-configuration> formControlName="provisionConfiguration">
</tb-device-profile-provision-configuration>
</ng-template>
</mat-expansion-panel> </mat-expansion-panel>
</mat-accordion> </mat-accordion>
</div> </div>

10
ui-ngx/src/app/modules/home/components/profile/tenant-profile-data.component.html

@ -22,9 +22,11 @@
<div translate>tenant-profile.profile-configuration</div> <div translate>tenant-profile.profile-configuration</div>
</mat-panel-title> </mat-panel-title>
</mat-expansion-panel-header> </mat-expansion-panel-header>
<tb-tenant-profile-configuration <ng-template matExpansionPanelContent>
formControlName="configuration" <tb-tenant-profile-configuration
required> formControlName="configuration"
</tb-tenant-profile-configuration> required>
</tb-tenant-profile-configuration>
</ng-template>
</mat-expansion-panel> </mat-expansion-panel>
</form> </form>

34
ui-ngx/src/app/modules/home/components/widget/data-key-config.component.ts

@ -36,7 +36,7 @@ import { EntityService } from '@core/http/entity.service';
import { DataKeysCallbacks } from '@home/components/widget/data-keys.component.models'; import { DataKeysCallbacks } from '@home/components/widget/data-keys.component.models';
import { DataKeyType } from '@shared/models/telemetry/telemetry.models'; import { DataKeyType } from '@shared/models/telemetry/telemetry.models';
import { Observable, of } from 'rxjs'; import { Observable, of } from 'rxjs';
import { map, mergeMap, tap } from 'rxjs/operators'; import { map, mergeMap, publishReplay, refCount, tap } from 'rxjs/operators';
import { alarmFields } from '@shared/models/alarm.models'; import { alarmFields } from '@shared/models/alarm.models';
import { JsFuncComponent } from '@shared/components/js-func.component'; import { JsFuncComponent } from '@shared/components/js-func.component';
import { JsonFormComponentData } from '@shared/components/json-form/json-form-component.models'; import { JsonFormComponentData } from '@shared/components/json-form/json-form-component.models';
@ -95,6 +95,7 @@ export class DataKeyConfigComponent extends PageComponent implements OnInit, Con
filteredKeys: Observable<Array<string>>; filteredKeys: Observable<Array<string>>;
private latestKeySearchResult: Array<string> = null; private latestKeySearchResult: Array<string> = null;
private fetchObservable$: Observable<Array<string>> = null;
keySearchText = ''; keySearchText = '';
@ -205,31 +206,42 @@ export class DataKeyConfigComponent extends PageComponent implements OnInit, Con
} }
private fetchKeys(searchText?: string): Observable<Array<string>> { private fetchKeys(searchText?: string): Observable<Array<string>> {
if (this.latestKeySearchResult === null || this.keySearchText !== searchText) { if (this.keySearchText !== searchText || this.latestKeySearchResult === null) {
this.keySearchText = searchText; this.keySearchText = searchText;
let fetchObservable: Observable<Array<DataKey>> = null; const dataKeyFilter = this.createKeyFilter(this.keySearchText);
return this.getKeys().pipe(
map(name => name.filter(dataKeyFilter)),
tap(res => this.latestKeySearchResult = res)
);
}
return of(this.latestKeySearchResult);
}
private getKeys() {
if (this.fetchObservable$ === null) {
let fetchObservable: Observable<Array<DataKey>>;
if (this.modelValue.type === DataKeyType.alarm) { if (this.modelValue.type === DataKeyType.alarm) {
const dataKeyFilter = this.createDataKeyFilter(this.keySearchText); fetchObservable = of(this.alarmKeys);
fetchObservable = of(this.alarmKeys.filter(dataKeyFilter));
} else { } else {
if (this.entityAliasId) { if (this.entityAliasId) {
const dataKeyTypes = [this.modelValue.type]; const dataKeyTypes = [this.modelValue.type];
fetchObservable = this.callbacks.fetchEntityKeys(this.entityAliasId, this.keySearchText, dataKeyTypes); fetchObservable = this.callbacks.fetchEntityKeys(this.entityAliasId, dataKeyTypes);
} else { } else {
fetchObservable = of([]); fetchObservable = of([]);
} }
} }
return fetchObservable.pipe( this.fetchObservable$ = fetchObservable.pipe(
map((dataKeys) => dataKeys.map((dataKey) => dataKey.name)), map((dataKeys) => dataKeys.map((dataKey) => dataKey.name)),
tap(res => this.latestKeySearchResult = res) publishReplay(1),
refCount()
); );
} }
return of(this.latestKeySearchResult); return this.fetchObservable$;
} }
private createDataKeyFilter(query: string): (key: DataKey) => boolean { private createKeyFilter(query: string): (key: string) => boolean {
const lowercaseQuery = query.toLowerCase(); const lowercaseQuery = query.toLowerCase();
return key => key.name.toLowerCase().indexOf(lowercaseQuery) === 0; return key => key.toLowerCase().startsWith(lowercaseQuery);
} }
public validateOnSubmit() { public validateOnSubmit() {

2
ui-ngx/src/app/modules/home/components/widget/data-keys.component.models.ts

@ -20,5 +20,5 @@ import { Observable } from 'rxjs';
export interface DataKeysCallbacks { export interface DataKeysCallbacks {
generateDataKey: (chip: any, type: DataKeyType) => DataKey; generateDataKey: (chip: any, type: DataKeyType) => DataKey;
fetchEntityKeys: (entityAliasId: string, query: string, types: Array<DataKeyType>) => Observable<Array<DataKey>>; fetchEntityKeys: (entityAliasId: string, types: Array<DataKeyType>) => Observable<Array<DataKey>>;
} }

37
ui-ngx/src/app/modules/home/components/widget/data-keys.component.ts

@ -38,7 +38,7 @@ import {
Validators Validators
} from '@angular/forms'; } from '@angular/forms';
import { Observable, of } from 'rxjs'; import { Observable, of } from 'rxjs';
import { filter, map, mergeMap, share, tap } from 'rxjs/operators'; import { filter, map, mergeMap, publishReplay, refCount, share, tap } from 'rxjs/operators';
import { Store } from '@ngrx/store'; import { Store } from '@ngrx/store';
import { AppState } from '@app/core/core.state'; import { AppState } from '@app/core/core.state';
import { TranslateService } from '@ngx-translate/core'; import { TranslateService } from '@ngx-translate/core';
@ -142,6 +142,7 @@ export class DataKeysComponent implements ControlValueAccessor, OnInit, AfterVie
searchText = ''; searchText = '';
private latestSearchTextResult: Array<DataKey> = null; private latestSearchTextResult: Array<DataKey> = null;
private fetchObservable$: Observable<Array<DataKey>> = null;
private dirty = false; private dirty = false;
@ -260,6 +261,7 @@ export class DataKeysComponent implements ControlValueAccessor, OnInit, AfterVie
if (!change.firstChange && change.currentValue !== change.previousValue) { if (!change.firstChange && change.currentValue !== change.previousValue) {
if (propName === 'entityAliasId') { if (propName === 'entityAliasId') {
this.searchText = ''; this.searchText = '';
this.fetchObservable$ = null;
this.latestSearchTextResult = null; this.latestSearchTextResult = null;
this.dirty = true; this.dirty = true;
} else if (['widgetType', 'datasourceType'].includes(propName)) { } else if (['widgetType', 'datasourceType'].includes(propName)) {
@ -405,14 +407,24 @@ export class DataKeysComponent implements ControlValueAccessor, OnInit, AfterVie
return key ? key.name : undefined; return key ? key.name : undefined;
} }
fetchKeys(searchText?: string): Observable<Array<DataKey>> { private fetchKeys(searchText?: string): Observable<Array<DataKey>> {
if (this.latestSearchTextResult === null || this.searchText !== searchText) { if (this.searchText !== searchText || this.latestSearchTextResult === null) {
this.searchText = searchText; this.searchText = searchText;
let fetchObservable: Observable<Array<DataKey>> = null; const dataKeyFilter = this.createDataKeyFilter(this.searchText);
return this.getKeys().pipe(
map(name => name.filter(dataKeyFilter)),
tap(res => this.latestSearchTextResult = res)
);
}
return of(this.latestSearchTextResult);
}
private getKeys(): Observable<Array<DataKey>> {
if (this.fetchObservable$ === null) {
let fetchObservable: Observable<Array<DataKey>>;
if (this.datasourceType === DatasourceType.function) { if (this.datasourceType === DatasourceType.function) {
const dataKeyFilter = this.createDataKeyFilter(this.searchText);
const targetKeysList = this.widgetType === widgetType.alarm ? this.alarmKeys : this.functionTypeKeys; const targetKeysList = this.widgetType === widgetType.alarm ? this.alarmKeys : this.functionTypeKeys;
fetchObservable = of(targetKeysList.filter(dataKeyFilter)); fetchObservable = of(targetKeysList);
} else { } else {
if (this.entityAliasId) { if (this.entityAliasId) {
const dataKeyTypes = [DataKeyType.timeseries]; const dataKeyTypes = [DataKeyType.timeseries];
@ -420,24 +432,25 @@ export class DataKeysComponent implements ControlValueAccessor, OnInit, AfterVie
dataKeyTypes.push(DataKeyType.attribute); dataKeyTypes.push(DataKeyType.attribute);
dataKeyTypes.push(DataKeyType.entityField); dataKeyTypes.push(DataKeyType.entityField);
if (this.widgetType === widgetType.alarm) { if (this.widgetType === widgetType.alarm) {
dataKeyTypes.push(DataKeyType.alarm); dataKeyTypes.push(DataKeyType.alarm);
} }
} }
fetchObservable = this.callbacks.fetchEntityKeys(this.entityAliasId, this.searchText, dataKeyTypes); fetchObservable = this.callbacks.fetchEntityKeys(this.entityAliasId, dataKeyTypes);
} else { } else {
fetchObservable = of([]); fetchObservable = of([]);
} }
} }
return fetchObservable.pipe( this.fetchObservable$ = fetchObservable.pipe(
tap(res => this.latestSearchTextResult = res) publishReplay(1),
refCount()
); );
} }
return of(this.latestSearchTextResult); return this.fetchObservable$;
} }
private createDataKeyFilter(query: string): (key: DataKey) => boolean { private createDataKeyFilter(query: string): (key: DataKey) => boolean {
const lowercaseQuery = query.toLowerCase(); const lowercaseQuery = query.toLowerCase();
return key => key.name.toLowerCase().indexOf(lowercaseQuery) === 0; return key => key.name.toLowerCase().startsWith(lowercaseQuery);
} }
textIsNotEmpty(text: string): boolean { textIsNotEmpty(text: string): boolean {

34
ui-ngx/src/app/modules/home/components/widget/lib/maps/leaflet-map.ts

@ -608,26 +608,32 @@ export default abstract class LeafletMap {
return polygon; return polygon;
} }
updatePoints(pointsData: FormattedData[], getTooltip: (point: FormattedData, setTooltip?: boolean) => string) { updatePoints(pointsData: FormattedData[][], getTooltip: (point: FormattedData) => string) {
if(pointsData.length) {
if (this.points) { if (this.points) {
this.map.removeLayer(this.points); this.map.removeLayer(this.points);
} }
this.points = new FeatureGroup(); this.points = new FeatureGroup();
pointsData.filter(pdata => !!this.convertPosition(pdata)).forEach(data => { }
const point = L.circleMarker(this.convertPosition(data), { for(let i = 0; i < pointsData.length; i++) {
color: this.options.pointColor, const pointsList = pointsData[i];
radius: this.options.pointSize pointsList.filter(pdata => !!this.convertPosition(pdata)).forEach(data => {
}); const point = L.circleMarker(this.convertPosition(data), {
if (!this.options.pointTooltipOnRightPanel) { color: this.options.pointColor,
point.on('click', () => getTooltip(data)); radius: this.options.pointSize
} });
else { if (!this.options.pointTooltipOnRightPanel) {
createTooltip(point, this.options, data.$datasource, getTooltip(data, false)); point.on('click', () => getTooltip(data));
} } else {
this.points.addLayer(point); createTooltip(point, this.options, data.$datasource, getTooltip(data));
}
this.points.addLayer(point);
}); });
}
if(pointsData.length) {
this.map.addLayer(this.points); this.map.addLayer(this.points);
} }
}
// Polyline // Polyline

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

@ -28,8 +28,12 @@
</button> </button>
</div> </div>
<div class="trip-animation-tooltip md-whiteframe-z4" fxLayout="column" <div class="trip-animation-tooltip md-whiteframe-z4" fxLayout="column"
[ngClass]="{'trip-animation-tooltip-hidden':!visibleTooltip}" [innerHTML]="mainTooltip" [ngClass]="{'trip-animation-tooltip-hidden':!visibleTooltip}"
[ngStyle]="{'background-color': settings.tooltipColor, 'opacity': settings.tooltipOpacity, 'color': settings.tooltipFontColor}"> [ngStyle]="{'background-color': settings.tooltipColor, 'opacity': settings.tooltipOpacity, 'color': settings.tooltipFontColor}">
<div *ngFor="let mainTooltip of mainTooltips"
[innerHTML]="mainTooltip"
style="padding: 10px 0">
</div>
</div> </div>
</div> </div>
<tb-history-selector *ngIf="historicalData" <tb-history-selector *ngIf="historicalData"

5
ui-ngx/src/app/modules/home/components/widget/trip-animation/trip-animation.component.scss

@ -73,6 +73,9 @@
} }
.trip-animation-tooltip { .trip-animation-tooltip {
display: flex;
overflow: auto;
max-height: 90%;
position: absolute; position: absolute;
top: 30px; top: 30px;
right: 0; right: 0;
@ -86,4 +89,4 @@
} }
} }
} }
} }

80
ui-ngx/src/app/modules/home/components/widget/trip-animation/trip-animation.component.ts

@ -47,6 +47,9 @@ import moment from 'moment';
import { isUndefined } from '@core/utils'; import { isUndefined } from '@core/utils';
import { ResizeObserver } from '@juggle/resize-observer'; import { ResizeObserver } from '@juggle/resize-observer';
interface dataMap {
[key: string] : FormattedData
}
@Component({ @Component({
// tslint:disable-next-line:component-selector // tslint:disable-next-line:component-selector
@ -70,7 +73,7 @@ export class TripAnimationComponent implements OnInit, AfterViewInit, OnDestroy
interpolatedTimeData = []; interpolatedTimeData = [];
widgetConfig: WidgetConfig; widgetConfig: WidgetConfig;
settings: TripAnimationSettings; settings: TripAnimationSettings;
mainTooltip = ''; mainTooltips = [];
visibleTooltip = false; visibleTooltip = false;
activeTrip: FormattedData; activeTrip: FormattedData;
label: string; label: string;
@ -115,7 +118,7 @@ export class TripAnimationComponent implements OnInit, AfterViewInit, OnDestroy
this.historicalData = parseArray(this.ctx.data).filter(arr => arr.length); this.historicalData = parseArray(this.ctx.data).filter(arr => arr.length);
if (this.historicalData.length) { if (this.historicalData.length) {
this.calculateIntervals(); this.calculateIntervals();
this.timeUpdated(this.currentTime && this.currentTime > this.minTime ? this.currentTime : this.minTime); this.timeUpdated(this.minTime);
} }
this.mapWidget.map.map?.invalidateSize(); this.mapWidget.map.map?.invalidateSize();
this.cd.detectChanges(); this.cd.detectChanges();
@ -140,32 +143,39 @@ export class TripAnimationComponent implements OnInit, AfterViewInit, OnDestroy
this.currentTime = time; this.currentTime = time;
const currentPosition = this.interpolatedTimeData const currentPosition = this.interpolatedTimeData
.map(dataSource => dataSource[time]) .map(dataSource => dataSource[time])
.filter(ds => ds); for(let j = 0; j < this.interpolatedTimeData.length; j++) {
if (isUndefined(currentPosition[0])) { if (isUndefined(currentPosition[j])) {
const timePoints = Object.keys(this.interpolatedTimeData[0]).map(item => parseInt(item, 10)); const timePoints = Object.keys(this.interpolatedTimeData[j]).map(item => parseInt(item, 10));
for (let i = 1; i < timePoints.length; i++) { for (let i = 1; i < timePoints.length; i++) {
if (timePoints[i - 1] < time && timePoints[i] > time) { if (timePoints[i - 1] < time && timePoints[i] > time) {
const beforePosition = this.interpolatedTimeData[0][timePoints[i - 1]]; const beforePosition = this.interpolatedTimeData[j][timePoints[i - 1]];
const afterPosition = this.interpolatedTimeData[0][timePoints[i]]; const afterPosition = this.interpolatedTimeData[j][timePoints[i]];
const ratio = getRatio(timePoints[i - 1], timePoints[i], time); const ratio = getRatio(timePoints[i - 1], timePoints[i], time);
currentPosition[0] = { currentPosition[j] = {
...beforePosition, ...beforePosition,
time, time,
...interpolateOnLineSegment(beforePosition, afterPosition, this.settings.latKeyName, this.settings.lngKeyName, ratio) ...interpolateOnLineSegment(beforePosition, afterPosition, this.settings.latKeyName, this.settings.lngKeyName, ratio)
}
break;
} }
break;
} }
} }
} }
for(let j = 0; j < this.interpolatedTimeData.length; j++) {
if (isUndefined(currentPosition[j])) {
currentPosition[j] = this.calculateLastPoints(this.interpolatedTimeData[j], time);
}
}
this.calcLabel(); this.calcLabel();
this.calcTooltip(currentPosition.find(position => position.entityName === this.activeTrip.entityName)); this.calcMainTooltip(currentPosition);
if (this.mapWidget && this.mapWidget.map && this.mapWidget.map.map) { if (this.mapWidget && this.mapWidget.map && this.mapWidget.map.map) {
this.mapWidget.map.updatePolylines(this.interpolatedTimeData.map(ds => _.values(ds)), true, this.activeTrip); const formattedInterpolatedTimeData = this.interpolatedTimeData.map(ds => _.values(ds));
this.mapWidget.map.updatePolylines(formattedInterpolatedTimeData, true);
if (this.settings.showPolygon) { if (this.settings.showPolygon) {
this.mapWidget.map.updatePolygons(this.interpolatedTimeData); this.mapWidget.map.updatePolygons(this.interpolatedTimeData);
} }
if (this.settings.showPoints) { if (this.settings.showPoints) {
this.mapWidget.map.updatePoints(_.values(_.union(this.interpolatedTimeData)[0]), this.calcTooltip); this.mapWidget.map.updatePoints(formattedInterpolatedTimeData.map(ds => _.union(ds)), this.calcTooltip);
} }
this.mapWidget.map.updateMarkers(currentPosition, true, (trip) => { this.mapWidget.map.updateMarkers(currentPosition, true, (trip) => {
this.activeTrip = trip; this.activeTrip = trip;
@ -177,6 +187,23 @@ export class TripAnimationComponent implements OnInit, AfterViewInit, OnDestroy
setActiveTrip() { setActiveTrip() {
} }
private calculateLastPoints(dataSource: dataMap, time: number): FormattedData {
const timeArr = Object.keys(dataSource);
let index = timeArr.findIndex((dtime, index) => {
return Number(dtime) >= time;
});
if(index !== -1) {
if(Number(timeArr[index]) !== time && index !== 0) {
index--;
}
} else {
index = timeArr.length - 1;
}
return dataSource[timeArr[index]];
}
calculateIntervals() { calculateIntervals() {
this.historicalData.forEach((dataSource, index) => { this.historicalData.forEach((dataSource, index) => {
this.minTime = dataSource[0]?.time || Infinity; this.minTime = dataSource[0]?.time || Infinity;
@ -194,16 +221,19 @@ export class TripAnimationComponent implements OnInit, AfterViewInit, OnDestroy
} }
} }
calcTooltip = (point?: FormattedData): string => { calcTooltip = (point: FormattedData): string => {
const data = point ? point : this.activeTrip; const data = point ? point : this.activeTrip;
const tooltipPattern: string = this.settings.useTooltipFunction ? const tooltipPattern: string = this.settings.useTooltipFunction ?
safeExecute(this.settings.tooltipFunction, [data, this.historicalData, point.dsIndex]) : this.settings.tooltipPattern; safeExecute(this.settings.tooltipFunction, [data, this.historicalData, point.dsIndex]) : this.settings.tooltipPattern;
const tooltipText = parseWithTranslation.parseTemplate(tooltipPattern, data, true); return parseWithTranslation.parseTemplate(tooltipPattern, data, true);
this.mainTooltip = this.sanitizer.sanitize( }
SecurityContext.HTML, tooltipText);
this.cd.detectChanges(); private calcMainTooltip(points: FormattedData[]): void {
this.activeTrip = point; const tooltips = [];
return tooltipText; for (let point of points) {
tooltips.push(this.sanitizer.sanitize(SecurityContext.HTML, this.calcTooltip(point)));
}
this.mainTooltips = tooltips;
} }
calcLabel() { calcLabel() {

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

@ -54,13 +54,13 @@ import { UtilsService } from '@core/services/utils.service';
import { DataKeyType } from '@shared/models/telemetry/telemetry.models'; import { DataKeyType } from '@shared/models/telemetry/telemetry.models';
import { TranslateService } from '@ngx-translate/core'; import { TranslateService } from '@ngx-translate/core';
import { EntityType } from '@shared/models/entity-type.models'; import { EntityType } from '@shared/models/entity-type.models';
import { forkJoin, Observable, of, Subscription } from 'rxjs'; import { Observable, of, Subscription } from 'rxjs';
import { WidgetConfigCallbacks } from '@home/components/widget/widget-config.component.models'; import { WidgetConfigCallbacks } from '@home/components/widget/widget-config.component.models';
import { import {
EntityAliasDialogComponent, EntityAliasDialogComponent,
EntityAliasDialogData EntityAliasDialogData
} from '@home/components/alias/entity-alias-dialog.component'; } from '@home/components/alias/entity-alias-dialog.component';
import { catchError, map, mergeMap, tap } from 'rxjs/operators'; import { catchError, mergeMap, tap } from 'rxjs/operators';
import { MatDialog } from '@angular/material/dialog'; import { MatDialog } from '@angular/material/dialog';
import { EntityService } from '@core/http/entity.service'; import { EntityService } from '@core/http/entity.service';
import { JsonFormComponentData } from '@shared/components/json-form/json-form-component.models'; import { JsonFormComponentData } from '@shared/components/json-form/json-form-component.models';
@ -792,54 +792,16 @@ export class WidgetConfigComponent extends PageComponent implements OnInit, Cont
); );
} }
private fetchEntityKeys(entityAliasId: string, query: string, dataKeyTypes: Array<DataKeyType>): Observable<Array<DataKey>> { private fetchEntityKeys(entityAliasId: string, dataKeyTypes: Array<DataKeyType>): Observable<Array<DataKey>> {
return this.aliasController.resolveSingleEntityInfo(entityAliasId).pipe( return this.aliasController.getAliasInfo(entityAliasId).pipe(
mergeMap((entity) => { mergeMap((aliasInfo) => {
if (entity) { return this.entityService.getEntityKeysByEntityFilter(
const fetchEntityTasks: Array<Observable<Array<DataKey>>> = []; aliasInfo.entityFilter,
for (const dataKeyType of dataKeyTypes) { dataKeyTypes,
fetchEntityTasks.push( {ignoreLoading: true, ignoreErrors: true}
this.entityService.getEntityKeys( ).pipe(
{entityType: entity.entityType, id: entity.id}, catchError(() => of([]))
query, );
dataKeyType,
{ignoreLoading: true, ignoreErrors: true}
).pipe(
map((keys) => {
const dataKeys: Array<DataKey> = [];
for (const key of keys) {
dataKeys.push({name: key, type: dataKeyType});
}
return dataKeys;
}
),
catchError(() => of([]))
));
}
return forkJoin(fetchEntityTasks).pipe(
map(arrayOfDataKeys => {
const result = new Array<DataKey>();
arrayOfDataKeys.forEach((dataKeyArray) => {
result.push(...dataKeyArray);
});
return result;
}
));
} else if (dataKeyTypes.includes(DataKeyType.alarm)) {
return this.entityService.getEntityKeys(null, query, DataKeyType.alarm).pipe(
map((keys) => {
const dataKeys: Array<DataKey> = [];
for (const key of keys) {
dataKeys.push({name: key, type: DataKeyType.alarm});
}
return dataKeys;
}
),
catchError(() => of([]))
);
} else {
return of([]);
}
}), }),
catchError(() => of([] as Array<DataKey>)) catchError(() => of([] as Array<DataKey>))
); );

6
ui-ngx/src/app/shared/models/entity.models.ts

@ -64,6 +64,12 @@ export interface EntityField {
time?: boolean; time?: boolean;
} }
export interface EntitiesKeysByQuery {
attribute: Array<string>;
timeseries: Array<string>;
entityTypes: EntityType[];
}
export const entityFields: {[fieldName: string]: EntityField} = { export const entityFields: {[fieldName: string]: EntityField} = {
createdTime: { createdTime: {
keyName: 'createdTime', keyName: 'createdTime',

Loading…
Cancel
Save