diff --git a/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java b/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java index 94417886d6..2df8298d23 100644 --- a/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java +++ b/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java @@ -16,18 +16,23 @@ package org.thingsboard.server.controller; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.http.ResponseEntity; import org.springframework.security.access.prepost.PreAuthorize; import org.springframework.web.bind.annotation.RequestBody; import org.springframework.web.bind.annotation.RequestMapping; 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.RestController; +import org.springframework.web.context.request.async.DeferredResult; 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.query.AlarmData; import org.thingsboard.server.common.data.query.AlarmDataQuery; import org.thingsboard.server.common.data.query.EntityCountQuery; 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.queue.util.TbCoreComponent; import org.thingsboard.server.service.query.EntityQueryService; @@ -40,6 +45,7 @@ public class EntityQueryController extends BaseController { @Autowired private EntityQueryService entityQueryService; + private static final int MAX_PAGE_SIZE = 100; @PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") @RequestMapping(value = "/entitiesQuery/count", method = RequestMethod.POST) @@ -76,4 +82,24 @@ public class EntityQueryController extends BaseController { throw handleException(e); } } + + @PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") + @RequestMapping(value = "/entitiesQuery/find/keys", method = RequestMethod.POST) + @ResponseBody + public DeferredResult 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); + } + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java b/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java index b8ccad466d..c9726693bf 100644 --- a/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java +++ b/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java @@ -15,11 +15,23 @@ */ 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 org.checkerframework.checker.nullness.qual.Nullable; import org.springframework.beans.factory.annotation.Autowired; 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.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.TenantId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.query.AlarmData; 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.EntityKeyType; 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.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.service.executors.DbCallbackExecutorService; +import org.thingsboard.server.service.security.AccessValidator; 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.List; +import java.util.Map; +import java.util.Set; +import java.util.function.Consumer; +import java.util.stream.Collectors; @Service @Slf4j @@ -52,6 +77,15 @@ public class DefaultEntityQueryService implements EntityQueryService { @Value("${server.ws.max_entities_per_alarm_subscription:1000}") private int maxEntitiesPerAlarmSubscription; + @Autowired + private DbCallbackExecutorService dbCallbackExecutor; + + @Autowired + private TimeseriesService timeseriesService; + + @Autowired + private AttributesService attributesService; + @Override public long countEntitiesByQuery(SecurityUser securityUser, EntityCountQuery 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); return new EntityDataQuery(query.getEntityFilter(), edpl, query.getEntityFields(), query.getLatestValues(), query.getKeyFilters()); } + + @Override + public DeferredResult getKeysByQuery(SecurityUser securityUser, TenantId tenantId, EntityDataQuery query, + boolean isTimeseries, boolean isAttributes) { + final DeferredResult response = new DeferredResult<>(); + if (!isAttributes && !isTimeseries) { + replyWithEmptyResponse(response); + return response; + } + + List ids = this.findEntityDataByQuery(securityUser, query).getData().stream() + .map(EntityData::getEntityId) + .collect(Collectors.toList()); + if (ids.isEmpty()) { + replyWithEmptyResponse(response); + return response; + } + + Set types = ids.stream().map(EntityId::getEntityType).collect(Collectors.toSet()); + final ListenableFuture> timeseriesKeysFuture; + final ListenableFuture> attributesKeysFuture; + + if (isTimeseries) { + timeseriesKeysFuture = dbCallbackExecutor.submit(() -> timeseriesService.findAllKeysByEntityIds(tenantId, ids)); + } else { + timeseriesKeysFuture = null; + } + + if (isAttributes) { + Map> typesMap = ids.stream().collect(Collectors.groupingBy(EntityId::getEntityType)); + List>> 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 response, Set types, List timeseriesKeys, List 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 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> future, Consumer> success, Consumer error) { + Futures.addCallback(future, new FutureCallback>() { + @Override + public void onSuccess(@Nullable List keys) { + success.accept(keys); + } + + @Override + public void onFailure(Throwable t) { + error.accept(t); + } + }, dbCallbackExecutor); + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/query/EntityQueryService.java b/application/src/main/java/org/thingsboard/server/service/query/EntityQueryService.java index 15f7d86252..763453c815 100644 --- a/application/src/main/java/org/thingsboard/server/service/query/EntityQueryService.java +++ b/application/src/main/java/org/thingsboard/server/service/query/EntityQueryService.java @@ -15,6 +15,9 @@ */ 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.query.AlarmData; import org.thingsboard.server.common.data.query.AlarmDataQuery; @@ -31,4 +34,7 @@ public interface EntityQueryService { PageData findAlarmDataByQuery(SecurityUser securityUser, AlarmDataQuery query); + DeferredResult getKeysByQuery(SecurityUser securityUser, TenantId tenantId, EntityDataQuery query, + boolean isTimeseries, boolean isAttributes); + } diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 791b032352..4a7c58e4db 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -192,8 +192,9 @@ cassandra: read_consistency_level: "${CASSANDRA_READ_CONSISTENCY_LEVEL:ONE}" write_consistency_level: "${CASSANDRA_WRITE_CONSISTENCY_LEVEL:ONE}" 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_partitions_max_cache_size: "${TS_KV_PARTITIONS_MAX_CACHE_SIZE:100000}" ts_key_value_ttl: "${TS_KV_TTL:0}" events_ttl: "${TS_EVENTS_TTL:0}" # Specify TTL of debug log in seconds. The current value corresponds to one week diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/attributes/AttributesService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/attributes/AttributesService.java index 2e5c895f84..47f1f91376 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/attributes/AttributesService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/attributes/AttributesService.java @@ -16,6 +16,7 @@ package org.thingsboard.server.dao.attributes; 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.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -42,4 +43,6 @@ public interface AttributesService { List findAllKeysByDeviceProfileId(TenantId tenantId, DeviceProfileId deviceProfileId); + List findAllKeysByEntityIds(TenantId tenantId, EntityType entityType, List entityIds); + } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java index dbfbe3d0f5..ccf933e03d 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java @@ -50,4 +50,6 @@ public interface TimeseriesService { ListenableFuture> removeAllLatest(TenantId tenantId, EntityId entityId); List findAllKeysByDeviceProfileId(TenantId tenantId, DeviceProfileId deviceProfileId); + + List findAllKeysByEntityIds(TenantId tenantId, List entityIds); } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index 0df14d4655..59a78570ad 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/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(); if (DataConstants.PROVISION.equals(userName) || DataConstants.PROVISION.equals(clientId)) { deviceSessionCtx.setProvisionOnly(true); - ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED)); + ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED, msg)); } else { X509Certificate cert; if (sslHandler != null && (cert = getX509Certificate()) != null) { - processX509CertConnect(ctx, cert); + processX509CertConnect(ctx, cert, msg); } else { processAuthTokenConnect(ctx, msg); } } } - private void processAuthTokenConnect(ChannelHandlerContext ctx, MqttConnectMessage msg) { - String userName = msg.payload().userName(); + private void processAuthTokenConnect(ChannelHandlerContext ctx, MqttConnectMessage connectMessage) { + String userName = connectMessage.payload().userName(); log.info("[{}] Processing connect msg for client with user name: {}!", sessionId, userName); TransportProtos.ValidateBasicMqttCredRequestMsg.Builder request = TransportProtos.ValidateBasicMqttCredRequestMsg.newBuilder() - .setClientId(msg.payload().clientIdentifier()); + .setClientId(connectMessage.payload().clientIdentifier()); if (userName != null) { request.setUserName(userName); } - String password = msg.payload().password(); + String password = connectMessage.payload().password(); if (password != null) { request.setPassword(password); } @@ -489,19 +489,19 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement new TransportServiceCallback() { @Override public void onSuccess(ValidateDeviceCredentialsResponse msg) { - onValidateDeviceResponse(msg, ctx); + onValidateDeviceResponse(msg, ctx, connectMessage); } @Override public void onError(Throwable 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(); } }); } - private void processX509CertConnect(ChannelHandlerContext ctx, X509Certificate cert) { + private void processX509CertConnect(ChannelHandlerContext ctx, X509Certificate cert, MqttConnectMessage connectMessage) { try { if (!context.isSkipValidityCheckForClientCert()) { cert.checkValidity(); @@ -512,18 +512,18 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement new TransportServiceCallback() { @Override public void onSuccess(ValidateDeviceCredentialsResponse msg) { - onValidateDeviceResponse(msg, ctx); + onValidateDeviceResponse(msg, ctx, connectMessage); } @Override public void onError(Throwable 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(); } }); } catch (Exception e) { - ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED)); + ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED, connectMessage)); ctx.close(); } } @@ -547,11 +547,11 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement doDisconnect(); } - private MqttConnAckMessage createMqttConnAckMsg(MqttConnectReturnCode returnCode) { + private MqttConnAckMessage createMqttConnAckMsg(MqttConnectReturnCode returnCode, MqttConnectMessage msg) { MqttFixedHeader mqttFixedHeader = new MqttFixedHeader(CONNACK, false, AT_MOST_ONCE, false, 0); MqttConnAckVariableHeader mqttConnAckVariableHeader = - new MqttConnAckVariableHeader(returnCode, true); + new MqttConnAckVariableHeader(returnCode, !msg.variableHeader().isCleanSession()); 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()) { - ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED)); + ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED, connectMessage)); ctx.close(); } else { deviceSessionCtx.setDeviceInfo(msg.getDeviceInfo()); @@ -640,7 +641,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement public void onSuccess(Void msg) { transportService.registerAsyncSession(deviceSessionCtx.getSessionInfo(), MqttTransportHandler.this); checkGatewaySession(); - ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED)); + ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED, connectMessage)); log.info("[{}] Client connected!", sessionId); } @@ -651,7 +652,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } else { 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(); } }); diff --git a/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesDao.java b/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesDao.java index 458a251f8a..f1af5a1b1b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesDao.java @@ -16,6 +16,7 @@ package org.thingsboard.server.dao.attributes; 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.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -41,4 +42,6 @@ public interface AttributesDao { ListenableFuture> removeAll(TenantId tenantId, EntityId entityId, String attributeType, List keys); List findAllKeysByDeviceProfileId(TenantId tenantId, DeviceProfileId deviceProfileId); + + List findAllKeysByEntityIds(TenantId tenantId, EntityType entityType, List entityIds); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesService.java b/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesService.java index 2e18006036..24f988c79d 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesService.java +++ b/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 org.springframework.beans.factory.annotation.Autowired; 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.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -65,6 +66,11 @@ public class BaseAttributesService implements AttributesService { return attributesDao.findAllKeysByDeviceProfileId(tenantId, deviceProfileId); } + @Override + public List findAllKeysByEntityIds(TenantId tenantId, EntityType entityType, List entityIds) { + return attributesDao.findAllKeysByEntityIds(tenantId, entityType, entityIds); + } + @Override public ListenableFuture> save(TenantId tenantId, EntityId entityId, String scope, List attributes) { validate(entityId, scope); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/AttributeKvRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/AttributeKvRepository.java index f6c3e195e7..3d14730b54 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/AttributeKvRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/AttributeKvRepository.java @@ -56,5 +56,8 @@ public interface AttributeKvRepository extends CrudRepository 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 findAllKeysByEntityIds(@Param("entityType") String entityType, @Param("entityIds") List entityIds); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java index a95bd4b612..a6a9348e8f 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java +++ b/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.Value; 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.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -145,6 +146,12 @@ public class JpaAttributeDao extends JpaAbstractDaoListeningExecutorService impl } } + @Override + public List findAllKeysByEntityIds(TenantId tenantId, EntityType entityType, List entityIds) { + return attributeKvRepository + .findAllKeysByEntityIds(entityType.name(), entityIds.stream().map(EntityId::getId).collect(Collectors.toList())); + } + @Override public ListenableFuture save(TenantId tenantId, EntityId entityId, String attributeType, AttributeKvEntry attribute) { AttributeKvEntity entity = new AttributeKvEntity(); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/SqlTimeseriesLatestDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/SqlTimeseriesLatestDao.java index ff602e6eed..a97fa97a5c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/SqlTimeseriesLatestDao.java +++ b/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.concurrent.ExecutionException; import java.util.function.Function; +import java.util.stream.Collectors; @Slf4j @Component @@ -169,6 +170,11 @@ public class SqlTimeseriesLatestDao extends BaseAbstractSqlTimeseriesDao impleme } } + @Override + public List findAllKeysByEntityIds(TenantId tenantId, List entityIds) { + return tsKvLatestRepository.findAllKeysByEntityIds(entityIds.stream().map(EntityId::getId).collect(Collectors.toList())); + } + private ListenableFuture getNewLatestEntryFuture(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query) { ListenableFuture> future = findNewLatestEntryFuture(tenantId, entityId, query); return Futures.transformAsync(future, entryList -> { diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/TsKvLatestRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/TsKvLatestRepository.java index cd3db69d70..da8c921487 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/TsKvLatestRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/latest/TsKvLatestRepository.java @@ -36,4 +36,9 @@ public interface TsKvLatestRepository extends CrudRepository 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 findAllKeysByEntityIds(@Param("entityIds") List entityIds); + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java index 7160baf939..12e081d234 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java +++ b/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); } + @Override + public List findAllKeysByEntityIds(TenantId tenantId, List entityIds) { + return timeseriesLatestDao.findAllKeysByEntityIds(tenantId, entityIds); + } + @Override public ListenableFuture save(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry) { validate(entityId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java index b9d8c62833..db5f8f8684 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java @@ -79,12 +79,17 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD protected static List FIXED_PARTITION = Arrays.asList(new Long[]{0L}); + private CassandraTsPartitionsCache cassandraTsPartitionsCache; + @Autowired private Environment environment; @Value("${cassandra.query.ts_key_value_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}") private long systemTtl; @@ -111,13 +116,16 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD super.startExecutor(); if (!isInstall()) { getFetchStmt(Aggregation.NONE, DESC_ORDER); - } - Optional partition = NoSqlTsPartitionDate.parse(partitioning); - if (partition.isPresent()) { - tsFormat = partition.get(); - } else { - log.warn("Incorrect configuration of partitioning {}", partitioning); - throw new RuntimeException("Failed to parse partitioning property: " + partitioning + "!"); + Optional partition = NoSqlTsPartitionDate.parse(partitioning); + if (partition.isPresent()) { + tsFormat = partition.get(); + if (!isFixedPartitioning() && partitionsCacheSize > 0) { + cassandraTsPartitionsCache = new CassandraTsPartitionsCache(partitionsCacheSize); + } + } 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); long partition = toPartitionTs(tsKvEntryTs); - log.debug("Saving partition {} for the entity [{}-{}] and key {}", partition, entityId.getEntityType(), entityId.getId(), key); - BoundStatementBuilder stmtBuilder = new BoundStatementBuilder((ttl == 0 ? getPartitionInsertStmt() : getPartitionInsertTtlStmt()).bind()); - stmtBuilder.setString(0, entityId.getEntityType().name()) - .setUuid(1, entityId.getId()) - .setLong(2, partition) - .setString(3, key); - if (ttl > 0) { - stmtBuilder.setInt(4, (int) ttl); + if (cassandraTsPartitionsCache == null) { + return doSavePartition(tenantId, entityId, key, ttl, partition); + } else { + CassandraPartitionCacheKey partitionSearchKey = new CassandraPartitionCacheKey(entityId, key, partition); + if (!cassandraTsPartitionsCache.has(partitionSearchKey)) { + ListenableFuture result = doSavePartition(tenantId, entityId, key, ttl, partition); + Futures.addCallback(result, new CacheCallback<>(partitionSearchKey), MoreExecutors.directExecutor()); + return result; + } else { + return Futures.immediateFuture(0); + } } - BoundStatement stmt = stmtBuilder.build(); - return getFuture(executeAsyncWrite(tenantId, stmt), rs -> 0); } @Override @@ -461,6 +470,38 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD return getFuture(executeAsyncWrite(tenantId, stmt), rs -> null); } + private ListenableFuture 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 implements FutureCallback { + 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) { if (systemTtl > 0) { if (ttl == 0) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java index f086f913c0..2cb67f1260 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java @@ -86,6 +86,11 @@ public class CassandraBaseTimeseriesLatestDao extends AbstractCassandraBaseTimes return Collections.emptyList(); } + @Override + public List findAllKeysByEntityIds(TenantId tenantId, List entityIds) { + return Collections.emptyList(); + } + @Override public ListenableFuture saveLatest(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry) { BoundStatementBuilder stmtBuilder = new BoundStatementBuilder(getLatestStmt().bind()); diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraPartitionCacheKey.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraPartitionCacheKey.java new file mode 100644 index 0000000000..791ce84113 --- /dev/null +++ b/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; + +} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraTsPartitionsCache.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraTsPartitionsCache.java new file mode 100644 index 0000000000..c167fb28cb --- /dev/null +++ b/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 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)); + } +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesLatestDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesLatestDao.java index 0b7156e03e..b2cd346277 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesLatestDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesLatestDao.java @@ -35,4 +35,6 @@ public interface TimeseriesLatestDao { ListenableFuture removeLatest(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query); List findAllKeysByDeviceProfileId(TenantId tenantId, DeviceProfileId deviceProfileId); + + List findAllKeysByEntityIds(TenantId tenantId, List entityIds); } diff --git a/dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java b/dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java new file mode 100644 index 0000000000..83975f46e8 --- /dev/null +++ b/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)); + } +} diff --git a/dao/src/test/resources/cassandra-test.properties b/dao/src/test/resources/cassandra-test.properties index 43a78abac4..4b3ea0a74d 100644 --- a/dao/src/test/resources/cassandra-test.properties +++ b/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_partitions_max_cache_size=100000 + cassandra.query.ts_key_value_ttl=0 cassandra.query.debug_events_ttl=604800 diff --git a/ui-ngx/src/app/core/http/entity.service.ts b/ui-ngx/src/app/core/http/entity.service.ts index 6a1f0191c4..953e67549c 100644 --- a/ui-ngx/src/app/core/http/entity.service.ts +++ b/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 { RuleChainService } from '@core/http/rule-chain.service'; 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 { 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 { deepClone, isDefined, isDefinedAndNotNull } from '@core/utils'; import { Asset } from '@shared/models/asset.models'; @@ -376,6 +382,13 @@ export class EntityService { return this.http.post>('/api/entitiesQuery/find', query, defaultHttpOptionsFromConfig(config)); } + public findEntityKeysByQuery(query: EntityDataQuery, attributes = true, timeseries = true, + config?: RequestConfig): Observable { + return this.http.post( + `/api/entitiesQuery/find/keys?attributes=${attributes}×eries=${timeseries}`, + query, defaultHttpOptionsFromConfig(config)); + } + public findAlarmDataByQuery(query: AlarmDataQuery, config?: RequestConfig): Observable> { return this.http.post>('/api/alarmsQuery/find', query, defaultHttpOptionsFromConfig(config)); } @@ -595,7 +608,7 @@ export class EntityService { return entityTypes; } - private getEntityFieldKeys(entityType: EntityType, searchText: string): Array { + private getEntityFieldKeys(entityType: EntityType, searchText: string = ''): Array { const entityFieldKeys: string[] = [entityFields.createdTime.keyName]; const query = searchText.toLowerCase(); switch (entityType) { @@ -637,7 +650,7 @@ export class EntityService { return query ? entityFieldKeys.filter((entityField) => entityField.toLowerCase().indexOf(query) === 0) : entityFieldKeys; } - private getAlarmKeys(searchText: string): Array { + private getAlarmKeys(searchText: string = ''): Array { const alarmKeys: string[] = Object.keys(alarmFields); const query = searchText.toLowerCase(); 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> { + if (!types.length) { + return of([]); + } + let entitiesKeysByQuery$: Observable; + 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 = []; + types.forEach(type => { + let keys: Array; + 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): Array { const datasources = subscriptionsInfo.map(subscriptionInfo => this.createDatasourceFromSubscriptionInfo(subscriptionInfo)); this.utils.generateColors(datasources); diff --git a/ui-ngx/src/app/modules/home/components/profile/alarm/device-profile-alarm.component.html b/ui-ngx/src/app/modules/home/components/profile/alarm/device-profile-alarm.component.html index cec0f48259..79fa9b7383 100644 --- a/ui-ngx/src/app/modules/home/components/profile/alarm/device-profile-alarm.component.html +++ b/ui-ngx/src/app/modules/home/components/profile/alarm/device-profile-alarm.component.html @@ -33,87 +33,91 @@ -
- - - {{'device-profile.alarm-type' | translate}} - - - {{ 'device-profile.alarm-type-required' | translate }} - - - {{ 'device-profile.alarm-type-unique' | translate }} - - -
- - - -
-
device-profile.advanced-settings
-
-
-
- - {{ 'device-profile.propagate-alarm' | translate }} - -
- - device-profile.alarm-rule-relation-types-list - - - {{key}} - close - - - - + +
+ + + {{'device-profile.alarm-type' | translate}} + + + {{ 'device-profile.alarm-type-required' | translate }} + + + {{ 'device-profile.alarm-type-unique' | translate }} + -
-
-
-
device-profile.create-alarm-rules
- - -
device-profile.clear-alarm-rule
-
-
- - -
-
-
- device-profile.no-clear-alarm-rule -
-
- + + + +
+
device-profile.advanced-settings
+
+
+
+ + + {{ 'device-profile.propagate-alarm' | translate }} + +
+ + device-profile.alarm-rule-relation-types-list + + + {{key}} + close + + + + + +
+
+
+
+
device-profile.create-alarm-rules
+ + +
device-profile.clear-alarm-rule
+
+
+ + +
+ +
+
+ device-profile.no-clear-alarm-rule +
+
+ +
-
+ diff --git a/ui-ngx/src/app/modules/home/components/profile/device-profile.component.html b/ui-ngx/src/app/modules/home/components/profile/device-profile.component.html index 4cba236d2a..bb3d6e7e43 100644 --- a/ui-ngx/src/app/modules/home/components/profile/device-profile.component.html +++ b/ui-ngx/src/app/modules/home/components/profile/device-profile.component.html @@ -91,10 +91,12 @@
device-profile.profile-configuration
- - + + + + @@ -102,10 +104,12 @@
device-profile.transport-configuration
- - + + + +
@@ -115,10 +119,12 @@ entityForm.get('profileData.alarms').value.length : 0} }}
- - + + + + @@ -126,9 +132,11 @@
device-profile.device-provisioning
- - + + + +
diff --git a/ui-ngx/src/app/modules/home/components/profile/tenant-profile-data.component.html b/ui-ngx/src/app/modules/home/components/profile/tenant-profile-data.component.html index 3d853940f4..ae7a9540db 100644 --- a/ui-ngx/src/app/modules/home/components/profile/tenant-profile-data.component.html +++ b/ui-ngx/src/app/modules/home/components/profile/tenant-profile-data.component.html @@ -22,9 +22,11 @@
tenant-profile.profile-configuration
- - + + + + diff --git a/ui-ngx/src/app/modules/home/components/widget/data-key-config.component.ts b/ui-ngx/src/app/modules/home/components/widget/data-key-config.component.ts index fde4b1a32f..a4187d9b6c 100644 --- a/ui-ngx/src/app/modules/home/components/widget/data-key-config.component.ts +++ b/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 { DataKeyType } from '@shared/models/telemetry/telemetry.models'; 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 { JsFuncComponent } from '@shared/components/js-func.component'; 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>; private latestKeySearchResult: Array = null; + private fetchObservable$: Observable> = null; keySearchText = ''; @@ -205,31 +206,42 @@ export class DataKeyConfigComponent extends PageComponent implements OnInit, Con } private fetchKeys(searchText?: string): Observable> { - if (this.latestKeySearchResult === null || this.keySearchText !== searchText) { + if (this.keySearchText !== searchText || this.latestKeySearchResult === null) { this.keySearchText = searchText; - let fetchObservable: Observable> = 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>; if (this.modelValue.type === DataKeyType.alarm) { - const dataKeyFilter = this.createDataKeyFilter(this.keySearchText); - fetchObservable = of(this.alarmKeys.filter(dataKeyFilter)); + fetchObservable = of(this.alarmKeys); } else { if (this.entityAliasId) { const dataKeyTypes = [this.modelValue.type]; - fetchObservable = this.callbacks.fetchEntityKeys(this.entityAliasId, this.keySearchText, dataKeyTypes); + fetchObservable = this.callbacks.fetchEntityKeys(this.entityAliasId, dataKeyTypes); } else { fetchObservable = of([]); } } - return fetchObservable.pipe( + this.fetchObservable$ = fetchObservable.pipe( 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(); - return key => key.name.toLowerCase().indexOf(lowercaseQuery) === 0; + return key => key.toLowerCase().startsWith(lowercaseQuery); } public validateOnSubmit() { diff --git a/ui-ngx/src/app/modules/home/components/widget/data-keys.component.models.ts b/ui-ngx/src/app/modules/home/components/widget/data-keys.component.models.ts index 452e390faa..f5a937a4be 100644 --- a/ui-ngx/src/app/modules/home/components/widget/data-keys.component.models.ts +++ b/ui-ngx/src/app/modules/home/components/widget/data-keys.component.models.ts @@ -20,5 +20,5 @@ import { Observable } from 'rxjs'; export interface DataKeysCallbacks { generateDataKey: (chip: any, type: DataKeyType) => DataKey; - fetchEntityKeys: (entityAliasId: string, query: string, types: Array) => Observable>; + fetchEntityKeys: (entityAliasId: string, types: Array) => Observable>; } diff --git a/ui-ngx/src/app/modules/home/components/widget/data-keys.component.ts b/ui-ngx/src/app/modules/home/components/widget/data-keys.component.ts index 22c9663be8..e31ad036c9 100644 --- a/ui-ngx/src/app/modules/home/components/widget/data-keys.component.ts +++ b/ui-ngx/src/app/modules/home/components/widget/data-keys.component.ts @@ -38,7 +38,7 @@ import { Validators } from '@angular/forms'; 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 { AppState } from '@app/core/core.state'; import { TranslateService } from '@ngx-translate/core'; @@ -142,6 +142,7 @@ export class DataKeysComponent implements ControlValueAccessor, OnInit, AfterVie searchText = ''; private latestSearchTextResult: Array = null; + private fetchObservable$: Observable> = null; private dirty = false; @@ -260,6 +261,7 @@ export class DataKeysComponent implements ControlValueAccessor, OnInit, AfterVie if (!change.firstChange && change.currentValue !== change.previousValue) { if (propName === 'entityAliasId') { this.searchText = ''; + this.fetchObservable$ = null; this.latestSearchTextResult = null; this.dirty = true; } else if (['widgetType', 'datasourceType'].includes(propName)) { @@ -405,14 +407,24 @@ export class DataKeysComponent implements ControlValueAccessor, OnInit, AfterVie return key ? key.name : undefined; } - fetchKeys(searchText?: string): Observable> { - if (this.latestSearchTextResult === null || this.searchText !== searchText) { + private fetchKeys(searchText?: string): Observable> { + if (this.searchText !== searchText || this.latestSearchTextResult === null) { this.searchText = searchText; - let fetchObservable: Observable> = 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> { + if (this.fetchObservable$ === null) { + let fetchObservable: Observable>; if (this.datasourceType === DatasourceType.function) { - const dataKeyFilter = this.createDataKeyFilter(this.searchText); const targetKeysList = this.widgetType === widgetType.alarm ? this.alarmKeys : this.functionTypeKeys; - fetchObservable = of(targetKeysList.filter(dataKeyFilter)); + fetchObservable = of(targetKeysList); } else { if (this.entityAliasId) { const dataKeyTypes = [DataKeyType.timeseries]; @@ -420,24 +432,25 @@ export class DataKeysComponent implements ControlValueAccessor, OnInit, AfterVie dataKeyTypes.push(DataKeyType.attribute); dataKeyTypes.push(DataKeyType.entityField); 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 { fetchObservable = of([]); } } - return fetchObservable.pipe( - tap(res => this.latestSearchTextResult = res) + this.fetchObservable$ = fetchObservable.pipe( + publishReplay(1), + refCount() ); } - return of(this.latestSearchTextResult); + return this.fetchObservable$; } private createDataKeyFilter(query: string): (key: DataKey) => boolean { const lowercaseQuery = query.toLowerCase(); - return key => key.name.toLowerCase().indexOf(lowercaseQuery) === 0; + return key => key.name.toLowerCase().startsWith(lowercaseQuery); } textIsNotEmpty(text: string): boolean { diff --git a/ui-ngx/src/app/modules/home/components/widget/lib/maps/leaflet-map.ts b/ui-ngx/src/app/modules/home/components/widget/lib/maps/leaflet-map.ts index 6b903101ad..5430aa2e69 100644 --- a/ui-ngx/src/app/modules/home/components/widget/lib/maps/leaflet-map.ts +++ b/ui-ngx/src/app/modules/home/components/widget/lib/maps/leaflet-map.ts @@ -608,26 +608,32 @@ export default abstract class LeafletMap { return polygon; } - updatePoints(pointsData: FormattedData[], getTooltip: (point: FormattedData, setTooltip?: boolean) => string) { + updatePoints(pointsData: FormattedData[][], getTooltip: (point: FormattedData) => string) { + if(pointsData.length) { if (this.points) { - this.map.removeLayer(this.points); + this.map.removeLayer(this.points); } this.points = new FeatureGroup(); - pointsData.filter(pdata => !!this.convertPosition(pdata)).forEach(data => { - const point = L.circleMarker(this.convertPosition(data), { - color: this.options.pointColor, - radius: this.options.pointSize - }); - if (!this.options.pointTooltipOnRightPanel) { - point.on('click', () => getTooltip(data)); - } - else { - createTooltip(point, this.options, data.$datasource, getTooltip(data, false)); - } - this.points.addLayer(point); + } + for(let i = 0; i < pointsData.length; i++) { + const pointsList = pointsData[i]; + pointsList.filter(pdata => !!this.convertPosition(pdata)).forEach(data => { + const point = L.circleMarker(this.convertPosition(data), { + color: this.options.pointColor, + radius: this.options.pointSize + }); + if (!this.options.pointTooltipOnRightPanel) { + point.on('click', () => getTooltip(data)); + } else { + createTooltip(point, this.options, data.$datasource, getTooltip(data)); + } + this.points.addLayer(point); }); + } + if(pointsData.length) { this.map.addLayer(this.points); } + } // Polyline diff --git a/ui-ngx/src/app/modules/home/components/widget/trip-animation/trip-animation.component.html b/ui-ngx/src/app/modules/home/components/widget/trip-animation/trip-animation.component.html index 2ea80cf2af..6264231b68 100644 --- a/ui-ngx/src/app/modules/home/components/widget/trip-animation/trip-animation.component.html +++ b/ui-ngx/src/app/modules/home/components/widget/trip-animation/trip-animation.component.html @@ -28,8 +28,12 @@
+ [ngClass]="{'trip-animation-tooltip-hidden':!visibleTooltip}" + [ngStyle]="{'background-color': settings.tooltipColor, 'opacity': settings.tooltipOpacity, 'color': settings.tooltipFontColor}"> +
+
arr.length); if (this.historicalData.length) { this.calculateIntervals(); - this.timeUpdated(this.currentTime && this.currentTime > this.minTime ? this.currentTime : this.minTime); + this.timeUpdated(this.minTime); } this.mapWidget.map.map?.invalidateSize(); this.cd.detectChanges(); @@ -140,32 +143,39 @@ export class TripAnimationComponent implements OnInit, AfterViewInit, OnDestroy this.currentTime = time; const currentPosition = this.interpolatedTimeData .map(dataSource => dataSource[time]) - .filter(ds => ds); - if (isUndefined(currentPosition[0])) { - const timePoints = Object.keys(this.interpolatedTimeData[0]).map(item => parseInt(item, 10)); - for (let i = 1; i < timePoints.length; i++) { - if (timePoints[i - 1] < time && timePoints[i] > time) { - const beforePosition = this.interpolatedTimeData[0][timePoints[i - 1]]; - const afterPosition = this.interpolatedTimeData[0][timePoints[i]]; - const ratio = getRatio(timePoints[i - 1], timePoints[i], time); - currentPosition[0] = { - ...beforePosition, - time, - ...interpolateOnLineSegment(beforePosition, afterPosition, this.settings.latKeyName, this.settings.lngKeyName, ratio) + for(let j = 0; j < this.interpolatedTimeData.length; j++) { + if (isUndefined(currentPosition[j])) { + const timePoints = Object.keys(this.interpolatedTimeData[j]).map(item => parseInt(item, 10)); + for (let i = 1; i < timePoints.length; i++) { + if (timePoints[i - 1] < time && timePoints[i] > time) { + const beforePosition = this.interpolatedTimeData[j][timePoints[i - 1]]; + const afterPosition = this.interpolatedTimeData[j][timePoints[i]]; + const ratio = getRatio(timePoints[i - 1], timePoints[i], time); + currentPosition[j] = { + ...beforePosition, + time, + ...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.calcTooltip(currentPosition.find(position => position.entityName === this.activeTrip.entityName)); + this.calcMainTooltip(currentPosition); 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) { this.mapWidget.map.updatePolygons(this.interpolatedTimeData); } 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.activeTrip = trip; @@ -177,6 +187,23 @@ export class TripAnimationComponent implements OnInit, AfterViewInit, OnDestroy 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() { this.historicalData.forEach((dataSource, index) => { 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 tooltipPattern: string = this.settings.useTooltipFunction ? safeExecute(this.settings.tooltipFunction, [data, this.historicalData, point.dsIndex]) : this.settings.tooltipPattern; - const tooltipText = parseWithTranslation.parseTemplate(tooltipPattern, data, true); - this.mainTooltip = this.sanitizer.sanitize( - SecurityContext.HTML, tooltipText); - this.cd.detectChanges(); - this.activeTrip = point; - return tooltipText; + return parseWithTranslation.parseTemplate(tooltipPattern, data, true); + } + + private calcMainTooltip(points: FormattedData[]): void { + const tooltips = []; + for (let point of points) { + tooltips.push(this.sanitizer.sanitize(SecurityContext.HTML, this.calcTooltip(point))); + } + this.mainTooltips = tooltips; } calcLabel() { diff --git a/ui-ngx/src/app/modules/home/components/widget/widget-config.component.ts b/ui-ngx/src/app/modules/home/components/widget/widget-config.component.ts index ba8cdaf5d2..b1984f997a 100644 --- a/ui-ngx/src/app/modules/home/components/widget/widget-config.component.ts +++ b/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 { TranslateService } from '@ngx-translate/core'; 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 { EntityAliasDialogComponent, EntityAliasDialogData } 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 { EntityService } from '@core/http/entity.service'; 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): Observable> { - return this.aliasController.resolveSingleEntityInfo(entityAliasId).pipe( - mergeMap((entity) => { - if (entity) { - const fetchEntityTasks: Array>> = []; - for (const dataKeyType of dataKeyTypes) { - fetchEntityTasks.push( - this.entityService.getEntityKeys( - {entityType: entity.entityType, id: entity.id}, - query, - dataKeyType, - {ignoreLoading: true, ignoreErrors: true} - ).pipe( - map((keys) => { - const dataKeys: Array = []; - 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(); - 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 = []; - for (const key of keys) { - dataKeys.push({name: key, type: DataKeyType.alarm}); - } - return dataKeys; - } - ), - catchError(() => of([])) - ); - } else { - return of([]); - } + private fetchEntityKeys(entityAliasId: string, dataKeyTypes: Array): Observable> { + return this.aliasController.getAliasInfo(entityAliasId).pipe( + mergeMap((aliasInfo) => { + return this.entityService.getEntityKeysByEntityFilter( + aliasInfo.entityFilter, + dataKeyTypes, + {ignoreLoading: true, ignoreErrors: true} + ).pipe( + catchError(() => of([])) + ); }), catchError(() => of([] as Array)) ); diff --git a/ui-ngx/src/app/shared/models/entity.models.ts b/ui-ngx/src/app/shared/models/entity.models.ts index 8f1d4f417e..300de484c1 100644 --- a/ui-ngx/src/app/shared/models/entity.models.ts +++ b/ui-ngx/src/app/shared/models/entity.models.ts @@ -64,6 +64,12 @@ export interface EntityField { time?: boolean; } +export interface EntitiesKeysByQuery { + attribute: Array; + timeseries: Array; + entityTypes: EntityType[]; +} + export const entityFields: {[fieldName: string]: EntityField} = { createdTime: { keyName: 'createdTime',