Browse Source

Added relation update msg. Push relation to the edge on sync

pull/2436/head
Volodymyr Babak 6 years ago
parent
commit
3a416131c6
  1. 9
      application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java
  2. 17
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  3. 45
      application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RelationUpdateMsgConstructor.java
  4. 118
      application/src/main/java/org/thingsboard/server/service/edge/rpc/init/DefaultSyncEdgeService.java
  5. 6
      application/src/main/java/org/thingsboard/server/service/edge/rpc/init/SyncEdgeService.java
  6. 2
      common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeQueueEntityType.java
  7. 14
      common/edge-api/src/main/proto/edge.proto

9
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;

17
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) {

45
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();
}
}

118
application/src/main/java/org/thingsboard/server/service/edge/rpc/init/DefaultInitEdgeService.java → 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<ResponseMsg> outputStream) {
initRuleChains(edge, outputStream);
initDevices(edge, outputStream);
initAssets(edge, outputStream);
initEntityViews(edge, outputStream);
initDashboards(edge, outputStream);
public void sync(Edge edge, StreamObserver<ResponseMsg> outputStream) {
Set<EntityId> 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<EntityId> pushedEntityIds, StreamObserver<ResponseMsg> outputStream) {
if (!pushedEntityIds.isEmpty()) {
List<ListenableFuture<List<EntityRelation>>> futures = new ArrayList<>();
for (EntityId entityId : pushedEntityIds) {
futures.add(syncRelations(edge, entityId, EntitySearchDirection.FROM));
futures.add(syncRelations(edge, entityId, EntitySearchDirection.TO));
}
ListenableFuture<List<List<EntityRelation>>> relationsListFuture = Futures.allAsList(futures);
Futures.transform(relationsListFuture, relationsList -> {
try {
Set<EntityRelation> uniqueEntityRelations = new HashSet<>();
if (!relationsList.isEmpty()) {
for (List<EntityRelation> 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<List<EntityRelation>> 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<ResponseMsg> outputStream) {
private void syncDevices(Edge edge, Set<EntityId> pushedEntityIds, StreamObserver<ResponseMsg> outputStream) {
try {
TimePageLink pageLink = new TimePageLink(100);
TimePageData<Device> 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<ResponseMsg> outputStream) {
private void syncAssets(Edge edge, Set<EntityId> pushedEntityIds, StreamObserver<ResponseMsg> outputStream) {
try {
TimePageLink pageLink = new TimePageLink(100);
TimePageData<Asset> 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<ResponseMsg> outputStream) {
private void syncEntityViews(Edge edge, Set<EntityId> pushedEntityIds, StreamObserver<ResponseMsg> outputStream) {
try {
TimePageLink pageLink = new TimePageLink(100);
TimePageData<EntityView> 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<ResponseMsg> outputStream) {
private void syncDashboards(Edge edge, Set<EntityId> pushedEntityIds, StreamObserver<ResponseMsg> outputStream) {
try {
TimePageLink pageLink = new TimePageLink(100);
TimePageData<DashboardInfo> 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<ResponseMsg> outputStream) {
private void syncRuleChains(Edge edge, Set<EntityId> pushedEntityIds, StreamObserver<ResponseMsg> outputStream) {
try {
TimePageLink pageLink = new TimePageLink(100);
TimePageData<RuleChain> 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<ResponseMsg> outputStream) {
public void syncRuleChainMetadata(Edge edge, RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg, StreamObserver<ResponseMsg> 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);

6
application/src/main/java/org/thingsboard/server/service/edge/rpc/init/InitEdgeService.java → 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<ResponseMsg> outputStream);
void sync(Edge edge, StreamObserver<ResponseMsg> outputStream);
void initRuleChainMetadata(Edge edge, RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg, StreamObserver<ResponseMsg> outputStream);
void syncRuleChainMetadata(Edge edge, RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg, StreamObserver<ResponseMsg> outputStream);
}

2
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
}

14
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;

Loading…
Cancel
Save