diff --git a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java index 42b66a73e8..5cbef17deb 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java @@ -54,6 +54,7 @@ import org.thingsboard.server.dao.cassandra.CassandraCluster; import org.thingsboard.server.dao.customer.CustomerService; import org.thingsboard.server.dao.dashboard.DashboardService; import org.thingsboard.server.dao.device.DeviceService; +import org.thingsboard.server.dao.edge.EdgeService; import org.thingsboard.server.dao.entityview.EntityViewService; import org.thingsboard.server.dao.event.EventService; import org.thingsboard.server.dao.nosql.CassandraBufferedRateExecutor; @@ -245,6 +246,11 @@ public class ActorSystemContext { @Getter private RuleChainTransactionService ruleChainTransactionService; + @Lazy + @Autowired + @Getter + private EdgeService edgeService; + @Value("${cluster.partition_id}") @Getter private long queuePartitionId; diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java index 44b6f3b6c5..d8468f6765 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java @@ -28,7 +28,9 @@ import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.actors.device.DeviceActorToRuleEngineMsg; import org.thingsboard.server.actors.service.DefaultActorService; import org.thingsboard.server.actors.shared.ComponentMsgProcessor; +import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.ShortEdgeInfo; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.RuleNodeId; @@ -43,6 +45,7 @@ import org.thingsboard.server.common.msg.cluster.ClusterEventMsg; import org.thingsboard.server.common.msg.cluster.ServerAddress; import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.system.ServiceToRuleEngineMsg; +import org.thingsboard.server.dao.edge.EdgeService; import org.thingsboard.server.dao.rule.RuleChainService; import java.util.ArrayList; @@ -65,6 +68,7 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor nodeActors; private final Map> nodeRoutes; private final RuleChainService service; + private final EdgeService edgeService; private RuleNodeId firstId; private RuleNodeCtx firstNode; @@ -79,6 +83,7 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor(); this.nodeRoutes = new HashMap<>(); this.service = systemContext.getRuleChainService(); + this.edgeService = systemContext.getEdgeService(); this.ruleChainName = ruleChainId.toString(); } @@ -326,6 +331,19 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor sessions = new ConcurrentHashMap<>(); + private static final ObjectMapper objectMapper = new ObjectMapper(); @Value("${edges.rpc.port}") private int rpcPort; @@ -56,8 +72,22 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase { @Autowired private EdgeContextComponent ctx; + @Autowired + private EdgeService edgeService; + + @Autowired + private AssetService assetService; + + @Autowired + private DeviceService deviceService; + + @Autowired + private AttributesService attributesService; + private Server server; + private ExecutorService executor; + @PostConstruct public void init() { log.info("Initializing Edge RPC service!"); @@ -81,9 +111,10 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase { throw new RuntimeException("Failed to start Edge RPC server!"); } log.info("Edge RPC service initialized!"); + executor = Executors.newSingleThreadExecutor(); + processHandleMessages(); } - @PreDestroy public void destroy() { if (server != null) { @@ -92,14 +123,28 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase { } @Override - public StreamObserver handleMsgs(StreamObserver responseObserver) { - return new EdgeGrpcSession(ctx, responseObserver, this::onEdgeConnect, this::onEdgeDisconnect).getInputStream(); + public StreamObserver handleMsgs(StreamObserver outputStream) { + return new EdgeGrpcSession(ctx, outputStream, this::onEdgeConnect, this::onEdgeDisconnect, edgeService, assetService, deviceService, attributesService, objectMapper).getInputStream(); } private void onEdgeConnect(EdgeId edgeId, EdgeGrpcSession edgeGrpcSession) { sessions.put(edgeId, edgeGrpcSession); } + private void processHandleMessages() { + executor.submit(() -> { + while (!Thread.interrupted()) { + try { + for (EdgeGrpcSession session : sessions.values()) { + session.processHandleMessages(); + } + } catch (Exception e) { + log.warn("Failed to process messages handling!", e); + } + } + }); + } + private void onEdgeDisconnect(EdgeId edgeId) { sessions.remove(edgeId); } 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 db395cbe83..5c683c3be9 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 @@ -1,25 +1,75 @@ +/** + * Copyright © 2016-2019 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; +import com.datastax.driver.core.utils.UUIDs; import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; import io.grpc.stub.StreamObserver; import lombok.Data; import lombok.extern.slf4j.Slf4j; +import org.thingsboard.server.common.data.Dashboard; +import org.thingsboard.server.common.data.DataConstants; +import org.thingsboard.server.common.data.Device; +import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.EntityView; +import org.thingsboard.server.common.data.Event; +import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.edge.Edge; +import org.thingsboard.server.common.data.edge.EdgeQueueEntry; +import org.thingsboard.server.common.data.id.AssetId; +import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.kv.AttributeKvEntry; +import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; +import org.thingsboard.server.common.data.kv.LongDataEntry; +import org.thingsboard.server.common.data.page.TimePageData; +import org.thingsboard.server.common.data.page.TimePageLink; +import org.thingsboard.server.common.data.rule.RuleChain; +import org.thingsboard.server.dao.asset.AssetService; +import org.thingsboard.server.dao.attributes.AttributesService; +import org.thingsboard.server.dao.device.DeviceService; +import org.thingsboard.server.dao.edge.EdgeService; +import org.thingsboard.server.dao.util.mapping.JacksonUtil; +import org.thingsboard.server.gen.edge.AssetUpdateMsg; import org.thingsboard.server.gen.edge.ConnectRequestMsg; import org.thingsboard.server.gen.edge.ConnectResponseCode; import org.thingsboard.server.gen.edge.ConnectResponseMsg; -import org.thingsboard.server.gen.edge.EdgeConfigurationProto; +import org.thingsboard.server.gen.edge.DashboardUpdateMsg; +import org.thingsboard.server.gen.edge.DeviceUpdateMsg; +import org.thingsboard.server.gen.edge.EdgeConfiguration; +import org.thingsboard.server.gen.edge.EntityViewUpdateMsg; import org.thingsboard.server.gen.edge.RequestMsg; import org.thingsboard.server.gen.edge.RequestMsgType; import org.thingsboard.server.gen.edge.ResponseMsg; +import org.thingsboard.server.gen.edge.RuleChainUpdateMsg; +import org.thingsboard.server.gen.edge.UpdateMsgType; import org.thingsboard.server.gen.edge.UplinkMsg; import org.thingsboard.server.gen.edge.UplinkResponseMsg; import org.thingsboard.server.service.edge.EdgeContextComponent; +import java.util.Arrays; +import java.util.Collections; +import java.util.List; import java.util.Optional; import java.util.UUID; +import java.util.concurrent.ExecutionException; import java.util.function.BiConsumer; import java.util.function.Consumer; @@ -30,6 +80,7 @@ public final class EdgeGrpcSession implements Cloneable { private final UUID sessionId; private final BiConsumer sessionOpenListener; private final Consumer sessionCloseListener; + private final ObjectMapper objectMapper; private EdgeContextComponent ctx; private Edge edge; @@ -37,14 +88,24 @@ public final class EdgeGrpcSession implements Cloneable { private StreamObserver outputStream; private boolean connected; - EdgeGrpcSession(EdgeContextComponent ctx, StreamObserver outputStream - , BiConsumer sessionOpenListener - , Consumer sessionCloseListener) { + private EdgeService edgeService; + private AssetService assetService; + private DeviceService deviceService; + private AttributesService attributesService; + + EdgeGrpcSession(EdgeContextComponent ctx, StreamObserver outputStream, + BiConsumer sessionOpenListener, Consumer sessionCloseListener, + EdgeService edgeService, AssetService assetService, DeviceService deviceService, AttributesService attributesService, ObjectMapper objectMapper) { this.sessionId = UUID.randomUUID(); this.ctx = ctx; this.outputStream = outputStream; this.sessionOpenListener = sessionOpenListener; this.sessionCloseListener = sessionCloseListener; + this.objectMapper = objectMapper; + this.edgeService = edgeService; + this.assetService = assetService; + this.deviceService = deviceService; + this.attributesService = attributesService; initInputStream(); } @@ -83,6 +144,193 @@ public final class EdgeGrpcSession implements Cloneable { }; } + void processHandleMessages() throws ExecutionException, InterruptedException { + Long queueStartTs = getQueueStartTs().get(); + // TODO: this 100 value must be chagned properly + TimePageLink pageLink = new TimePageLink(30, queueStartTs + 1000); + TimePageData pageData; + UUID ifOffset = null; + do { + pageData = edgeService.findQueueEvents(edge.getTenantId(), edge.getId(), pageLink); + if (!pageData.getData().isEmpty()) { + for (Event event : pageData.getData()) { + EdgeQueueEntry entry; + try { + entry = objectMapper.treeToValue(event.getBody(), EdgeQueueEntry.class); + UpdateMsgType msgType = getResponseMsgType(entry.getType()); + switch (entry.getEntityType()) { + case DEVICE: + Device device = objectMapper.readValue(entry.getData(), Device.class); + onDeviceUpdated(msgType, device); + break; + case ASSET: + Asset asset = objectMapper.readValue(entry.getData(), Asset.class); + onAssetUpdated(msgType, asset); + break; + case ENTITY_VIEW: + EntityView entityView = objectMapper.readValue(entry.getData(), EntityView.class); + onEntityViewUpdated(msgType, entityView); + break; + case DASHBOARD: + Dashboard dashboard = objectMapper.readValue(entry.getData(), Dashboard.class); + onDashboardUpdated(msgType, dashboard); + break; + case RULE_CHAIN: + RuleChain ruleChain = objectMapper.readValue(entry.getData(), RuleChain.class); + onRuleChainUpdated(msgType, ruleChain); + break; + } + } catch (Exception e) { + log.error("Exception during processing records from queue", e); + } + ifOffset = event.getUuidId(); + } + } + if (pageData.hasNext()) { + pageLink = pageData.getNextPageLink(); + } + } while (pageData.hasNext()); + + if (ifOffset != null) { + Long newStartTs = UUIDs.unixTimestamp(ifOffset); + updateQueueStartTs(newStartTs); + } + try { + Thread.sleep(1000); + } catch (InterruptedException e) { + log.error("Error during sleep", e); + } + } + + private void updateQueueStartTs(Long newStartTs) { + List attributes = Collections.singletonList(new BaseAttributeKvEntry(new LongDataEntry("queueStartTs", newStartTs), System.currentTimeMillis())); + attributesService.save(edge.getTenantId(), edge.getId(), DataConstants.SERVER_SCOPE, attributes); + } + + private ListenableFuture getQueueStartTs() { + ListenableFuture> future = + attributesService.find(edge.getTenantId(), edge.getId(), DataConstants.SERVER_SCOPE, "queueStartTs"); + return Futures.transform(future, attributeKvEntryOpt -> { + if (attributeKvEntryOpt != null && attributeKvEntryOpt.isPresent()) { + AttributeKvEntry attributeKvEntry = attributeKvEntryOpt.get(); + return attributeKvEntry.getLongValue().isPresent() ? attributeKvEntry.getLongValue().get() : 0L; + } else { + return 0L; + } + } ); + } + + private void onDeviceUpdated(UpdateMsgType msgType, Device device) { + outputStream.onNext(ResponseMsg.newBuilder() + .setDeviceUpdateMsg(constructDeviceUpdatedMsg(msgType, device)) + .build()); + } + + private void onAssetUpdated(UpdateMsgType msgType, Asset asset) { + outputStream.onNext(ResponseMsg.newBuilder() + .setAssetUpdateMsg(constructAssetUpdatedMsg(msgType, asset)) + .build()); + } + + private void onEntityViewUpdated(UpdateMsgType msgType, EntityView entityView) { + outputStream.onNext(ResponseMsg.newBuilder() + .setEntityViewUpdateMsg(constructEntityViewUpdatedMsg(msgType, entityView)) + .build()); + } + + private void onRuleChainUpdated(UpdateMsgType msgType, RuleChain ruleChain) { + outputStream.onNext(ResponseMsg.newBuilder() + .setRuleChainUpdateMsg(constructRuleChainUpdatedMsg(msgType, ruleChain)) + .build()); + } + + private void onDashboardUpdated(UpdateMsgType msgType, Dashboard dashboard) { + outputStream.onNext(ResponseMsg.newBuilder() + .setDashboardUpdateMsg(constructDashboardUpdatedMsg(msgType, dashboard)) + .build()); + } + + private UpdateMsgType getResponseMsgType(String msgType) { + switch (msgType) { + case DataConstants.ENTITY_UPDATED: + return UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE; + case DataConstants.ENTITY_CREATED: + case DataConstants.ENTITY_ASSIGNED_TO_EDGE: + return UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE; + case DataConstants.ENTITY_DELETED: + case DataConstants.ENTITY_UNASSIGNED_FROM_EDGE: + return UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE; + default: + throw new RuntimeException("Unsupported mstType [" + msgType + "]"); + } + } + + private RuleChainUpdateMsg constructRuleChainUpdatedMsg(UpdateMsgType msgType, RuleChain ruleChain) { + RuleChainUpdateMsg.Builder builder = RuleChainUpdateMsg.newBuilder() + .setMsgType(msgType) + .setIdMSB(ruleChain.getId().getId().getMostSignificantBits()) + .setIdLSB(ruleChain.getId().getId().getLeastSignificantBits()) + .setName(ruleChain.getName()) + .setRoot(ruleChain.isRoot()) + .setDebugMode(ruleChain.isDebugMode()) + .setConfiguration(JacksonUtil.toString(ruleChain.getConfiguration())); + if (ruleChain.getFirstRuleNodeId() != null) { + builder.setFirstRuleNodeIdMSB(ruleChain.getFirstRuleNodeId().getId().getMostSignificantBits()) + .setFirstRuleNodeIdLSB(ruleChain.getFirstRuleNodeId().getId().getLeastSignificantBits()); + } + return builder.build(); + } + + private DashboardUpdateMsg constructDashboardUpdatedMsg(UpdateMsgType msgType, Dashboard dashboard) { + DashboardUpdateMsg.Builder builder = DashboardUpdateMsg.newBuilder() + .setMsgType(msgType) + .setIdMSB(dashboard.getId().getId().getMostSignificantBits()) + .setIdLSB(dashboard.getId().getId().getLeastSignificantBits()) + .setName(dashboard.getName()); + return builder.build(); + } + + private DeviceUpdateMsg constructDeviceUpdatedMsg(UpdateMsgType msgType, Device device) { + DeviceUpdateMsg.Builder builder = DeviceUpdateMsg.newBuilder() + .setMsgType(msgType) + .setName(device.getName()) + .setType(device.getName()); + return builder.build(); + } + + private AssetUpdateMsg constructAssetUpdatedMsg(UpdateMsgType msgType, Asset asset) { + AssetUpdateMsg.Builder builder = AssetUpdateMsg.newBuilder() + .setMsgType(msgType) + .setName(asset.getName()) + .setType(asset.getName()); + return builder.build(); + } + + private EntityViewUpdateMsg constructEntityViewUpdatedMsg(UpdateMsgType msgType, EntityView entityView) { + String relatedName; + String relatedType; + org.thingsboard.server.gen.edge.EntityType relatedEntityType; + if (entityView.getEntityId().getEntityType().equals(EntityType.DEVICE)) { + Device device = deviceService.findDeviceById(entityView.getTenantId(), new DeviceId(entityView.getEntityId().getId())); + relatedName = device.getName(); + relatedType = device.getType(); + relatedEntityType = org.thingsboard.server.gen.edge.EntityType.DEVICE; + } else { + Asset asset = assetService.findAssetById(entityView.getTenantId(), new AssetId(entityView.getEntityId().getId())); + relatedName = asset.getName(); + relatedType = asset.getType(); + relatedEntityType = org.thingsboard.server.gen.edge.EntityType.ASSET; + } + EntityViewUpdateMsg.Builder builder = EntityViewUpdateMsg.newBuilder() + .setMsgType(msgType) + .setName(entityView.getName()) + .setType(entityView.getName()) + .setRelatedName(relatedName) + .setRelatedType(relatedType) + .setRelatedEntityType(relatedEntityType); + return builder.build(); + } + private UplinkResponseMsg processUplinkMsg(UplinkMsg uplinkMsg) { return null; } @@ -103,23 +351,23 @@ public final class EdgeGrpcSession implements Cloneable { return ConnectResponseMsg.newBuilder() .setResponseCode(ConnectResponseCode.BAD_CREDENTIALS) .setErrorMsg("Failed to validate the edge!") - .setConfiguration(EdgeConfigurationProto.getDefaultInstance()).build(); + .setConfiguration(EdgeConfiguration.getDefaultInstance()).build(); } catch (Exception e) { log.error("[{}] Failed to process edge connection!", request.getEdgeRoutingKey(), e); return ConnectResponseMsg.newBuilder() .setResponseCode(ConnectResponseCode.SERVER_UNAVAILABLE) .setErrorMsg("Failed to process edge connection!") - .setConfiguration(EdgeConfigurationProto.getDefaultInstance()).build(); + .setConfiguration(EdgeConfiguration.getDefaultInstance()).build(); } } return ConnectResponseMsg.newBuilder() .setResponseCode(ConnectResponseCode.BAD_CREDENTIALS) .setErrorMsg("Failed to find the edge! Routing key: " + request.getEdgeRoutingKey()) - .setConfiguration(EdgeConfigurationProto.getDefaultInstance()).build(); + .setConfiguration(EdgeConfiguration.getDefaultInstance()).build(); } - private EdgeConfigurationProto constructEdgeConfigProto(Edge edge) throws JsonProcessingException { - return EdgeConfigurationProto.newBuilder() + private EdgeConfiguration constructEdgeConfigProto(Edge edge) throws JsonProcessingException { + return EdgeConfiguration.newBuilder() .setTenantIdMSB(edge.getTenantId().getId().getMostSignificantBits()) .setTenantIdLSB(edge.getTenantId().getId().getLeastSignificantBits()) .setName(edge.getName()) 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 5cc093fcad..b9cafde6b6 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 @@ -15,15 +15,22 @@ */ package org.thingsboard.server.dao.edge; +import com.fasterxml.jackson.databind.JsonNode; import com.google.common.util.concurrent.ListenableFuture; import org.thingsboard.server.common.data.EntitySubtype; +import org.thingsboard.server.common.data.Event; +import org.thingsboard.server.common.data.Tenant; 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.EdgeId; +import org.thingsboard.server.common.data.id.EntityId; 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 org.thingsboard.server.common.msg.TbMsg; import java.util.List; import java.util.Optional; @@ -66,6 +73,9 @@ public interface EdgeService { ListenableFuture> findEdgeTypesByTenantId(TenantId tenantId); + void pushEventToEdge(TenantId tenantId, TbMsg tbMsg); + + TimePageData findQueueEvents(TenantId tenantId, EdgeId edgeId, TimePageLink pageLink); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java index 2e2130c4d4..7090039d6d 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java @@ -56,6 +56,8 @@ public class DataConstants { public static final String ATTRIBUTES_DELETED = "ATTRIBUTES_DELETED"; public static final String ALARM_ACK = "ALARM_ACK"; public static final String ALARM_CLEAR = "ALARM_CLEAR"; + public static final String ENTITY_ASSIGNED_TO_EDGE = "ENTITY_ASSIGNED_TO_EDGE"; + public static final String ENTITY_UNASSIGNED_FROM_EDGE = "ENTITY_UNASSIGNED_FROM_EDGE"; public static final String RPC_CALL_FROM_SERVER_TO_DEVICE = "RPC_CALL_FROM_SERVER_TO_DEVICE"; @@ -63,4 +65,6 @@ public class DataConstants { public static final String SECRET_KEY_FIELD_NAME = "secretKey"; public static final String DURATION_MS_FIELD_NAME = "durationMs"; + public static final String EDGE_QUEUE_EVENT_TYPE = "EDGE_QUEUE"; + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeQueueEntry.java b/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeQueueEntry.java new file mode 100644 index 0000000000..4cfc210267 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeQueueEntry.java @@ -0,0 +1,11 @@ +package org.thingsboard.server.common.data.edge; + +import lombok.Data; +import org.thingsboard.server.common.data.EntityType; + +@Data +public class EdgeQueueEntry { + private String type; + private EntityType entityType; + private String data; +} diff --git a/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java b/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java index fd18407b5c..53c443ef1f 100644 --- a/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java +++ b/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java @@ -24,15 +24,20 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; import org.thingsboard.edge.exception.EdgeConnectionException; -import org.thingsboard.server.gen.edge.CloudDownlinkDataProto; +import org.thingsboard.server.gen.edge.AssetUpdateMsg; import org.thingsboard.server.gen.edge.ConnectRequestMsg; import org.thingsboard.server.gen.edge.ConnectResponseCode; import org.thingsboard.server.gen.edge.ConnectResponseMsg; -import org.thingsboard.server.gen.edge.EdgeConfigurationProto; +import org.thingsboard.server.gen.edge.DashboardUpdateMsg; +import org.thingsboard.server.gen.edge.DeviceUpdateMsg; +import org.thingsboard.server.gen.edge.DownlinkMsg; +import org.thingsboard.server.gen.edge.EdgeConfiguration; import org.thingsboard.server.gen.edge.EdgeRpcServiceGrpc; +import org.thingsboard.server.gen.edge.EntityViewUpdateMsg; import org.thingsboard.server.gen.edge.RequestMsg; import org.thingsboard.server.gen.edge.RequestMsgType; import org.thingsboard.server.gen.edge.ResponseMsg; +import org.thingsboard.server.gen.edge.RuleChainUpdateMsg; import org.thingsboard.server.gen.edge.UplinkMsg; import org.thingsboard.server.gen.edge.UplinkResponseMsg; @@ -65,8 +70,13 @@ public class EdgeGrpcClient implements EdgeRpcClient { public void connect(String edgeKey, String edgeSecret, Consumer onUplinkResponse, - Consumer onEdgeUpdate, - Consumer onDownlink, + Consumer onEdgeUpdate, + Consumer onDeviceUpdate, + Consumer onAssetUpdate, + Consumer onEntityViewUpdate, + Consumer onRuleChainUpdate, + Consumer onDashboardUpdate, + Consumer onDownlink, Consumer onError) { NettyChannelBuilder builder = NettyChannelBuilder.forAddress(rpcHost, rpcPort).usePlaintext(); if (sslEnabled) { @@ -80,7 +90,7 @@ public class EdgeGrpcClient implements EdgeRpcClient { channel = builder.build(); EdgeRpcServiceGrpc.EdgeRpcServiceStub stub = EdgeRpcServiceGrpc.newStub(channel); log.info("[{}] Sending a connect request to the TB!", edgeKey); - this.inputStream = stub.handleMsgs(initOutputStream(edgeKey, onUplinkResponse, onEdgeUpdate, onDownlink, onError)); + this.inputStream = stub.handleMsgs(initOutputStream(edgeKey, onUplinkResponse, onEdgeUpdate, onDeviceUpdate, onAssetUpdate, onEntityViewUpdate, onRuleChainUpdate, onDashboardUpdate, onDownlink, onError)); this.inputStream.onNext(RequestMsg.newBuilder() .setMsgType(RequestMsgType.CONNECT_RPC_MESSAGE) .setConnectRequestMsg(ConnectRequestMsg.newBuilder().setEdgeRoutingKey(edgeKey).setEdgeSecret(edgeSecret).build()) @@ -103,7 +113,16 @@ public class EdgeGrpcClient implements EdgeRpcClient { .build()); } - private StreamObserver initOutputStream(String edgeKey, Consumer onUplinkResponse, Consumer onEdgeUpdate, Consumer onDownlink, Consumer onError) { + private StreamObserver initOutputStream(String edgeKey, + Consumer onUplinkResponse, + Consumer onEdgeUpdate, + Consumer onDeviceUpdate, + Consumer onAssetUpdate, + Consumer onEntityViewUpdate, + Consumer onRuleChainUpdate, + Consumer onDashboardUpdate, + Consumer onDownlink, + Consumer onError) { return new StreamObserver() { @Override public void onNext(ResponseMsg responseMsg) { @@ -119,9 +138,24 @@ public class EdgeGrpcClient implements EdgeRpcClient { } else if (responseMsg.hasUplinkResponseMsg()) { log.debug("[{}] Uplink response message received {}", edgeKey, responseMsg.getUplinkResponseMsg()); onUplinkResponse.accept(responseMsg.getUplinkResponseMsg()); + } else if (responseMsg.hasDeviceUpdateMsg()) { + log.debug("[{}] Device update message received {}", edgeKey, responseMsg.getDeviceUpdateMsg()); + onDeviceUpdate.accept(responseMsg.getDeviceUpdateMsg()); + } else if (responseMsg.hasAssetUpdateMsg()) { + log.debug("[{}] Asset update message received {}", edgeKey, responseMsg.getAssetUpdateMsg()); + onAssetUpdate.accept(responseMsg.getAssetUpdateMsg()); + } else if (responseMsg.hasEntityViewUpdateMsg()) { + log.debug("[{}] EntityView update message received {}", edgeKey, responseMsg.getEntityViewUpdateMsg()); + onEntityViewUpdate.accept(responseMsg.getEntityViewUpdateMsg()); + } else if (responseMsg.hasRuleChainUpdateMsg()) { + log.debug("[{}] Rule Chain udpate message received {}", edgeKey, responseMsg.getRuleChainUpdateMsg()); + onRuleChainUpdate.accept(responseMsg.getRuleChainUpdateMsg()); + } else if (responseMsg.hasDashboardUpdateMsg()) { + log.debug("[{}] Dashboard message received {}", edgeKey, responseMsg.getDashboardUpdateMsg()); + onDashboardUpdate.accept(responseMsg.getDashboardUpdateMsg()); } else if (responseMsg.hasDownlinkMsg()) { - log.debug("[{}] Downlink message received for device {}", edgeKey, responseMsg.getDownlinkMsg().getCloudData().getDeviceName()); - onDownlink.accept(responseMsg.getDownlinkMsg().getCloudData()); + log.debug("[{}] Downlink message received for rule chain {}", edgeKey, responseMsg.getDownlinkMsg()); + onDownlink.accept(responseMsg.getDownlinkMsg()); } } diff --git a/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java b/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java index 95e0dbdfd6..aa390a0d94 100644 --- a/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java +++ b/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java @@ -15,8 +15,13 @@ */ package org.thingsboard.edge.rpc; -import org.thingsboard.server.gen.edge.CloudDownlinkDataProto; -import org.thingsboard.server.gen.edge.EdgeConfigurationProto; +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.DownlinkMsg; +import org.thingsboard.server.gen.edge.EdgeConfiguration; +import org.thingsboard.server.gen.edge.EntityViewUpdateMsg; +import org.thingsboard.server.gen.edge.RuleChainUpdateMsg; import org.thingsboard.server.gen.edge.UplinkMsg; import org.thingsboard.server.gen.edge.UplinkResponseMsg; @@ -27,11 +32,15 @@ public interface EdgeRpcClient { void connect(String integrationKey, String integrationSecret, Consumer onUplinkResponse, - Consumer onEdgeUpdate, - Consumer onDownlink, + Consumer onEdgeUpdate, + Consumer onDeviceUpdate, + Consumer onAssetUpdate, + Consumer onEntityViewUpdate, + Consumer onRuleChainUpdate, + Consumer onDashboardUpdate, + Consumer onDownlink, Consumer onError); - void disconnect() throws InterruptedException; void sendUplinkMsg(UplinkMsg uplinkMsg) throws InterruptedException; diff --git a/common/edge-api/src/main/proto/edge.proto b/common/edge-api/src/main/proto/edge.proto index a2d3ce82f3..73fd216c70 100644 --- a/common/edge-api/src/main/proto/edge.proto +++ b/common/edge-api/src/main/proto/edge.proto @@ -34,17 +34,21 @@ service EdgeRpcService { message RequestMsg { RequestMsgType msgType = 1; ConnectRequestMsg connectRequestMsg = 2; - UplinkMsg uplinkMsg = 3; + DeviceUpdateMsg deviceUpdateMsg = 3; + UplinkMsg uplinkMsg = 4; } message ResponseMsg { - ResponseMsgType msgType = 1; - ConnectResponseMsg connectResponseMsg = 2; - UplinkResponseMsg uplinkResponseMsg = 3; - DownlinkMsg downlinkMsg = 4; + ConnectResponseMsg connectResponseMsg = 1; + UplinkResponseMsg uplinkResponseMsg = 2; + DeviceUpdateMsg deviceUpdateMsg = 3; + RuleChainUpdateMsg ruleChainUpdateMsg = 4; + DashboardUpdateMsg dashboardUpdateMsg = 5; + AssetUpdateMsg assetUpdateMsg = 6; + EntityViewUpdateMsg entityViewUpdateMsg = 7; + DownlinkMsg downlinkMsg = 8; } - enum RequestMsgType { CONNECT_RPC_MESSAGE = 0; UPLINK_RPC_MESSAGE = 1; @@ -64,10 +68,10 @@ enum ConnectResponseCode { message ConnectResponseMsg { ConnectResponseCode responseCode = 1; string errorMsg = 2; - EdgeConfigurationProto configuration = 3; + EdgeConfiguration configuration = 3; } -message EdgeConfigurationProto { +message EdgeConfiguration { int64 tenantIdMSB = 1; int64 tenantIdLSB = 2; string name = 5; @@ -75,23 +79,98 @@ message EdgeConfigurationProto { string type = 7; } -enum ResponseMsgType { - SAVE_ENTITY_MESSAGE = 0; - DELETE_ENTITY_MESSAGE = 1; +enum UpdateMsgType { + ENTITY_CREATED_RPC_MESSAGE = 0; + ENTITY_UPDATED_RPC_MESSAGE = 1; + ENTITY_DELETED_RPC_MESSAGE = 2; } -message CloudDownlinkDataProto { +message DeviceData { string deviceName = 1; string deviceType = 2; bytes tbMsg = 3; } +message AssetData { + string assetName = 1; + string assetType = 2; + bytes tbMsg = 3; +} + +message EntityViewData { + string entityViewName = 1; + string entityViewType = 2; + bytes tbMsg = 3; +} + +message RuleChainData { + string ruleChainName = 1; + string ruleChainType = 2; + bytes tbMsg = 3; +} + +message DashboardData { + string dashboardName = 1; + string dashboardType = 2; + bytes tbMsg = 3; +} + +message RuleChainUpdateMsg { + UpdateMsgType msgType = 1; + int64 idMSB = 2; + int64 idLSB = 3; + string name = 4; + int64 firstRuleNodeIdMSB = 5; + int64 firstRuleNodeIdLSB = 6; + bool root = 7; + bool debugMode = 8; + string configuration = 9; +} + +message DashboardUpdateMsg { + UpdateMsgType msgType = 1; + int64 idMSB = 2; + int64 idLSB = 3; + string name = 4; +} + +message DeviceUpdateMsg { + UpdateMsgType msgType = 1; + string name = 2; + string type = 3; +} + +message AssetUpdateMsg { + UpdateMsgType msgType = 1; + string name = 2; + string type = 3; +} + +message EntityViewUpdateMsg { + UpdateMsgType msgType = 1; + string name = 2; + string type = 3; + string relatedName = 4; + string relatedType = 5; + EntityType relatedEntityType = 6; +} + +enum EntityType { + DEVICE = 0; + ASSET = 1; +} + /** * Main Messages; */ message UplinkMsg { int32 uplinkMsgId = 1; + repeated DeviceData deviceData = 2; + repeated AssetData assetData = 3; + repeated EntityViewData entityViewData = 4; + repeated RuleChainData ruleChainData = 5; + repeated DashboardData dashboardData = 6; } message UplinkResponseMsg { @@ -100,5 +179,11 @@ message UplinkResponseMsg { } message DownlinkMsg { - CloudDownlinkDataProto cloudData = 1; + int32 downlinkMsgId = 1; + repeated DeviceData deviceData = 2; + repeated AssetData assetData = 3; + repeated EntityViewData entityViewData = 4; + repeated RuleChainData ruleChainData = 5; + repeated DashboardData dashboardData = 6; } + diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/BaseEdgeService.java b/dao/src/main/java/org/thingsboard/server/dao/edge/BaseEdgeService.java index 7deb6844b0..95b61d4053 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/BaseEdgeService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/BaseEdgeService.java @@ -5,7 +5,7 @@ * 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 + * 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, @@ -15,9 +15,12 @@ */ package org.thingsboard.server.dao.edge; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; import com.google.common.base.Function; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; +import lombok.Data; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.cache.Cache; @@ -27,10 +30,14 @@ import org.springframework.cache.annotation.Cacheable; import org.springframework.stereotype.Service; import org.springframework.util.StringUtils; import org.thingsboard.server.common.data.Customer; +import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.EntitySubtype; import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.Event; +import org.thingsboard.server.common.data.ShortEdgeInfo; import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.edge.Edge; +import org.thingsboard.server.common.data.edge.EdgeQueueEntry; import org.thingsboard.server.common.data.edge.EdgeSearchQuery; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.EdgeId; @@ -38,23 +45,31 @@ import org.thingsboard.server.common.data.id.EntityId; 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 org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.EntitySearchDirection; +import org.thingsboard.server.common.data.rule.RuleChain; +import org.thingsboard.server.common.msg.TbMsg; 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.event.EventService; import org.thingsboard.server.dao.exception.DataValidationException; +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 javax.annotation.Nullable; +import java.io.IOException; import java.util.ArrayList; import java.util.Collections; import java.util.Comparator; import java.util.List; import java.util.Optional; +import java.util.UUID; import java.util.stream.Collectors; import static org.thingsboard.server.common.data.CacheConstants.EDGE_CACHE; @@ -69,6 +84,8 @@ import static org.thingsboard.server.dao.service.Validator.validateString; @Slf4j public class BaseEdgeService extends AbstractEntityService implements EdgeService { + private static final ObjectMapper mapper = new ObjectMapper(); + public static final String INCORRECT_TENANT_ID = "Incorrect tenantId "; public static final String INCORRECT_PAGE_LINK = "Incorrect page link "; public static final String INCORRECT_CUSTOMER_ID = "Incorrect customerId "; @@ -86,9 +103,15 @@ public class BaseEdgeService extends AbstractEntityService implements EdgeServic @Autowired private CacheManager cacheManager; + @Autowired + private EventService eventService; + @Autowired private DashboardService dashboardService; + @Autowired + private RuleChainService ruleChainService; + @Override public Edge findEdgeById(TenantId tenantId, EdgeId edgeId) { log.trace("Executing findEdgeById [{}]", edgeId); @@ -150,7 +173,8 @@ public class BaseEdgeService extends AbstractEntityService implements EdgeServic Edge edge = edgeDao.findById(tenantId, edgeId.getId()); - deleteEntityRelations(tenantId, edgeId); + dashboardService.unassignEdgeDashboards(tenantId, edgeId); + ruleChainService.unassignEdgeRuleChains(tenantId, edgeId); List list = new ArrayList<>(); list.add(edge.getTenantId()); @@ -158,7 +182,7 @@ public class BaseEdgeService extends AbstractEntityService implements EdgeServic Cache cache = cacheManager.getCache(EDGE_CACHE); cache.evict(list); - dashboardService.unassignEdgeDashboards(tenantId, edgeId); + deleteEntityRelations(tenantId, edgeId); edgeDao.removeById(tenantId, edgeId.getId()); } @@ -275,6 +299,106 @@ public class BaseEdgeService extends AbstractEntityService implements EdgeServic }); } + @Override + public void pushEventToEdge(TenantId tenantId, TbMsg tbMsg) { + try { + switch (tbMsg.getOriginator().getEntityType()) { + case ASSET: + processAsset(tenantId, tbMsg); + break; + case DEVICE: + processDevice(tenantId, tbMsg); + break; + case DASHBOARD: + processDashboard(tenantId, tbMsg); + break; + case RULE_CHAIN: + processRuleChain(tenantId, tbMsg); + break; + case ENTITY_VIEW: + processEntityView(tenantId, tbMsg); + break; + default: + log.debug("Entity type [{}] is not designed to be pushed to edge", tbMsg.getOriginator().getEntityType()); + } + } catch (IOException e) { + log.error("Can't push to edge updates, entity type [{}], data [{}]", tbMsg.getOriginator().getEntityType(), tbMsg.getData(), e); + } + + + } + + private void processDevice(TenantId tenantId, TbMsg tbMsg) { + // TODO + } + + private void processDashboard(TenantId tenantId, TbMsg tbMsg) { + processAssignedEntity(tenantId, tbMsg, EntityType.DASHBOARD); + } + + private void processEntityView(TenantId tenantId, TbMsg tbMsg) { + // TODO + } + + private void processAsset(TenantId tenantId, TbMsg tbMsg) { + // TODO + } + + private void processAssignedEntity(TenantId tenantId, TbMsg tbMsg, EntityType entityType) { + EdgeId edgeId; + switch (tbMsg.getType()) { + case DataConstants.ENTITY_ASSIGNED_TO_EDGE: + edgeId = new EdgeId(UUID.fromString(tbMsg.getMetaData().getValue("assignedEdgeId"))); + pushEventToEdge(tenantId, edgeId, tbMsg.getType(), entityType, tbMsg.getData()); + break; + case DataConstants.ENTITY_UNASSIGNED_FROM_EDGE: + edgeId = new EdgeId(UUID.fromString(tbMsg.getMetaData().getValue("unassignedEdgeId"))); + pushEventToEdge(tenantId, edgeId, tbMsg.getType(), entityType, tbMsg.getData()); + break; + } + } + + private void processRuleChain(TenantId tenantId, TbMsg tbMsg) throws IOException { + switch (tbMsg.getType()) { + case DataConstants.ENTITY_ASSIGNED_TO_EDGE: + case DataConstants.ENTITY_UNASSIGNED_FROM_EDGE: + processAssignedEntity(tenantId, tbMsg, EntityType.RULE_CHAIN); + break; + case DataConstants.ENTITY_DELETED: + case DataConstants.ENTITY_CREATED: + case DataConstants.ENTITY_UPDATED: + RuleChain ruleChain = mapper.readValue(tbMsg.getData(), RuleChain.class); + for (ShortEdgeInfo assignedEdge : ruleChain.getAssignedEdges()) { + pushEventToEdge(tenantId, assignedEdge.getEdgeId(), tbMsg.getType(), EntityType.RULE_CHAIN, tbMsg.getData()); + } + break; + default: + log.warn("Unsupported message type " + tbMsg.getType()); + } + + } + + private void pushEventToEdge(TenantId tenantId, EdgeId edgeId, String type, EntityType entityType, String data) { + log.debug("Pushing event to edge queue. tenantId [{}], edgeId [{}], type [{}], data [{}]", tenantId, edgeId, type, data); + + EdgeQueueEntry queueEntry = new EdgeQueueEntry(); + queueEntry.setType(type); + queueEntry.setEntityType(entityType); + queueEntry.setData(data); + + Event event = new Event(); + event.setEntityId(edgeId); + event.setTenantId(tenantId); + event.setType(DataConstants.EDGE_QUEUE_EVENT_TYPE); + event.setBody(mapper.valueToTree(queueEntry)); + eventService.saveAsync(event); + } + + @Override + public TimePageData findQueueEvents(TenantId tenantId, EdgeId edgeId, TimePageLink pageLink) { + return eventService.findEvents(tenantId, edgeId, DataConstants.EDGE_QUEUE_EVENT_TYPE, pageLink); + } + private DataValidator edgeValidator = new DataValidator() { diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/nosql/EdgeEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/nosql/EdgeEntity.java index de8d81f659..7d1deff8b6 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/nosql/EdgeEntity.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/nosql/EdgeEntity.java @@ -5,7 +5,7 @@ * 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 + * 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, diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/EdgeEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sql/EdgeEntity.java index 6fce28077f..8f1d37d1c3 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sql/EdgeEntity.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sql/EdgeEntity.java @@ -5,7 +5,7 @@ * 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 + * 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, diff --git a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java index 28af5707ba..2da4203339 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.rule; +import com.fasterxml.jackson.databind.ObjectMapper; import com.google.common.base.Function; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; @@ -43,6 +44,7 @@ import org.thingsboard.server.common.data.rule.RuleChainConnectionInfo; import org.thingsboard.server.common.data.rule.RuleChainMetaData; import org.thingsboard.server.common.data.rule.RuleNode; import org.thingsboard.server.dao.edge.EdgeDao; +import org.thingsboard.server.dao.edge.EdgeService; import org.thingsboard.server.dao.entity.AbstractEntityService; import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.dao.service.DataValidator; @@ -66,6 +68,8 @@ import java.util.stream.Collectors; @Slf4j public class BaseRuleChainService extends AbstractEntityService implements RuleChainService { + private static final ObjectMapper objectMapper = new ObjectMapper(); + @Autowired private RuleChainDao ruleChainDao; @@ -78,6 +82,9 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC @Autowired private EdgeDao edgeDao; + @Autowired + private EdgeService edgeService; + @Override public RuleChain saveRuleChain(RuleChain ruleChain) { ruleChainValidator.validate(ruleChain, RuleChain::getTenantId); @@ -381,10 +388,9 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC log.warn("[{}] Failed to create ruleChain relation. Edge Id: [{}]", ruleChainId, edgeId); throw new RuntimeException(e); } - return saveRuleChain(ruleChain); - } else { - return ruleChain; + ruleChain = saveRuleChain(ruleChain); } + return ruleChain; } @Override diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaEdgeDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaEdgeDao.java index 64d5412169..14ba59d796 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaEdgeDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaEdgeDao.java @@ -5,7 +5,7 @@ * 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 + * 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,