From 3a416131c6f93194746f970d82eebf1bd465785e Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Thu, 11 Jun 2020 19:25:31 +0300 Subject: [PATCH] Added relation update msg. Push relation to the edge on sync --- .../service/edge/EdgeContextComponent.java | 9 +- .../service/edge/rpc/EdgeGrpcSession.java | 17 ++- .../RelationUpdateMsgConstructor.java | 45 +++++++ ...rvice.java => DefaultSyncEdgeService.java} | 118 ++++++++++++++---- ...tEdgeService.java => SyncEdgeService.java} | 6 +- .../common/data/edge/EdgeQueueEntityType.java | 2 +- common/edge-api/src/main/proto/edge.proto | 14 +++ 7 files changed, 182 insertions(+), 29 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RelationUpdateMsgConstructor.java rename application/src/main/java/org/thingsboard/server/service/edge/rpc/init/{DefaultInitEdgeService.java => DefaultSyncEdgeService.java} (69%) rename application/src/main/java/org/thingsboard/server/service/edge/rpc/init/{InitEdgeService.java => SyncEdgeService.java} (85%) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java index f936ee82a5..6959a522bf 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java @@ -36,8 +36,9 @@ import org.thingsboard.server.service.edge.rpc.constructor.AssetUpdateMsgConstru import org.thingsboard.server.service.edge.rpc.constructor.DashboardUpdateMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.DeviceUpdateMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.EntityViewUpdateMsgConstructor; +import org.thingsboard.server.service.edge.rpc.constructor.RelationUpdateMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.UserUpdateMsgConstructor; -import org.thingsboard.server.service.edge.rpc.init.InitEdgeService; +import org.thingsboard.server.service.edge.rpc.init.SyncEdgeService; import org.thingsboard.server.service.edge.rpc.constructor.RuleChainUpdateMsgConstructor; import org.thingsboard.server.service.queue.TbClusterService; import org.thingsboard.server.service.state.DeviceStateService; @@ -100,7 +101,7 @@ public class EdgeContextComponent { @Lazy @Autowired - private InitEdgeService initEdgeService; + private SyncEdgeService syncEdgeService; @Lazy @Autowired @@ -130,6 +131,10 @@ public class EdgeContextComponent { @Autowired private UserUpdateMsgConstructor userUpdateMsgConstructor; + @Lazy + @Autowired + private RelationUpdateMsgConstructor relationUpdateMsgConstructor; + @Lazy @Autowired private EdgeEventStorageSettings edgeEventStorageSettings; 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 e6b68f486a..d029b5ff20 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 @@ -141,7 +141,7 @@ public final class EdgeGrpcSession implements Closeable { outputStream.onError(new RuntimeException(responseMsg.getErrorMsg())); } if (ConnectResponseCode.ACCEPTED == responseMsg.getResponseCode()) { - ctx.getInitEdgeService().init(edge, outputStream); + ctx.getSyncEdgeService().sync(edge, outputStream); } } if (connected) { @@ -360,6 +360,10 @@ public final class EdgeGrpcSession implements Closeable { User user = objectMapper.readValue(entry.getData(), User.class); onUserUpdated(msgType, user); break; + case RELATION: + EntityRelation entityRelation = objectMapper.readValue(entry.getData(), EntityRelation.class); + onEntityRelationUpdated(msgType, entityRelation); + break; } } @@ -463,6 +467,15 @@ public final class EdgeGrpcSession implements Closeable { .build()); } + private void onEntityRelationUpdated(UpdateMsgType msgType, EntityRelation entityRelation) { + EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg.newBuilder() + .setRelationUpdateMsg(ctx.getRelationUpdateMsgConstructor().constructRelationUpdatedMsg(msgType, entityRelation)) + .build(); + outputStream.onNext(ResponseMsg.newBuilder() + .setEntityUpdateMsg(entityUpdateMsg) + .build()); + } + private UpdateMsgType getResponseMsgType(String msgType) { if (msgType.equals(SessionMsgType.POST_TELEMETRY_REQUEST.name()) || msgType.equals(SessionMsgType.POST_ATTRIBUTES_REQUEST.name()) || @@ -548,7 +561,7 @@ public final class EdgeGrpcSession implements Closeable { } if (uplinkMsg.getRuleChainMetadataRequestMsgList() != null && !uplinkMsg.getRuleChainMetadataRequestMsgList().isEmpty()) { for (RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg : uplinkMsg.getRuleChainMetadataRequestMsgList()) { - ctx.getInitEdgeService().initRuleChainMetadata(edge, ruleChainMetadataRequestMsg, outputStream); + ctx.getSyncEdgeService().syncRuleChainMetadata(edge, ruleChainMetadataRequestMsg, outputStream); } } } catch (Exception e) { diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RelationUpdateMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RelationUpdateMsgConstructor.java new file mode 100644 index 0000000000..1c0af319f3 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RelationUpdateMsgConstructor.java @@ -0,0 +1,45 @@ +/** + * 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.service.edge.rpc.constructor; + +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.relation.EntityRelation; +import org.thingsboard.server.dao.util.mapping.JacksonUtil; +import org.thingsboard.server.gen.edge.RelationUpdateMsg; +import org.thingsboard.server.gen.edge.UpdateMsgType; + +@Component +@Slf4j +public class RelationUpdateMsgConstructor { + + public RelationUpdateMsg constructRelationUpdatedMsg(UpdateMsgType msgType, EntityRelation entityRelation) { + RelationUpdateMsg.Builder builder = RelationUpdateMsg.newBuilder() + .setMsgType(msgType) + .setFromIdMSB(entityRelation.getFrom().getId().getMostSignificantBits()) + .setFromIdLSB(entityRelation.getFrom().getId().getLeastSignificantBits()) + .setFromEntityType(entityRelation.getFrom().getEntityType().name()) + .setToIdMSB(entityRelation.getTo().getId().getMostSignificantBits()) + .setToIdLSB(entityRelation.getTo().getId().getLeastSignificantBits()) + .setToEntityType(entityRelation.getTo().getEntityType().name()) + .setType(entityRelation.getType()) + .setAdditionalInfo(JacksonUtil.toString(entityRelation.getAdditionalInfo())); + if (entityRelation.getTypeGroup() != null) { + builder.setTypeGroup(entityRelation.getTypeGroup().name()); + } + return builder.build(); + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/init/DefaultInitEdgeService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/init/DefaultSyncEdgeService.java similarity index 69% rename from application/src/main/java/org/thingsboard/server/service/edge/rpc/init/DefaultInitEdgeService.java rename to application/src/main/java/org/thingsboard/server/service/edge/rpc/init/DefaultSyncEdgeService.java index 192fb8395b..d51474046e 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/init/DefaultInitEdgeService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/init/DefaultSyncEdgeService.java @@ -16,6 +16,8 @@ package org.thingsboard.server.service.edge.rpc.init; import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; +import com.google.common.util.concurrent.MoreExecutors; import io.grpc.stub.StreamObserver; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; @@ -24,25 +26,31 @@ import org.thingsboard.server.common.data.Dashboard; import org.thingsboard.server.common.data.DashboardInfo; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.EntityView; +import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.edge.Edge; +import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.RuleChainId; -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.EntityRelationsQuery; +import org.thingsboard.server.common.data.relation.EntitySearchDirection; +import org.thingsboard.server.common.data.relation.RelationsSearchParameters; import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.common.data.rule.RuleChainMetaData; import org.thingsboard.server.dao.asset.AssetService; import org.thingsboard.server.dao.dashboard.DashboardService; import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.entityview.EntityViewService; +import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.dao.rule.RuleChainService; import org.thingsboard.server.gen.edge.AssetUpdateMsg; import org.thingsboard.server.gen.edge.DashboardUpdateMsg; import org.thingsboard.server.gen.edge.DeviceUpdateMsg; import org.thingsboard.server.gen.edge.EntityUpdateMsg; import org.thingsboard.server.gen.edge.EntityViewUpdateMsg; +import org.thingsboard.server.gen.edge.RelationUpdateMsg; import org.thingsboard.server.gen.edge.ResponseMsg; import org.thingsboard.server.gen.edge.RuleChainMetadataRequestMsg; import org.thingsboard.server.gen.edge.RuleChainMetadataUpdateMsg; @@ -52,18 +60,25 @@ import org.thingsboard.server.service.edge.rpc.constructor.AssetUpdateMsgConstru import org.thingsboard.server.service.edge.rpc.constructor.DashboardUpdateMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.DeviceUpdateMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.EntityViewUpdateMsgConstructor; +import org.thingsboard.server.service.edge.rpc.constructor.RelationUpdateMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.RuleChainUpdateMsgConstructor; +import java.util.ArrayList; +import java.util.HashSet; +import java.util.List; +import java.util.Set; import java.util.UUID; -import java.util.concurrent.Future; @Service @Slf4j -public class DefaultInitEdgeService implements InitEdgeService { +public class DefaultSyncEdgeService implements SyncEdgeService { @Autowired private RuleChainService ruleChainService; + @Autowired + private RelationService relationService; + @Autowired private DeviceService deviceService; @@ -79,6 +94,9 @@ public class DefaultInitEdgeService implements InitEdgeService { @Autowired private RuleChainUpdateMsgConstructor ruleChainUpdateMsgConstructor; + @Autowired + private RelationUpdateMsgConstructor relationUpdateMsgConstructor; + @Autowired private DeviceUpdateMsgConstructor deviceUpdateMsgConstructor; @@ -92,15 +110,68 @@ public class DefaultInitEdgeService implements InitEdgeService { private DashboardUpdateMsgConstructor dashboardUpdateMsgConstructor; @Override - public void init(Edge edge, StreamObserver outputStream) { - initRuleChains(edge, outputStream); - initDevices(edge, outputStream); - initAssets(edge, outputStream); - initEntityViews(edge, outputStream); - initDashboards(edge, outputStream); + public void sync(Edge edge, StreamObserver outputStream) { + Set pushedEntityIds = new HashSet<>(); + syncRuleChains(edge, pushedEntityIds, outputStream); + syncDevices(edge, pushedEntityIds, outputStream); + syncAssets(edge, pushedEntityIds, outputStream); + syncEntityViews(edge, pushedEntityIds, outputStream); + syncDashboards(edge, pushedEntityIds, outputStream); + syncRelations(edge, pushedEntityIds, outputStream); + } + + private void syncRelations(Edge edge, Set pushedEntityIds, StreamObserver outputStream) { + if (!pushedEntityIds.isEmpty()) { + List>> futures = new ArrayList<>(); + for (EntityId entityId : pushedEntityIds) { + futures.add(syncRelations(edge, entityId, EntitySearchDirection.FROM)); + futures.add(syncRelations(edge, entityId, EntitySearchDirection.TO)); + } + ListenableFuture>> relationsListFuture = Futures.allAsList(futures); + Futures.transform(relationsListFuture, relationsList -> { + try { + Set uniqueEntityRelations = new HashSet<>(); + if (!relationsList.isEmpty()) { + for (List entityRelations : relationsList) { + if (!entityRelations.isEmpty()) { + uniqueEntityRelations.addAll(entityRelations); + } + } + } + if (!uniqueEntityRelations.isEmpty()) { + log.trace("[{}] [{}] relation(s) are going to be pushed to edge.", edge.getId(), uniqueEntityRelations.size()); + for (EntityRelation relation : uniqueEntityRelations) { + try { + RelationUpdateMsg relationUpdateMsg = + relationUpdateMsgConstructor.constructRelationUpdatedMsg( + UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, + relation); + EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg.newBuilder() + .setRelationUpdateMsg(relationUpdateMsg) + .build(); + outputStream.onNext(ResponseMsg.newBuilder() + .setEntityUpdateMsg(entityUpdateMsg) + .build()); + } catch (Exception e) { + log.error("Exception during loading relation [{}] to edge on init!", relation, e); + } + } + } + } catch (Exception e) { + log.error("Exception during loading relation(s) to edge on init!", e); + } + return null; + }, MoreExecutors.directExecutor()); + } + } + + private ListenableFuture> syncRelations(Edge edge, EntityId entityId, EntitySearchDirection direction) { + EntityRelationsQuery query = new EntityRelationsQuery(); + query.setParameters(new RelationsSearchParameters(entityId, direction, -1, false)); + return relationService.findByQuery(edge.getTenantId(), query); } - private void initDevices(Edge edge, StreamObserver outputStream) { + private void syncDevices(Edge edge, Set pushedEntityIds, StreamObserver outputStream) { try { TimePageLink pageLink = new TimePageLink(100); TimePageData pageData; @@ -119,6 +190,7 @@ public class DefaultInitEdgeService implements InitEdgeService { outputStream.onNext(ResponseMsg.newBuilder() .setEntityUpdateMsg(entityUpdateMsg) .build()); + pushedEntityIds.add(device.getId()); } } if (pageData.hasNext()) { @@ -126,11 +198,11 @@ public class DefaultInitEdgeService implements InitEdgeService { } } while (pageData.hasNext()); } catch (Exception e) { - log.error("Exception during loading edge device(s) on init!"); + log.error("Exception during loading edge device(s) on init!", e); } } - private void initAssets(Edge edge, StreamObserver outputStream) { + private void syncAssets(Edge edge, Set pushedEntityIds, StreamObserver outputStream) { try { TimePageLink pageLink = new TimePageLink(100); TimePageData pageData; @@ -149,6 +221,7 @@ public class DefaultInitEdgeService implements InitEdgeService { outputStream.onNext(ResponseMsg.newBuilder() .setEntityUpdateMsg(entityUpdateMsg) .build()); + pushedEntityIds.add(asset.getId()); } } if (pageData.hasNext()) { @@ -156,11 +229,11 @@ public class DefaultInitEdgeService implements InitEdgeService { } } while (pageData.hasNext()); } catch (Exception e) { - log.error("Exception during loading edge asset(s) on init!"); + log.error("Exception during loading edge asset(s) on init!", e); } } - private void initEntityViews(Edge edge, StreamObserver outputStream) { + private void syncEntityViews(Edge edge, Set pushedEntityIds, StreamObserver outputStream) { try { TimePageLink pageLink = new TimePageLink(100); TimePageData pageData; @@ -179,6 +252,7 @@ public class DefaultInitEdgeService implements InitEdgeService { outputStream.onNext(ResponseMsg.newBuilder() .setEntityUpdateMsg(entityUpdateMsg) .build()); + pushedEntityIds.add(entityView.getId()); } } if (pageData.hasNext()) { @@ -186,11 +260,11 @@ public class DefaultInitEdgeService implements InitEdgeService { } } while (pageData.hasNext()); } catch (Exception e) { - log.error("Exception during loading edge entity view(s) on init!"); + log.error("Exception during loading edge entity view(s) on init!", e); } } - private void initDashboards(Edge edge, StreamObserver outputStream) { + private void syncDashboards(Edge edge, Set pushedEntityIds, StreamObserver outputStream) { try { TimePageLink pageLink = new TimePageLink(100); TimePageData pageData; @@ -210,6 +284,7 @@ public class DefaultInitEdgeService implements InitEdgeService { outputStream.onNext(ResponseMsg.newBuilder() .setEntityUpdateMsg(entityUpdateMsg) .build()); + pushedEntityIds.add(dashboard.getId()); } } if (pageData.hasNext()) { @@ -217,11 +292,11 @@ public class DefaultInitEdgeService implements InitEdgeService { } } while (pageData.hasNext()); } catch (Exception e) { - log.error("Exception during loading edge dashboard(s) on init!"); + log.error("Exception during loading edge dashboard(s) on init!", e); } } - private void initRuleChains(Edge edge, StreamObserver outputStream) { + private void syncRuleChains(Edge edge, Set pushedEntityIds, StreamObserver outputStream) { try { TimePageLink pageLink = new TimePageLink(100); TimePageData pageData; @@ -241,6 +316,7 @@ public class DefaultInitEdgeService implements InitEdgeService { outputStream.onNext(ResponseMsg.newBuilder() .setEntityUpdateMsg(entityUpdateMsg) .build()); + pushedEntityIds.add(ruleChain.getId()); } } if (pageData.hasNext()) { @@ -248,12 +324,12 @@ public class DefaultInitEdgeService implements InitEdgeService { } } while (pageData.hasNext()); } catch (Exception e) { - log.error("Exception during loading edge rule chain(s) on init!"); + log.error("Exception during loading edge rule chain(s) on init!", e); } } @Override - public void initRuleChainMetadata(Edge edge, RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg, StreamObserver outputStream) { + public void syncRuleChainMetadata(Edge edge, RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg, StreamObserver outputStream) { if (ruleChainMetadataRequestMsg.getRuleChainIdMSB() != 0 && ruleChainMetadataRequestMsg.getRuleChainIdLSB() != 0) { RuleChainId ruleChainId = new RuleChainId(new UUID(ruleChainMetadataRequestMsg.getRuleChainIdMSB(), ruleChainMetadataRequestMsg.getRuleChainIdLSB())); RuleChainMetaData ruleChainMetaData = ruleChainService.loadRuleChainMetaData(edge.getTenantId(), ruleChainId); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/init/InitEdgeService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/init/SyncEdgeService.java similarity index 85% rename from application/src/main/java/org/thingsboard/server/service/edge/rpc/init/InitEdgeService.java rename to application/src/main/java/org/thingsboard/server/service/edge/rpc/init/SyncEdgeService.java index 8aeb89bf23..c83a9ec3b0 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/init/InitEdgeService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/init/SyncEdgeService.java @@ -20,9 +20,9 @@ import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.gen.edge.ResponseMsg; import org.thingsboard.server.gen.edge.RuleChainMetadataRequestMsg; -public interface InitEdgeService { +public interface SyncEdgeService { - void init(Edge edge, StreamObserver outputStream); + void sync(Edge edge, StreamObserver outputStream); - void initRuleChainMetadata(Edge edge, RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg, StreamObserver outputStream); + void syncRuleChainMetadata(Edge edge, RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg, StreamObserver outputStream); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeQueueEntityType.java b/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeQueueEntityType.java index accae613c1..7ba316a529 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeQueueEntityType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeQueueEntityType.java @@ -16,5 +16,5 @@ package org.thingsboard.server.common.data.edge; public enum EdgeQueueEntityType { - DASHBOARD, ASSET, DEVICE, ENTITY_VIEW, ALARM, RULE_CHAIN, RULE_CHAIN_METADATA, EDGE, USER, CUSTOMER + DASHBOARD, ASSET, DEVICE, ENTITY_VIEW, ALARM, RULE_CHAIN, RULE_CHAIN_METADATA, EDGE, USER, CUSTOMER, RELATION } diff --git a/common/edge-api/src/main/proto/edge.proto b/common/edge-api/src/main/proto/edge.proto index 9c38c71baa..2256acccf8 100644 --- a/common/edge-api/src/main/proto/edge.proto +++ b/common/edge-api/src/main/proto/edge.proto @@ -54,6 +54,7 @@ message EntityUpdateMsg { AlarmUpdateMsg alarmUpdateMsg = 7; UserUpdateMsg userUpdateMsg = 8; CustomerUpdateMsg customerUpdateMsg = 9; + RelationUpdateMsg relationUpdateMsg = 10; } enum RequestMsgType { @@ -222,6 +223,19 @@ message CustomerUpdateMsg { string additionalInfo = 13; } +message RelationUpdateMsg { + UpdateMsgType msgType = 1; + int64 fromIdMSB = 2; + int64 fromIdLSB = 3; + string fromEntityType = 4; + int64 toIdMSB = 5; + int64 toIdLSB = 6; + string toEntityType = 7; + string type = 8; + string typeGroup = 9; + string additionalInfo = 10; +} + message UserUpdateMsg { UpdateMsgType msgType = 1; int64 idMSB = 2;