From cc645a938e9a44f582a82f699967d1deed695d66 Mon Sep 17 00:00:00 2001 From: deaflynx Date: Thu, 13 Aug 2020 16:53:10 +0300 Subject: [PATCH 1/7] Code cleaning from PE comparison of master vs. develop/2.5.3 --- ui/src/app/api/dashboard.service.js | 8 -------- 1 file changed, 8 deletions(-) diff --git a/ui/src/app/api/dashboard.service.js b/ui/src/app/api/dashboard.service.js index 7280c9d06d..9d59e24ebc 100644 --- a/ui/src/app/api/dashboard.service.js +++ b/ui/src/app/api/dashboard.service.js @@ -284,13 +284,6 @@ function DashboardService($rootScope, $http, $q, $location, $filter) { } dashboard.assignedCustomersText = assignedCustomersTitles.join(', '); } - dashboard.assignedEdgesIds = []; - if (dashboard.assignedEdges && dashboard.assignedEdges.length) { - for (var j = 0; j < dashboard.assignedEdges.length; j++) { - var assignedEdge = dashboard.assignedEdges[j]; - dashboard.assignedEdgesIds.push(assignedEdge.edgeId.id); - } - } return dashboard; } @@ -298,7 +291,6 @@ function DashboardService($rootScope, $http, $q, $location, $filter) { delete dashboard.publicCustomerId; delete dashboard.assignedCustomersText; delete dashboard.assignedCustomersIds; - delete dashboard.assignedEdgesIds; return dashboard; } From 581f23b0b067bdfb06c9a9955eb4dd0331db4187 Mon Sep 17 00:00:00 2001 From: Bohdan Smetaniuk Date: Fri, 21 Aug 2020 18:50:11 +0300 Subject: [PATCH 2/7] fixes + processing attributes delete msg from edge --- .../edge/DefaultEdgeNotificationService.java | 50 ++--------------- .../service/edge/rpc/EdgeGrpcSession.java | 31 ++++++++++- .../constructor/EntityDataMsgConstructor.java | 6 +- .../server/dao/edge/EdgeService.java | 6 +- .../server/dao/edge/EdgeServiceImpl.java | 55 ++++++++++++++++++- .../rule/engine/edge/TbMsgPushToEdgeNode.java | 9 +-- 6 files changed, 95 insertions(+), 62 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java b/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java index 279ec96064..27489c98cc 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java @@ -242,7 +242,7 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService { } } } else { - ListenableFuture> edgeIdsFuture = findRelatedEdgeIdsByEntityId(tenantId, entityId); + ListenableFuture> edgeIdsFuture = edgeService.findRelatedEdgeIdsByEntityId(tenantId, entityId, dbCallbackExecutorService); Futures.transform(edgeIdsFuture, edgeIds -> { if (edgeIds != null && !edgeIds.isEmpty()) { for (EdgeId edgeId : edgeIds) { @@ -321,7 +321,7 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService { if (alarm != null) { EdgeEventType edgeEventType = getEdgeQueueTypeByEntityType(alarm.getOriginator().getEntityType()); if (edgeEventType != null) { - ListenableFuture> relatedEdgeIdsByEntityIdFuture = findRelatedEdgeIdsByEntityId(tenantId, alarm.getOriginator()); + ListenableFuture> relatedEdgeIdsByEntityIdFuture = edgeService.findRelatedEdgeIdsByEntityId(tenantId, alarm.getOriginator(), dbCallbackExecutorService); Futures.transform(relatedEdgeIdsByEntityIdFuture, relatedEdgeIdsByEntityId -> { if (relatedEdgeIdsByEntityId != null) { for (EdgeId edgeId : relatedEdgeIdsByEntityId) { @@ -346,8 +346,8 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService { if (!relation.getFrom().getEntityType().equals(EntityType.EDGE) && !relation.getTo().getEntityType().equals(EntityType.EDGE)) { List>> futures = new ArrayList<>(); - futures.add(findRelatedEdgeIdsByEntityId(tenantId, relation.getTo())); - futures.add(findRelatedEdgeIdsByEntityId(tenantId, relation.getFrom())); + futures.add(edgeService.findRelatedEdgeIdsByEntityId(tenantId, relation.getTo(), dbCallbackExecutorService)); + futures.add(edgeService.findRelatedEdgeIdsByEntityId(tenantId, relation.getFrom(), dbCallbackExecutorService)); ListenableFuture>> combinedFuture = Futures.allAsList(futures); Futures.transform(combinedFuture, listOfListsEdgeIds -> { Set uniqueEdgeIds = new HashSet<>(); @@ -373,48 +373,6 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService { } } - private ListenableFuture> findRelatedEdgeIdsByEntityId(TenantId tenantId, EntityId entityId) { - switch (entityId.getEntityType()) { - case DEVICE: - case ASSET: - case ENTITY_VIEW: - ListenableFuture> originatorEdgeRelationsFuture = - relationService.findByToAndTypeAsync(tenantId, entityId, EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE); - return Futures.transform(originatorEdgeRelationsFuture, originatorEdgeRelations -> { - if (originatorEdgeRelations != null && originatorEdgeRelations.size() > 0) { - return Collections.singletonList(new EdgeId(originatorEdgeRelations.get(0).getFrom().getId())); - } else { - return Collections.emptyList(); - } - }, dbCallbackExecutorService); - case DASHBOARD: - return convertToEdgeIds(edgeService.findEdgesByTenantIdAndDashboardId(tenantId, new DashboardId(entityId.getId()))); - case RULE_CHAIN: - return convertToEdgeIds(edgeService.findEdgesByTenantIdAndRuleChainId(tenantId, new RuleChainId(entityId.getId()))); - case USER: - User userById = userService.findUserById(tenantId, new UserId(entityId.getId())); - TextPageData edges; - if (userById.getCustomerId() == null || userById.getCustomerId().getId().equals(ModelConstants.NULL_UUID)) { - edges = edgeService.findEdgesByTenantId(tenantId, new TextPageLink(Integer.MAX_VALUE)); - } else { - edges = edgeService.findEdgesByTenantIdAndCustomerId(tenantId, new CustomerId(entityId.getId()), new TextPageLink(Integer.MAX_VALUE)); - } - return convertToEdgeIds(Futures.immediateFuture(edges.getData())); - default: - return Futures.immediateFuture(Collections.emptyList()); - } - } - - private ListenableFuture> convertToEdgeIds(ListenableFuture> future) { - return Futures.transform(future, edges -> { - if (edges != null && !edges.isEmpty()) { - return edges.stream().map(IdBased::getId).collect(Collectors.toList()); - } else { - return Collections.emptyList(); - } - }, dbCallbackExecutorService); - } - private EdgeEventType getEdgeQueueTypeByEntityType(EntityType entityType) { switch (entityType) { case DEVICE: diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java index 121772a01b..86c45c78ba 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java @@ -33,6 +33,7 @@ import lombok.Data; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang.RandomStringUtils; import org.checkerframework.checker.nullness.qual.Nullable; +import org.thingsboard.rule.engine.api.msg.DeviceAttributesEventNotificationMsg; import org.thingsboard.server.common.data.AdminSettings; import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.Dashboard; @@ -64,6 +65,7 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.id.WidgetTypeId; import org.thingsboard.server.common.data.id.WidgetsBundleId; +import org.thingsboard.server.common.data.kv.AttributeKey; import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; import org.thingsboard.server.common.data.kv.LongDataEntry; @@ -88,6 +90,7 @@ import org.thingsboard.server.common.transport.util.JsonUtils; import org.thingsboard.server.gen.edge.AdminSettingsUpdateMsg; import org.thingsboard.server.gen.edge.AlarmUpdateMsg; import org.thingsboard.server.gen.edge.AssetUpdateMsg; +import org.thingsboard.server.gen.edge.AttributeDeleteMsg; import org.thingsboard.server.gen.edge.AttributesRequestMsg; import org.thingsboard.server.gen.edge.ConnectRequestMsg; import org.thingsboard.server.gen.edge.ConnectResponseCode; @@ -125,11 +128,12 @@ import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.service.edge.EdgeContextComponent; import java.io.Closeable; -import java.io.IOException; import java.util.ArrayList; import java.util.Collections; +import java.util.HashSet; import java.util.List; import java.util.Optional; +import java.util.Set; import java.util.UUID; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; @@ -382,7 +386,7 @@ public final class EdgeGrpcSession implements Closeable { ctx.getAttributesService().save(edge.getTenantId(), edge.getId(), DataConstants.SERVER_SCOPE, attributes); } - private DownlinkMsg processTelemetryMessage(EdgeEvent edgeEvent) throws IOException { + private DownlinkMsg processTelemetryMessage(EdgeEvent edgeEvent) { log.trace("Executing processTelemetryMessage, edgeEvent [{}]", edgeEvent); EntityId entityId = null; switch (edgeEvent.getEdgeEventType()) { @@ -823,6 +827,9 @@ public final class EdgeGrpcSession implements Closeable { result.add(processPostTelemetry(entityId, entityData.getPostTelemetryMsg(), metaData)); } } + if (entityData.hasAttributeDeleteMsg()) { + result.add(processAttributeDeleteMsg(entityId, entityData.getAttributeDeleteMsg(), entityData.getEntityType())); + } } } @@ -948,6 +955,26 @@ public final class EdgeGrpcSession implements Closeable { return Futures.immediateFuture(null); } + private ListenableFuture processAttributeDeleteMsg(EntityId entityId, AttributeDeleteMsg attributeDeleteMsg, String entityType) { + try { + String scope = attributeDeleteMsg.getScope(); + List attributeNames = attributeDeleteMsg.getAttributeNamesList(); + ctx.getAttributesService().removeAll(edge.getTenantId(), entityId, scope, attributeNames); + if (EntityType.DEVICE.name().equals(entityType)) { + Set attributeKeys = new HashSet<>(); + for (String attributeName : attributeNames) { + attributeKeys.add(new AttributeKey(scope, attributeName)); + } + ctx.getTbClusterService().pushMsgToCore(DeviceAttributesEventNotificationMsg.onDelete( + edge.getTenantId(), (DeviceId) entityId, attributeKeys), null); + } + } catch (Exception e) { + log.error("Can't process attribute delete msg [{}]", attributeDeleteMsg, e); + return Futures.immediateFailedFuture(new RuntimeException("Can't process attribute delete msg " + attributeDeleteMsg, e)); + } + return Futures.immediateFuture(null); + } + private ListenableFuture onDeviceUpdate(DeviceUpdateMsg deviceUpdateMsg) { DeviceId edgeDeviceId = new DeviceId(new UUID(deviceUpdateMsg.getIdMSB(), deviceUpdateMsg.getIdLSB())); switch (deviceUpdateMsg.getMsgType()) { diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java index c1d211af58..0f7fb0523d 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java @@ -43,9 +43,11 @@ public class EntityDataMsgConstructor { case TIMESERIES_UPDATED: try { JsonObject data = entityData.getAsJsonObject(); - long ts = System.currentTimeMillis(); + long ts; if (data.get("ts") != null && !data.get("ts").isJsonNull()) { - ts = data.getAsJsonObject("ts").getAsLong(); + ts = data.getAsJsonPrimitive("ts").getAsLong(); + } else { + ts = System.currentTimeMillis(); } builder.setPostTelemetryMsg(JsonConverter.convertToTelemetryProto(data.getAsJsonObject("data"), ts)); } catch (Exception e) { diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java index 5c6e5d7559..71c56068d1 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java @@ -22,15 +22,15 @@ import org.thingsboard.server.common.data.edge.EdgeSearchQuery; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.DashboardId; import org.thingsboard.server.common.data.id.EdgeId; +import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.TextPageData; import org.thingsboard.server.common.data.page.TextPageLink; -import org.thingsboard.server.common.data.page.TimePageData; -import org.thingsboard.server.common.data.page.TimePageLink; import java.util.List; import java.util.Optional; +import java.util.concurrent.Executor; public interface EdgeService { @@ -75,6 +75,8 @@ public interface EdgeService { ListenableFuture> findEdgesByTenantIdAndRuleChainId(TenantId tenantId, RuleChainId ruleChainId); ListenableFuture> findEdgesByTenantIdAndDashboardId(TenantId tenantId, DashboardId dashboardId); + + ListenableFuture> findRelatedEdgeIdsByEntityId(TenantId tenantId, EntityId entityId, Executor executorService); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java index 01f82faa17..4ce86beea5 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java @@ -31,31 +31,35 @@ import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.EntitySubtype; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.Tenant; +import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.EdgeSearchQuery; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.DashboardId; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.IdBased; import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.page.TextPageData; import org.thingsboard.server.common.data.page.TextPageLink; -import org.thingsboard.server.common.data.page.TimePageData; -import org.thingsboard.server.common.data.page.TimePageLink; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.EntitySearchDirection; +import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.dao.customer.CustomerDao; import org.thingsboard.server.dao.dashboard.DashboardService; import org.thingsboard.server.dao.entity.AbstractEntityService; import org.thingsboard.server.dao.exception.DataValidationException; +import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.dao.rule.RuleChainService; import org.thingsboard.server.dao.service.DataValidator; import org.thingsboard.server.dao.service.PaginatedRemover; import org.thingsboard.server.dao.service.Validator; import org.thingsboard.server.dao.tenant.TenantDao; +import org.thingsboard.server.dao.user.UserService; import javax.annotation.Nullable; import java.util.ArrayList; @@ -63,6 +67,7 @@ import java.util.Collections; import java.util.Comparator; import java.util.List; import java.util.Optional; +import java.util.concurrent.Executor; import java.util.stream.Collectors; import static org.thingsboard.server.common.data.CacheConstants.EDGE_CACHE; @@ -91,6 +96,9 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic @Autowired private CustomerDao customerDao; + @Autowired + private UserService userService; + @Autowired private CacheManager cacheManager; @@ -420,4 +428,47 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic } }; + @Override + public ListenableFuture> findRelatedEdgeIdsByEntityId(TenantId tenantId, EntityId entityId, Executor executorService) { + switch (entityId.getEntityType()) { + case DEVICE: + case ASSET: + case ENTITY_VIEW: + ListenableFuture> originatorEdgeRelationsFuture = + relationService.findByToAndTypeAsync(tenantId, entityId, EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE); + return Futures.transform(originatorEdgeRelationsFuture, originatorEdgeRelations -> { + if (originatorEdgeRelations != null && originatorEdgeRelations.size() > 0) { + return Collections.singletonList(new EdgeId(originatorEdgeRelations.get(0).getFrom().getId())); + } else { + return Collections.emptyList(); + } + }, executorService); + case DASHBOARD: + return convertToEdgeIds(findEdgesByTenantIdAndDashboardId(tenantId, new DashboardId(entityId.getId())), executorService); + case RULE_CHAIN: + return convertToEdgeIds(findEdgesByTenantIdAndRuleChainId(tenantId, new RuleChainId(entityId.getId())), executorService); + case USER: + User userById = userService.findUserById(tenantId, new UserId(entityId.getId())); + TextPageData edges; + if (userById.getCustomerId() == null || userById.getCustomerId().getId().equals(ModelConstants.NULL_UUID)) { + edges = findEdgesByTenantId(tenantId, new TextPageLink(Integer.MAX_VALUE)); + } else { + edges = findEdgesByTenantIdAndCustomerId(tenantId, new CustomerId(entityId.getId()), new TextPageLink(Integer.MAX_VALUE)); + } + return convertToEdgeIds(Futures.immediateFuture(edges.getData()), executorService); + default: + return Futures.immediateFuture(Collections.emptyList()); + } + } + + private ListenableFuture> convertToEdgeIds(ListenableFuture> future, Executor executorService) { + return Futures.transform(future, edges -> { + if (edges != null && !edges.isEmpty()) { + return edges.stream().map(IdBased::getId).collect(Collectors.toList()); + } else { + return Collections.emptyList(); + } + }, executorService); + } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java index be265ebc19..b531af7ed3 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java @@ -232,14 +232,7 @@ public class TbMsgPushToEdgeNode implements TbNode { TextPageData edgesByTenantId = ctx.getEdgeService().findEdgesByTenantId(tenantId, new TextPageLink(Integer.MAX_VALUE)); return Futures.immediateFuture(edgesByTenantId.getData().stream().map(IdBased::getId).collect(Collectors.toList())); } else { - ListenableFuture> future = ctx.getRelationService().findByToAndTypeAsync(tenantId, originatorId, EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE); - return Futures.transform(future, relations -> { - List result = new ArrayList<>(); - if (relations != null && relations.size() > 0) { - result.add(new EdgeId(relations.get(0).getFrom().getId())); - } - return result; - }, ctx.getDbCallbackExecutor()); + return ctx.getEdgeService().findRelatedEdgeIdsByEntityId(tenantId, originatorId, ctx.getDbCallbackExecutor()); } } From 1d5f5c5d8037c5cc7dfa6ee70cb15ca554d0c327 Mon Sep 17 00:00:00 2001 From: Bohdan Smetaniuk Date: Tue, 25 Aug 2020 14:04:47 +0300 Subject: [PATCH 3/7] code fixes --- .../edge/DefaultEdgeNotificationService.java | 8 ++-- .../service/edge/rpc/EdgeGrpcSession.java | 48 +++++++++++++++---- .../server/dao/edge/EdgeService.java | 3 +- .../server/dao/edge/EdgeServiceImpl.java | 15 +++--- .../rule/engine/edge/TbMsgPushToEdgeNode.java | 2 +- 5 files changed, 53 insertions(+), 23 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java b/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java index 27489c98cc..fbc5bcf376 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java @@ -242,7 +242,7 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService { } } } else { - ListenableFuture> edgeIdsFuture = edgeService.findRelatedEdgeIdsByEntityId(tenantId, entityId, dbCallbackExecutorService); + ListenableFuture> edgeIdsFuture = edgeService.findRelatedEdgeIdsByEntityId(tenantId, entityId); Futures.transform(edgeIdsFuture, edgeIds -> { if (edgeIds != null && !edgeIds.isEmpty()) { for (EdgeId edgeId : edgeIds) { @@ -321,7 +321,7 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService { if (alarm != null) { EdgeEventType edgeEventType = getEdgeQueueTypeByEntityType(alarm.getOriginator().getEntityType()); if (edgeEventType != null) { - ListenableFuture> relatedEdgeIdsByEntityIdFuture = edgeService.findRelatedEdgeIdsByEntityId(tenantId, alarm.getOriginator(), dbCallbackExecutorService); + ListenableFuture> relatedEdgeIdsByEntityIdFuture = edgeService.findRelatedEdgeIdsByEntityId(tenantId, alarm.getOriginator()); Futures.transform(relatedEdgeIdsByEntityIdFuture, relatedEdgeIdsByEntityId -> { if (relatedEdgeIdsByEntityId != null) { for (EdgeId edgeId : relatedEdgeIdsByEntityId) { @@ -346,8 +346,8 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService { if (!relation.getFrom().getEntityType().equals(EntityType.EDGE) && !relation.getTo().getEntityType().equals(EntityType.EDGE)) { List>> futures = new ArrayList<>(); - futures.add(edgeService.findRelatedEdgeIdsByEntityId(tenantId, relation.getTo(), dbCallbackExecutorService)); - futures.add(edgeService.findRelatedEdgeIdsByEntityId(tenantId, relation.getFrom(), dbCallbackExecutorService)); + futures.add(edgeService.findRelatedEdgeIdsByEntityId(tenantId, relation.getTo())); + futures.add(edgeService.findRelatedEdgeIdsByEntityId(tenantId, relation.getFrom())); ListenableFuture>> combinedFuture = Futures.allAsList(futures); Futures.transform(combinedFuture, listOfListsEdgeIds -> { Set uniqueEdgeIds = new HashSet<>(); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java index 86c45c78ba..99ffe79d9e 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java @@ -25,6 +25,7 @@ import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; +import com.google.common.util.concurrent.SettableFuture; import com.google.gson.Gson; import com.google.gson.JsonElement; import com.google.gson.JsonObject; @@ -937,25 +938,46 @@ public final class EdgeGrpcSession implements Closeable { } private ListenableFuture processPostTelemetry(EntityId entityId, TransportProtos.PostTelemetryMsg msg, TbMsgMetaData metaData) { + SettableFuture futureToSet = SettableFuture.create(); for (TransportProtos.TsKvListProto tsKv : msg.getTsKvListList()) { JsonObject json = JsonUtils.getJsonObject(tsKv.getKvList()); metaData.putValue("ts", tsKv.getTs() + ""); TbMsg tbMsg = TbMsg.newMsg(SessionMsgType.POST_TELEMETRY_REQUEST.name(), entityId, metaData, gson.toJson(json)); - // TODO: voba - verify that null callback is OK - ctx.getTbClusterService().pushMsgToRuleEngine(edge.getTenantId(), tbMsg.getOriginator(), tbMsg, null); + ctx.getTbClusterService().pushMsgToRuleEngine(edge.getTenantId(), tbMsg.getOriginator(), tbMsg, new TbQueueCallback() { + @Override + public void onSuccess(TbQueueMsgMetadata metadata) { + futureToSet.set(null); + } + + @Override + public void onFailure(Throwable t) { + futureToSet.setException(t); + } + }); } - return Futures.immediateFuture(null); + return futureToSet; } private ListenableFuture processPostAttributes(EntityId entityId, TransportProtos.PostAttributeMsg msg, TbMsgMetaData metaData) { + SettableFuture futureToSet = SettableFuture.create(); JsonObject json = JsonUtils.getJsonObject(msg.getKvList()); TbMsg tbMsg = TbMsg.newMsg(SessionMsgType.POST_ATTRIBUTES_REQUEST.name(), entityId, metaData, gson.toJson(json)); - // TODO: voba - verify that null callback is OK - ctx.getTbClusterService().pushMsgToRuleEngine(edge.getTenantId(), tbMsg.getOriginator(), tbMsg, null); - return Futures.immediateFuture(null); + ctx.getTbClusterService().pushMsgToRuleEngine(edge.getTenantId(), tbMsg.getOriginator(), tbMsg, new TbQueueCallback() { + @Override + public void onSuccess(TbQueueMsgMetadata metadata) { + futureToSet.set(null); + } + + @Override + public void onFailure(Throwable t) { + futureToSet.setException(t); + } + }); + return futureToSet; } private ListenableFuture processAttributeDeleteMsg(EntityId entityId, AttributeDeleteMsg attributeDeleteMsg, String entityType) { + SettableFuture futureToSet = SettableFuture.create(); try { String scope = attributeDeleteMsg.getScope(); List attributeNames = attributeDeleteMsg.getAttributeNamesList(); @@ -966,13 +988,23 @@ public final class EdgeGrpcSession implements Closeable { attributeKeys.add(new AttributeKey(scope, attributeName)); } ctx.getTbClusterService().pushMsgToCore(DeviceAttributesEventNotificationMsg.onDelete( - edge.getTenantId(), (DeviceId) entityId, attributeKeys), null); + edge.getTenantId(), (DeviceId) entityId, attributeKeys), new TbQueueCallback() { + @Override + public void onSuccess(TbQueueMsgMetadata metadata) { + futureToSet.set(null); + } + + @Override + public void onFailure(Throwable t) { + futureToSet.setException(t); + } + }); } } catch (Exception e) { log.error("Can't process attribute delete msg [{}]", attributeDeleteMsg, e); return Futures.immediateFailedFuture(new RuntimeException("Can't process attribute delete msg " + attributeDeleteMsg, e)); } - return Futures.immediateFuture(null); + return futureToSet; } private ListenableFuture onDeviceUpdate(DeviceUpdateMsg deviceUpdateMsg) { diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java index 71c56068d1..13702a13ef 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java @@ -30,7 +30,6 @@ import org.thingsboard.server.common.data.page.TextPageLink; import java.util.List; import java.util.Optional; -import java.util.concurrent.Executor; public interface EdgeService { @@ -76,7 +75,7 @@ public interface EdgeService { ListenableFuture> findEdgesByTenantIdAndDashboardId(TenantId tenantId, DashboardId dashboardId); - ListenableFuture> findRelatedEdgeIdsByEntityId(TenantId tenantId, EntityId entityId, Executor executorService); + ListenableFuture> findRelatedEdgeIdsByEntityId(TenantId tenantId, EntityId entityId); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java index 4ce86beea5..6d3aa3e901 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java @@ -67,7 +67,6 @@ import java.util.Collections; import java.util.Comparator; import java.util.List; import java.util.Optional; -import java.util.concurrent.Executor; import java.util.stream.Collectors; import static org.thingsboard.server.common.data.CacheConstants.EDGE_CACHE; @@ -429,7 +428,7 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic }; @Override - public ListenableFuture> findRelatedEdgeIdsByEntityId(TenantId tenantId, EntityId entityId, Executor executorService) { + public ListenableFuture> findRelatedEdgeIdsByEntityId(TenantId tenantId, EntityId entityId) { switch (entityId.getEntityType()) { case DEVICE: case ASSET: @@ -442,11 +441,11 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic } else { return Collections.emptyList(); } - }, executorService); + }, MoreExecutors.directExecutor()); case DASHBOARD: - return convertToEdgeIds(findEdgesByTenantIdAndDashboardId(tenantId, new DashboardId(entityId.getId())), executorService); + return convertToEdgeIds(findEdgesByTenantIdAndDashboardId(tenantId, new DashboardId(entityId.getId()))); case RULE_CHAIN: - return convertToEdgeIds(findEdgesByTenantIdAndRuleChainId(tenantId, new RuleChainId(entityId.getId())), executorService); + return convertToEdgeIds(findEdgesByTenantIdAndRuleChainId(tenantId, new RuleChainId(entityId.getId()))); case USER: User userById = userService.findUserById(tenantId, new UserId(entityId.getId())); TextPageData edges; @@ -455,20 +454,20 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic } else { edges = findEdgesByTenantIdAndCustomerId(tenantId, new CustomerId(entityId.getId()), new TextPageLink(Integer.MAX_VALUE)); } - return convertToEdgeIds(Futures.immediateFuture(edges.getData()), executorService); + return convertToEdgeIds(Futures.immediateFuture(edges.getData())); default: return Futures.immediateFuture(Collections.emptyList()); } } - private ListenableFuture> convertToEdgeIds(ListenableFuture> future, Executor executorService) { + private ListenableFuture> convertToEdgeIds(ListenableFuture> future) { return Futures.transform(future, edges -> { if (edges != null && !edges.isEmpty()) { return edges.stream().map(IdBased::getId).collect(Collectors.toList()); } else { return Collections.emptyList(); } - }, executorService); + }, MoreExecutors.directExecutor()); } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java index b531af7ed3..1c0e44050d 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java @@ -232,7 +232,7 @@ public class TbMsgPushToEdgeNode implements TbNode { TextPageData edgesByTenantId = ctx.getEdgeService().findEdgesByTenantId(tenantId, new TextPageLink(Integer.MAX_VALUE)); return Futures.immediateFuture(edgesByTenantId.getData().stream().map(IdBased::getId).collect(Collectors.toList())); } else { - return ctx.getEdgeService().findRelatedEdgeIdsByEntityId(tenantId, originatorId, ctx.getDbCallbackExecutor()); + return ctx.getEdgeService().findRelatedEdgeIdsByEntityId(tenantId, originatorId); } } From 67fc2f5a197bbe53bf9bd3eb578193961a23e8a5 Mon Sep 17 00:00:00 2001 From: Bohdan Smetaniuk Date: Tue, 25 Aug 2020 14:22:58 +0300 Subject: [PATCH 4/7] added logs --- .../service/edge/rpc/EdgeGrpcSession.java | 44 +++++++++---------- 1 file changed, 21 insertions(+), 23 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java index 99ffe79d9e..670c595e72 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java @@ -951,6 +951,7 @@ public final class EdgeGrpcSession implements Closeable { @Override public void onFailure(Throwable t) { + log.error("Can't process post telemetry [{}]", msg, t); futureToSet.setException(t); } }); @@ -970,6 +971,7 @@ public final class EdgeGrpcSession implements Closeable { @Override public void onFailure(Throwable t) { + log.error("Can't process post attributes [{}]", msg, t); futureToSet.setException(t); } }); @@ -978,31 +980,27 @@ public final class EdgeGrpcSession implements Closeable { private ListenableFuture processAttributeDeleteMsg(EntityId entityId, AttributeDeleteMsg attributeDeleteMsg, String entityType) { SettableFuture futureToSet = SettableFuture.create(); - try { - String scope = attributeDeleteMsg.getScope(); - List attributeNames = attributeDeleteMsg.getAttributeNamesList(); - ctx.getAttributesService().removeAll(edge.getTenantId(), entityId, scope, attributeNames); - if (EntityType.DEVICE.name().equals(entityType)) { - Set attributeKeys = new HashSet<>(); - for (String attributeName : attributeNames) { - attributeKeys.add(new AttributeKey(scope, attributeName)); + String scope = attributeDeleteMsg.getScope(); + List attributeNames = attributeDeleteMsg.getAttributeNamesList(); + ctx.getAttributesService().removeAll(edge.getTenantId(), entityId, scope, attributeNames); + if (EntityType.DEVICE.name().equals(entityType)) { + Set attributeKeys = new HashSet<>(); + for (String attributeName : attributeNames) { + attributeKeys.add(new AttributeKey(scope, attributeName)); + } + ctx.getTbClusterService().pushMsgToCore(DeviceAttributesEventNotificationMsg.onDelete( + edge.getTenantId(), (DeviceId) entityId, attributeKeys), new TbQueueCallback() { + @Override + public void onSuccess(TbQueueMsgMetadata metadata) { + futureToSet.set(null); } - ctx.getTbClusterService().pushMsgToCore(DeviceAttributesEventNotificationMsg.onDelete( - edge.getTenantId(), (DeviceId) entityId, attributeKeys), new TbQueueCallback() { - @Override - public void onSuccess(TbQueueMsgMetadata metadata) { - futureToSet.set(null); - } - @Override - public void onFailure(Throwable t) { - futureToSet.setException(t); - } - }); - } - } catch (Exception e) { - log.error("Can't process attribute delete msg [{}]", attributeDeleteMsg, e); - return Futures.immediateFailedFuture(new RuntimeException("Can't process attribute delete msg " + attributeDeleteMsg, e)); + @Override + public void onFailure(Throwable t) { + log.error("Can't process attribute delete msg [{}]", attributeDeleteMsg, t); + futureToSet.setException(t); + } + }); } return futureToSet; } From 883f8fff4c0db557ea7ef10743164c28026750fd Mon Sep 17 00:00:00 2001 From: Bohdan Smetaniuk Date: Tue, 25 Aug 2020 19:23:49 +0300 Subject: [PATCH 5/7] fixed bug with connection to first node --- .../edge/rpc/constructor/RuleChainUpdateMsgConstructor.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RuleChainUpdateMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RuleChainUpdateMsgConstructor.java index 52e487d973..1d81763a1f 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RuleChainUpdateMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RuleChainUpdateMsgConstructor.java @@ -68,6 +68,8 @@ public class RuleChainUpdateMsgConstructor { .addAllRuleChainConnections(constructRuleChainConnections(ruleChainMetaData.getRuleChainConnections())); if (ruleChainMetaData.getFirstNodeIndex() != null) { builder.setFirstNodeIndex(ruleChainMetaData.getFirstNodeIndex()); + } else { + builder.setFirstNodeIndex(-1); } builder.setMsgType(msgType); return builder.build(); From 4ac8c0675b61256926f4e2cad49378a6e52467ea Mon Sep 17 00:00:00 2001 From: deaflynx Date: Thu, 27 Aug 2020 08:46:32 +0300 Subject: [PATCH 6/7] Fixed fetch edge rule chains --- ui/src/app/edge/edge.controller.js | 2 +- ui/src/app/edge/set-root-rule-chain-to-edges.controller.js | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/ui/src/app/edge/edge.controller.js b/ui/src/app/edge/edge.controller.js index 7e57fce6a6..94f076c1af 100644 --- a/ui/src/app/edge/edge.controller.js +++ b/ui/src/app/edge/edge.controller.js @@ -552,7 +552,7 @@ export function EdgeController($rootScope, userService, edgeService, customerSer $event.stopPropagation(); } var pageSize = 10; - ruleChainService.getRuleChains({limit: pageSize, textSearch: ''}).then( + ruleChainService.getEdgesRuleChains({limit: pageSize, textSearch: ''}).then( function success(_ruleChains) { var ruleChains = { pageSize: pageSize, diff --git a/ui/src/app/edge/set-root-rule-chain-to-edges.controller.js b/ui/src/app/edge/set-root-rule-chain-to-edges.controller.js index a474fde6a7..a8d134b20e 100644 --- a/ui/src/app/edge/set-root-rule-chain-to-edges.controller.js +++ b/ui/src/app/edge/set-root-rule-chain-to-edges.controller.js @@ -53,7 +53,7 @@ export default function SetRootRuleChainToEdgesController(ruleChainService, edge fetchMoreItems_: function () { if (vm.ruleChains.hasNext && !vm.ruleChains.pending) { vm.ruleChains.pending = true; - ruleChainService.getRuleChains(vm.ruleChains.nextPageLink).then( + ruleChainService.getEdgesRuleChains(vm.ruleChains.nextPageLink).then( function success(ruleChains) { vm.ruleChains.data = vm.ruleChains.data.concat(ruleChains.data); vm.ruleChains.nextPageLink = ruleChains.nextPageLink; From aeaf2ed6c6f167ae37b5eca22d3a41cdd08663d6 Mon Sep 17 00:00:00 2001 From: deaflynx Date: Thu, 27 Aug 2020 11:05:13 +0300 Subject: [PATCH 7/7] Added manage customer edges --- .../app/customer/customer-fieldset.tpl.html | 1 + ui/src/app/customer/customer.controller.js | 22 ++++++++++++++++++ ui/src/app/customer/customer.directive.js | 1 + ui/src/app/customer/customers.tpl.html | 1 + ui/src/app/edge/edge.routes.js | 23 +++++++++++++++++++ ui/src/app/locale/locale.constant-en_US.json | 3 +++ 6 files changed, 51 insertions(+) diff --git a/ui/src/app/customer/customer-fieldset.tpl.html b/ui/src/app/customer/customer-fieldset.tpl.html index 11b2d58fb8..a23959ef96 100644 --- a/ui/src/app/customer/customer-fieldset.tpl.html +++ b/ui/src/app/customer/customer-fieldset.tpl.html @@ -19,6 +19,7 @@ {{ 'customer.manage-assets' | translate }} {{ 'customer.manage-devices' | translate }} {{ 'customer.manage-dashboards' | translate }} +{{ 'customer.manage-edges' | translate }} {{ 'customer.delete' | translate }}
diff --git a/ui/src/app/customer/customer.controller.js b/ui/src/app/customer/customer.controller.js index 42a3801efa..c61962f9f3 100644 --- a/ui/src/app/customer/customer.controller.js +++ b/ui/src/app/customer/customer.controller.js @@ -77,6 +77,20 @@ export default function CustomerController(customerService, $state, $stateParams }, icon: "dashboard" }, + { + onAction: function ($event, item) { + openCustomerEdges($event, item); + }, + name: function() { return $translate.instant('edge.edges') }, + details: function(customer) { + if (customer && customer.additionalInfo && customer.additionalInfo.isPublic) { + return $translate.instant('customer.manage-public-edges') + } else { + return $translate.instant('customer.manage-customer-edges') + } + }, + icon: "router" + }, { onAction: function ($event, item) { vm.grid.deleteItem($event, item); @@ -147,6 +161,7 @@ export default function CustomerController(customerService, $state, $stateParams vm.openCustomerAssets = openCustomerAssets; vm.openCustomerDevices = openCustomerDevices; vm.openCustomerDashboards = openCustomerDashboards; + vm.openCustomerEdges = openCustomerEdges; function deleteCustomerTitle(customer) { return $translate.instant('customer.delete-customer-title', {customerTitle: customer.title}); @@ -216,4 +231,11 @@ export default function CustomerController(customerService, $state, $stateParams $state.go('home.customers.dashboards', {customerId: customer.id.id}); } + function openCustomerEdges($event, customer) { + if ($event) { + $event.stopPropagation(); + } + $state.go('home.customers.edges', {customerId: customer.id.id}); + } + } diff --git a/ui/src/app/customer/customer.directive.js b/ui/src/app/customer/customer.directive.js index ff1befaf9e..33245654bb 100644 --- a/ui/src/app/customer/customer.directive.js +++ b/ui/src/app/customer/customer.directive.js @@ -55,6 +55,7 @@ export default function CustomerDirective($compile, $templateCache, $translate, onManageAssets: '&', onManageDevices: '&', onManageDashboards: '&', + onManageEdges: '&', onDeleteCustomer: '&' } }; diff --git a/ui/src/app/customer/customers.tpl.html b/ui/src/app/customer/customers.tpl.html index 396fc6734b..755425a2ff 100644 --- a/ui/src/app/customer/customers.tpl.html +++ b/ui/src/app/customer/customers.tpl.html @@ -29,6 +29,7 @@ on-manage-assets="vm.openCustomerAssets(event, vm.grid.detailsConfig.currentItem)" on-manage-devices="vm.openCustomerDevices(event, vm.grid.detailsConfig.currentItem)" on-manage-dashboards="vm.openCustomerDashboards(event, vm.grid.detailsConfig.currentItem)" + on-manage-edges="vm.openCustomerEdges(event, vm.grid.detailsConfig.currentItem)" on-delete-customer="vm.grid.deleteItem(event, vm.grid.detailsConfig.currentItem)"> diff --git a/ui/src/app/edge/edge.routes.js b/ui/src/app/edge/edge.routes.js index 1afd319609..e3443d1211 100644 --- a/ui/src/app/edge/edge.routes.js +++ b/ui/src/app/edge/edge.routes.js @@ -161,5 +161,28 @@ export default function EdgeRoutes($stateProvider, types) { ncyBreadcrumb: { label: '{"icon": "dashboard", "label": "{{ vm.dashboard.title }}", "translate": "false"}' } + }) + .state('home.customers.edges', { + url: '/:customerId/edges', + params: {'topIndex': 0}, + module: 'private', + auth: ['TENANT_ADMIN'], + views: { + "content@home": { + templateUrl: edgesTemplate, + controllerAs: 'vm', + controller: 'EdgeController' + } + }, + data: { + edgesType: 'customer', + searchEnabled: true, + searchByEntitySubtype: true, + searchEntityType: types.entityType.edge, + pageTitle: 'customer.edges' + }, + ncyBreadcrumb: { + label: '{"icon": "router", "label": "{{ vm.customerEdgesTitle }}", "translate": "false"}' + } }); } diff --git a/ui/src/app/locale/locale.constant-en_US.json b/ui/src/app/locale/locale.constant-en_US.json index f8d1c8568a..a86d4b2ffd 100644 --- a/ui/src/app/locale/locale.constant-en_US.json +++ b/ui/src/app/locale/locale.constant-en_US.json @@ -441,6 +441,7 @@ "manage-assets": "Manage assets", "manage-devices": "Manage devices", "manage-dashboards": "Manage dashboards", + "manage-edges": "Manage edges", "title": "Title", "title-required": "Title is required.", "description": "Description", @@ -812,6 +813,8 @@ "assign-edges-text": "Assign { count, plural, 1 {1 edge} other {# edges} } to customer", "unassign-edge-title": "Are you sure you want to unassign the edge '{{edgeName}}'?", "unassign-edge-text": "After the confirmation the edge will be unassigned and won't be accessible by the customer.", + "unassign-edges-title": "Are you sure you want to unassign { count, plural, 1 {1 edge} other {# edges} }?", + "unassign-edges-text": "After the confirmation all selected edges will be unassigned and won't be accessible by the customer.", "make-public": "Make edge public", "make-public-edge-title": "Are you sure you want to make the edge '{{edgeName}}' public?", "make-public-edge-text": "After the confirmation the edge and all its data will be made public and accessible by others.",